En esta serie, PySpark para principiantes: Dominar los conceptos básicos, entonces ya comprenderá el corazón de Spark: datos distribuidos, DataFrames y ejecución diferida. Instalaste PySpark, estableciste una SparkSession, leíste un CSV y realizaste manipulaciones simples de los datos en un Dataframe. Dejaré un enlace a esa historia al final de este.
Una cosa que vale la pena repetir de ese artículo original es que a menudo uso los términos PySpark y Spark indistintamente, pero estrictamente hablando, Spark es el marco informático distribuido general (escrito en Scala) y PySpark es una API de Python dedicada a Spark.
Más allá de lo básico
Ahora bien, sucede algo interesante cuando superas esa etapa de principiante. Rápidamente te das cuenta de que tu segundo proyecto PySpark requiere una mentalidad ligeramente diferente:
Quiere leer/escribir datos de una manera más segura, rápida y predecible. Quiere combinar conjuntos de datos sin sentirse inseguro acerca de las uniones. Quiere comprender por qué Spark se comporta como lo hace y cómo empujarlo suavemente en la dirección correcta.
Este artículo lo guía a través de los siguientes pasos. Es deliberadamente lento y práctico. Sin partes internas profundas. Sin ajuste de clúster. Sin optimizaciones complicadas de Spark. Justo lo que los verdaderos principiantes necesitan saber cuando pasan de ejemplos de juguetes a trabajos pequeños del mundo real.
Estamos usando Spark de código abierto, ejecutándolo localmente, como antes.
1. Dar el siguiente paso: leer los datos correctamente
En mi primer artículo, utilizamos el cargador CSV más simple posible:
df = spark.read.csv("ventas.csv", encabezado=Verdadero, inferSchema=Verdadero)
Funciona, y está bien para los primeros experimentos, pero esconde un problema sutil.
Spark está adivinando tus tipos de datos
Cuando usa la directiva inferSchema=True, Spark analiza una pequeña muestra de su archivo y usa esa información para adivinar si una columna es un número entero, una cadena, un booleano o un doble. Eso significa:
Si 99 filas parecen ser numéricas y la fila número 100 está en blanco, Spark podría interpretar la columna como una cadena. Si alguien edita el archivo la próxima semana y accidentalmente agrega £23,50 en lugar de 23,50, Spark podría tratar toda la columna de manera diferente. Si su archivo es grande, la muestra que utiliza Spark no representará todo el conjunto de datos.
Esto puede provocar comportamientos misteriosos más adelante, el tipo de errores que los principiantes encuentran más difíciles de diagnosticar.
Un mejor hábito para principiantes: definir un esquema para sus datos
Piense en un esquema como la versión de Spark de un modelo para leer datos. Antes de construir algo, le dices a Spark cosas como:
Los nombres de las columnas.
¿Qué tipo de datos deberían ser?
Si un valor de columna es opcional o no.
Así es como se ve en nuestro ejemplo de datos de ventas. Recuerde que los datos se veían así:
ID_transacción,nombre_cliente,monto_neto,monto_impuesto, es_miembro 101,Alice,250.50,25.05,true 102,Bob,120.00,6.00, false 103,Charlie,450.75,25.07,true 104,David,89.99,5.73,false
Para especificar los tipos de los campos anteriores en Spark, definimos nuestro esquema usando un código como el siguiente.
desde pyspark.sql tipos de importación como T esquema = T.StructType([ T.StructField("transaction_id", T.IntegerType(), False), T.StructField("customer_name", T.StringType(), False), T.StructField("net_amount", T.DoubleType(), True), T.StructField("tax_amount", T.DoubleType(), Verdadero), T.StructField("is_member", T.BooleanType(), Verdadero), ])
Los nombres de las columnas y los parámetros de tipo se explican por sí solos. El parámetro True[False] indica que puede [no] haber valores NULL en la columna. Tenga en cuenta que el indicador de nulidad Verdadero/Falso es principalmente metadatos de esquema e información de optimización. No siempre se aplica estrictamente para cada fuente de datos como lo es una restricción NOT NULL de una base de datos.
Opciones más útiles al leer datos CSV
Hay un montón de opciones útiles de lectura de CSV que puedes combinar con la directiva de esquema que hacen que la carga de datos CSV sea aún más confiable.
Las opciones más comunes incluyen:
mode=”PERMISSIVE”: mantiene las filas incorrectas tanto como sea posible mode=”DROPMALFORMED”: elimina las filas con formato incorrecto mode=”FAILFAST”: errores inmediatamente header= True[False]: ¿El archivo contiene [o no] un registro de encabezado nullValue: qué texto debe reemplazar los valores nulos en el formato de fecha/marca de tiempo de entrada?
Ahora podemos cargar sales_data en un marco de datos como este:
df = ( spark.read .option("header", True) # Otros modos: "PERMISSIVE" y "DROPMALFORMED". .option("mode", "FAILFAST") .option("nullValue", "N/A") .schema(schema) .csv("sales_data.csv") )
¿Por qué es esto importante para los principiantes?
Sabes cuáles son los tipos de datos antes de empezar a trabajar. Si se especifica, Spark rechazará filas extrañas en lugar de interpretarlas silenciosamente. Tus transformaciones se vuelven más predecibles. Si une dos conjuntos de datos más adelante, las discrepancias de tipos no le sorprenderán.
2. Comprender las transformaciones de datos
Recuerde, en mi artículo anterior, en nuestros primeros pasos en la manipulación de marcos de datos con PySpark, agregamos una columna derivada adicional a nuestro marco de datos usando un código como este:
df2 = df.withColumn("monto_bruto", df.monto_neto + df.monto_impuesto)
Le expliqué que esta línea no calcula nada todavía. Simplemente agrega un paso al plan interno de Spark:
1. Lea el CSV 2. Agregue una nueva columna (monto_bruto = neto + impuestos)
Entonces podrías agregar más pasos como este:
df3 = df2.withColumn("porcentaje_impuesto", df2.monto_impuesto / df2.monto_bruto * 100)
Aún así, no se ha realizado ningún cálculo. Sólo cuando realizas una acción como…
df3.mostrar()
… Spark dice:
"Está bien, ahora necesito ejecutar todos estos pasos".
Esto es lo que significa "ejecución diferida", pero lo importante para los principiantes no es el nombre. Es el efecto, y significa,
Puedes encadenar muchas transformaciones sin “pagar” por ellas hasta que necesites el resultado. Spark puede reorganizar el orden internamente para ejecutar las cosas de manera eficiente. No pierde tiempo realizando pasos intermedios sobre datos que podría filtrar más adelante.
Piensa en ello como si fuera una tarea cotidiana, como preparar un sándwich:
Reúnes todos los ingredientes. Lo ensamblas en tu mente. En realidad, sólo empiezas a cortar y preparar una vez que sabes exactamente lo que estás haciendo.
3. Limpiar datos antes de que causen problemas
Los datos reales suelen ser confusos y a menudo contienen valores faltantes, cadenas en blanco, registros duplicados o valores de marcador de posición como "N/A" y "desconocido".
En PySpark, el objetivo es detectar y solucionar problemas obvios de manera temprana para que el resto de su flujo de trabajo se comporte de manera predecible. PySpark tiene una serie de funciones útiles que le permiten hacer esto.
Eliminar filas con valores faltantes
La función de limpieza más sencilla es dropna().
df_clean = df.dropna()
Esto elimina cualquier fila que contenga un valor nulo en cualquier columna. Esto puede resultar útil, pero a menudo resulta demasiado agresivo.
Más comúnmente, solo elimina filas donde faltan columnas importantes en esa fila en particular:
df_clean = df.dropna(subset=["monto_neto", "monto_impuesto"])
Esto significa:
Mantenga la fila mientras importe_neto y importe_impuesto estén presentes.
Es posible que otras columnas aún contengan valores nulos, y eso podría estar bien.
Llenando valores faltantes
A veces no desea eliminar filas. Sólo desea reemplazar los valores faltantes con algo sensato.
Ahí es donde fillna() es útil.
df_clean = df.fillna({"ciudad": "Desconocido"})
También puedes completar columnas numéricas:
df_clean = df.fillna({"importe_impuesto": 0.0})
Esto resulta útil cuando un valor faltante tiene un significado claro. Por ejemplo, un importe de descuento faltante podría razonablemente convertirse en 0,0. Pero ten cuidado. Completar los valores faltantes puede cambiar el significado de sus datos si elige el valor predeterminado incorrecto.
Cambiar tipos de columnas con cast()
A veces, Spark lee una columna como del tipo incorrecto, especialmente cuando se trabaja con archivos CSV. Si es así, puedes convertir una columna usando el operador cast():
desde pyspark.sql funciones de importación como F df_clean = df.withColumn("net_amount",F.col("net_amount").cast("double"))
Esto es especialmente común cuando fechas, números o valores booleanos se han leído como cadenas.
Eliminar filas duplicadas
Pueden aparecer filas duplicadas cuando los archivos se exportan más de una vez, se unen incorrectamente o se combinan desde múltiples fuentes. Puede eliminar duplicados exactos como este:
df_clean = df.dropDuplicates()
O elimine duplicados según una o más columnas seleccionadas.
df_clean = df.dropDuplicates(["transaction_id"])
Esa segunda versión suele ser más útil porque dice:
Cada ID de transacción solo debe aparecer una vez.
Un pequeño ejemplo de limpieza de datos
Juntando esas ideas:
de pyspark.sql funciones de importación como F df_clean = ( df # Eliminar transacciones a las que les faltan valores requeridos. .dropna(subset=["transaction_id", "net_amount"]) # Proporcionar valores predeterminados para valores opcionales. .fillna( { "city": "Unknown", "tax_amount": 0.0, } ) # Aplicar los tipos numéricos esperados. .withColumn( "net_amount", F.col("net_amount").cast("double"), ) .withColumn( "tax_amount", F.col("tax_amount").cast("double"), ) # Mantenga una fila para cada transacción .dropDuplicates(["transaction_id"]) )
4. Unir conjuntos de datos en PySpark sin perderse
Si ha trabajado con bases de datos antes, probablemente haya escrito declaraciones SQL que unen dos o más tablas. Las uniones en Spark funcionan de la misma manera, pero en Dataframes.
¿Qué es una unión?
Si el concepto de unión es nuevo para usted, es una forma de hacer coincidir filas de un DataFrame con filas relacionadas de otro DataFrame. En otras palabras, responde a una pregunta como:
"¿Qué filas de este DataFrame corresponden a filas de ese DataFrame?"
Esa es la idea principal detrás de cada unión en PySpark. Una vez que esa parte está clara, la sintaxis y los tipos de unión se vuelven mucho más fáciles de entender.
Si tiene dos marcos de datos como este:
datos_ventas.csv
id_transacción, nombre_cliente, importe_neto, importe_impuesto 101, Alice, 250,50, 25,05 102, Bob, 120,00, 6,00
clientes.csv
nombre_cliente, ciudad, nivel_lealtad Alice, Nueva York, Gold Bob, Londres, Plata
Puedes unirte a ellos en su campo común nombre_cliente de esta manera:
df_sales = spark.read.csv("sales_data.csv", header=True) df_customers = spark.read.csv("customers.csv", header=True) df_joined = df_sales.join(df_customers, on="customer_name", how="inner") df_joined.show() # Salida +————-+————–+———-+———-+——–+————-+ |nombre_cliente|id_transacción|monto_neto|importe_impuesto|ciudad |nivel_lealtad| +————-+————–+———-+———-+——–+————-+ |Alice |101 |250.50 |25.05 |Nueva York|Oro | |Bob |102 |120.00 |6.00 |Londres |Plata | +————-+————–+———-+———-+——–+————-+
¿Qué combinación deberían utilizar los principiantes?
Hay varios tipos diferentes de uniones disponibles en Spark. Para el 99 % de los casos de uso para principiantes, utilizará uno de los siguientes:
interior: muestra solo las filas coincidentes restantes: muestra todo lo que está en la tabla de la izquierda, además de las coincidencias exterior: muestra todas las filas de ambas tablas
Y de estos, la combinación interna será de lejos el tipo de combinación más común que utilizará en su trabajo diario.
No se preocupe todavía por la "difusión", la "ordenación-fusión", el "shuffle-hash" ni ninguna otra estrategia de unión avanzada. A medida que crezca su experiencia con Spark, podrá leer sobre ellos cuando lo desee.
Sólo recuerda:
Las uniones son computacionalmente más costosas que las simples operaciones de columnas, así que úselas cuando sea necesario, pero no de manera casual.
5. Lectura y escritura de datos a la “manera Spark”: Parquet
La mayoría de los principiantes se quedan con CSV porque les resulta familiar. Pero CSV es lento, rígido y carece de soporte para tipos de datos y, en la vida real, Parquet es el formato de datos nativo de Spark. Parquet es un formato de datos comprimidos en columnas ideal para análisis de datos, informes de datos y cargas de trabajo de lectura intensa.
Cuando Spark lee un conjunto de datos de Parquet:
Solo carga las columnas que realmente necesitas. Entiende cada tipo de datos. Se carga significativamente más rápido que CSV.
Escribe el contenido del marco de datos en archivos de formato Parquet como este:
df_joined.write.mode("sobrescribir").parquet("salida/ventas_enriquecidas")
Luego podrás volver a leerlo instantáneamente así,
df_fast = spark.read.parquet("salida/ventas_enriquecidas") df_fast.show()
NÓTESE BIEN. Usar Parquet para la entrada y salida de archivos es la “actualización” de rendimiento más sencilla para cualquier principiante de Spark.
6. Pensando en los flujos de trabajo de PySpark
Una vez que comprenda cómo leer datos, limpiarlos, transformarlos, unirlos y volver a escribirlos, el siguiente paso es aprender a organizar esas acciones en un flujo de trabajo simple. Un proyecto principiante de PySpark suele seguir esta secuencia:
Leer datos -> verificarlos y limpiarlos -> agregar columnas útiles -> combinar con otros datos -> escribir el resultado
Puede parecer obvio, pero es un cambio importante. Ya no estás simplemente experimentando con un DataFrame a la vez. Estás construyendo un proceso repetible.
Mantenga cada etapa simple
Un hábito útil para principiantes es darle a cada etapa de su flujo de trabajo un propósito claro. Por ejemplo:
df_raw = spark.read.schema(schema).csv("sales_data.csv", header=True) df_clean = df_raw.dropna(subset=["net_amount", "tax_amount"]) df_enriched = df_clean.withColumn( "bruss_amount", F.col("net_amount") + F.col("importe_impuesto") ) df_final = df_enriched.join(df_customers, on="nombre_cliente", how="left") df_final.write.mode("overwrite").parquet("output/final_dataset")
Este estilo es un poco más detallado que encadenar todo en una expresión larga, pero es mucho más fácil de leer cuando estás aprendiendo.
Cada nombre de DataFrame le indica dónde se encuentra en el flujo de trabajo:
df_raw -> los datos tal como llegaron df_clean -> los datos después de la limpieza básica df_enriched -> los datos después de agregar un nuevo significado df_final -> el conjunto de datos listo para guardar
Por qué esto importa
Cuando algo sale mal, esta estructura facilita mucho la depuración.
Puede inspeccionar cada etapa mirando los datos:
df_raw.show() df_clean.show() df_enriched.show()
Puede comprobar el recuento de filas:
df_raw.count() df_clean.count() df_final.count()
Esto ayuda a responder preguntas útiles como:
¿Las filas desaparecieron inesperadamente durante la limpieza? ¿La combinación creó más filas de las esperadas? ¿Una columna calculada produjo nulos?
El modelo mental simple de: Entradas → preparación → combinación → salida lo llevará sorprendentemente lejos en su viaje con PySpark.
7. Una suave introducción a la interfaz de usuario de Spark.
Spark tiene una pequeña interfaz de usuario web que se activa cuando ejecutas una acción como .count() o .write(). Cuando su trabajo de Spark se esté ejecutando localmente, visite:
http://localhost:4040
Deberías ver algo como esto en la pantalla.
Parece un poco abrumador, pero no es necesario que comprendas todas las pestañas. En esta etapa, sólo necesita saber que la interfaz de usuario existe y por qué es útil. Y es útil porque le ayuda a ver qué trabajos de Spark se han ejecutado o se están ejecutando actualmente.
Y, a medida que crece su experiencia en Spark, la interfaz de usuario puede ayudarlo a comprender por qué los trabajos fallaron o tardan más de lo esperado en ejecutarse. Pero eso viene mucho más tarde. Por ahora, trate la interfaz de usuario de Spark como el tablero de su automóvil: no necesita entender el motor para darse cuenta cuando algo parece extraño.
Resumen: ahora está listo para su primer proyecto real de PySpark
En este punto, ha pasado de "Puedo ejecutar Spark" a "Puedo crear una canalización de Spark limpia y sencilla".
Ahora sabes cómo:
lea datos de forma segura, límpielos y prepárelos, enriquézcalos con nuevas columnas, combine múltiples conjuntos de datos, guarde el resultado de manera eficiente y observe Spark lo suficiente para mantenerse seguro.
Nada en este artículo requiere un clúster. Nada requirió ajustes avanzados. Así es exactamente como comienzan muchos proyectos reales de PySpark.
Cuando tenga más experiencia, es posible que desee ampliar sus conocimientos investigando algunos de estos temas.
leer planes de ejecución entender barajados administrar particiones otros tipos de uniones ajuste simple del rendimiento
Estos son algunos de los temas que espero cubrir en un artículo futuro, pero por ahora, usted ha dominado su próximo hito importante y puede crear algo significativo y útil con PySpark.
Por cierto, aquí está el enlace al primer artículo de esta serie.
PySpark para principiantes: Dominar los conceptos básicos, que mencioné al principio.