Ir al contenido principal

Tutorial de Apache Spark: ML con PySpark

Este tutorial de Apache Spark te introduce al procesamiento de big data, el análisis y el ML con PySpark.
Actualizado 17 sept 2026  · 15 min leer

Explorar con IA

ChatGPTClaudePerplexity

Apache Spark y Python para big data y aprendizaje automático

Apache Spark es un motor rápido, fácil de usar y generalista para el procesamiento de big data que incluye módulos integrados para streaming, SQL, aprendizaje automático (ML) y procesamiento de grafos. Esta tecnología es una habilidad muy demandada en ingeniería de datos, y también resulta útil para científicos de datos al hacer análisis exploratorio (EDA), extracción de características y, por supuesto, ML.

En este tutorial conectarás Spark con Python a través de PySpark, la API de Python que expone el modelo de programación de Spark. Concretamente, te centrarás en:

Tutorial de Apache Spark

Si te interesa más usar Spark con R, echa un vistazo al curso gratuito de DataCamp Introducción a Spark en R con sparklyr o descarga la chuleta de PySpark SQL.

Instalar Apache Spark

Instalar Spark y hacerlo funcionar puede ser un reto. En esta sección verás los pasos para instalarlo en tu PC.

Lo primero es comprobar si cumples los prerrequisitos. Spark está escrito en Scala y se ejecuta en la máquina virtual de Java (JVM). Por eso, necesitas comprobar si tienes instalado el Java Development Kit (JDK). Lo haces porque el JDK te proporciona una o más implementaciones de la JVM. Lo ideal es instalar la versión más reciente; en el momento de escribir esto, JDK8.

¡Después ya puedes descargar Spark!

Descargar pyspark con pip

También puedes descargar e instalar PySpark con pip. Es bastante sencillo, como instalar cualquier otro paquete. Ejecuta el comando habitual y se hará el trabajo pesado por ti:

$ pip install pyspark

Como alternativa, puedes ir a la página de descargas de Spark. Deja las opciones por defecto en los tres primeros pasos y encontrarás un enlace de descarga en el paso 4. Haz clic en ese enlace para descargarlo. Para este tutorial, descargarás la versión 2.2.0 de Spark y el paquete “Pre-built for Apache Hadoop 2.7 and later”.

Nota: la descarga puede tardar un rato en completarse.

Descargar Spark con Homebrew

También puedes instalar Spark con Homebrew, un gestor de paquetes libre y de código abierto. Es especialmente útil si trabajas con macOS.

Simplemente ejecuta los siguientes comandos para buscar Spark, obtener más información e instalarlo finalmente en tu ordenador:

# Buscar spark
$ brew search spark

# Obtener más información sobre apache-spark
$ brew info apache-spark

# Instalar apache-spark
$ brew install apache-spark

Descargar y configurar Spark

Después, descomprime el archivo que aparece en tu carpeta Downloads. Puede hacerse automáticamente haciendo doble clic en el archivo spark-2.2.0-bin-hadoop2.7.tgz o abriendo la Terminal y ejecutando el siguiente comando:

$ tar xvf spark-2.2.0-bin-hadoop2.7.tgz

Luego mueve la carpeta descomprimida a /usr/local/spark con esta línea:

$ mv spark-2.1.0-bin-hadoop2.7 /usr/local/spark

Nota: si te aparece un error de permiso denegado al mover la carpeta a la nueva ubicación, añade sudo delante del comando. La línea quedará como $ sudo mv spark-2.1.0-bin-hadoop2.7 /usr/local/spark. Se te pedirá la contraseña, que suele ser la misma que usas para desbloquear tu equipo al iniciarlo :)

Ahora que está todo listo, abre el archivo README en la ruta /usr/local/spark. Puedes hacerlo ejecutando

$ cd /usr/local/spark

Esto te lleva a la carpeta que necesitas. A continuación, puedes inspeccionarla y leer el README que incluye.

Primero, usa $ ls para listar los archivos y carpetas que hay en spark. Verás un archivo README.md. Puedes abrirlo ejecutando uno de estos comandos:

# Abrir y editar el archivo
$ nano README.md

# Solo leer el archivo 
$ cat README.md

Consejo: usa la tecla tabulador para autocompletar mientras escribes el nombre del archivo :) Te ahorrará tiempo.

Verás que el README incluye información general sobre Spark, documentación online, compilación de Spark, las shells interactivas de Scala y Python, programas de ejemplo y mucho más.

Lo más interesante aquí puede ser la sección sobre cómo compilar Spark, pero solo es relevante si no has descargado una versión precompilada. En este tutorial has descargado una versión precompilada. Pulsa CTRL + X para salir del README y volver a la carpeta de Spark.

Si elegiste una versión que no está compilada, ejecuta el comando que aparece en el README. En el momento de escribir esto, es el siguiente:

$ build/mvn -DskipTests clean package run

Ten en cuenta que este comando puede tardar en ejecutarse.

Fundamentos de PySpark: RDDs

