En el ecosistema del procesamiento de datos masivos, Apache Spark se ha consolidado como el motor unificado de computación más rápido y versátil del mercado.
Su capacidad para procesar información en memoria a través de clústeres distribuidos se apoya en dos conceptos fundamentales: las transformaciones y acciones en Spark.
Comprender cómo interactúan estas dos categorías de operaciones sobre las estructuras de datos (como los RDDs, DataFrames y Datasets) es un requisito primordial para cualquier profesional que busque optimizar flujos de trabajo en Big Data.
A lo largo de esta guía técnica analizaremos la arquitectura del Resilient Distributed Dataset (RDD), la estrategia de evaluación perezosa (lazy evaluation), la clasificación de transformaciones estrechas y anchas, y repasaremos ejemplos prácticos en Scala para cada comando esencial.
¿Qué es un RDD y cómo funciona el modelo de ejecución en Spark?
Para entender las transformaciones y acciones en Spark, primero debemos definir la abstracción sobre la que operan: el RDD (Resilient Distributed Dataset).
Un RDD es una colección de objetos inmutable, dividida en particiones y distribuida entre los diferentes nodos de un clúster gestionado por un Cluster Manager.
Los RDDs poseen tres características operativas clave:
- Inmutabilidad: Un RDD no se puede modificar una vez creado. Cualquier operación aplicada genera un nuevo RDD.
- Tolerancia a fallos (Resiliencia): Si una partición de datos se pierde debido al fallo de un nodo, Spark es capaz de reconstruirla reconstruyendo el historial de transformaciones.
- Evaluación perezosa (Lazy Evaluation): Spark no ejecuta los cálculos en el momento en que se define una transformación. En su lugar, registra la instrucción en un Grafo Acíclico Dirigido (DAG) y pospone la ejecución real hasta que se invoca una acción.
| Criterio técnico | Transformaciones en Spark | Acciones en Spark |
|---|---|---|
| Resultado devuelto | Generan un nuevo RDD a partir de un dataset existente. | Devuelven un valor escalar, un array o guardan datos en disco. |
| Modo de ejecución | Evaluación perezosa (Lazy): Se registran en el DAG. | Eager / Inmediata: Desencadenan el cálculo del clúster. |
| Impacto en el DAG | Construyen la secuencia de linaje y el plan de ejecucion. | Inician la ejecución de los jobs, stages y tasks. |
| Ejemplos representativos | map, filter, flatMap, union, intersection. |
collect, count, reduce, first, saveAsTextFile. |
Transformaciones en Spark: tipos y ejemplos en Scala

