El Módulo 4 terminó con una tabla que decía dónde vive cada dato de Kilómetro Cero: los pedidos en Cassandra, el stock en PostgreSQL, el catálogo en Redis, las fotos en MinIO y, en el lago de datos de HDFS, un fichero por día con los cientos de miles de eventos de pedidos.eventos y los millones de clics de la web. Repartir los datos resolvía el problema de guardarlos; ahora aparece el siguiente: calcular algo con ellos. Sumar las ventas de la Semana de la Vendimia por productor y mercado sobre los eventos de siete días, o entrenar recomendaciones con los clics de un año, no cabe en una máquina, y aunque cupiera tardaría horas. Esta lección explica qué cambia cuando lo que se distribuye es un cálculo y no solo un dato: por qué conviene llevar el código hasta donde están los datos, qué formas hay de repartir el trabajo (scatter/gather, divide y vencerás, colas de trabajo, BSP, dataflow, actores), en qué se diferencian los lotes, los flujos y las consultas interactivas, y qué problemas nuevos aparecen (el coste de mover datos entre nodos, el sesgo, los rezagados, la reejecución tras un fallo, el límite del escalado). Todo lo que viene después en el módulo (MapReduce, Spark, Flink, Airflow) es una implementación concreta de estas ideas, así que merece la pena entenderlas primero con código Python de una sola máquina, que es lo que haremos en simulaciones/.
Contenido
- Distribuir un cálculo: llevar el código a los datos
- Paralelismo de datos y paralelismo de tareas
- Patrones de computación distribuida
- Lotes, flujos e interactivo
- El coste de la comunicación: el shuffle
- Sesgo de datos
- Rezagados y ejecución especulativa
- Tolerancia a fallos por reejecución determinista
- Escalabilidad: Amdahl y Gustafson
- Práctica: scatter/gather, cola de trabajo y sesgo en
simulaciones/ - Errores Comunes y Consejos
- Ejercicios
- Conclusión
- Distribuir un cálculo: llevar el código a los datos
Cuando el catálogo cabía en una tabla de PostgreSQL, "calcular las ventas por productor" era una consulta SQL: el motor leía las filas de su disco local, las agrupaba y devolvía veinte números. El dato y el cálculo estaban en la misma máquina. En el lago de datos de 04-02 ya no es así: /km0/eventos/2026-09-14/pedidos.jsonl tiene 150 MB en dos bloques que viven en tres DataNodes distintos, y los siete días de la Semana de la Vendimia suman más de un gigabyte repartido por todo el clúster. Hay dos formas de calcular sobre eso:
- Llevar los datos al código. Un proceso en
analiticadescarga el gigabyte por la red, lo recorre y suma. Funciona, pero la red es el recurso más lento que tenemos (la falacia 3 de 01-04, "el ancho de banda es infinito"): a 1 Gbit/s, mover 1 GB son unos 10 segundos solo de transferencia, y con un año de clics (terabytes) el enfoque es sencillamente inviable. Además, el proceso que suma es uno solo, así que el cálculo no escala. - Llevar el código a los datos (data locality). Enviar a cada DataNode un programa pequeño (unos kilobytes) que lea el bloque que ya tiene en su disco local, calcule un resultado parcial (unos kilobytes: ventas por productor de ese bloque) y devuelva solo eso. Se mueve el código y los resultados, que son minúsculos; los datos, que son enormes, no se mueven.
Esta inversión es la idea central de toda la computación distribuida moderna y la razón por la que HDFS y MapReduce nacieron juntos: el sistema de archivos expone dónde está cada bloque (hdfs fsck -locations lo mostraba en 04-02) precisamente para que el planificador de cálculo pueda ejecutar cada tarea en el nodo que tiene su bloque, o al menos en el mismo rack. Cuando la localidad no es posible (el nodo está ocupado, o los datos están en un almacenamiento de objetos como MinIO que no expone localidad), se paga la red, y por eso las plataformas modernas separan cómputo y almacenamiento pero compensan con redes de 25–100 Gbit/s y formatos columnares que leen solo las columnas necesarias (05-03).
El segundo cambio es que el cálculo se convierte en muchas tareas independientes más una fase de combinación. En vez de un programa que recorre todo, escribimos una función que procesa un trozo y otra que combina resultados parciales. Esa descomposición no es gratuita: hay cálculos que se descomponen de forma natural (sumar ventas por productor) y otros que no (ordenar todos los eventos por importe, calcular una mediana exacta, recorrer un grafo de "clientes que compraron lo mismo"). Los patrones del apartado 3 son las formas conocidas de descomponer.
- Paralelismo de datos y paralelismo de tareas
Hay dos maneras de repartir un trabajo entre varios nodos, y conviene distinguirlas porque exigen infraestructuras distintas:
| Paralelismo de datos | Paralelismo de tareas | |
|---|---|---|
| Qué se reparte | Los datos: cada nodo ejecuta el mismo código sobre un trozo distinto | Las tareas: cada nodo ejecuta código distinto sobre los mismos datos o sobre datos relacionados |
| Ejemplo en Kilómetro Cero | Sumar ventas por productor: 8 nodos, cada uno con un octavo de pedidos.jsonl |
Para un pedido: uno calcula el importe, otro valida el stock, otro estima el reparto; los tres a la vez |
| Cómo escala | Con el tamaño de los datos: el doble de datos, el doble de nodos, el mismo tiempo | Con el número de tareas distintas, que suele ser pequeño y fijo |
| Coordinación | Al final, para combinar resultados parciales | Entre tareas, por dependencias (B necesita el resultado de A) |
| Dificultad principal | Repartir de forma equilibrada; combinar sin perder información | Sincronizar y gestionar dependencias; el camino crítico limita |
| Modelos que lo explotan | MapReduce, Spark, Flink, SQL distribuido | Pipelines de tareas (Airflow, 05-05), microservicios que colaboran (Módulo 8) |
La analítica de Kilómetro Cero es casi toda paralelismo de datos, y por eso este módulo se centra en él. Pero los dos se combinan: el pipeline diario de 05-05 es paralelismo de tareas (validar, agregar, cargar, notificar) en el que la tarea "agregar" es, por dentro, paralelismo de datos en Spark. Y la ley que gobierna el escalado (apartado 9) es distinta en cada caso: el paralelismo de datos se acerca al escalado lineal porque la fracción secuencial es pequeña; el de tareas está limitado por la cadena de dependencias más larga.
- Patrones de computación distribuida
3.1 Scatter/gather (dispersar y reunir)
Es el patrón más simple y el que ya apareció, sin nombre, en 04-01 al hablar de índices secundarios en Cassandra: un coordinador reparte el trabajo entre N trabajadores (scatter), cada uno calcula sobre su parte, y el coordinador reúne y combina los resultados parciales (gather).
flowchart LR
C[Coordinador<br/>analitica] -- trozo 1 --> T1[Trabajador 1<br/>ventas parciales]
C -- trozo 2 --> T2[Trabajador 2<br/>ventas parciales]
C -- trozo 3 --> T3[Trabajador 3<br/>ventas parciales]
C -- trozo 4 --> T4[Trabajador 4<br/>ventas parciales]
T1 --> G[Gather:<br/>sumar parciales]
T2 --> G
T3 --> G
T4 --> G
G --> R[ventas por productor]
Funciona cuando el cálculo es asociativo y conmutativo: sumar, contar, máximo, mínimo, unión de conjuntos. Da igual en qué orden y en qué agrupación se combinen los parciales, el resultado es el mismo. Una media no es directamente combinable (la media de medias es incorrecta si los trozos tienen distinto tamaño), pero se convierte en combinable llevando (suma, cuenta) en lugar de la media; una mediana exacta no se convierte, y hay que aproximarla o pagar un ordenamiento global. Es la primera pregunta que hay que hacerse ante cualquier cálculo distribuido: ¿qué resultado parcial devuelve cada trozo, y cómo se combinan dos parciales?
El límite del scatter/gather es que el coordinador es único: reparte, espera a todos y combina. Si combinar es costoso (millones de claves distintas) o los parciales son grandes, el coordinador se convierte en el cuello de botella. MapReduce (05-02) resuelve exactamente esto distribuyendo también la fase de combinación.
3.2 Divide y vencerás distribuido
Es el scatter/gather aplicado recursivamente: un problema se parte en subproblemas, cada subproblema se vuelve a partir hasta que cabe en un nodo, y los resultados se combinan subiendo por el árbol. Un ordenamiento distribuido funciona así: cada nodo ordena su trozo (mergesort local), y las combinaciones sucesivas mezclan listas ordenadas. Los frameworks lo usan para operaciones de reducción con muchos nodos: en vez de que un coordinador sume 1 000 parciales, se suman en árbol (treeReduce en Spark), con 1 000 → 32 → 1 combinaciones en paralelo.
3.3 Cola de trabajo con trabajadores (work queue)
En el scatter/gather el coordinador decide de antemano qué trozo va a cada trabajador. Con una cola de trabajo no decide: publica todas las tareas en una cola y cada trabajador toma la siguiente cuando termina la anterior. Es el patrón de los consumidores competidores de 02-04 aplicado al cálculo, y tiene dos ventajas grandes:
- Equilibrio dinámico. Si un trabajador es lento (máquina vieja, tarea grande), simplemente toma menos tareas; los rápidos absorben el resto. Nadie espera ocioso.
- Tolerancia a fallos natural. Si un trabajador muere a mitad de una tarea, la tarea vuelve a la cola (por el lease o el ack pendiente) y otro la ejecuta. Para que eso sea correcto, la tarea debe ser idempotente: ejecutarla dos veces debe dar el mismo resultado que una (apartado 8).
Es el modelo interno de casi todos los planificadores: YARN, Spark y Flink mantienen colas de tareas pendientes y las asignan a los ejecutores que quedan libres, con preferencia por el nodo que tiene los datos. Lo veremos en simulaciones/cola_trabajo.py.
3.4 Bulk Synchronous Parallel (BSP)
Hay cálculos que no se hacen en una pasada, sino en iteraciones donde cada paso depende del anterior: el PageRank de un grafo, el entrenamiento de un modelo por descenso de gradiente, la propagación de "clientes que compraron lo mismo" por el grafo de pedidos para las recomendaciones de Kilómetro Cero. El modelo BSP (Valiant, 1990) organiza esos cálculos en supersteps (superpasos):
- Cálculo local: cada nodo trabaja solo con sus datos y los mensajes que recibió en el superpaso anterior.
- Comunicación: cada nodo envía mensajes a los demás (por ejemplo, un vértice del grafo envía su valor a sus vecinos).
- Barrera: nadie empieza el siguiente superpaso hasta que todos han terminado el actual y todos los mensajes han llegado.
flowchart TB
subgraph S1[Superpaso 1]
direction LR
A1[Nodo A<br/>calcula] --> M1[mensajes]
B1[Nodo B<br/>calcula] --> M1
C1[Nodo C<br/>calcula] --> M1
end
M1 --> BAR1{{Barrera: todos han terminado}}
BAR1 --> S2
subgraph S2[Superpaso 2]
direction LR
A2[Nodo A<br/>calcula] --> M2[mensajes]
B2[Nodo B<br/>calcula] --> M2
C2[Nodo C<br/>calcula] --> M2
end
M2 --> BAR2{{Barrera}}
BAR2 --> FIN[... hasta converger]
La barrera es lo que hace el modelo fácil de razonar (dentro de un superpaso no hay carreras: cada nodo solo ve mensajes del superpaso anterior) y también lo que lo hace sensible a los rezagados (apartado 7): el superpaso dura lo que dure el nodo más lento. Pregel de Google, y sus descendientes Apache Giraph y GraphX de Spark, son BSP "pensando como un vértice": cada vértice del grafo recibe mensajes, actualiza su valor y envía mensajes a sus vecinos, superpaso tras superpaso hasta que ningún vértice cambia. En Kilómetro Cero, "productos que suelen comprarse juntos" es un grafo donde los vértices son productos y las aristas pesan por el número de pedidos compartidos; dos o tres superpasos de propagación bastan para encontrar que queso-curado y vino-crianza están más cerca de lo que sugieren sus categorías.
3.5 Pipeline / dataflow: el DAG de operadores
El modelo de flujo de datos (dataflow) describe el cálculo como un grafo dirigido acíclico (DAG) de operadores: leer, filtrar, transformar, agrupar, unir, escribir. Cada operador recibe datos de los anteriores y emite datos a los siguientes; los datos fluyen por las aristas. El sistema decide cómo paralelizar cada operador (cuántas instancias, en qué nodos), cómo encadenar operadores que no necesitan redistribuir datos (filtrar y transformar pueden ir en el mismo proceso, fila a fila), y dónde hay que redistribuir (agrupar por productor obliga a que todas las filas de un productor lleguen a la misma instancia: es el shuffle del apartado 5).
flowchart LR
L[leer pedidos.jsonl] --> F[filtrar tipo = pedido.creado]
F --> E[explotar líneas del pedido]
E --> S[[shuffle por productor y mercado]]
S --> A[sumar importe]
C[leer catálogo] --> J
A --> J[unir con catálogo]
J --> W[escribir Parquet]
Es el modelo de Spark (05-03) y de Flink (05-04), y también, con otro vocabulario, el de los motores SQL distribuidos y el de los planificadores de pipelines (05-05, donde los nodos del DAG son trabajos enteros en vez de operadores). Frente a MapReduce, que obliga a expresar todo como parejas map/reduce encadenadas, el dataflow deja al programador escribir la transformación completa y al optimizador decidir las fases. Frente a BSP, el dataflow no tiene barreras globales: un operador procesa en cuanto tiene datos, y solo el shuffle sincroniza.
3.6 Actores
El modelo de actores (Hewitt, 1973) toma otro camino: no reparte datos ni describe un grafo, sino que modela el sistema como muchos objetos pequeños (actores) que solo se comunican por mensajes asíncronos, cada uno con su estado privado y su buzón, procesando un mensaje cada vez. No hay memoria compartida ni bloqueos: un actor recibe un mensaje, cambia su estado, envía mensajes a otros actores o crea actores nuevos. Erlang lo lleva en el lenguaje desde los años 80 (las centralitas de Ericsson, RabbitMQ de 02-04 está escrito en Erlang) y Akka lo trajo a la JVM. Encaja de forma natural con estado por entidad: un actor por repartidor de furgoneta-3 que recibe sus posiciones y mantiene su ruta, o un actor por pedido que ejecuta su saga (03-05). Es menos adecuado para el cálculo masivo sobre datos históricos, que es lo que ocupa este módulo, así que lo dejamos nombrado: su lugar vuelve a aparecer cuando el estado por clave es el protagonista, y de hecho los operadores con estado de Flink (05-04) se parecen mucho a actores particionados por clave.
3.7 Resumen de patrones
| Patrón | Reparto | Coordinación | Encaja cuando | Ejemplo |
|---|---|---|---|---|
| Scatter/gather | Estático, por el coordinador | Una vez, al final | Cálculo asociativo, pocos parciales | Ventas por productor de un día |
| Divide y vencerás | Recursivo | En árbol | Combinación costosa, muchos nodos | Ordenar todos los eventos por importe |
| Cola de trabajo | Dinámico, por demanda | Ninguna entre trabajadores | Tareas de duración desigual, fallos frecuentes | Redimensionar 100 000 fotos de MinIO |
| BSP | Por vértice/partición | Barrera por superpaso | Iterativo, grafos | Productos comprados juntos |
| Dataflow (DAG) | Por operador y partición | Solo en el shuffle | Transformaciones encadenadas, lotes o flujos | Pipeline de ventas diarias, panel de reparto |
| Actores | Por entidad | Mensajes asíncronos | Estado por entidad, concurrencia | Un actor por repartidor |
- Lotes, flujos e interactivo
Independientemente del patrón, hay tres modos de procesar según cuándo llegan los datos y cuándo se necesita la respuesta:
| Por lotes (batch) | Por flujos (streaming) | Interactivo (ad hoc) | |
|---|---|---|---|
| Entrada | Un conjunto acotado y completo: "los eventos del 14 de septiembre" | Un flujo no acotado que no termina: pedidos.eventos en Kafka |
Un conjunto acotado, pero la pregunta se decide en el momento |
| Cuándo se calcula | Programado: cada noche, cada hora | Continuamente, evento a evento o en microlotes | Cuando alguien pregunta |
| Latencia esperada | Minutos a horas | Milisegundos a segundos | Segundos |
| Tamaño de datos por ejecución | Gigabytes a petabytes | Kilobytes por evento; millones de eventos por hora | Gigabytes, con índices o formatos columnares para ir rápido |
| Resultado | Completo y exacto sobre el conjunto | Aproximado o provisional, se refina al llegar más datos (05-04) | Exacto sobre lo que hay |
| Tolerancia a fallos | Reejecutar el lote entero | Checkpoints de estado y offsets | Reejecutar la consulta |
| Kilómetro Cero | Ventas por productor/mercado/día; recomendaciones entrenadas cada noche con los clics | Panel de reparto con posiciones de furgoneta-3; alerta de stock bajo |
"¿Cuánto vendió Bodega Roble Alto en Lleida durante la Semana de la Vendimia?" desde la consola de analitica |
| Herramientas | MapReduce (05-02), Spark (05-03) | Kafka Streams, Flink, Spark Structured Streaming (05-04) | Spark SQL, Presto/Trino, Hive nombrado (05-02) |
Los tres modos comparten los patrones del apartado 3 y los problemas de los apartados 5–8, pero cada uno los sufre de forma distinta: el sesgo en un lote alarga una noche, en un flujo atasca una partición para siempre. Y la frontera entre lotes y flujos es más difusa de lo que parece: un lote de "los eventos del 14 de septiembre" es un flujo al que se ha puesto principio y fin; un flujo procesado en microlotes de un segundo son lotes muy pequeños. Esa idea, que un motor puede tratar ambos modos con el mismo DAG, es la que Spark y Flink explotan y la que 05-04 desarrollará.
- El coste de la comunicación: el shuffle
Distribuir un cálculo tiene un coste que no existe en una máquina: mover datos entre nodos cuando el siguiente paso necesita agruparlos de otra manera. Sumar ventas por productor exige que todas las líneas de Quesería Montblanc acaben en el mismo nodo, y esas líneas están repartidas por todos los bloques de todos los nodos. La operación que las junta se llama shuffle (barajar), y es con diferencia la fase más cara de cualquier trabajo distribuido, porque:
- Implica todos-a-todos: cada nodo envía una parte de sus datos a cada uno de los demás. Con N nodos son N² flujos de red.
- Suele pasar por disco: los datos se serializan, se escriben ordenados por clave de destino, se transfieren y se vuelven a leer. En MapReduce cada shuffle es una escritura completa en disco (05-02); Spark lo mantiene en memoria cuando puede (05-03).
- No se puede solapar del todo con el cálculo: el receptor necesita todos sus datos antes de agrupar, así que el shuffle es una barrera implícita.
La regla práctica es minimizar lo que cruza el shuffle: reducir antes de mover. Si cada nodo suma localmente sus líneas por productor antes de enviarlas (un combiner, en vocabulario de 05-02; una agregación parcial, en Spark), en lugar de mover 250 000 líneas mueve 20 parciales por nodo. Esa es la razón de insistir en que el cálculo sea asociativo: solo entonces se puede prerreducir. Y explica el consejo de 04-01 sobre las tablas materializadas por mercado y productor: alguien había pagado el shuffle una vez, en el momento de escribir, para no pagarlo en cada consulta.
Un cálculo que no necesita shuffle (filtrar, transformar fila a fila, sumar un total global que se combina en un solo número) es embarazosamente paralelo y escala casi linealmente. Uno que necesita varios shuffles encadenados (agrupar por productor, unir con el catálogo, reagrupar por mercado) está limitado por la red y por el peor de los apartados siguientes.
- Sesgo de datos
El reparto ideal da a cada nodo la misma cantidad de trabajo. El reparto real depende de la distribución de las claves, y las distribuciones reales son desiguales. Durante la "Semana del Queso Artesano" Quesería Montblanc concentra la mitad de las líneas de pedido de Kilómetro Cero; si el shuffle agrupa por productor, el nodo que recibe queseria-montblanc procesa la mitad de los datos mientras los otros siete se reparten la otra mitad. El trabajo tarda lo que tarda ese nodo: con 8 nodos y una clave que pesa el 50 %, el speedup máximo es 2, no 8. Es sesgo de datos (data skew), y es la causa más frecuente de trabajos distribuidos que "no escalan aunque añada máquinas".
Se detecta mirando la duración de las tareas de una misma fase: si la mayoría acaba en 20 s y una tarda 3 min, hay una clave caliente. Las soluciones, todas ellas formas de romper la clave gorda, se tratan en la práctica (apartado 10.3) y se retomarán en 05-03 con el nombre que les da Spark, salting:
- Prerreducir antes del shuffle: si cada nodo suma sus líneas de Quesería Montblanc, lo que cruza la red es un parcial por nodo, no la mitad de los datos.
- Repartir la clave caliente añadiendo un sufijo aleatorio (
queseria-montblanc#0…#7), agrupar por la clave con sufijo y volver a agrupar los ocho parciales en un segundo paso mucho más pequeño. - Elegir otra clave de partición cuando el cálculo lo permite: agrupar por
(productor, mercado)reparte a Quesería Montblanc entre cuatro mercados. - Tratar aparte las claves calientes conocidas (filtrarlas, procesarlas con más paralelismo) y unir después.
El sesgo también existe en la entrada: si el fichero del 14 de septiembre son dos bloques y el del 15 son diez, los trabajos por día tendrán duraciones muy distintas. Y existe en el tiempo: los eventos de Kafka se concentran a mediodía y a última hora de la tarde, así que un flujo particionado por hora tiene horas gordas.
- Rezagados y ejecución especulativa
Aunque el reparto sea perfecto, alguna tarea acabará tardando mucho más que las demás sin que los datos lo justifiquen: la máquina tiene un disco degradado, otro trabajo compite por su CPU, la red de su rack está saturada, la JVM está en una pausa de recolección de basura. Son los rezagados (stragglers). En un trabajo de 1 000 tareas la probabilidad de que alguna caiga en una máquina con problemas es alta, y como el trabajo termina cuando termina la última tarea, un solo rezagado alarga todo el trabajo. En BSP el efecto se multiplica por el número de superpasos.
La solución clásica, introducida por MapReduce y presente en Spark (spark.speculation) y en Hadoop, es la ejecución especulativa: cuando una fase está casi terminada y una tarea lleva mucho más tiempo que la mediana de las demás, el planificador lanza una copia de esa tarea en otro nodo; la primera que termina gana, y la otra se cancela. Es un gasto (se ejecuta trabajo redundante) que compensa porque el coste de una tarea extra es mucho menor que el de todo el clúster esperando. Solo es posible porque las tareas son deterministas e idempotentes, que es el tema del apartado siguiente: si la copia y la original escribieran ambas su resultado, habría que garantizar que se escribe uno solo (salida atómica, 05-02).
- Tolerancia a fallos por reejecución determinista
En una máquina, si el programa falla a la mitad, se relanza desde el principio. Con 1 000 nodos y un trabajo de tres horas, la probabilidad de que algún nodo falle durante el trabajo es prácticamente 1 (falacia 1 de 01-04), y relanzar todo cada vez es inaceptable. La estrategia de los frameworks de este módulo es la reejecución de grano fino: si falla una tarea, se reejecuta esa tarea, no el trabajo. Para que eso sea correcto hacen falta tres propiedades que ya conocemos de 02-05:
- Entrada inmutable. La tarea lee un trozo de datos que no cambia (un bloque de HDFS, un rango de offsets de Kafka, una partición de un RDD). Reejecutar lee lo mismo.
- Cálculo determinista. La misma entrada produce la misma salida. Sin números aleatorios sin semilla, sin depender de la hora actual ni del orden de llegada por la red. Si hace falta aleatoriedad, se deriva de la clave o de una semilla fija.
- Salida idempotente o atómica. O bien escribir dos veces es inocuo (un
PUTde un objeto con el mismo nombre en MinIO, unINSERT ... ON CONFLICT DO NOTHINGcon el id de la tarea), o bien la salida se escribe en un lugar temporal y se publica de golpe al terminar (renombrar un fichero, confirmar una transacción). Así una tarea que murió a medias no deja una salida parcial que confunda a su reejecución ni a la ejecución especulativa.
Con esas tres propiedades, el fallo de un nodo se reduce a "sus tareas vuelven a la cola", que es lo que ya hacía la cola de trabajo del apartado 3.3. Lo que cambia con cada framework es cómo reconstruye la entrada de una tarea cuando esa entrada era el resultado de otra tarea anterior: MapReduce la escribe siempre en disco (HDFS o local), Spark recuerda cómo se calculó y la recomputa (el linaje de 05-03), Flink guarda checkpoints periódicos del estado (05-04). Y lo que no cambia es que todo descansa sobre tareas idempotentes: el consumidor idempotente de 02-05 y la tarea de Spark son la misma idea a distinta escala.
- Escalabilidad: Amdahl y Gustafson
En 01-03 vimos la ley de Amdahl: si una fracción 1 − p del trabajo es secuencial, el speedup con N nodos está acotado por 1 / (1 − p) por muchos nodos que se añadan. En un trabajo distribuido la parte secuencial es la que no se reparte: leer la lista de bloques, planificar tareas, el gather final, escribir el resultado en un solo fichero, y sobre todo el shuffle y la espera a los rezagados, que aunque se ejecuten en paralelo se comportan como una barrera. Con p = 0,95 el techo es 20 aunque se usen 1 000 nodos, lo que parece decir que la computación distribuida a gran escala no compensa.
La respuesta de John Gustafson (1988) es que Amdahl supone un problema de tamaño fijo, y no es así como se usan los clústeres: nadie compra 1 000 nodos para calcular más rápido las ventas de un día, sino para calcular las ventas de un año, o de todos los clics, en el mismo tiempo. La ley de Gustafson mide el speedup escalado: si con N nodos el trabajo tarda un tiempo T del que una fracción s es secuencial y 1 − s es paralela, ese mismo trabajo en un solo nodo habría tardado s + (1 − s) · N veces T, así que:
Con s = 0,05 y N = 1 000, el speedup escalado es 950: casi lineal, porque al crecer el problema la fracción secuencial (planificar, reunir 20 números) se mantiene mientras la paralela (leer terabytes) crece con N. Las dos leyes son ciertas, miden cosas distintas, y juntas dan la regla de diseño de este módulo:
| Amdahl | Gustafson | |
|---|---|---|
| Supone | Tamaño del problema fijo | Tiempo fijo, problema que crece con N |
| Pregunta | ¿Cuánto más rápido acabo lo mismo? | ¿Cuánto más proceso en el mismo tiempo? |
| Fórmula | 1 / ((1 − p) + p / N) |
N − s · (N − 1) |
| Lección para Kilómetro Cero | Un día de ventas no se acelera con 100 nodos: el gather y el arranque dominan | Un año de clics se procesa en la misma noche con 100 nodos que un día con 1 |
| Cómo mejorar | Reducir la parte secuencial: prerreducir, evitar el gather único | Mantener la parte secuencial constante al crecer los datos |
La consecuencia práctica: el paralelismo compensa cuando el trabajo por nodo es grande frente al coste fijo de arrancar y coordinar. Lanzar 1 000 tareas de 100 ms cada una en un clúster con 2 s de latencia de planificación es peor que 10 tareas de 10 s. Lo comprobaremos midiendo en la práctica.
- Práctica: scatter/gather, cola de trabajo y sesgo en
simulaciones/
simulaciones/Toda la práctica se ejecuta en una máquina con multiprocessing, que reparte trabajo entre procesos igual que un clúster lo reparte entre nodos, con la diferencia de que la "red" es memoria local. Es suficiente para ver los patrones, medir speedups y reproducir el sesgo. El fichero de entrada es una versión pequeña de /km0/eventos/2026-09-14/pedidos.jsonl, con la envoltura de eventos de 02-05 y un evento pedido.creado por línea:
{"id_evento":"e-000123-1","tipo":"pedido.creado","version":1,"fecha_ms":1789380000000,"origen":"pedidos","datos":{"pedido_id":"P-2026-000123","cliente":"ana","mercado":"girona","lineas":[{"producto":"tomate-rosa","productor":"huerta-la-vega","cantidad":2,"precio":3.90},{"producto":"queso-curado","productor":"queseria-montblanc","cantidad":1,"precio":12.50}]}}
{"id_evento":"e-000124-1","tipo":"pedido.creado","version":1,"fecha_ms":1789380045000,"origen":"pedidos","datos":{"pedido_id":"P-2026-000124","cliente":"marc","mercado":"lleida","lineas":[{"producto":"vino-crianza","productor":"bodega-roble-alto","cantidad":6,"precio":9.80}]}}
{"id_evento":"e-000125-1","tipo":"pedido.creado","version":1,"fecha_ms":1789380090000,"origen":"pedidos","datos":{"pedido_id":"P-2026-000125","cliente":"lucia","mercado":"valencia","lineas":[{"producto":"queso-fresco","productor":"queseria-montblanc","cantidad":3,"precio":4.20},{"producto":"calabacin","productor":"huerta-la-vega","cantidad":4,"precio":1.60}]}}Estas tres líneas bastan para leer el código; para medir algo hace falta más volumen, así que el primer script incluye un generador que crea cientos de miles de eventos con la misma forma y una distribución sesgada hacia Quesería Montblanc.
10.1 simulaciones/ventas_scatter_gather.py
# km0/simulaciones/ventas_scatter_gather.py
"""Scatter/gather local: reparte pedidos.jsonl entre N trabajadores y suma ventas por productor.
Uso: python ventas_scatter_gather.py generar 400000 # crea eventos/2026-09-14/pedidos.jsonl
python ventas_scatter_gather.py calcular 1 2 4 8 # mide con 1, 2, 4 y 8 trabajadores
"""
import json, os, random, sys, time
from collections import Counter
from multiprocessing import Pool
RUTA = "eventos/2026-09-14/pedidos.jsonl"
PRODUCTOS = [ # (producto, productor, precio, peso en la distribución)
("tomate-rosa", "huerta-la-vega", 3.90, 15),
("calabacin", "huerta-la-vega", 1.60, 10),
("queso-curado", "queseria-montblanc", 12.50, 35), # Semana del Queso Artesano:
("queso-fresco", "queseria-montblanc", 4.20, 15), # Montblanc concentra el 50 %
("vino-crianza", "bodega-roble-alto", 9.80, 25),
]
MERCADOS = ["girona", "lleida", "tarragona", "valencia"]
CLIENTES = ["ana", "marc", "lucia"]
def generar(n_pedidos: int) -> None:
"""Escribe n_pedidos eventos pedido.creado con una distribución sesgada de productores."""
random.seed(42) # determinista: mismo fichero en cada ejecución
os.makedirs(os.path.dirname(RUTA), exist_ok=True)
pesos = [p[3] for p in PRODUCTOS]
with open(RUTA, "w", encoding="utf-8") as f:
for i in range(n_pedidos):
lineas = [
{"producto": prod, "productor": productor, "cantidad": random.randint(1, 6), "precio": precio}
for prod, productor, precio, _ in random.choices(PRODUCTOS, weights=pesos, k=random.randint(1, 3))
]
evento = {
"id_evento": f"e-{i:06d}-1", "tipo": "pedido.creado", "version": 1,
"fecha_ms": 1789344000000 + i * 200, "origen": "pedidos",
"datos": {"pedido_id": f"P-2026-{i:06d}", "cliente": random.choice(CLIENTES),
"mercado": random.choice(MERCADOS), "lineas": lineas},
}
f.write(json.dumps(evento) + "\n")
def trozos_por_bytes(ruta: str, n: int) -> list[tuple[int, int]]:
"""Divide el fichero en n rangos [inicio, fin) alineados a saltos de línea.
Imita lo que hace HDFS con los bloques: cada trabajador recibe un rango de bytes,
no una lista de líneas, para no tener que leer todo el fichero en el coordinador.
"""
tamano = os.path.getsize(ruta)
cortes = [0]
with open(ruta, "rb") as f:
for k in range(1, n):
f.seek(tamano * k // n) # salto aproximado
f.readline() # avanza hasta el final de la línea partida
cortes.append(f.tell())
cortes.append(tamano)
return [(cortes[i], cortes[i + 1]) for i in range(n)]
def mapear_trozo(rango: tuple[int, int]) -> Counter:
"""Fase 'scatter': un trabajador lee su rango y devuelve ventas parciales por productor."""
inicio, fin = rango
parcial = Counter()
with open(RUTA, "rb") as f:
f.seek(inicio)
while f.tell() < fin:
linea = f.readline()
if not linea:
break
ev = json.loads(linea)
if ev["tipo"] != "pedido.creado":
continue
for ln in ev["datos"]["lineas"]:
parcial[ln["productor"]] += ln["cantidad"] * ln["precio"]
return parcial
def reducir(parciales: list[Counter]) -> Counter:
"""Fase 'gather': combinar parciales. Sumar es asociativo y conmutativo, el orden no importa."""
total = Counter()
for p in parciales:
total.update(p)
return total
def calcular(n_trabajadores: int) -> tuple[Counter, float]:
t0 = time.perf_counter()
rangos = trozos_por_bytes(RUTA, n_trabajadores)
if n_trabajadores == 1:
parciales = [mapear_trozo(rangos[0])] # sin Pool: evita el coste de arrancar procesos
else:
with Pool(n_trabajadores) as pool:
parciales = pool.map(mapear_trozo, rangos)
total = reducir(parciales)
return total, time.perf_counter() - t0
if __name__ == "__main__":
if sys.argv[1] == "generar":
generar(int(sys.argv[2]))
print(f"Generado {RUTA}: {os.path.getsize(RUTA) / 1e6:.1f} MB")
else:
base = None
for n in map(int, sys.argv[2:]):
total, seg = calcular(n)
base = base or seg
print(f"{n:2d} trabajadores: {seg:6.2f} s speedup {base / seg:4.2f}x "
f"Montblanc = {total['queseria-montblanc']:,.2f} €")Puntos que conviene entender del código:
- El coordinador no lee los datos.
trozos_por_bytessolo calcula rangos de bytes, con unseeky unreadlinepara alinear cada corte al principio de una línea (sin eso, una línea quedaría partida entre dos trabajadores y ambos la descartarían o la contarían mal). Es exactamente lo que hace Hadoop con los input splits y por qué JSON Lines o CSV son "divisibles" mientras que un JSON con un array gigante no lo es. - Cada trabajador abre el fichero por su cuenta y lee solo su rango. En un clúster, ese trabajador correría en el nodo que tiene el bloque (localidad); aquí todos comparten disco.
- El parcial es pequeño: un
Countercon tres claves, independientemente de que el trozo tenga mil o un millón de líneas. Lo que cruza la "red" (la cola interna dePool) son tres números por trabajador. reducires trivial porque sumar es asociativo. Si el cálculo fuera "importe medio por pedido", el parcial debería ser(suma, cuenta)por productor y la reducción dividir al final.
Una ejecución en un portátil de 8 núcleos con 400 000 pedidos (unos 130 MB):
$ python ventas_scatter_gather.py generar 400000 Generado eventos/2026-09-14/pedidos.jsonl: 131.6 MB $ python ventas_scatter_gather.py calcular 1 2 4 8 16 1 trabajadores: 6.84 s speedup 1.00x Montblanc = 1,893,412.30 € 2 trabajadores: 3.61 s speedup 1.89x Montblanc = 1,893,412.30 € 4 trabajadores: 1.96 s speedup 3.49x Montblanc = 1,893,412.30 € 8 trabajadores: 1.18 s speedup 5.80x Montblanc = 1,893,412.30 € 16 trabajadores: 1.14 s speedup 6.00x Montblanc = 1,893,412.30 €
El resultado es idéntico con cualquier número de trabajadores (la reducción no depende del reparto) y el speedup se aleja del ideal conforme crecen los procesos: con 8 obtenemos 5,8, y con 16 (más procesos que núcleos) nada. Es Amdahl en acción: arrancar el Pool cuesta unos 100 ms fijos, el gather y la alineación de trozos son secuenciales, y el disco es compartido. Con un fichero diez veces mayor la fracción secuencial se diluye y el speedup con 8 se acerca a 7,5: Gustafson.
10.2 simulaciones/cola_trabajo.py
El segundo script convierte el cálculo en tareas idempotentes en una cola, con trabajadores que las toman por demanda y uno que muere a mitad de una tarea. Cada tarea es "calcular las ventas por productor de un mercado y un día", y su salida es un fichero salida/<dia>-<mercado>.json escrito de forma atómica.
# km0/simulaciones/cola_trabajo.py
"""Cola de trabajo con trabajadores que compiten por tareas idempotentes.
Un trabajador muere a propósito a mitad de una tarea; el coordinador detecta la muerte,
devuelve la tarea a la cola y lanza un trabajador de repuesto. El resultado final es el mismo.
"""
import json, os, sys, time
from collections import Counter
from multiprocessing import Process, Queue
RUTA = "eventos/2026-09-14/pedidos.jsonl"
SALIDA = "salida"
MERCADOS = ["girona", "lleida", "tarragona", "valencia"]
def ejecutar_tarea(tarea: dict) -> Counter:
"""Ventas por productor de un mercado. Determinista: misma entrada, mismo resultado."""
parcial = Counter()
with open(RUTA, encoding="utf-8") as f:
for linea in f:
ev = json.loads(linea)
if ev["tipo"] == "pedido.creado" and ev["datos"]["mercado"] == tarea["mercado"]:
for ln in ev["datos"]["lineas"]:
parcial[ln["productor"]] += ln["cantidad"] * ln["precio"]
return parcial
def escribir_atomico(ruta: str, contenido: dict) -> None:
"""Escribe en un temporal y renombra: nadie ve nunca un fichero a medias."""
tmp = f"{ruta}.tmp-{os.getpid()}"
with open(tmp, "w", encoding="utf-8") as f:
json.dump(contenido, f)
os.replace(tmp, ruta) # rename atómico en POSIX
def trabajador(nombre: str, pendientes: Queue, eventos: Queue, morir_en: str | None) -> None:
while True:
tarea = pendientes.get()
if tarea is None: # señal de fin
return
eventos.put(("inicio", nombre, tarea["id"]))
if morir_en == tarea["id"]:
time.sleep(0.2)
os._exit(1) # muerte súbita: sin excepción, sin limpieza
resultado = ejecutar_tarea(tarea)
escribir_atomico(f"{SALIDA}/{tarea['dia']}-{tarea['mercado']}.json", dict(resultado))
eventos.put(("fin", nombre, tarea["id"]))
def coordinador(n_trabajadores: int) -> None:
os.makedirs(SALIDA, exist_ok=True)
pendientes, eventos = Queue(), Queue()
tareas = {f"2026-09-14/{m}": {"id": f"2026-09-14/{m}", "dia": "2026-09-14", "mercado": m} for m in MERCADOS}
for t in tareas.values():
pendientes.put(t)
en_curso: dict[str, str] = {} # trabajador -> id de tarea
terminadas: set[str] = set()
procesos: dict[str, Process] = {}
def lanzar(nombre, morir_en=None):
p = Process(target=trabajador, args=(nombre, pendientes, eventos, morir_en), daemon=True)
p.start()
procesos[nombre] = p
for i in range(n_trabajadores):
lanzar(f"w{i}", morir_en="2026-09-14/lleida" if i == 1 else None) # w1 morirá en lleida
while len(terminadas) < len(tareas):
while not eventos.empty():
tipo, nombre, id_tarea = eventos.get()
if tipo == "inicio":
en_curso[nombre] = id_tarea
print(f"[coord] {nombre} empieza {id_tarea}")
else:
terminadas.add(id_tarea); en_curso.pop(nombre, None)
print(f"[coord] {nombre} termina {id_tarea}")
for nombre, p in list(procesos.items()): # detección de fallos: el proceso ya no vive
if not p.is_alive() and nombre in en_curso:
perdida = en_curso.pop(nombre)
print(f"[coord] {nombre} ha muerto con {perdida}: reencolo y lanzo repuesto")
pendientes.put(tareas[perdida]) # la tarea vuelve a la cola, intacta
del procesos[nombre]
lanzar(nombre + "'")
time.sleep(0.05)
for _ in procesos:
pendientes.put(None)
total = Counter()
for m in MERCADOS:
with open(f"{SALIDA}/2026-09-14-{m}.json", encoding="utf-8") as f:
total.update(json.load(f))
print("Total por productor:", {k: round(v, 2) for k, v in total.items()})
if __name__ == "__main__":
coordinador(int(sys.argv[1]) if len(sys.argv) > 1 else 3)$ python cola_trabajo.py 3
[coord] w0 empieza 2026-09-14/girona
[coord] w1 empieza 2026-09-14/lleida
[coord] w2 empieza 2026-09-14/tarragona
[coord] w1 ha muerto con 2026-09-14/lleida: reencolo y lanzo repuesto
[coord] w0 termina 2026-09-14/girona
[coord] w0 empieza 2026-09-14/valencia
[coord] w1' empieza 2026-09-14/lleida
[coord] w2 termina 2026-09-14/tarragona
[coord] w0 termina 2026-09-14/valencia
[coord] w1' termina 2026-09-14/lleida
Total por productor: {'huerta-la-vega': 712338.1, 'queseria-montblanc': 1893412.3, 'bodega-roble-alto': 984501.6}Lo que enseña la ejecución:
- Equilibrio dinámico:
w0termina Girona y toma Valencia sin que nadie se lo asigne; con mercados de tamaño desigual, los trabajadores rápidos absorben más tareas. - Detección y reejecución: el coordinador no recibe ningún mensaje de error de
w1(murió conos._exit, como una máquina que se apaga); lo detecta porque el proceso ya no está vivo, igual que YARN detecta un NodeManager por heartbeats perdidos. Reencola la tarea tal cual y lanza un repuesto (en otra ejecución puede serw2, si queda libre antes, quien tome Lleida: la cola no asigna, los trabajadores compiten). - Idempotencia:
w1murió después de empezar y quizá con el fichero temporal a medias;w1'ejecuta la misma tarea, escribe su propio temporal y lo renombra. El temporal huérfano dew1queda ensalida/sin que nadie lo lea (un trabajo real lo limpiaría). Si la salida hubiera sido unINSERTen PostgreSQL sin clave única, la reejecución habría duplicado filas: es el mismo problema del consumidor de 02-05. - El resultado es el mismo que el del scatter/gather, lo cual es la definición de que el fallo ha sido tolerado.
10.3 Sesgo: cuando una partición tarda el doble
El tercer experimento reutiliza ventas_scatter_gather.py pero cambia el reparto: en lugar de trozos de bytes, reparte por productor, que es lo que haría un shuffle ingenuo "agrupar por productor". Añade al script esta función y la llamada:
# A nivel de módulo: Pool necesita poder serializar (pickle) la función, y una función anidada no lo es.
def mapear_productor(productor):
t0 = time.perf_counter(); total = 0.0
with open(RUTA, encoding="utf-8") as f:
for linea in f:
ev = json.loads(linea)
for ln in ev["datos"]["lineas"]:
if ln["productor"] == productor:
total += ln["cantidad"] * ln["precio"]
return productor, total, time.perf_counter() - t0
def calcular_por_productor() -> None:
"""Reparto por clave: cada trabajador procesa un productor. Reproduce el sesgo de un shuffle."""
productores = ["huerta-la-vega", "queseria-montblanc", "bodega-roble-alto"]
with Pool(3) as pool:
for productor, total, seg in pool.map(mapear_productor, productores):
print(f"{productor:20s} {total:14,.2f} € {seg:5.2f} s")(Para que el sesgo se vea en el tiempo, y no solo en la cantidad de datos, cada trabajador en un shuffle real solo recibiría sus líneas; en esta simulación cada uno filtra el fichero completo, así que añade en el bucle interior un pequeño trabajo proporcional a las líneas propias, por ejemplo hashlib.md5(linea).hexdigest() solo cuando ln["productor"] == productor.) La salida muestra el problema:
huerta-la-vega 712,338.10 € 1.71 s queseria-montblanc 1,893,412.30 € 3.52 s bodega-roble-alto 984,501.60 € 1.93 s
Tres trabajadores, y el trabajo dura 3,52 s: lo que tarda Quesería Montblanc. Los otros dos están ociosos la mitad del tiempo. La corrección es repartir la clave caliente en subclaves y hacer una segunda reducción:
# A nivel de módulo: Pool necesita poder serializar (pickle) la función, y una función anidada no lo es.
def mapear_subclave(clave):
productor, sub, n_sub = clave; total = 0.0
with open(RUTA, encoding="utf-8") as f:
for linea in f:
ev = json.loads(linea)
# la subclave se deriva del id de pedido: determinista, reparte uniformemente
if productor == "queseria-montblanc" and hash(ev["datos"]["pedido_id"]) % n_sub != sub:
continue
for ln in ev["datos"]["lineas"]:
if ln["productor"] == productor:
total += ln["cantidad"] * ln["precio"]
return productor, total
def calcular_con_salting(n_sub: int = 4) -> None:
"""Rompe la clave caliente en n_sub subclaves y vuelve a combinar. Dos fases de reducción."""
claves = [("huerta-la-vega", 0, 1), ("bodega-roble-alto", 0, 1)] + \
[("queseria-montblanc", k, n_sub) for k in range(n_sub)] # 2 + 4 = 6 tareas
with Pool(6) as pool:
parciales = pool.map(mapear_subclave, claves)
total = Counter()
for productor, importe in parciales: # segunda reducción: 6 números, trivial
total[productor] += importe
print(dict(total))Ahora la tarea más larga procesa un octavo de los datos en vez de la mitad, y el trabajo entero baja a poco más de 1 s con seis procesos. Nótese que hash() de Python está aleatorizado por proceso para cadenas (PYTHONHASHSEED), así que para que la subclave sea determinista entre ejecuciones y reejecuciones hay que fijar la semilla o usar zlib.crc32(pedido_id.encode()) % n_sub; es un ejemplo pequeño de cómo se cuela el no determinismo del apartado 8. Spark hará esto mismo con dos groupBy y una columna de sal en 05-03.
Errores Comunes y Consejos
- Mover los datos al cálculo por costumbre. El primer instinto de quien viene del monolito es "descargo el fichero y lo proceso". Con gigabytes funciona mal y con terabytes no funciona. Pregunta siempre dónde están los datos y si el cálculo puede ir allí.
- Diseñar un resultado parcial que no se combina. Medias, percentiles, distintos (count distinct) y medianas no se suman. Lleva
(suma, cuenta), usa estructuras aproximadas (HyperLogLog para distintos, t-digest para percentiles) o acepta un shuffle completo. - Ignorar el shuffle. Un
groupBysobre una clave de alta cardinalidad es una operación todos-a-todos. Prerreduce antes de mover, elige claves con cardinalidad razonable y mira los bytes de shuffle en la interfaz del framework (05-02 y 05-03 los muestran). - Culpar al clúster del sesgo. Cuando "no escala", mira la duración de las tareas de la fase lenta. Si una tarda cinco veces la mediana, no faltan máquinas: sobra una clave.
- Tareas no idempotentes. Un
INSERTsin clave, un contador incrementado, unappenda un fichero: cualquiera de ellos convierte la reejecución (y la ejecución especulativa) en un duplicado. Escribe a temporal y renombra, usa claves naturales, o escribe por partición completa con sobreescritura (05-05). - No determinismo escondido.
hash()de Python con semilla aleatoria,datetime.now(),randomsin semilla, el orden de undictde una versión antigua, el orden de llegada de mensajes. La reejecución de una tarea debe producir bytes idénticos. - Muchas tareas diminutas. El planificador tiene un coste por tarea (milisegundos a segundos). 100 000 tareas de 50 ms son peores que 1 000 de 5 s. Como orientación, entre 2 y 4 tareas por núcleo y fase, con duraciones de segundos a minutos.
- Extrapolar Amdahl a Gustafson (o al revés). Si el problema es fijo, los nodos extra no ayudan; si el problema crece, sí. Antes de pedir más máquinas, decide cuál de los dos casos es el tuyo.
Ejercicios
Ejercicio 1: ¿Qué se combina y qué no?
Para cada uno de estos cálculos sobre los eventos pedido.creado de la Semana de la Vendimia, indica (a) qué devuelve cada trabajador como resultado parcial y (b) cómo se combinan dos parciales, o por qué no es posible combinarlos y qué alternativa hay:
- Importe total vendido por Bodega Roble Alto.
- Importe medio por pedido en cada mercado.
- Número de clientes distintos que compraron
vino-crianza. - Los 10 productos más vendidos por cantidad.
- La mediana del importe de los pedidos.
Ejercicio 2: Reejecución y salida atómica
En cola_trabajo.py, sustituye escribir_atomico por una escritura directa (open(ruta, "w") y json.dump) y haz que el trabajador muera durante la escritura (por ejemplo, escribe la mitad del JSON, haz f.flush() y os._exit(1)). Describe qué ocurre en el coordinador al final, cómo lo detectarías en producción y por qué el os.replace lo evita. Después, propón cómo escribirías la salida si en lugar de un fichero fuera una tabla PostgreSQL ventas_mercado_dia(dia, mercado, productor, importe) de modo que la reejecución siguiera siendo segura.
Ejercicio 3: Amdahl, Gustafson y el tamaño de la tarea
Con los tiempos de la ejecución del apartado 10.1 (1 trabajador: 6,84 s; 8 trabajadores: 1,18 s), estima la fracción secuencial del script según Amdahl. Con esa fracción, ¿qué speedup obtendrías con 64 trabajadores sobre el mismo fichero? ¿Y cuál sería el speedup escalado de Gustafson con 64 trabajadores si el fichero creciera 64 veces? Por último: el planificador de un clúster real añade 1,5 s por tarea entre asignación y arranque; si el fichero de 130 MB se divide en 1 000 trozos, ¿cuánto tarda el trabajo con 8 nodos, y con cuántos trozos deberías dividirlo?
Soluciones
Ejercicio 1.
- Parcial: un número (suma de
cantidad × preciode las líneas de Bodega Roble Alto). Combinación: suma. Asociativo y conmutativo; el caso ideal. - Parcial: por mercado, la pareja
(suma_importes, numero_pedidos). Combinación: sumar componente a componente; la media se calcula solo al final,suma / cuenta. Combinar medias directamente sería incorrecto salvo que todos los trozos tuvieran el mismo número de pedidos. - Parcial: el conjunto de ids de cliente que compraron
vino-crianzaen ese trozo. Combinación: unión de conjuntos, y al final el tamaño. Es combinable pero el parcial puede ser grande (cientos de miles de ids); si eso es un problema, un HyperLogLog por trozo (unos KB) se combina con una unión y da el cardinal con un error del 1–2 %. - Parcial: cantidades por producto (un
Counter), no "los 10 mejores del trozo": el undécimo de un trozo puede ser el primero global. Combinación: sumar losCountery elegir los 10 al final. Si el número de productos fuera enorme, se podría guardar el top-K por trozo con K grande como aproximación, aceptando error. - No es combinable: la mediana de medianas no es la mediana. Alternativas: ordenar globalmente (un shuffle por rangos de importe: caro pero exacto), o un t-digest/percentil aproximado por trozo, que sí se combina y da la mediana con error acotado.
Ejercicio 2.
Con escritura directa, w1 deja salida/2026-09-14-lleida.json con medio JSON. El coordinador reencola y w1' vuelve a abrir el fichero en modo "w", que lo trunca, así que en este caso concreto el resultado final es correcto; pero entre la muerte y la reejecución (segundos, o minutos en un clúster) el fichero existe y está corrupto: cualquier lector (el pipeline de 05-05, un hdfs dfs -cat) fallaría con un JSONDecodeError, o peor, leería un parcial si el formato fuera CSV. Si el trabajador muriera después de escribir pero antes de enviar fin, la reejecución también sobrescribiría, sin daño. El problema real aparece si la reejecución no trunca (modo "a") o si la salida es un sistema sin truncado. En producción se detecta por lectores que fallan o por ficheros con tamaño inesperado; os.replace lo evita porque el fichero con el nombre definitivo aparece de golpe y completo o no aparece: es la salida atómica del apartado 8 y la que MapReduce implementa con directorios _temporary (05-02).
Para PostgreSQL: una clave primaria (dia, mercado, productor) y escritura con INSERT ... ON CONFLICT (dia, mercado, productor) DO UPDATE SET importe = EXCLUDED.importe, todo dentro de una transacción que empieza con DELETE FROM ventas_mercado_dia WHERE dia = %s AND mercado = %s y termina con COMMIT: la tarea reemplaza su partición entera de forma atómica, y ejecutarla dos veces deja exactamente las mismas filas. Es el "reprocesar un día sin duplicar" de 05-05.
Ejercicio 3.
Amdahl: S(8) = 1 / ((1 − p) + p / 8) = 6,84 / 1,18 = 5,80. Despejando, (1 − p) + p / 8 = 1 / 5,80 = 0,1724 → 1 − 0,875 p = 0,1724 → p = 0,946. Fracción secuencial 1 − p ≈ 5,4 %. Con 64 trabajadores: S(64) = 1 / (0,054 + 0,946 / 64) = 1 / 0,0688 = 14,5: en la práctica, menos, porque el portátil tiene 8 núcleos y el disco es uno. Gustafson con s = 0,054 y N = 64: 64 − 0,054 × 63 = 60,6: procesar 64 veces más datos con 64 nodos en casi el mismo tiempo.
Con 1 000 trozos y 8 nodos, cada trozo son 130 KB (unos 7 ms de proceso) más 1,5 s de coste de planificación: 1 000 tareas × 1,507 s / 8 nodos ≈ 188 s. Peor que los 6,84 s de un solo proceso. Con 8 trozos: 8 × (0,86 s + 1,5 s) / 8 ≈ 2,4 s. Con 16 o 24 trozos (2–3 por nodo) el equilibrio dinámico ayuda a los rezagados sin disparar el coste fijo: unos 2,5–3 s. La regla: el trabajo por tarea debe ser al menos un orden de magnitud mayor que el coste de planificarla.
Conclusión
Distribuir un cálculo es más que ejecutarlo en varias máquinas: es llevar el código a donde están los datos, descomponer el trabajo en tareas cuyos resultados parciales se puedan combinar, y aceptar que la combinación (el shuffle) es la parte cara. Hemos separado el paralelismo de datos, que es el de este módulo, del paralelismo de tareas, que es el de los pipelines; y hemos recorrido los patrones con los que se organiza el reparto: scatter/gather para lo asociativo, divide y vencerás para combinar en árbol, la cola de trabajo para equilibrar y tolerar fallos por demanda, BSP con sus superpasos y barreras para lo iterativo y los grafos, el DAG de operadores del dataflow que Spark y Flink implementan, y los actores para el estado por entidad. Los tres modos (lotes, flujos, interactivo) comparten patrones y problemas: el sesgo que hace que Quesería Montblanc alargue todo el trabajo, los rezagados que la ejecución especulativa esquiva, y la tolerancia a fallos que solo funciona si las tareas son deterministas e idempotentes con salida atómica, la misma regla que gobernaba a los consumidores de 02-05. Amdahl nos recordó que un problema fijo tiene un techo, y Gustafson que los clústeres existen para problemas que crecen. En simulaciones/ hemos medido un speedup de 5,8 con 8 procesos, visto morir y renacer a un trabajador sin alterar el resultado, y roto la clave caliente de Quesería Montblanc en subclaves.
Todo esto lo hicimos a mano, con multiprocessing y un coordinador de cuarenta líneas que reparte, detecta muertes y reencola. Un framework de computación distribuida es exactamente eso, pero para miles de nodos, con localidad de datos, shuffle distribuido, ejecución especulativa y salida atómica resueltos de una vez para todos los trabajos. El primero que lo consiguió, y el que fijó el vocabulario que seguimos usando, fue MapReduce, y con él Hadoop, que es la siguiente lección: cómo el mismo cálculo de ventas por productor se expresa como map, shuffle & sort y reduce, y cómo YARN reparte las tareas por el clúster.
Curso de Arquitecturas Distribuidas
Módulo 1: Introducción a los Sistemas Distribuidos
- Conceptos Básicos de Sistemas Distribuidos
- Modelos de Sistemas Distribuidos
- Ventajas y Desafíos de los Sistemas Distribuidos
- Las Falacias de la Computación Distribuida
- Tiempo, Relojes y Ordenación de Eventos
- Del Monolito a la Plataforma Distribuida: el Caso Kilómetro Cero
Módulo 2: Comunicación en Sistemas Distribuidos
- Protocolos de Comunicación
- RPC y RMI
- gRPC y Serialización de Datos
- Mensajería y Colas de Mensajes
- Patrones de Comunicación Asíncrona
Módulo 3: Consistencia y Replicación
- Modelos de Consistencia
- El Teorema CAP y PACELC
- Algoritmos de Consenso
- Replicación de Datos
- Transacciones Distribuidas y Sagas
Módulo 4: Almacenamiento Distribuido
- Particionado de Datos y Hashing Consistente
- Sistemas de Archivos Distribuidos
- Almacenamiento de Objetos
- Bases de Datos Distribuidas
- Cachés Distribuidos
Módulo 5: Computación Distribuida
- Modelos de Computación Distribuida
- MapReduce y Hadoop
- Spark y Computación en Memoria
- Procesamiento de Flujos de Datos
- Planificación de Trabajos y Pipelines de Datos
Módulo 6: Seguridad en Sistemas Distribuidos
- Autenticación y Autorización
- Cifrado y Protección de Datos
- Gestión de Identidades
- Seguridad entre Servicios: mTLS y Gestión de Secretos
- Puertas de Enlace, Limitación de Tasa y Auditoría
Módulo 7: Monitoreo y Mantenimiento
- Monitoreo de Sistemas Distribuidos
- Logs Centralizados y Trazabilidad Distribuida
- Gestión de Fallos y Recuperación
- Patrones de Resiliencia: Timeouts, Reintentos y Circuit Breaker
- Automatización y Orquestación
- Pruebas en Sistemas Distribuidos e Ingeniería del Caos