Ahora que has instalado Spark y PySpark, empecemos explorando la shell interactiva de Spark y repasando algunos básicos que necesitarás para ponerte en marcha. En el resto del tutorial trabajarás con PySpark en un cuaderno de Jupyter.

Aplicaciones de Spark vs. Spark shell

La shell interactiva es un entorno de tipo Read-Eval(uate)-Print-Loop (REPL): lo que escribes se lee, se evalúa y se imprime para que continúes con tu análisis. Te puede recordar a IPython, una potente shell interactiva de Python que conocerás por Jupyter. Si quieres saber más, lee el artículo de DataCamp IPython o Jupyter.

Puedes usar la shell, disponible tanto para Python como para Scala, para cualquier trabajo interactivo que necesites.

Además de la shell, también puedes escribir y desplegar aplicaciones de Spark. A diferencia de escribir aplicaciones, en la shell la SparkSession ya está creada para que puedas ponerte a trabajar sin perder tiempo creándola.

Quizá te preguntes: ¿qué es SparkSession?

Es el punto de entrada principal a la funcionalidad de Spark: representa la conexión a un clúster de Spark y te permite crear RDDs y difundir variables en ese clúster. Cuando trabajas con Spark, todo empieza y termina con esta SparkSession. Nota: antes de Spark 2.0.0 los tres objetos de conexión principales eran SparkContext, SqlContext y HiveContext.

Verás más sobre esto más adelante. Por ahora, centrémonos en la shell.

La shell de Spark en Python

Desde la carpeta spark en /usr/local/spark, puedes ejecutar

$ ./bin/pyspark

Al principio verás aparecer algo de texto. Luego verás “Spark”, así:

Python 2.7.13 (v2.7.13:a06454b1afa1, Dec 17 2016, 12:39:47) 
[GCC 4.2.1 (Apple Inc. build 5666) (dot 3)] on darwin
Type "help", "copyright", "credits" or "license" for more information.
Using Spark's default log4j profile: org/apache/spark/log4j-defaults.properties
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
17/07/26 11:41:26 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
17/07/26 11:41:47 WARN ObjectStore: Failed to get database global_temp, returning NoSuchObjectException
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /__ / .__/\_,_/_/ /_/\_\   version 2.2.0
      /_/

Using Python version 2.7.13 (v2.7.13:a06454b1afa1, Dec 17 2016 12:39:47)
SparkSession available as 'spark'.
>>>

Cuando veas esto, ya puedes empezar a experimentar en la shell interactiva.

Consejo: si prefieres usar la shell de IPython en lugar de la de Spark, puedes hacerlo definiendo esta variable de entorno:

export PYSPARK_DRIVER_PYTHON="/usr/local/ipython/bin/ipython"

Crear RDDs

Vamos a empezar en pequeño y crear un RDD, la unidad básica de Spark. Un RDD representa datos, pero no es un único objeto, colección de registros, conjunto de resultados o conjunto de datos. Está pensado para datos que residen en múltiples máquinas: un solo RDD puede estar repartido por miles de JVM, porque Spark particiona automáticamente los datos para lograr el paralelismo. Por supuesto, puedes ajustar el paralelismo para obtener más particiones. Por eso un RDD es realmente una colección de particiones.

Puedes crear fácilmente un RDD simple con la función parallelize() pasando unos datos (un iterable, como una lista o una colección):

>>> rdd1 = spark.sparkContext.parallelize([('a',7),('a',2),('b',2)])
>>> rdd2 = spark.sparkContext.parallelize([("a",["x","y","z"]), ("b",["p", "r"])])
>>> rdd3 = spark.sparkContext.parallelize(range(100))

Nota: el objeto SparkSession contiene el objeto SparkContext, al que accedes con spark.sparkContext. Por compatibilidad retroactiva, aún puedes usar sc, como en rdd1 = sc.parallelize(['a',7),('a',2),('b',2)]).

Operaciones sobre RDD

Ahora que has creado los RDDs, puedes operar en paralelo sobre los datos distribuidos en rdd1 y rdd2. Hay dos tipos de operaciones: transformaciones y acciones.

Para entender intuitivamente la diferencia, algunas de las transformaciones más comunes son map(), filter(), flatMap(), sample(), randomSplit(), coalesce() y repartition(), y algunas acciones habituales son reduce(), collect(), first(), take(), count(), saveAsHadoopFile().

Las transformaciones son perezosas y crean uno o varios RDDs nuevos; las acciones producen valores que no son RDD: devuelven un conjunto de resultados, un número, un archivo, …

Por ejemplo, puedes agregar todos los elementos de rdd1 con esta lambda sencilla y devolver los resultados al driver:

>>> rdd1.reduce(lambda a,b: a+b)

Esto te dará como resultado: ('a', 7, 'a', 2, 'b', 2). Otra transformación es flatMapValues(), que se aplica a RDDs de pares clave-valor como rdd2. En este caso, pasas cada valor por una función flatMap sin cambiar las claves (la lambda de abajo) y luego ejecutas una acción recopilando los resultados con collect().

