A menudo comienza con herramientas como pandas. Son intuitivos, potentes y perfectos para conjuntos de datos pequeños y medianos. Pero tan pronto como sus datos crecen más allá de lo que cabe cómodamente en la memoria, comienzan a surgir problemas de rendimiento. Aquí es donde entra en juego PySpark.
Tenga en cuenta que en este artículo a menudo usaré los términos Spark y PySpark indistintamente. Para nuestros propósitos, no importa, pero debes recordar que son diferentes. Spark es el marco informático distribuido general (escrito en Scala) y PySpark es una API de Python dedicada a Spark.
¿Qué es PySpark?
PySpark es la API de Python para Apache Spark, un marco informático distribuido para procesar de manera eficiente grandes volúmenes de datos. En lugar de ejecutar todos los cálculos en una sola máquina, Spark distribuye el trabajo entre varias máquinas (un clúster), lo que le permite procesar datos a escala mientras escribe código que aún resulta familiar para los usuarios de Python.
Una de las ventajas clave de PySpark es que abstrae gran parte de la complejidad de los sistemas distribuidos. No es necesario administrar manualmente los subprocesos, la memoria o la comunicación de red. Spark maneja estas inquietudes por usted, mientras usted se concentra en describir lo que desea hacer con los datos en lugar de cómo se deben ejecutar.
Si es un completo recién llegado a Spark, hay tres ideas centrales clave que debe aprender antes de usarlo. Estos son:
1. Clústeres
Cuando la gente escucha que Spark se ejecuta en un "clúster", puede parecer intimidante. En la práctica, no es necesario un conocimiento profundo de los sistemas distribuidos para comenzar. Un clúster es simplemente un grupo de servidores conectados en red que pueden colaborar. En una aplicación Spark que se ejecuta en un clúster, una máquina actúa como controlador, coordinando el trabajo, mientras que las otras actúan como ejecutores y realizan cálculos en fragmentos de datos. Cuando los nodos Ejecutores han terminado su trabajo, envían una señal al nodo Controlador, y el Controlador puede realizar lo que sea necesario con el conjunto de resultados final.
┌───────────────────┐ │ Conductor │ │(tu aplicación PySpark) │ └─────────┬─────────┘ │ | El Conductor alquila trabajo | a uno o más ejecutores ┌────────────────────┼─── ────────────────────────┐ │ │ │ ┌───────▼────────┐ ┌───────▼────────┐ ┌───────▼────────┐ │ Ejecutor 1 │ │ Ejecutor 2 │ │ Ejecutor N │ │ procesa parte│ │ procesa parte│ …… │ procesa parte│ │ de los datos │ | de los datos │ │ de los datos │ └────────────────┘ └────────────────┘ └────────────────┘
Solo recuerde que no necesita ejecutar Spark en un clúster de computación físico. Cuando ejecuta PySpark localmente, Spark simula un clúster en su computadora portátil o PC usando múltiples núcleos. Uno de los puntos fuertes de PySpark es que el mismo código se puede implementar posteriormente en un clúster real, ya sea en la nube o localmente, con solo cambios muy pequeños.
Esta separación de coordinación y ejecución permite que Spark crezca. A medida que crecen los conjuntos de datos, se pueden agregar más ejecutores para procesar datos en paralelo, lo que reduce el tiempo de ejecución sin necesidad de realizar cambios en el código.
2. El marco de datos de Spark
En el corazón de PySpark se encuentra la API DataFrame, que es la forma principal de trabajar con datos en Spark. Un DataFrame es simplemente una tabla de datos, formada por filas y columnas, muy similar a una tabla en una base de datos o un DataFrame en pandas. Si ha utilizado SQL o pandas antes, las ideas básicas le resultarán familiares.
Con Spark DataFrames, puede realizar tareas de datos comunes, como filtrar filas, seleccionar columnas, agrupar datos, unir tablas y calcular resúmenes como recuentos o promedios. Estas operaciones son fáciles de leer y escribir, lo que le permite centrarse en lo que desea hacer con los datos en lugar de en los detalles técnicos de cómo se ejecutan.
Lo que hace especial a Spark es lo que sucede detrás de escena. Spark determina automáticamente la forma más eficiente de ejecutar sus operaciones de DataFrame y luego las ejecuta en paralelo en varias computadoras en un clúster. No es necesario que lo administres tú mismo: Spark maneja cosas como dividir los datos, coordinar el trabajo y recuperarse de fallas si algo sale mal.
Debido a esto, Spark DataFrames puede manejar conjuntos de datos muy grandes, incluso aquellos demasiado grandes para caber en la memoria de una sola máquina. Al mismo tiempo, proporcionan una interfaz sencilla y familiar, lo que convierte a PySpark en una herramienta potente pero accesible para trabajar con big data.
3. Evaluación perezosa versus entusiasta
Otra fortaleza de PySpark que vale la pena conocer es su enfoque de ejecución perezosa versus entusiasta.
La mayoría de las bibliotecas de datos de Python, como Pandas, utilizan una ejecución entusiasta. Esto significa que cuando ejecuta una operación, se ejecuta inmediatamente, seguida de la siguiente operación, y así sucesivamente.
PySpark aborda esto de manera diferente mediante el uso de una técnica llamada ejecución diferida. Cuando escribe transformaciones de datos, como seleccionar columnas o filtrar filas, Spark no las ejecuta inmediatamente. En cambio, crea un plan de ejecución optimizado y ejecuta el cálculo solo cuando se activa una acción (como mostrar resultados o escribir datos en el disco). Esto permite a Spark optimizar el flujo de trabajo antes de la ejecución, haciendo que su código sea más eficiente sin esfuerzo adicional de su parte.
Datos de ejecución ansiosa (por ejemplo, pandas) ──filtro── ► resultado (calculado inmediatamente) En pandas, cada operación se ejecuta tan pronto como se llama. Esto es intuitivo pero puede resultar ineficiente para grandes conjuntos de datos. PySpark utiliza ejecución diferida. Datos de ejecución diferida (PySpark) ──filtro── ► │ └─groupby── ► (el plan se compila aquí) │ └─agg── ► (aún no hay ejecución) │ acción ── ► se ejecuta aquí
Para aclarar este punto, considere el siguiente escenario. Digamos que tenemos un marco de datos de 10 millones de registros que queremos…
a) Agregue una nueva columna vacía llamada X
b) Filtrar los datos de alguna forma que nos haga eliminar el 50% de los registros.
c) Realizar una agregación en los registros restantes para que la nueva columna X contenga el valor MAX de otro valor en esa fila
d) Imprima la fila con el valor más alto de X
En un sistema que realiza una ejecución entusiasta, como Pandas, cada paso se realiza exactamente como lo describimos anteriormente. Para 10 millones de registros, se vería así:
Agregar columna: el sistema crea una nueva versión del conjunto de datos de 10 millones de filas en la memoria, agregando la columna X. Filtro: el sistema filtra los 10 millones de filas, lo que genera 5 millones de eliminaciones, y escribe un nuevo conjunto de datos de 5 millones de filas en la memoria. Agregación: calcula el valor MAX para cada fila y actualiza la columna. Imprimir: encuentra la fila superior y te la muestra.
El problema es que hemos hecho una enorme cantidad de “trabajo pesado” (agregando una columna a 10 millones de filas) solo para inmediatamente desperdiciar la mitad de ese trabajo en el siguiente paso.
Spark, por otro lado, debido a su modelo de ejecución diferida, no realiza ningún trabajo cuando define los pasos (a), (b) o (c). En cambio, construye un plan lógico (también llamado DAG (gráfico acíclico dirigido) para hacer el trabajo.
Cuando finalmente activa el paso (d), la Acción, el optimizador de Spark analiza todo el plan y se da cuenta de que puede funcionar de manera mucho más inteligente:
Predicate Pushdown: Spark ve el filtro (elimina el 50% de los registros). En lugar de agregar la columna X a 10 millones de filas, mueve el filtrado hasta el principio. Optimización: solo agrega la columna X y agrega los 5 millones de filas restantes. Resultado: Evita procesar 5 millones de registros, ahorrando un 50% de memoria y tiempo de CPU.
Configurar el entorno de desarrollo
Ok, ya es suficiente teoría. Veamos cómo puede instalar PySpark en su sistema y ejecutar algunos fragmentos de código de ejemplo. Ahora, como texto introductorio para principiantes, la creación de un clúster de múltiples nodos del mundo real está más allá del alcance de este artículo. Pero como mencioné antes, Spark puede crear un clúster sintético en su PC o computadora portátil si es de múltiples núcleos, lo cual será si su sistema tiene menos de 10 años.
Lo primero que haremos es configurar un entorno de desarrollo separado para este trabajo, asegurándonos de que nuestros proyectos estén aislados y no interfieran entre sí. Estoy usando WSL2 Ubuntu para Windows y Conda para esta parte, pero siéntete libre de usar cualquier entorno y método al que estés acostumbrado.
Instale PySpark, etc.
# 1. Cree un nuevo entorno con Python 3.11 (muy estable para Spark) conda create -n spark_env python=3.11 -y # 2. Actívelo conda enable spark_env # 3. Instale PySpark y PyArrow (necesarios para archivos Parquet) pip install pyspark pyarrow jupyter
Para verificar que PySpark se haya instalado correctamente, escriba el comando pyspark en una ventana de terminal.
$ pyspark Python 3.11.14 | empaquetado por conda-forge | (principal, 22 de octubre de 2025, 22:46:25) [GCC 14.3.0] en Linux Escriba "ayuda", "derechos de autor", "créditos" o "licencia" para obtener más información. ADVERTENCIA: Uso de módulos de incubadora: jdk.incubator.vector ADVERTENCIA: el paquete sun.security.action no está en java.base Uso del perfil log4j predeterminado de Spark: org/apache/spark/log4j2-defaults.properties 26/01/15 16:15:21 ADVERTENCIA Utilidades: su nombre de host, tpr-desktop, se resuelve en una dirección de bucle invertido: 127.0.1.1; usando 10.255.255.254 en su lugar (en la interfaz baja) 26/01/15 16:15:21 Utilidades WARN: configure SPARK_LOCAL_IP si necesita vincularse a otra dirección Usando el perfil log4j predeterminado de Spark: org/apache/spark/log4j2-defaults.properties Configuración del nivel de registro predeterminado en "WARN". Para ajustar el nivel de registro, utilice sc.setLogLevel(newLevel). Para SparkR, utilice setLogLevel(newLevel). 26/01/15 16:15:22 ADVERTENCIA NativeCodeLoader: No se puede cargar la biblioteca nativa-hadoop para su plataforma… usando clases integradas de Java cuando corresponda ADVERTENCIA: Se ha llamado a un método obsoleto terminal en sun.misc.Unsafe ADVERTENCIA: sun.misc.Unsafe::arrayBaseOffset ha sido llamado por org.apache.spark.unsafe.Platform (archivo:/home/tom/miniconda3/envs/pandas_to_pyspark/lib/python3.11/site-packages/pyspark/jars/spark-unsafe_2.13-4.1.1.jar) ADVERTENCIA: Considere informar esto a los mantenedores de la clase org.apache.spark.unsafe.Platform ADVERTENCIA: sun.misc.Unsafe::arrayBaseOffset se eliminará en una versión futura Bienvenido a ____ __ / __/__ ___ _____/ /__ _ / _ / _ `/ __/ '_/ /__ / .__/_,_/_/ /_/_ versión 4.1.1 /_/ Uso de Python versión 3.11.14 (principal, 22 de octubre de 2025 22:46:25) Interfaz de usuario web de contexto Spark disponible en http://10.255.255.254:4040 Contexto de Spark disponible como 'sc' (master = local[*], ID de aplicación = local-1768493723158). SparkSession disponible como 'spark'. >>>
Si no ve el banner de bienvenida de Spark, entonces algo salió mal y debe volver a verificar su instalación.
Ejemplo 1: creación de un clúster local
En realidad, esto es bastante fácil. Simplemente escriba lo siguiente en su cuaderno.
from pyspark.sql import SparkSession # Inicializa la sesión de Spark spark = SparkSession.builder .master("local[*]") .appName("MyLocalCluster") .config("spark.driver.memory", "2g") .getOrCreate() # Verifica que el clúster se esté ejecutando print(f"Spark está ejecutando la versión: {spark.version}") print(f"Master URL: {spark.sparkContext.master}") # # La salida # Spark está ejecutando la versión: 4.1.1 URL maestra: local[*]
El concepto SparkSession es importante. En los primeros días de Spark, los usuarios tenían que hacer malabares con múltiples "puntos de entrada" (como SparkContext para funciones principales, SQLContext para marcos de datos y HiveContext para bases de datos). Fue confuso para los principiantes.
SparkSession se introdujo en Spark 2.0 como la "ventanilla única" para todo. Es el único punto de entrada para interactuar con la funcionalidad Spark.
Ejemplo 2: creación de un marco de datos
Crear marcos de datos y manipular los datos que contienen en PySpark será lo que harás la mayor parte del tiempo. Y es bastante sencillo de hacer. Aquí, definimos que nuestro marco de datos contendrá tres registros y tres columnas con nombre.
# 1. Defina sus datos como una lista de tuplas data = [ ("Alice", 34, "New York"), ("Bob", 45, "London"), ("Catherine", 29, "Paris") ] # 2. Defina los nombres de sus columnas columns = ["Name", "Ege", "City"] # 3. Cree el DataFrame df = spark.createDataFrame(data, columns) # 4. Mostrar el resultado df.show() # # La salida # +———+—+——–+ | Nombre|Edad| Ciudad| +———+—+——–+ | Alicia| 34|Nueva York| | Bob| 45| Londres| |Catalina| 29| París| +———+—+——–+
Lo más probable es que cualquier marco de datos que utilice se cree inicialmente leyendo datos de un archivo o base de datos. Cree un archivo CSV llamado sales_data.csv en su sistema con el siguiente contenido.
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
Crear un marco de datos a partir de un archivo como este es simple,
# Cargar el archivo CSV df = spark.read.format("csv") .option("header", "true") .option("inferSchema", "true") .load("sales_data.csv") # Mostrar los datos print("Dataframe Content:") df.show() # Mostrar los tipos de datos (Esquema) print("Data Schema:") df.printSchema() # # La salida # Contenido del marco de datos: +————–+————-+———-+———-+———-+ |id_transacción|nombre_cliente|monto_neto|monto_impuesto| es_miembro| +————-+————-+———-+———-+———-+ | 101| Alicia| 250,5| 25.05| verdadero| | 102| Bob| 120,0| 6.0| falso| | 103| charlie| 450,75| 25.07| verdadero| | 104| David| 89,99| 5.73| falso| +————–+————-+———-+———-+———-+ Esquema de datos: raíz |– ID_transacción: entero (anulable = verdadero) |– nombre_cliente: cadena (anulable = verdadero) |– monto_neto: doble (anulable = verdadero) |– monto_impuesto: doble (anulable = verdadero) |– is_member: cadena (anulable = verdadero)
Ejemplo 3: procesamiento de datos
Por supuesto, una vez que tenga sus datos de entrada en un marco de datos, lo siguiente que querrá hacer es procesarlos o manipularlos de alguna manera. Eso también es fácil. Refiriéndose a los datos_ventas que acabamos de cargar, digamos que queremos calcular el monto bruto (neto + impuestos) y la tasa impositiva como porcentaje del monto bruto para cada registro y agregarlos a nuestro marco de datos inicial.
desde pyspark.sql la importación funciona como F # 1. Agregue 'monto_bruto' sumando el neto y el impuesto # 2. Agregue 'porcentaje_impuesto' dividiendo el impuesto por el nuevo monto bruto df_extended = df.withColumn("monto_bruto", F.col("monto_neto") + F.col("monto_impuesto")) .withColumn("percentaje_impuesto", (F.col("tax_amount") / (F.col("net_amount") + F.col("tax_amount"))) * 100) # 3. Opcional: redondear el porcentaje a 2 decimales para facilitar la lectura df_extended = df_extended.withColumn("tax_percentage", F.round(F.col("tax_percentage"), 2)) # Mostrar las nuevas columnas junto con las antiguas unos df_extended.show() # # La salida # +————–+————-+———-+———-+———-+————+————–+ |transaction_id|customer_name|net_amount|tax_amount| es_miembro|monto_bruto|porcentaje_impuesto| +————–+————-+———-+———-+———-+————+————–+ | 101| Alicia| 250,5| 25.05| verdadero| 275,55| 9.09| | 102| Bob| 120,0| 6.0| falso| 126,0| 4.76| | 103| charlie| 450,75| 25.07| verdadero| 475,82| 5.27| | 104| David| 89,99| 5.73| falso| 95,72| 5,99| +————-+————-+———-+———-+———-+————+————–+
Resumen
Con esto concluye nuestra breve estancia en el mundo de la informática distribuida con PySpark. Le expliqué qué es PySpark y por qué debería considerar usarlo si los datos que está procesando exceden sus límites de memoria. En resumen, la capacidad de PySpark para escalar a grandes clústeres de múltiples nodos, su modelo de ejecución diferida y la estructura de datos del marco de datos lo convierten en una potencia de procesamiento de datos ideal.
PySpark se utiliza ampliamente en ingeniería de datos, análisis y procesos de aprendizaje automático. Se integra bien con plataformas en la nube, admite una variedad de fuentes de datos (como CSV, Parquet y bases de datos) y escala desde una computadora portátil hasta grandes clústeres de producción.
Si se siente cómodo con Python y desea trabajar con grandes conjuntos de datos sin abandonar la sintaxis familiar, PySpark es un excelente siguiente paso. Cierra la brecha entre el análisis de datos simple y el procesamiento de datos a gran escala, lo que lo convierte en una herramienta valiosa para cualquiera que ingrese al mundo del big data.
Con suerte, puede utilizar mis simples ejemplos de codificación y explicaciones para dar el siguiente paso hacia el uso de PySpark en el mundo real, en un clúster real, y realizar un procesamiento adecuado de big data.