Las transformaciones se dividen conceptualmente en dos categorías según la necesidad de mover información a través de la red:
- Transformaciones estrechas (Narrow): Cada partición del RDD de origen es utilizada por como máximo una partición del RDD resultante (por ejemplo,
mapofilter). No requieren redistribución de datos (shuffle). - Transformaciones anchas (Wide): Múltiples particiones de origen contribuyen a generar una partición de destino (por ejemplo,
groupByKeyoreduceByKey). Requieren un proceso de shuffle a través de los nodos del clúster.
Repasemos los métodos de transformación más comunes mediante la inicialización de colecciones con la opción sc.parallelize en Scala:
1. map(func)
Devuelve un nuevo RDD tras pasar cada elemento del conjunto original a través de una función de mapeo individual:
// Inicializar un RDD de números enteros
val v1 = sc.parallelize(List(2, 4, 8))
# Multiplicar cada elemento por 2
val v2 = v1.map(_ * 2)
# Mostrar el resultado
v2.collect
// res0: Array[Int] = Array(4, 8, 16)
2. filter(func)
Realiza un filtrado sobre los elementos del RDD original, devolviendo un nuevo RDD únicamente con aquellos valores que cumplen la condición booleana especificada:
// Crear un RDD con cadenas de texto
val v1 = sc.parallelize(List("ABC", "BCD", "DEF"))
// Filtrar los elementos que contienen la letra "A"
val v2 = v1.filter(_.contains("A"))
v2.collect
// res0: Array[String] = Array(ABC)
3. flatMap(func)
Similar a la operación map, pero cada elemento de entrada se puede mapear a cero o más elementos de salida, aplanando el resultado en una secuencia continua de valores:
val x = sc.parallelize(List("Ejemplo proyecto Alejandro", "Hola mundo"), 2)
// Uso de map (devuelve un array de arrays)
val yMap = x.map(x => x.split(" "))
yMap.collect
// res0: Array[Array[String]] = Array(Array(Ejemplo, proyecto, Alejandro), Array(Hola, mundo))
// Uso de flatMap (devuelve una secuencia aplanada)
val yFlat = x.flatMap(x => x.split(" "))
yFlat.collect
// res1: Array[String] = Array(Ejemplo, proyecto, Alejandro, Hola, mundo)
4. mapPartitions(func)
Similar a map, pero se ejecuta de forma independiente sobre cada partición completa del RDD en lugar de procesar elemento por elemento, optimizando la creación de objetos pesados o conexiones a **bases de datos**:
val a = sc.parallelize(1 to 9, 3)
def myfunc[T](iter: Iterator[T]): Iterator[(T, T)] = {
var res = List[(T, T)]()
var pre = iter.next
while (iter.hasNext) {
val cur = iter.next
res .::= (pre, cur)
pre = cur
}
res.iterator
}
a.mapPartitions(myfunc).collect
// res0: Array[(Int, Int)] = Array((2,3), (1,2), (5,6), (4,5), (8,9), (7,8))
5. sample(withReplacement, fraction, seed)
Extrae una muestra aleatoria o fracción de los datos con o sin reemplazo, utilizando opcionalmente una semilla generadora para asegurar la reproducibilidad del experimento:
val randRDD = sc.parallelize(List(
(7, "cat"), (6, "mouse"), (7, "cup"),
(6, "book"), (7, "tv"), (6, "screen"), (7, "heater")
))
val sampleMap = List((7, 0.4), (6, 0.6)).toMap
randRDD.sampleByKey(false, sampleMap, 42).collect
// res0: Array[(Int, String)] = Array((6,book), (7,tv), (7,heater))
6. union(otherDataset)
Combina los elementos de dos RDDs distintos, devolviendo un nuevo RDD que contiene la unión de ambos conjuntos:
val a = sc.parallelize(1 to 3, 1)
val b = sc.parallelize(5 to 7, 1)
a.union(b).collect()
// res0: Array[Int] = Array(1, 2, 3, 5, 6, 7)
7. intersection(otherDataset)
Evalúa dos conjuntos de datos distribuidos y devuelve únicamente aquellos elementos que están presentes en ambos RDDs:
val x = sc.parallelize(1 to 20)
val y = sc.parallelize(10 to 30)
val z = x.intersection(y)
z.collect
// res0: Array[Int] = Array(16, 14, 12, 18, 20, 10, 13, 19, 15, 11, 17)
Acciones en Apache Spark: desencadenando la ejecución real
Las acciones son las funciones encargadas de romper la evaluación perezosa.
Al invocar una acción, Spark compila el DAG registrado, evalúa el plan de ejecución y envía las tareas a los workers para devolver el resultado al programa controlador (driver) o escribir los datos en un sistema de archivos distribuido.
1. reduce(func)
Agrega todos los elementos del dataset utilizando una función conmutativa y asociativa para que pueda calcularse correctamente en paralelo en los **diferentes nodos**:
val a = sc.parallelize(1 to 100, 3)
// Sumar recursivamente todos los elementos
a.reduce(_ + _)
// res0: Int = 5050
2. collect()
Recupera todos los elementos de un RDD distribuido, los convierte en un array local y los traslada a la memoria del driver:
val c = sc.parallelize(List("Gnu", "Cat", "Rat", "Dog", "Gnu", "Rat"), 2)
c.collect
// res0: Array[String] = Array(Gnu, Cat, Rat, Dog, Gnu, Rat)
Precaución de producción: Utilizar collect() en datasets masivos que superen la capacidad de la memoria RAM del driver provocará una excepción de tipo OutOfMemoryError. Para inspeccionar conjuntos de datos grandes, utiliza take(n).
3. count()
Devuelve el número total de elementos contenidos en el dataset:
val a = sc.parallelize(1 to 4)
a.count
// res0: Long = 4
4. first()
Retorna únicamente el primer elemento del conjunto de datos distribuido:
val c = sc.parallelize(List("Gnu", "Cat", "Rat", "Dog"), 2)
c.first
// res0: String = Gnu
5. take(n)
Devuelve un array con los primeros n elementos del dataset. Es la alternativa segura a collect() para inspeccionar datos:
val b = sc.parallelize(List("dog", "cat", "ape", "salmon", "gnu"), 2)
b.take(2)
// res0: Array[String] = Array(dog, cat)
6. takeSample(withReplacement, num, [seed])
Devuelve un array local con una muestra aleatoria de elementos numéricos del RDD, con o sin sustitución:
val x = sc.parallelize(1 to 200, 3)
x.takeSample(true, 20, 1)
// res0: Array[Int] = Array(74, 164, 160, 41, 123, 27, 134, 5, 22, 185, 129, 107, 140, 191, 187, 26, 55, 186, 181, 60)
7. takeOrdered(n, [ordering])
Devuelve los primeros n elementos del RDD utilizando su ordenamiento implícito natural o un comparador personalizado:
val b = sc.parallelize(List("dog", "cat", "ape", "salmon", "gnu"), 2)
b.takeOrdered(2)
// res0: Array[String] = Array(ape, cat)
8. saveAsTextFile(path)
Escribe los elementos del RDD en el sistema de archivos local o distribuido (como HDFS o Amazon S3) guardando cada partición como un archivo de texto independiente dentro de la ruta configurada:
val a = sc.parallelize(1 to 10000, 3)
a.saveAsTextFile("/home/usuario/datos")
// Estructura de carpetas generada en disco:
// /home/usuario/datos/part-00000
// /home/usuario/datos/part-00001
// /home/usuario/datos/part-00002
// /home/usuario/datos/_SUCCESS
Buenas prácticas de optimización y persistencia en memoria
Para exprimir el máximo rendimiento de Apache Spark en proyectos reales de ingeniería de datos, conviene aplicar las siguientes recomendaciones:
- Persistencia de RDDs reutilizados (cache/persist): Si vas a invocar múltiples acciones sobre un mismo RDD transformado, utiliza
rdd.cache()ordd.persist(). De lo contrario, Spark reevaluará toda la cadena de transformaciones desde el origen en cada acción. - Minimizar el uso de Shuffle: Las transformaciones anchas requieren mover datos a través de la red entre nodos. Prioriza filtrados previos con
filterantes de ejecutar operaciones comojoinogroupByKey. - Evolución hacia DataFrames y Spark SQL: Aunque los RDDs proporcionan control a bajo nivel, en la actualidad el motor de Spark está optimizado para operar mediante DataFrames y Spark SQL, apoyándose en el optimizador Catalyst para mejorar el rendimiento de las consultas.
Cómo conectar el ecosistema Big Data con tu carrera profesional
Dominar la arquitectura distribuida, la optimización de memoria y la programación con transformaciones y acciones en Spark es una competencia imprescindible para consolidarse como Data Engineer o Data Scientist.
Si estás organizando tu plan de estudio desde cero, te invitamos a consultar nuestra guía sobre qué aprender primero en programación.
Comprender cómo se estructuran los proyectos evaluando la diferencia entre frontend, backend y full stack te ayudará a entender el flujo de datos hacia la interfaz.
Durante el desarrollo de tus scripts y canalizaciones en Spark, administrarás repositorios sabiendo qué es Git y por qué es tan importante para trabajar en equipo.
A nivel de almacenamiento de información distribuida, verificarás transacciones seguras comprobando qué es ACID en bases de datos.
En organizaciones avanzadas que automatizan el despliegue de modelos analíticos en la nube, este conocimiento conecta con entender qué es MLOps y por qué es clave en ingeniería de software.
Para aquellos profesionales que buscan especializarse en la gestión de infraestructura en la nube y automatización de sistemas distribuidos, recomendamos explorar el Programa Técnico Avanzado en DevOps con IA y LLMops.
Conclusión
Saber utilizar las transformaciones y acciones en Spark es el pilar para construir tuberías de procesamiento de información escalables en Big Data.
Comprender la evaluación perezosa, evitar el consumo excesivo de memoria con collect() y estructurar transformaciones eficientes te capacitará para procesar datos masivos con total solvencia.
Aprender estas tecnologías con una metodología práctica y guiada por mentores en activo te abrirá las puertas de las principales empresas del sector IT.

Si quieres dominar Apache Spark, Scala, Python, algoritmos de aprendizaje automático y sistemas de almacenamiento masivo con proyectos reales respaldado por expertos en activo, descubre el Bootcamp Full Stack Big Data & Machine Learning de KeepCoding y transforma tu futuro profesional hoy mismo.