>>> rdd2.flatMapValues(lambda x: x).collect()
[('a', 'x'), ('a', 'y'), ('a', 'z'), ('b', 'p'), ('b', 'r')]

Los datos

Ahora que has cubierto lo básico con la shell interactiva, es hora de trabajar con datos reales. Para este tutorial usarás el conjunto de datos de California Housing. Ojo: en realidad es “small data” y usar Spark aquí puede ser excesivo; este tutorial es didáctico y busca que veas cómo usar PySpark para construir un modelo de aprendizaje automático.

Cargar y explorar tus datos

Aunque ya sabes un poco más sobre tus datos, conviene explorarlos a fondo. Antes, sin embargo, vas a preparar Jupyter Notebook con Spark y dar los primeros pasos para definir el SparkContext.

PySpark en Jupyter Notebook

En esta parte no usarás la shell, sino que crearás tu propia aplicación en un cuaderno de Jupyter. Ya tienes todo lo necesario instalado, así que no hay mucho que hacer para que PySpark funcione en Jupyter.

Arranca el cuaderno como siempre con $ jupyter notebook. Crea un cuaderno nuevo, importa la librería findspark y usa la función init(). En este caso indicarás la ruta /usr/local/spark a init() porque sabes que ahí instalaste Spark.

# Importa findspark 
import findspark

# Inicializa e indica la ruta
findspark.init("/usr/local/spark")

# O usa esta alternativa
#findspark.init()

Consejo: si no tienes claro si tu ruta está bien configurada o dónde instalaste Spark, usa findspark.find() para detectar automáticamente la ubicación.

Si buscas otras formas de trabajar con Spark en Jupyter, consulta nuestra guía para principiantes de Apache Spark en Python.

Con esto listo, ya puedes crear tu primer programa en Spark.

Crear tu primer programa en Spark

Lo primero es importar SparkContext desde el paquete pyspark e inicializarlo. Recuerda que antes no necesitabas hacerlo porque la shell interactiva lo creaba automáticamente. Aquí tendrás que trabajar un poquito más :)

Importa el módulo SparkSession desde pyspark.sql y construye una SparkSession con el método builder(). Luego puedes definir la URL del master, el nombre de la aplicación, añadir configuración adicional como la memoria del executor y, por último, usar getOrCreate() para obtener la sesión actual o crear una nueva si no hay ninguna.

# Importar SparkSession
from pyspark.sql import SparkSession

# Construir la SparkSession
spark = SparkSession.builder \
   .master("local") \
   .appName("Linear Regression Model") \
   .config("spark.executor.memory", "1gb") \
   .getOrCreate()
   
sc = spark.sparkContext

Nota: si aparece un FileNotFoundError del estilo “No such file or directory: ‘/User/YourName/Downloads/spark-2.1.0-bin-hadoop2.7/./bin/spark-submit’”, tendrás que (re)definir tu PATH de Spark. Ve a tu directorio home con $ cd y edita el archivo .bash_profile con $ nano .bash_profile.

Añade algo como esto al final del archivo

export SPARK_HOME="/usr/local/spark"

Usa CTRL + X para salir y guarda los cambios con Y. Luego no olvides aplicarlos ejecutando source .bash_profile.

Consejo: también puedes definir variables de entorno adicionales si lo necesitas. Probablemente no te harán falta, pero está bien saberlo. Por ejemplo:

# Fijar un valor para la semilla de hash
export PYTHONHASHSEED=0

# Usar un ejecutable de Python alternativo
export PYSPARK_PYTHON=/usr/local/ipython/bin/ipython

# Ampliar la ruta de búsqueda por defecto de librerías compartidas
export LD_LIBRARY_PATH=/usr/local/ipython/bin/ipython

# Ampliar la ruta de búsqueda por defecto de librerías privadas 
export PYTHONPATH=$SPARK_HOME/python/lib/py4j-*-src.zip:$PYTHONPATH:$SPARK_HOME/python/

Nota: ahora has inicializado una SparkSession por defecto. En la mayoría de los casos querrás configurarla más a fondo, algo imprescindible con big data. Si quieres saber más, consulta esta página.

Cargar tus datos

Este tutorial usa el conjunto de datos California Housing. Apareció en 1997 en el artículo Sparse Spatial Autoregressions, de Pace, R. Kelley y Ronald Barry, publicado en la revista Statistics and Probability Letters. Los autores lo construyeron a partir del censo de California de 1990.

Los datos contienen una fila por grupo censal. Un grupo censal es la unidad geográfica más pequeña para la que la Oficina del Censo de EE. UU. publica datos de muestra (suele tener entre 600 y 3.000 personas). En esta muestra, un grupo censal incluye de media 1425,5 personas en un área compacta. Encontrarás esta información en esta página o leyendo el artículo mencionado, disponible aquí.

