Published on

¿Qué son Spark, Scala y PySpark?

Authors

Veintidós millones de viajes en taxi. Una laptop con 8 GB de memoria. Veintitrés segundos.

Para eso sirven estas herramientas. Aquí va qué es cada una, en lenguaje simple.

Todo lo de abajo es reproducible. El notebook complementario descarga los datos, corre cada función e imprime cada número de este artículo. Toma como un minuto en una laptop normal.

El problema

Las herramientas normales de datos cargan el archivo completo en memoria. Abres un archivo de 200 MB y 200 MB se van a la RAM. Bien.

Ahora intenta un archivo de 50 GB en una laptop de 8 GB. No se pone lento. Se vuelve imposible.

Puedes comprar una computadora más grande, lo cual tiene un límite y un precio. O puedes repartir el trabajo entre muchas computadoras. Spark es la segunda opción.

Spark

Imagina contar una baraja de cartas. Solo, cuentas las 52. Con cuatro amigos, cada quien cuenta 13 y suman los totales.

Spark hace eso con datos. Corta tu archivo en pedazos, le da cada pedazo a un trabajador distinto y combina las respuestas. Tú nunca administras nada de eso. Dices qué quieres, Spark decide cómo dividirlo.

Scala

Spark está escrito en un lenguaje llamado Scala, que corre sobre la Máquina Virtual de Java.

Nunca vas a escribir Scala. Pero se cuela en tu vida de exactamente una forma: Spark necesita Java instalado, aunque tú nunca lo toques.

Cuando monté esto, todo se instaló correctamente y aun así falló al instante. No había Java. Ese error atrapa a mucha gente.

PySpark

PySpark te deja escribir Spark en Python. No es una versión de Spark hecha en Python. Es un control remoto para la de Scala.

Tu Python no procesa los números. Describe el trabajo, manda esa descripción a Spark y recibe respuestas.

Por eso se requiere Java, y por eso los errores de Spark están llenos de nombres de Java que tú nunca escribiste. También significa que Python no es la opción lenta. Python y Scala mandan instrucciones idénticas al mismo motor.

La idea que confunde a todos

Esta función convierte siete columnas a sus tipos correctos:

def clean_data(df):
    df2 = df.withColumn("passenger_count",df.passenger_count.cast('int'))\
    .withColumn("total_amount",df.total_amount.cast('float'))\
    .withColumn("tip_amount",df.tip_amount.cast('float'))\
    .withColumn("trip_distance",df.trip_distance.cast('float'))\
    .withColumn("fare_amount",df.fare_amount.cast('float'))\
    .withColumn("tpep_pickup_datetime",df.tpep_pickup_datetime.cast('timestamp'))\
    .withColumn("tpep_dropoff_datetime",df.tpep_dropoff_datetime.cast('timestamp'))
    return df2

Córrela sobre 22 millones de filas y termina al instante. Porque no convierte nada.

El notebook mide esto lado a lado. clean_data regresa en 0.1 segundos; el .count() que va después es lo que realmente hace el trabajo.

Spark tiene dos tipos de operaciones. Las que describen trabajo, como withColumn y groupBy, regresan de inmediato. Las que exigen respuestas, como show() y count(), hacen que todo se ejecute de verdad.

Es la diferencia entre escribir una lista del súper e ir al súper. Puedes agregar cincuenta cosas en segundos porque no has salido de la casa.

De ahí se desprenden dos cosas. Spark ve tu plan completo antes de correrlo, así que puede optimizarlo. Y tu error de la celda 3 no va a aparecer hasta la celda 9, donde se pida la primera respuesta real.

Lo que encontró

Velocidad promedio de los taxis, sobre 22.6 millones de viajes:

DíaVelocidad
Jueves9.6 mph
Viernes9.7 mph
Martes10.2 mph
Domingo11.7 mph

El jueves es el peor día para ir en taxi. El domingo el mejor. Toda la diferencia es menos de 2 mph, lo cual dice algo sobre Manhattan.

Las noches cuestan menos por milla que los días, 5.66contra5.66 contra 6.41, porque menos tráfico significa menos minutos de taxímetro parado.

Dónde Spark realmente se gana su lugar

Veintidós millones de filas es una demo. Así se ven los casos reales.

Los datos no caben en ningún lado. Un procesador de pagos guarda años de transacciones con tarjeta. Una telco registra cada llamada y conexión. Una planta transmite lecturas de sensores de mil máquinas. Son miles de millones de filas. No hay laptop, ni servidor individual, que las abra.

El trabajo tiene fecha de entrega. Muchos pipelines nocturnos tienen que estar listos antes de que abra el negocio. Cuando un proceso que toma nueve horas necesita tomar una, no llegas ahí optimizando en una sola máquina. Agregas máquinas.

Hay que unir dos cosas grandes. Cruzar un año de transacciones contra una tabla de clientes es fácil de describir y brutal de ejecutar, porque cada fila que hace match tiene que encontrarse físicamente con su pareja. Spark reparte eso en un clúster. Una sola máquina se queda sin espacio.

Los datos viven en la nube. La mayoría de las empresas guarda sus datos como archivos en S3 o Google Cloud Storage, no en una base de datos. Spark los lee directo, en paralelo, en el formato en que ya están. Esa es la realidad diaria de casi toda la ingeniería de datos.

El mismo código tiene que correr en ambos tamaños. Esta es la que la gente subestima. En el notebook de este artículo, Spark corre con local[4], o sea cuatro hilos en una laptop. Apuntarlo a un clúster de 200 máquinas cambia esa línea. Nada más. Prototipas con una muestra en tu escritorio y corres el mismo código sobre todo en producción.

Cuándo saltártelo

Si tus datos caben cómodamente en memoria, no uses Spark. Vas a pagar por arrancar una JVM, planificar un trabajo distribuido y coordinar trabajadores, y no obtienes nada a cambio. Con unos pocos miles de filas es más lento que no hacer nada especial.

Spark no es una forma más rápida de manejar datos pequeños. Es a lo que recurres cuando los datos dejaron de ser pequeños, y la respuesta honesta a "¿debería usar Spark?" normalmente es no.

Córrelo tú mismo

El notebook está aquí, o puedes descargarlo directo. Descarga los datos de taxi solo, así que no hay nada que buscar ni configurar.

python3 -m venv ~/sparkenv
~/sparkenv/bin/python -m pip install pyspark==3.5.3 jdk4py ipykernel
~/sparkenv/bin/python -m ipykernel install --user --name pyspark-env --display-name "Python (PySpark)"

jdk4py es el truco útil ahí. Instala un runtime de Java como paquete de Python, así te saltas la instalación de Java a nivel de sistema que normalmente pide contraseña de administrador.

Abre el notebook, elige el kernel Python (PySpark) y corre todo. Después prueba los experimentos del final: cambia local[4] por local[1] y observa cómo se mueven los tiempos, o agrega df.cache() y observa cómo Spark deja de releer del disco.

Recibe el siguiente

Artículos sobre datos, analítica y las decisiones que deciden si un modelo se usa o se ignora.

CompartirLinkedInXReddit