Estos datos espaciales contienen 20.640 observaciones sobre precios de vivienda con 9 variables económicas:

  • Longitude se refiere a la distancia angular de un lugar geográfico al norte o sur del ecuador terrestre para cada grupo censal;
  • Latitude se refiere a la distancia angular al este u oeste del ecuador para cada grupo censal;
  • Housing median age es la edad mediana de las personas de un grupo censal. Nota: la mediana es el valor que se sitúa en el punto medio de una distribución de frecuencias;
  • Total rooms es el número total de habitaciones en las viviendas por grupo censal;
  • Total bedrooms es el número total de dormitorios en las viviendas por grupo censal;
  • Population es el número de habitantes de un grupo censal;
  • Households se refiere a las unidades familiares y sus ocupantes por grupo censal;
  • Median income registra la renta mediana de las personas que pertenecen a un grupo censal; y
  • Median house value es la variable dependiente y se refiere al valor mediano de la vivienda por grupo censal.

Además, verás que se han excluido del conjunto los grupos con ceros en todas las variables independientes y dependientes.

Median house value es la variable dependiente y será la variable objetivo de tu modelo de ML.

Puedes descargar los datos aquí. Busca la carpeta houses.zip, descárgala y descomprímela para acceder a los archivos.

A continuación, usarás el método textFile() para leer los datos desde la carpeta a RDDs. Este método toma un URI del archivo, que en este caso es la ruta local de tu máquina, y lo lee como una colección de líneas. Para mayor comodidad, leerás no solo el archivo .data, sino también el .domain que contiene el encabezado. Así podrás comprobar el orden de las variables.

# Cargar los datos
rdd = sc.textFile('/Users/yourName/Downloads/CaliforniaHousing/cal_housing.data')

# Cargar el encabezado
header = sc.textFile('/Users/yourName/Downloads/CaliforniaHousing/cal_housing.domain')

Exploración de datos

Ya has recopilado bastante información mirando la página web del conjunto, pero siempre es mejor inspeccionar los datos directamente con Spark en Python.

Importante: como la ejecución en Spark es “perezosa”, todavía no se ha ejecutado nada. Tus datos aún no se han leído realmente. Las variables rdd y header son solo planes. Tienes que empujar a Spark a trabajar, así que usa collect() para ver el header:

header.collect()

El método collect() trae el RDD completo a una sola máquina, y verás este resultado:

[u'longitude: continuous.', u'latitude: continuous.', u'housingMedianAge: continuous. ', u'totalRooms: continuous. ', u'totalBedrooms: continuous. ', u'population: continuous. ', u'households: continuous. ', u'medianIncome: continuous. ', u'medianHouseValue: continuous. ']

Consejo: cuidado con collect(). Puede hacer que el driver se quede sin memoria. Una alternativa más segura para ojear unos pocos elementos del RDD es take(). En general, limita el conjunto de resultados siempre que puedas, igual que haces con SQL.

Aprendes que el orden de las variables coincide con el presentado antes y que todas las columnas deberían tener valores continuos. Fuerza a Spark a hacer algo más y echa un vistazo a los datos de California para confirmarlo.

Llama a take() sobre tu RDD:

rdd.take(2)

Al ejecutar la línea anterior tomas los 2 primeros elementos del RDD. El resultado es el esperado: como has leído los archivos con textFile(), las líneas se leen tal cual. Las entradas están separadas por comas y las filas también están separadas por comas:

[u'-122.230000,37.880000,41.000000,880.000000,129.000000,322.000000,126.000000,8.325200,452600.000000', u'-122.220000,37.860000,21.000000,7099.000000,1106.000000,2401.000000,1138.000000,8.301400,358500.000000']

Hay que solucionarlo. No necesitas dividir cada entrada, pero sí asegurarte de que cada fila sea un elemento separado. Para ello, usa map() pasando una lambda que divide cada línea por comas. Comprueba el resultado con take() como antes:

Recuerda: las funciones lambda son anónimas y se crean en tiempo de ejecución.

# Dividir las líneas por comas
rdd = rdd.map(lambda line: line.split(","))

# Inspeccionar las 2 primeras filas 
rdd.take(2)

Obtendrás este resultado:

[[u'-122.230000', u'37.880000', u'41.000000', u'880.000000', u'129.000000', u'322.000000', u'126.000000', u'8.325200', u'452600.000000'], [u'-122.220000', u'37.860000', u'21.000000', u'7099.000000', u'1106.000000', u'2401.000000', u'1138.000000', u'8.301400', u'358500.000000']]

También puedes usar estas funciones para inspeccionar tus datos:

# Ver la primera línea 
rdd.first()

# Tomar los elementos superiores
rdd.top(2)

Si vienes de Pandas o de data frames en R, quizá esperabas ver un encabezado, pero no lo hay. Para facilitarte la vida, pasarás del RDD a un DataFrame. Siempre que puedas, es preferible usar DataFrames. Especialmente en Python, su rendimiento es mejor que el de los RDDs.

¿Cuál es la diferencia entre ambos?

Usa RDDs cuando quieras realizar transformaciones y acciones de bajo nivel sobre datos no estructurados; es decir, cuando no te importe imponer un esquema ni acceder a los atributos por nombre o columna. En cuanto al rendimiento, con RDDs no buscas necesariamente las ventajas que ofrecen los DataFrames para datos (semi)estructurados. Los RDDs son útiles cuando quieres manipular datos con construcciones funcionales en lugar de expresiones de dominio específicas.

En resumen, ahora cambiarás a DataFrames para usar expresiones de alto nivel, hacer consultas SQL y acceder por columnas.

Vamos a ello.

El primer paso es crear un SchemaRDD o un RDD de objetos Row con un esquema. Es normal, porque como en un DataFrame, finalmente quieres filas y columnas. Cada entrada se vincula a una fila y una columna, y las columnas tienen tipos de datos.

Volverás a usar map() con una lambda que mapea cada entrada a un campo de una Row. Para visualizarlo, piensa en esta primera línea:

[u'-122.230000', u'37.880000', u'41.000000', u'880.000000', u'129.000000', u'322.000000', u'126.000000', u'8.325200', u'452600.000000']

La lambda indica que construirás una fila en un SchemaRDD y que el elemento en el índice 0 se llamará “longitude”, y así sucesivamente.

Con este SchemaRDD, puedes convertirlo fácilmente a DataFrame con toDF().

# Importar los módulos necesarios 
from pyspark.sql import Row

# Mapear el RDD a un DF
df = rdd.map(lambda line: Row(longitude=line[0], 
                              latitude=line[1], 
                              housingMedianAge=line[2],
                              totalRooms=line[3],
                              totalBedRooms=line[4],
                              population=line[5], 
                              households=line[6],
                              medianIncome=line[7],
                              medianHouseValue=line[8])).toDF()

Ahora que tienes el DataFrame df, puedes inspeccionarlo con first() y take(), y también con head() y show():

# Mostrar las 20 primeras filas 
df.show()

Verás inmediatamente que esto se ve muy distinto del RDD con el que trabajabas antes:

tutorial pyspark

Consejo: usa df.columns para obtener las columnas del DataFrame.

Los datos parecen ordenados en columnas, pero ¿qué pasa con los tipos? Al leer los datos, Spark intenta inferir el esquema, ¿lo ha logrado aquí? Usa df.dtypes o df.printSchema() para saber más sobre los tipos.

# Imprimir los tipos de las columnas de `df`
# df.dtypes

# Imprimir el esquema de `df`
df.printSchema()

Como no ejecutas la primera línea, solo verás este resultado:

root
 |-- households: string (nullable = true)
 |-- housingMedianAge: string (nullable = true)
 |-- latitude: string (nullable = true)
 |-- longitude: string (nullable = true)
 |-- medianHouseValue: string (nullable = true)
 |-- medianIncome: string (nullable = true)
 |-- population: string (nullable = true)
 |-- totalBedRooms: string (nullable = true)
 |-- totalRooms: string (nullable = true)

Todas las columnas siguen siendo string… ¡Decepcionante!

Si quieres continuar con este DataFrame, tendrás que arreglarlo y asignar tipos de datos más adecuados a todas las columnas. También mejorarás el rendimiento. Intuitivamente, podrías optar por algo así, donde declaras que cada columna de df se convierta a FloatType():

from pyspark.sql.types import *

df = df.withColumn("longitude", df["longitude"].cast(FloatType())) \
   .withColumn("latitude", df["latitude"].cast(FloatType())) \
   .withColumn("housingMedianAge",df["housingMedianAge"].cast(FloatType())) \
   .withColumn("totalRooms", df["totalRooms"].cast(FloatType())) \ 
   .withColumn("totalBedRooms", df["totalBedRooms"].cast(FloatType())) \ 
   .withColumn("population", df["population"].cast(FloatType())) \ 
   .withColumn("households", df["households"].cast(FloatType())) \ 
   .withColumn("medianIncome", df["medianIncome"].cast(FloatType())) \ 
   .withColumn("medianHouseValue", df["medianHouseValue"].cast(FloatType()))

Pero estas llamadas repetidas son farragosas, propensas a errores y poco elegantes. ¿Por qué no escribes una función que lo haga de forma limpia?

La siguiente función definida por el usuario (UDF) recibe un DataFrame, nombres de columnas y el nuevo tipo de datos. Para cada nombre de columna, toma la columna y la convierte al nuevo tipo. Luego devuelve el DataFrame:

# Importar todo desde `sql.types`
from pyspark.sql.types import *

# Función para convertir el tipo de columnas de un DataFrame
def convertColumn(df, names, newType):
  for name in names: 
     df = df.withColumn(name, df[name].cast(newType))
  return df 

# Asignar todos los nombres de columna a `columns`
columns = ['households', 'housingMedianAge', 'latitude', 'longitude', 'medianHouseValue', 'medianIncome', 'population', 'totalBedRooms', 'totalRooms']

# Convertir las columnas de `df` a `FloatType()`
df = convertColumn(df, columns, FloatType())

¡Mucho mejor! Puedes comprobar rápidamente los tipos con printSchema(), como antes.

Con esto listo, toca ponerse con la exploración real. Ya has visto que el acceso columnar y las consultas SQL son dos ventajas de los DataFrames. Vamos a profundizar un poco. Empieza seleccionando dos columnas de df y mostrando 10 filas:

df.select('population','totalBedRooms').show(10)

Esta consulta te devuelve:

+----------+-------------+
|population|totalBedRooms|
+----------+-------------+
|     322.0|        129.0|
|    2401.0|       1106.0|
|     496.0|        190.0|
|     558.0|        235.0|
|     565.0|        280.0|
|     413.0|        213.0|
|    1094.0|        489.0|
|    1157.0|        687.0|
|    1206.0|        665.0|
|    1551.0|        707.0|
+----------+-------------+
solo se muestran las 10 primeras filas

También puedes hacer consultas más complejas, como esta:

df.groupBy("housingMedianAge").count().sort("housingMedianAge",ascending=False).show()

Que te devuelve:

+----------------+-----+                                                        
|housingMedianAge|count|
+----------------+-----+
|            52.0| 1273|
|            51.0|   48|
|            50.0|  136|
|            49.0|  134|
|            48.0|  177|
|            47.0|  198|
|            46.0|  245|
|            45.0|  294|
|            44.0|  356|
|            43.0|  353|
|            42.0|  368|
|            41.0|  296|
|            40.0|  304|
|            39.0|  369|
|            38.0|  394|
|            37.0|  537|
|            36.0|  862|
|            35.0|  824|
|            34.0|  689|
|            33.0|  615|
+----------------+-----+
solo se muestran las 20 primeras filas

Además de consultar, puedes describir tus datos y obtener estadísticas de resumen. Esto te ayudará después.

df.describe().show()

PySpark Machine Learning

Fíjate en los valores mínimos y máximos de los atributos numéricos. Varias variables tienen un rango muy amplio: necesitarás normalizar el conjunto de datos.

Preprocesamiento de datos

Con toda la información que has obtenido en este pequeño análisis exploratorio, ya puedes preprocesar los datos para alimentar el modelo.

  • No deberías preocuparte por valores perdidos; los ceros se han excluido del conjunto.
  • Probablemente deberías estandarizar los datos, ya que has visto que los rangos min/máx son bastante grandes.
  • Posiblemente podrías añadir atributos derivados, como habitaciones por dormitorio o habitaciones por hogar.
  • Tu variable dependiente también es bastante grande; para facilitar el trabajo, ajustarás ligeramente sus valores.

Preprocesar los valores objetivo

Empecemos con medianHouseValue, tu variable dependiente. Para facilitar el trabajo con los valores objetivo, expresarás los valores de vivienda en unidades de 100.000. Es decir, un objetivo como 452600.000000 pasará a ser 4.526:

# Importar todo desde `sql.functions` 
from pyspark.sql.functions import *

# Ajustar los valores de `medianHouseValue`
df = df.withColumn("medianHouseValue", col("medianHouseValue")/100000)

# Mostrar las 2 primeras líneas de `df`
df.take(2)

Verás claramente que los valores se han ajustado correctamente al inspeccionar el resultado de take():

[Row(households=126.0, housingMedianAge=41.0, latitude=37.880001068115234, longitude=-122.2300033569336, medianHouseValue=4.526, medianIncome=8.325200080871582, population=322.0, totalBedRooms=129.0, totalRooms=880.0), Row(households=1138.0, housingMedianAge=21.0, latitude=37.86000061035156, longitude=-122.22000122070312, medianHouseValue=3.585, medianIncome=8.301400184631348, population=2401.0, totalBedRooms=1106.0, totalRooms=7099.0)]

Ingeniería de características

Ahora que has ajustado medianHouseValue, también puedes añadir variables derivadas que mencionábamos arriba. Vas a añadir estas columnas:

  • Rooms per household: número de habitaciones por hogar en cada grupo censal;
  • Population per household: te indica cuántas personas viven por hogar en cada grupo censal; y
  • Bedrooms per room: te da una idea de cuántas habitaciones son dormitorios por grupo censal;

Como trabajas con DataFrames, lo mejor es usar select() para seleccionar las columnas con las que vas a operar: totalRooms, households y population. Además, indica que trabajas con columnas usando col(); de lo contrario no podrás hacer operaciones elemento a elemento como las divisiones previstas:

# Importa todo desde `sql.functions` si no lo has hecho
from pyspark.sql.functions import *

# Dividir `totalRooms` entre `households`
roomsPerHousehold = df.select(col("totalRooms")/col("households"))

# Dividir `population` entre `households`
populationPerHousehold = df.select(col("population")/col("households"))

# Dividir `totalBedRooms` entre `totalRooms`
bedroomsPerRoom = df.select(col("totalBedRooms")/col("totalRooms"))

# Añadir las nuevas columnas a `df`
df = df.withColumn("roomsPerHousehold", col("totalRooms")/col("households")) \
   .withColumn("populationPerHousehold", col("population")/col("households")) \
   .withColumn("bedroomsPerRoom", col("totalBedRooms")/col("totalRooms"))
   
# Inspeccionar el resultado
df.first()

Verás que, para la primera fila, hay unas 6,98 habitaciones por hogar, los hogares del grupo tienen unas 2,5 personas y el porcentaje de dormitorios es bajo (0,14):

Row(households=126.0, housingMedianAge=41.0, latitude=37.880001068115234, longitude=-122.2300033569336, medianHouseValue=4.526, medianIncome=8.325200080871582, population=322.0, totalBedRooms=129.0, totalRooms=880.0, roomsPerHousehold=6.984126984126984, populationPerHousehold=2.5555555555555554, bedroomsPerRoom=0.14659090909090908)

A continuación, y anticipando un posible problema al estandarizar, vas a reordenar las columnas. Como no quieres estandarizar la variable objetivo, asegúrate de aislarla.

En este caso, usa select() y pasa los nombres de las columnas en el orden adecuado. La variable objetivo medianHouseValue irá primero, para que no se vea afectada por la estandarización.

Nota: este es el momento de excluir variables que quizá no quieras considerar. Aquí dejaremos fuera longitude, latitude, housingMedianAge y totalRooms.

# Reordenar y seleccionar columnas
df = df.select("medianHouseValue", 
              "totalBedRooms", 
              "population", 
              "households", 
              "medianIncome", 
              "roomsPerHousehold", 
              "populationPerHousehold", 
              "bedroomsPerRoom")

Estandarización

Ahora que has reordenado los datos, ya casi puedes normalizarlos. Falta un paso: separar las características de la variable objetivo. En esencia, se trata de aislar la primera columna del DataFrame del resto.

En este caso, usarás la función map() con RDDs para hacerlo. Verás también que usas DenseVector(). Un vector denso es un vector local respaldado por un array de dobles que representa sus valores; sirve para almacenar arrays de valores en PySpark.

Después, vuelves a crear un DataFrame a partir de input_data y renombras las columnas pasando una lista con "label" y "features":

# Importar `DenseVector`
from pyspark.ml.linalg import DenseVector

# Definir `input_data` 
input_data = df.rdd.map(lambda x: (x[0], DenseVector(x[1:])))

# Sustituir `df` por el nuevo DataFrame
df = spark.createDataFrame(input_data, ["label", "features"])

Ahora sí puedes escalar los datos. Puedes usar Spark ML para esto: la librería facilita y escala el ML en big data. Encontrarás algoritmos de ML y todo lo necesario para construir pipelines prácticos. En este caso no necesitas tanto preprocesado, así que quizá un pipeline sea excesivo, pero si quieres profundizar, consulta esta página.

La columna de entrada son las features y la columna de salida con las reescaladas, que se añadirá a scaled_df, se llamará "features_scaled":

# Importar `StandardScaler` 
from pyspark.ml.feature import StandardScaler

# Inicializar `standardScaler`
standardScaler = StandardScaler(inputCol="features", outputCol="features_scaled")

# Ajustar el scaler al DataFrame
scaler = standardScaler.fit(df)

# Transformar los datos de `df` con el scaler
scaled_df = scaler.transform(df)

# Inspeccionar el resultado
scaled_df.take(2)

Mira tu DataFrame y el resultado. Verás que, efectivamente, se ha añadido una tercera columna features_scaled, que puedes comparar con features:

[Row(label=4.526, features=DenseVector([129.0, 322.0, 126.0, 8.3252, 6.9841, 2.5556, 0.1466]), features_scaled=DenseVector([0.3062, 0.2843, 0.3296, 4.3821, 2.8228, 0.2461, 2.5264])), Row(label=3.585, features=DenseVector([1106.0, 2401.0, 1138.0, 8.3014, 6.2381, 2.1098, 0.1558]), features_scaled=DenseVector([2.6255, 2.1202, 2.9765, 4.3696, 2.5213, 0.2031, 2.6851]))]

Nota: estas líneas se parecen mucho a lo que harías en Scikit-Learn.

Construir un modelo de aprendizaje automático con Spark ML

Con el preprocesamiento hecho, por fin toca construir tu modelo de regresión lineal. Como siempre, primero hay que dividir los datos en entrenamiento y prueba. Por suerte, con randomSplit() es muy sencillo:

# Dividir los datos en entrenamiento y prueba
train_data, test_data = scaled_df.randomSplit([.8,.2],seed=1234)

Pasas una lista con dos números que representan el tamaño de los conjuntos de entrenamiento y prueba, y una semilla para la reproducibilidad. Si quieres saber más, echa un vistazo al tutorial de machine learning en Python de DataCamp.

Después, sin más, puedes crear tu modelo.

Nota: el argumento elasticNetParam corresponde a α y regParam (el parámetro de regularización) corresponde a λ. Más información aquí.

# Importar `LinearRegression`
from pyspark.ml.regression import LinearRegression

# Inicializar `lr`
lr = LinearRegression(labelCol="label", maxIter=10, regParam=0.3, elasticNetParam=0.8)

# Ajustar el modelo con los datos
linearModel = lr.fit(train_data)

Con el modelo listo, puedes generar predicciones para los datos de prueba: usa transform() para predecir las etiquetas de test_data. Luego, usa operaciones RDD para extraer las predicciones y las etiquetas reales y emparejarlas en una lista llamada predictionAndLabel.

Por último, inspecciona los valores predichos y reales accediendo a la lista con corchetes []:

# Generar predicciones
predicted = linearModel.transform(test_data)

# Extraer las predicciones y las etiquetas "correctas"
predictions = predicted.select("prediction").rdd.map(lambda x: x[0])
labels = predicted.select("label").rdd.map(lambda x: x[0])

# Combinar `predictions` y `labels`
predictionAndLabel = predictions.zip(labels).collect()

# Imprimir las 5 primeras parejas 
predictionAndLabel[:5]

Verás los siguientes valores reales y predichos (en ese orden):

[(1.4491508524918457, 0.14999), (1.5705029404692372, 0.14999), (2.148727956912464, 0.14999), (1.5831547768979277, 0.344), (1.5182107797955968, 0.398)]

Evaluar el modelo

Mirar valores predichos está bien, pero mejor aún es ver métricas para entender la calidad del modelo. Empieza imprimiendo los coeficientes y el intercepto:

# Coeficientes del modelo
linearModel.coefficients

# Intercepto del modelo
linearModel.intercept

Obtendrás este resultado:

# Los coeficientes
[0.0,0.0,0.0,0.276239709215,0.0,0.0,0.0]

# El intercepto
0.990399577462

Después, usa el atributo summary para obtener el rootMeanSquaredError y el r2:

# Obtener el RMSE
linearModel.summary.rootMeanSquaredError

# Obtener el R2
linearModel.summary.r2
  • El RMSE mide el error entre valores predichos y observados. Cuanto menor es, más cercanos están.

  • El R2 (coeficiente de determinación) muestra cuán cerca están los datos de la recta de regresión ajustada. Va de 0 a 1, donde 0 indica que el modelo no explica la variabilidad y 1 que la explica completamente. En general, cuanto mayor el R2, mejor se ajusta el modelo.

Obtendrás lo siguiente:

# RMSE
0.8692118678997669

# R2
0.4240895287218379

¡A tu modelo aún le vendrán bien mejoras! Si quieres seguir, puedes probar con los parámetros del modelo, las variables incluidas en el DataFrame, etc. Pero aquí termina el tutorial por ahora.

Antes de irte…

Antes de cerrar, detén la SparkSession con esta línea:

spark.stop()

Lleva el big data más lejos

¡Enhorabuena! Has llegado al final del tutorial, donde has aprendido a crear un modelo de regresión lineal con Spark ML.

Si te interesa aprender más sobre PySpark, plantéate hacer el curso Introduction to PySpark de DataCamp y echa un vistazo al tutorial de Apache Spark: ML con PySpark.

Temas
Python
Ciencia de datos
Aprendizaje automático

Aprende más sobre Python y PySpark

Curso

Fundamentos de PySpark

4 h
157.8K
Aprende a implementar la gestión de datos distribuidos y el machine learning en Spark utilizando el paquete PySpark.
Ver detallesRight Arrow
Iniciar Curso
Ver másRight Arrow
Relacionado

Tutorial

Tutorial de Pyspark: Primeros pasos con Pyspark

Descubre qué es Pyspark y cómo se puede utilizar, con ejemplos.
Natassha Selvaraj's photo

Natassha Selvaraj

10 min

Tutorial

Instalación de PySpark (Todos los sistemas operativos)

Este tutorial mostrará la instalación de PySpark y cómo gestionar las variables de entorno en los sistemas operativos Windows, Linux y Mac.

Olivia Smith

8 min

Tutorial

Multiprocesamiento en Python: Guía de hilos y procesos

Aprende a gestionar hilos y procesos con el módulo de multiprocesamiento de Python. Descubre las técnicas clave de la programación paralela. Mejora la eficacia de tu código con ejemplos.
Kurtis Pykes 's photo

Kurtis Pykes

7 min

Clustering k-means

Tutorial

Introducción a k-Means Clustering con scikit-learn en Python

En este tutorial, aprenda a aplicar k-Means Clustering con scikit-learn en Python

Kevin Babitz

8 min

Tutorial

Tutorial de Generación de nubes de palabras en Python

Aprende a realizar Análisis exploratorios de datos para el Procesamiento del lenguaje natural utilizando WordCloud en Python.
Duong Vu's photo

Duong Vu

11 min

Tutorial

Tutorial sobre el uso de XGBoost en Python

Descubre la potencia de XGBoost, uno de los marcos de machine learning más populares entre los científicos de datos, con este tutorial paso a paso en Python.
Ver MásVer Más