alpinashop_analitica ya existe y tiene dentro el histórico de pedidos. Pero ese histórico llegó allí de una forma que no se puede repetir todas las noches: Lucía ejecutó a mano una consulta federada contra la réplica de PostgreSQL. Funcionó una vez. La pregunta es qué pasa mañana, y pasado, y el día que alguien pida ver las ventas de la campaña hoy y no mañana.
La respuesta ingenua es escribir un script. Un fichero sincronizar.py que se conecte a Cloud SQL, lea los pedidos del día, los transforme y los inserte en BigQuery, lanzado por un cron en una máquina virtual. Ese script funcionará durante semanas, y luego fallará. Fallará porque los datos crecen y el script tarda cinco horas; porque la VM se reinicia a mitad de proceso y nadie sabe si escribió la mitad de las filas; porque un pedido con un carácter raro lanza una excepción y se pierde el resto del volcado; porque nadie sabe si se ejecutó ayer; y porque el día que alguien pida datos en tiempo real, todo el diseño hay que tirarlo.
Cloud Dataflow es el servicio gestionado de Google para ejecutar pipelines de datos escritos con Apache Beam. Su promesa concreta es doble: escribes la lógica de transformación una sola vez y Google se encarga del paralelismo, del reintento, del escalado y de la coherencia; y ese mismo código sirve tanto para procesar dos años de histórico como para procesar los pedidos según entran.
En esta lección vas a entender el modelo de Beam, escribirás un pipeline por lotes que carga el histórico de AlpinaShop, lo probarás en local, lo lanzarás en Dataflow, y luego entrarás en la parte que de verdad separa el procesamiento de datos serio del amateur: el tiempo. Qué significa "las ventas de las 10:00" cuando un evento generado a las 10:00 llega a las 10:07, y cómo se responde a eso sin mentir.
Contenido
- Qué es un pipeline de datos y por qué un script no basta
- Apache Beam: el modelo unificado
- Los cuatro conceptos:
Pipeline,PCollection,PTransform, runner - Las transformaciones que usarás el 90 % del tiempo
- Primer pipeline por lotes, explicado línea a línea
- Ejecución local con
DirectRunner - Ejecución gestionada con
DataflowRunner - El tiempo en streaming: evento frente a proceso
- Marcas de agua, ventanas, disparadores y datos tardíos
- El pipeline de streaming de
pedidos-nuevos - Plantillas de Dataflow: la vía práctica
- Escalado automático y Dataflow Prime
- Monitorización, paralelismo y sesgo de datos
- Coste: qué se paga exactamente
- Cuándo Dataflow no es la respuesta
- Qué es un pipeline de datos y por qué un script no basta
Un pipeline de datos es una secuencia declarada de operaciones que llevan datos de un origen a un destino transformándolos por el camino. El adjetivo importante es declarada: describes qué quieres que ocurra, no cómo se reparte el trabajo entre máquinas.
Comparemos con honestidad el script en una VM y el pipeline gestionado:
| Aspecto | Script en una VM | Pipeline en Dataflow |
|---|---|---|
| Paralelismo | El que programes tú, a mano, con hilos | Automático: reparte por trabajadores |
| Escalado | Cambiar la VM y reiniciar | Horizontal y automático durante la ejecución |
| Fallo de una máquina | Se pierde todo el proceso | Se reintenta el fragmento afectado |
| Fallo de un registro | Excepción que tumba el proceso | Se desvía a una salida de errores y sigue |
| Semántica de escritura | La que consigas | Exactamente una vez en las conexiones nativas |
| Estado si se reinicia | Desconocido | Gestionado por el servicio |
| Lote y streaming | Dos programas distintos | El mismo código |
| Coste en reposo | La VM encendida siempre | Cero: no hay nada encendido |
| Observabilidad | Los print que hayas puesto |
Grafo, métricas y logs integrados |
La fila que más duele en la práctica es la del fallo de un registro. Un volcado de 800.000 pedidos en el que el pedido 412.337 tiene un importe con coma en vez de punto no debe perder los 387.663 pedidos restantes. Un script mal hecho los pierde; uno bien hecho requiere un esfuerzo considerable de tratamiento de errores que en Beam viene de serie con las salidas etiquetadas.
Y la fila del coste en reposo tiene su matiz: Dataflow en modo lote cuesta cero cuando no se ejecuta, pero un pipeline de streaming está permanentemente encendido y factura sin parar. Volveremos a ello en el apartado 14, porque es la sorpresa más habitual.
- Apache Beam: el modelo unificado
Apache Beam es un modelo de programación open source —donado por Google a la Apache Software Foundation— para definir pipelines de datos que luego se ejecutan en distintos motores.
La palabra clave es unificado. Antes de Beam, procesar por lotes y procesar en streaming eran mundos separados con herramientas separadas, y las empresas mantenían dos implementaciones de la misma lógica de negocio: una para el histórico y otra para el tiempo real, que inevitablemente divergían y daban números distintos. Beam parte de la idea de que un lote es simplemente un flujo acotado y un stream es un flujo no acotado, y que la lógica de transformación es la misma en ambos casos.
flowchart LR
subgraph SDK["SDK de Apache Beam"]
P["Codigo del pipeline<br/>Python / Java / Go"]
end
subgraph Runners["Runners"]
D["DirectRunner<br/>local, pruebas"]
DF["DataflowRunner<br/>Google Cloud"]
FL["FlinkRunner"]
SP["SparkRunner"]
end
P --> D
P --> DF
P --> FL
P --> SP
El runner es el motor que ejecuta el pipeline. El mismo fichero Python puede correr en tu portátil con DirectRunner, en Dataflow con DataflowRunner, o en un clúster de Flink o Spark. Eso reduce el bloqueo con el proveedor: si AlpinaShop tuviera que salir de Google Cloud, la lógica de sus pipelines viajaría.
Beam tiene SDK para Java, Python y Go. Usaremos Python, coherente con el resto de la aplicación de AlpinaShop, que ya es Flask.
# Instalacion del SDK con las dependencias de Google Cloud
python -m venv venv-beam
source venv-beam/bin/activate
pip install 'apache-beam[gcp]==2.64.0'Fijar la versión es deliberado: los trabajadores de Dataflow usarán exactamente la versión del SDK con la que lanzas el pipeline, y las diferencias entre versiones son una fuente clásica de fallos que solo aparecen en la nube.
- Los cuatro conceptos:
Pipeline, PCollection, PTransform, runner
Pipeline, PCollection, PTransform, runnerPipeline es el objeto que contiene el grafo completo. Se construye, se declara y se ejecuta. Nada ocurre mientras lo escribes: estás dibujando un plano.
PCollection es un conjunto de datos distribuido e inmutable. No es una lista de Python: puede tener cero elementos o billones, puede estar repartida por cien máquinas, y no se puede modificar. Cada transformación produce una PCollection nueva. La inmutabilidad es lo que permite reintentar un fragmento fallido sin corromper nada.
Una PCollection puede ser:
- Acotada (bounded): tiene un final conocido. Un fichero, una tabla. Es el caso de lote.
- No acotada (unbounded): no termina nunca. Un topic de Pub/Sub. Es el caso de streaming.
PTransform es una operación que toma una o más PCollection y produce una o más PCollection. Se aplica con el operador |, que en Beam está sobrecargado para significar "aplica esta transformación".
Runner es el motor de ejecución, ya visto.
La sintaxis, que resulta extraña la primera vez:
El operador >> asocia un nombre al paso. No es decorativo: ese nombre es el que aparece en el grafo de la consola de Dataflow y en las métricas. Un pipeline con pasos llamados Map(<lambda at main.py:34>) es imposible de depurar en producción. Nombra todos los pasos, siempre.
flowchart TD
A["PCollection: lineas de texto crudas<br/>gs://alpinashop-catalogo/exportaciones/..."]
B["PTransform: ParseCSV<br/>ParDo"]
C["PCollection: diccionarios de pedido"]
D["PTransform: ValidarYLimpiar<br/>ParDo con salidas multiples"]
E["PCollection: pedidos validos"]
F["PCollection: pedidos rechazados"]
G["PTransform: WriteToBigQuery"]
H["PTransform: WriteToText<br/>cuarentena en el bucket"]
A --> B --> C --> D
D --> E --> G
D --> F --> H
Ese grafo es exactamente el que vas a escribir en el apartado 5. Fíjate en que la rama de errores es parte del diseño, no un añadido.
- Las transformaciones que usarás el 90 % del tiempo
| Transformación | Qué hace | Ejemplo en AlpinaShop |
|---|---|---|
beam.Map(f) |
Aplica f a cada elemento; devuelve uno |
Convertir una línea CSV en diccionario |
beam.FlatMap(f) |
Aplica f; devuelve 0, 1 o N elementos |
Desglosar un pedido en sus líneas |
beam.Filter(f) |
Se queda con los que cumplen f |
Descartar pedidos cancelados |
beam.ParDo(DoFn) |
La forma general: clase con estado, ciclo de vida y salidas múltiples | Validar y separar válidos de rechazados |
beam.GroupByKey() |
Agrupa pares (clave, valor) por clave |
Agrupar líneas por sku |
beam.CombinePerKey(f) |
Agrupa y reduce con una función asociativa | Sumar ventas por sku |
beam.CoGroupByKey() |
Une varias PCollection por clave (el JOIN de Beam) |
Cruzar pedidos con productos |
beam.Keys() / beam.Values() |
Extrae claves o valores | |
beam.Distinct() |
Elimina duplicados | Sesiones únicas |
beam.io.ReadFromText / WriteToText |
Lectura/escritura de ficheros, incluido Cloud Storage | |
beam.io.ReadFromPubSub |
Lee un topic o suscripción | El pipeline de pedidos-nuevos |
beam.io.WriteToBigQuery |
Escribe en una tabla | Destino de todo |
Dos aclaraciones que evitan errores conceptuales importantes:
GroupByKey frente a CombinePerKey. GroupByKey lleva todos los valores de una clave a una sola máquina y los materializa en memoria. Si un sku tiene tres millones de líneas, esa máquina puede reventar. CombinePerKey con una función asociativa y conmutativa (sum, max, min) hace agregación parcial en cada trabajador antes de mover nada por la red: cada worker suma lo suyo y solo viajan los resultados parciales. Siempre que puedas usar CombinePerKey, úsalo; es órdenes de magnitud más eficiente y no se rompe con claves calientes.
Map frente a ParDo. Map es azúcar sintáctico sobre ParDo. Usa Map para transformaciones simples y sin estado; usa ParDo con una clase DoFn cuando necesites inicialización costosa (abrir un cliente de API una vez por trabajador, no una vez por elemento), métricas propias, o varias salidas.
- Primer pipeline por lotes, explicado línea a línea
Objetivo concreto: leer los ficheros CSV de exportación de pedidos que hay en gs://alpinashop-catalogo/exportaciones/2026/03/14/, validarlos, limpiarlos, calcular un agregado de ventas por SKU y escribir dos cosas en alpinashop_analitica: las líneas limpias en lineas_pedido y el agregado en una tabla nueva ventas_diarias_sku. Los registros defectuosos van a un fichero de cuarentena en el bucket.
"""
Pipeline por lotes de AlpinaShop: exportaciones CSV -> BigQuery.
Ejecucion local:
python pipeline_pedidos.py --fecha 2026-03-14
Ejecucion en Dataflow:
python pipeline_pedidos.py --fecha 2026-03-14 --runner DataflowRunner ...
"""
import argparse
import csv
import io
import logging
from datetime import datetime
from decimal import Decimal, InvalidOperation
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions
PROYECTO = "alpinashop-datos"
DATASET = "alpinashop_analitica"
BUCKET = "alpinashop-catalogo"Todo lo de arriba es Python normal. Las constantes en mayúsculas van al principio para que cambiar de proyecto no obligue a buscar cadenas por el fichero.
class ParsearLineaCSV(beam.DoFn):
"""Convierte una linea de texto CSV en un diccionario.
Se implementa como DoFn y no como Map porque necesitamos
contadores propios y porque queremos dos salidas: validos y rechazados.
"""
SALIDA_RECHAZOS = "rechazos"
def __init__(self):
# Los contadores de Beam se agregan entre todos los trabajadores
# y se ven en la consola de Dataflow. Son la forma correcta
# de instrumentar un pipeline; los print no sirven de nada.
self.contador_ok = beam.metrics.Metrics.counter("parseo", "filas_ok")
self.contador_ko = beam.metrics.Metrics.counter("parseo", "filas_rechazadas")
def process(self, linea):
try:
campos = next(csv.reader(io.StringIO(linea), delimiter=","))
except Exception as exc:
self.contador_ko.inc()
yield beam.pvalue.TaggedOutput(
self.SALIDA_RECHAZOS,
{"linea": linea, "motivo": f"csv_ilegible: {exc}"},
)
return
if len(campos) != 8:
self.contador_ko.inc()
yield beam.pvalue.TaggedOutput(
self.SALIDA_RECHAZOS,
{"linea": linea, "motivo": f"esperadas 8 columnas, hay {len(campos)}"},
)
return
yield {
"pedido_id": campos[0].strip(),
"linea_num": campos[1].strip(),
"fecha_pedido": campos[2].strip(),
"sku": campos[3].strip().upper(),
"cantidad": campos[4].strip(),
"precio_unitario": campos[5].strip().replace(",", "."),
"descuento_linea": campos[6].strip().replace(",", "."),
"estado": campos[7].strip().lower(),
}Puntos importantes de esta clase:
yielden lugar dereturn. UnDoFnes un generador: puede emitir cero, uno o muchos elementos por entrada. Unreturncon valor no funciona como esperas.TaggedOutputmarca el elemento para que salga por una rama distinta del grafo. Es el mecanismo de Beam para "esto ha ido mal pero el pipeline sigue". Sin él, una excepción reintenta el paquete cuatro veces y luego mata el trabajo entero.replace(",", ".")normaliza los decimales del ERP español, que exporta89,90. Es exactamente el tipo de suciedad real que hay en cualquier exportación.- Los contadores (
Metrics.counter) suben a la consola de Dataflow y permiten responder "¿cuántas filas se rechazaron anoche?" sin abrir un log.
class ValidarPedido(beam.DoFn):
"""Convierte tipos y aplica reglas de negocio."""
SALIDA_RECHAZOS = "rechazos"
def process(self, fila):
motivos = []
# 1) Fecha
try:
fecha = datetime.strptime(fila["fecha_pedido"], "%Y-%m-%d").date()
except ValueError:
motivos.append("fecha invalida")
fecha = None
# 2) Cantidad: entero estrictamente positivo
try:
cantidad = int(fila["cantidad"])
if cantidad <= 0:
motivos.append("cantidad no positiva")
except ValueError:
motivos.append("cantidad no numerica")
cantidad = None
# 3) Importes en Decimal, NUNCA en float (ver 04-01)
try:
precio = Decimal(fila["precio_unitario"])
descuento = Decimal(fila["descuento_linea"] or "0")
if precio < 0 or descuento < 0:
motivos.append("importe negativo")
except InvalidOperation:
motivos.append("importe no numerico")
precio = descuento = None
# 4) SKU con el formato del catalogo de AlpinaShop
if not fila["sku"] or len(fila["sku"]) < 4:
motivos.append("sku ausente o demasiado corto")
if motivos:
yield beam.pvalue.TaggedOutput(
self.SALIDA_RECHAZOS,
{"linea": str(fila), "motivo": "; ".join(motivos)},
)
return
importe = (precio * cantidad) - descuento
yield {
"pedido_id": fila["pedido_id"],
"linea_num": int(fila["linea_num"]),
"fecha_pedido": fecha.isoformat(),
"sku": fila["sku"],
"cantidad": cantidad,
"precio_unitario": str(precio), # BigQuery acepta NUMERIC como cadena
"descuento_linea": str(descuento),
"importe_linea": str(importe),
}Detalle que cuesta caro descubrir en producción: los valores NUMERIC se pasan a WriteToBigQuery como cadena, no como float. Si conviertes a float para serializar, has reintroducido el error de coma flotante que evitamos con tanto cuidado en 04-01.
Ahora el pipeline completo:
def construir_pipeline(pipeline, fecha):
ruta_entrada = f"gs://{BUCKET}/exportaciones/{fecha.replace('-', '/')}/pedidos-*.csv"
ruta_cuarentena = f"gs://{BUCKET}/cuarentena/{fecha}/rechazos"
# 1) LEER: cada linea del CSV es un elemento de la PCollection
crudo = (
pipeline
| "LeerCSV" >> beam.io.ReadFromText(ruta_entrada, skip_header_lines=1)
)
# 2) PARSEAR con dos salidas
parseado = (
crudo
| "ParsearCSV" >> beam.ParDo(ParsearLineaCSV()).with_outputs(
ParsearLineaCSV.SALIDA_RECHAZOS, main="validos"
)
)
# 3) VALIDAR, tambien con dos salidas
validado = (
parseado.validos
| "ValidarPedido" >> beam.ParDo(ValidarPedido()).with_outputs(
ValidarPedido.SALIDA_RECHAZOS, main="limpios"
)
)
lineas_limpias = validado.limpios
# 4) ESCRIBIR el detalle en BigQuery
(
lineas_limpias
| "EscribirLineas" >> beam.io.WriteToBigQuery(
table=f"{PROYECTO}:{DATASET}.lineas_pedido",
schema="pedido_id:STRING,linea_num:INTEGER,fecha_pedido:DATE,"
"sku:STRING,cantidad:INTEGER,precio_unitario:NUMERIC,"
"descuento_linea:NUMERIC,importe_linea:NUMERIC",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
additional_bq_parameters={
"timePartitioning": {"type": "DAY", "field": "fecha_pedido"},
"clustering": {"fields": ["sku"]},
},
)
)
# 5) AGREGAR ventas por SKU con CombinePerKey (no GroupByKey)
(
lineas_limpias
| "ClaveSKU" >> beam.Map(
lambda f: ((f["fecha_pedido"], f["sku"]), Decimal(f["importe_linea"]))
)
| "SumarPorSKU" >> beam.CombinePerKey(sum)
| "FormatearAgregado" >> beam.Map(
lambda kv: {
"dia": kv[0][0],
"sku": kv[0][1],
"ventas_eur": str(kv[1]),
}
)
| "EscribirAgregado" >> beam.io.WriteToBigQuery(
table=f"{PROYECTO}:{DATASET}.ventas_diarias_sku",
schema="dia:DATE,sku:STRING,ventas_eur:NUMERIC",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE,
)
)
# 6) CUARENTENA: los rechazos de ambos pasos, juntos, a Cloud Storage
(
(parseado[ParsearLineaCSV.SALIDA_RECHAZOS],
validado[ValidarPedido.SALIDA_RECHAZOS])
| "UnirRechazos" >> beam.Flatten()
| "SerializarRechazos" >> beam.Map(
lambda r: f'{r["motivo"]}\t{r["linea"]}'
)
| "EscribirCuarentena" >> beam.io.WriteToText(
ruta_cuarentena, file_name_suffix=".tsv"
)
)
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--fecha", required=True, help="Fecha a procesar, yyyy-MM-dd")
conocidos, resto = parser.parse_known_args()
opciones = PipelineOptions(resto)
# Necesario para que los trabajadores instalen las dependencias del fichero
opciones.view_as(SetupOptions).save_main_session = True
with beam.Pipeline(options=opciones) as p:
construir_pipeline(p, conocidos.fecha)
if __name__ == "__main__":
logging.getLogger().setLevel(logging.INFO)
main()Cinco cosas que merecen comentario:
with beam.Pipeline(...) as p: al salir del bloquewith, Beam llama arun()y espera. Sin elwith, hay que llamar ap.run().wait_until_finish()explícitamente.WRITE_APPENDfrente aWRITE_TRUNCATE: el detalle se añade (queremos histórico acumulado); el agregado se reescribe entero (queremos la foto vigente). Elegir mal aquí duplica datos silenciosamente.additional_bq_parameterscrea la tabla ya particionada y clusterizada si no existía. Es la forma de no perder lo aprendido en 04-01 cuando la tabla la crea el pipeline.beam.Flatten()une variasPCollectiondel mismo tipo en una. Es la unión de ramas del grafo, el equivalente a unUNION ALL.save_main_session=Trueserializa el ámbito global del módulo para los trabajadores. Sin esto, un pipeline que funciona en local falla en Dataflow conNameErrorsobre las constantes o los imports. Es el error de novato número uno.
- Ejecución local con
DirectRunner
DirectRunnerAntes de gastar un céntimo, se prueba en local. El DirectRunner ejecuta el pipeline en tu máquina, con un subconjunto de datos.
# Datos de prueba locales
mkdir -p ./pruebas && cat > ./pruebas/pedidos-test.csv <<'EOF'
pedido_id,linea_num,fecha_pedido,sku,cantidad,precio_unitario,descuento_linea,estado
PED-2026-0042,1,2026-03-14,MOCH-40L-AZ,1,89,90,0,confirmado
PED-2026-0042,2,2026-03-14,FRON-300L,2,34.50,5.00,confirmado
PED-2026-0043,1,2026-03-14,CRAM-12P,-1,120.00,0,confirmado
PED-2026-0044,1,fecha-mala,TIEN-2P,1,240.00,0,confirmado
EOF
python pipeline_pedidos.py \
--fecha 2026-03-14 \
--runner DirectRunnerEse fichero de prueba está sucio a propósito, y así es como debe ser cualquier juego de pruebas de un pipeline:
- La primera línea tiene nueve campos porque
89,90mete una coma de más. El parser la rechazará con "esperadas 8 columnas". Es exactamente el fallo real de un ERP mal configurado. - La tercera tiene cantidad negativa: la rechaza el validador.
- La cuarta tiene una fecha inválida: la rechaza el validador.
- Solo la segunda pasa.
El DirectRunner es deliberadamente estricto: comprueba la inmutabilidad de los elementos, serializa y deserializa entre pasos, y desordena los datos a propósito. Si tu pipeline depende del orden de llegada o modifica un objeto in situ, el DirectRunner lo detecta y falla, mientras que en la nube produciría resultados incorrectos de forma intermitente. Que sea lento y quisquilloso es la funcionalidad, no un defecto.
Buena práctica: además de probar el pipeline entero, prueba las transformaciones con unittest y las utilidades de Beam:
import unittest
import apache_beam as beam
from apache_beam.testing.test_pipeline import TestPipeline
from apache_beam.testing.util import assert_that, equal_to
class TestValidarPedido(unittest.TestCase):
def test_rechaza_cantidad_negativa(self):
entrada = [{
"pedido_id": "PED-1", "linea_num": "1", "fecha_pedido": "2026-03-14",
"sku": "CRAM-12P", "cantidad": "-1",
"precio_unitario": "120.00", "descuento_linea": "0", "estado": "confirmado",
}]
with TestPipeline() as p:
salidas = (
p | beam.Create(entrada)
| beam.ParDo(ValidarPedido()).with_outputs(
ValidarPedido.SALIDA_RECHAZOS, main="limpios")
)
assert_that(salidas.limpios, equal_to([]), label="sin validos")beam.Create(...) fabrica una PCollection a partir de una lista de Python: es la forma de inyectar datos de prueba. assert_that con equal_to compara el contenido sin importar el orden, que es lo correcto en un sistema distribuido.
- Ejecución gestionada con
DataflowRunner
DataflowRunnerCon las pruebas en verde, al servicio real.
# Bucket propio para los artefactos de Dataflow (no mezclar con el catalogo)
gcloud storage buckets create gs://alpinashop-dataflow \
--project=alpinashop-datos --location=europe-west1 \
--uniform-bucket-level-access
# Cuenta de servicio dedicada al pipeline, con minimo privilegio (03-04)
gcloud iam service-accounts create sa-dataflow-pedidos \
--project=alpinashop-datos \
--display-name="Pipelines de Dataflow de pedidos"
SA="[email protected]"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/dataflow.worker"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/bigquery.dataEditor"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/bigquery.jobUser"
gcloud storage buckets add-iam-policy-binding gs://alpinashop-dataflow \
--member="serviceAccount:$SA" --role="roles/storage.objectAdmin"
gcloud storage buckets add-iam-policy-binding gs://alpinashop-catalogo \
--member="serviceAccount:$SA" --role="roles/storage.objectViewer"Y el lanzamiento:
python pipeline_pedidos.py \
--fecha 2026-03-14 \
--runner DataflowRunner \
--project alpinashop-datos \
--region europe-west1 \
--temp_location gs://alpinashop-dataflow/temp \
--staging_location gs://alpinashop-dataflow/staging \
--service_account_email "$SA" \
--subnetwork regions/europe-west1/subnetworks/sn-datos-euw1 \
--no_use_public_ips \
--machine_type n2-standard-2 \
--max_num_workers 10 \
--job_name alpinashop-pedidos-20260314 \
--labels entorno=produccion,equipo=datos,centro-coste=analiticaOpción por opción, porque cada una tiene consecuencias:
| Opción | Qué hace y por qué |
|---|---|
--region europe-west1 |
Dónde se ejecutan los trabajadores. Debe coincidir con la ubicación del dataset y con la del bucket, o pagarás salida entre regiones y añadirás latencia |
--temp_location |
Ficheros temporales y de shuffle intermedio. Obligatorio |
--staging_location |
Donde se sube el código del pipeline y sus dependencias |
--service_account_email |
La identidad de los trabajadores. Sin esto se usa la cuenta de servicio de Compute por defecto, que suele ser Editor del proyecto: un incumplimiento directo del mínimo privilegio de 03-04 |
--subnetwork |
Los trabajadores arrancan dentro de alpinashop-vpc, en la subred de datos. No en una red por defecto |
--no_use_public_ips |
Sin IP pública: salen a internet por el Cloud NAT que ya configuraste en 03-01, y alcanzan las API de Google por Private Google Access |
--machine_type |
Tipo de VM del trabajador |
--max_num_workers |
Techo del autoescalado: la red de seguridad contra una factura desbocada |
--job_name |
Nombre visible. Incluir la fecha ayuda a localizar reprocesos |
--labels |
Etiquetas de facturación, mismo esquema que el resto de AlpinaShop |
Comprobación del estado:
gcloud dataflow jobs list --region=europe-west1 \
--format="table(id, name, type, state, createTime)" --limit=5
# Detalle y metricas de un job concreto
gcloud dataflow jobs describe JOB_ID --region=europe-west1
gcloud dataflow metrics list JOB_ID --region=europe-west1 \
--format="table(name.name, scalar)" --filter="name.name~filas_"Ese último comando devuelve los contadores filas_ok y filas_rechazadas que instrumentamos. Ese es el retorno de haber usado métricas de Beam en lugar de print.
- El tiempo en streaming: evento frente a proceso
Aquí empieza la parte difícil, y también la que hace que merezca la pena aprender Beam en lugar de improvisar.
En lote, el tiempo es simple: tienes el fichero, tiene un final, procesas y terminas. En streaming no hay final, y aparece una distinción que lo cambia todo:
- Tiempo del evento (event time): cuándo ocurrió el hecho en el mundo real. El cliente pulsó "Comprar" a las 10:00:00.
- Tiempo de proceso (processing time): cuándo el dato llega a tu pipeline. 10:00:03, o 10:07:12 si el móvil estaba en un túnel, o 11:30 si la app guardó el evento en local hasta recuperar cobertura.
Con AlpinaShop es muy concreto. Un cliente en el metro de Barcelona navega el catálogo, añade una mochila al carrito a las 18:42, pierde cobertura, y la app envía los eventos acumulados a las 18:51.
flowchart LR
subgraph Real["Tiempo del evento (mundo real)"]
E1["18:41 ver_producto"]
E2["18:42 anadir_carrito"]
E3["18:44 iniciar_pago"]
end
subgraph Proceso["Tiempo de proceso (llegada al pipeline)"]
P1["18:51 los tres a la vez"]
end
E1 --> P1
E2 --> P1
E3 --> P1
Ahora la pregunta de negocio: ¿cuántos productos se vieron entre las 18:40 y las 18:45? Si agrupas por tiempo de proceso, la respuesta es cero, y es falsa. Si agrupas por tiempo del evento, la respuesta incluye a este cliente, y es la correcta.
Beam agrupa por tiempo del evento por defecto. Esa es su decisión de diseño más importante y la razón por la que sus números cuadran con los del informe por lotes del día siguiente. Un sistema que agrupa por tiempo de proceso da resultados que dependen de la red y que nunca son reproducibles: reprocesar el mismo día da un resultado distinto.
- Marcas de agua, ventanas, disparadores y datos tardíos
Si esperas por tiempo del evento, surge la pregunta inevitable: ¿cuándo dejas de esperar? Un evento de las 18:42 podría llegar mañana. ¿Cierras la ventana o esperas eternamente?
La respuesta de Beam son cuatro mecanismos que se combinan.
Ventanas (Windowing)
Trocean la PCollection no acotada en trozos finitos sobre los que sí se puede agregar.
| Tipo | Definición | Uso en AlpinaShop |
|---|---|---|
| Fija (fixed/tumbling) | Intervalos contiguos que no se solapan | Pedidos por hora para el panel de campaña |
| Deslizante (sliding) | Intervalos que se solapan | Media móvil de visitas: ventana de 30 min cada 5 min |
| De sesión (session) | Se agrupan eventos separados por menos de un hueco dado | Sesiones de navegación reales: todo lo que hace un usuario hasta estar 30 min inactivo |
| Global | Una sola ventana infinita | Solo con disparadores explícitos |
from apache_beam import window
# Ventana fija de 1 hora: para el contador de pedidos
por_hora = eventos | "VentanaHora" >> beam.WindowInto(
window.FixedWindows(60 * 60)
)
# Ventana deslizante: media movil de 30 min, actualizada cada 5
media_movil = eventos | "VentanaDeslizante" >> beam.WindowInto(
window.SlidingWindows(size=30 * 60, period=5 * 60)
)
# Ventana de sesion: agrupa la actividad de un usuario
sesiones = eventos | "VentanaSesion" >> beam.WindowInto(
window.Sessions(gap_size=30 * 60)
)La ventana de sesión es especialmente elegante y no tiene equivalente sencillo en SQL: su duración no se fija de antemano, la determinan los propios datos. Es exactamente la definición de "sesión de navegación" que necesita la tabla visitas de 04-01.
Marca de agua (watermark)
La marca de agua es la estimación que hace el sistema de "ya no espero eventos anteriores a este instante". Dataflow la calcula automáticamente observando las marcas de tiempo de los datos que van llegando: si de Pub/Sub llevan diez minutos llegando eventos de las 18:50 en adelante, la marca de agua avanza más allá de las 18:45 y las ventanas anteriores se pueden cerrar.
No es una garantía, es una heurística. Siempre puede llegar algo después. Por eso existen los otros dos mecanismos.
Disparadores (triggers)
El disparador decide cuándo emitir el resultado de una ventana. Por defecto, al pasar la marca de agua: un resultado por ventana, cuando se considera completa.
Pero el panel de la campaña de otoño no puede esperar una hora a ver el primer número. Con un disparador se emiten resultados parciales:
from apache_beam.transforms.trigger import (
AfterWatermark, AfterProcessingTime, AccumulationMode
)
pedidos_por_hora = (
eventos
| "Ventana" >> beam.WindowInto(
window.FixedWindows(60 * 60),
trigger=AfterWatermark(
early=AfterProcessingTime(60), # avance parcial cada minuto
late=AfterProcessingTime(10 * 60), # correcciones cada 10 min si llega tarde
),
allowed_lateness=2 * 60 * 60, # aceptamos hasta 2 h de retraso
accumulation_mode=AccumulationMode.ACCUMULATING,
)
| "Contar" >> beam.CombinePerKey(sum)
)Interpretación práctica de esta configuración para AlpinaShop:
- Cada minuto se emite un resultado provisional: el panel se mueve y la gente ve que el sistema está vivo.
- Cuando pasa la marca de agua, se emite el resultado considerado definitivo.
- Durante dos horas más (
allowed_lateness), si llegan eventos de esa hora —el cliente del metro—, se emite una corrección cada diez minutos. - Pasadas las dos horas, lo que llegue se descarta. Ese descarte es una decisión de negocio explícita, no un accidente, y hay que instrumentarlo con un contador para saber cuánto se pierde.
AccumulationMode.ACCUMULATING significa que cada emisión contiene el total acumulado de la ventana, así que el destino debe sobrescribir. La alternativa, DISCARDING, emite solo lo nuevo desde la última emisión, y el destino debe sumar. Confundirlas produce cifras dobladas o divididas, y es un error muy difícil de ver en un panel.
Datos tardíos
Todo lo que llega después de la marca de agua es un dato tardío. Con allowed_lateness decides cuánto tiempo lo aceptas. La regla es simple: cuanto más esperas, más correctos son los números y más recursos (estado en memoria) consume el pipeline. Dos horas para un comercio electrónico europeo es un valor razonable; dos días sería carísimo y no cambiaría las decisiones.
- El pipeline de streaming de
pedidos-nuevos
pedidos-nuevosEn la próxima lección, 04-04, crearemos el topic de Pub/Sub pedidos-nuevos, donde la aplicación Flask publicará un mensaje JSON cada vez que se confirme un pedido. Anticipemos el consumidor, porque es el caso de uso canónico de Dataflow.
"""
Pipeline de streaming: pedidos-nuevos (Pub/Sub) -> BigQuery.
Se ejecuta de forma continua, sin fin.
"""
import json
import logging
from datetime import datetime
import apache_beam as beam
from apache_beam import window
from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions
from apache_beam.transforms.trigger import AfterWatermark, AfterProcessingTime, AccumulationMode
PROYECTO = "alpinashop-datos"
SUSCRIPCION = f"projects/{PROYECTO}/subscriptions/sub-analitica"
class DecodificarPedido(beam.DoFn):
SALIDA_RECHAZOS = "rechazos"
def process(self, mensaje, momento=beam.DoFn.TimestampParam):
try:
datos = json.loads(mensaje.data.decode("utf-8"))
except Exception as exc:
yield beam.pvalue.TaggedOutput(
self.SALIDA_RECHAZOS,
{"payload": str(mensaje.data[:500]), "motivo": f"json invalido: {exc}"},
)
return
# Los atributos del mensaje de Pub/Sub viajan aparte del cuerpo
origen = mensaje.attributes.get("origen", "desconocido")
yield {
"pedido_id": datos["pedido_id"],
"fecha_pedido": datos["fecha"][:10],
"momento_pedido": datos["fecha"],
"cliente_id": datos.get("cliente_id"),
"canal": origen,
"estado": "confirmado",
"total_pedido": str(datos["total"]),
"momento_ingesta": datetime.utcnow().isoformat(),
}
def main():
opciones = PipelineOptions(streaming=True, save_main_session=True)
opciones.view_as(StandardOptions).streaming = True
with beam.Pipeline(options=opciones) as p:
mensajes = (
p
| "LeerPubSub" >> beam.io.ReadFromPubSub(
subscription=SUSCRIPCION,
with_attributes=True,
timestamp_attribute="momento_evento", # <-- clave
)
)
decodificados = (
mensajes
| "Decodificar" >> beam.ParDo(DecodificarPedido()).with_outputs(
DecodificarPedido.SALIDA_RECHAZOS, main="validos")
)
# A) Detalle en tiempo real, fila a fila
(
decodificados.validos
| "EscribirPedidos" >> beam.io.WriteToBigQuery(
table=f"{PROYECTO}:alpinashop_analitica.pedidos_streaming",
schema="pedido_id:STRING,fecha_pedido:DATE,momento_pedido:TIMESTAMP,"
"cliente_id:STRING,canal:STRING,estado:STRING,"
"total_pedido:NUMERIC,momento_ingesta:TIMESTAMP",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
method="STORAGE_WRITE_API",
)
)
# B) Agregado por hora y canal, para el panel de campana
(
decodificados.validos
| "VentanaHora" >> beam.WindowInto(
window.FixedWindows(3600),
trigger=AfterWatermark(early=AfterProcessingTime(60)),
allowed_lateness=7200,
accumulation_mode=AccumulationMode.ACCUMULATING,
)
| "ClaveCanal" >> beam.Map(lambda d: (d["canal"], float(d["total_pedido"])))
| "SumarPorCanal" >> beam.CombinePerKey(sum)
| "Formatear" >> beam.Map(
lambda kv, w=beam.DoFn.WindowParam: {
"hora_inicio": w.start.to_utc_datetime().isoformat(),
"canal": kv[0],
"ventas_eur": round(kv[1], 2),
}
)
| "EscribirAgregado" >> beam.io.WriteToBigQuery(
table=f"{PROYECTO}:alpinashop_analitica.ventas_por_hora",
schema="hora_inicio:TIMESTAMP,canal:STRING,ventas_eur:FLOAT",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
method="STORAGE_WRITE_API",
)
)
if __name__ == "__main__":
logging.getLogger().setLevel(logging.INFO)
main()Tres decisiones que hay que entender:
timestamp_attribute="momento_evento": le dice a Beam que use el atributo del mensaje como tiempo del evento, en lugar del momento de publicación en Pub/Sub. Sin esto, un mensaje retenido nueve minutos en el móvil se contabilizaría en la hora equivocada. Es la línea que hace que todo el apartado 9 funcione de verdad, y se olvida constantemente.method="STORAGE_WRITE_API": usa la API moderna de escritura, más barata y con semántica de exactamente una vez, en lugar de las inserciones en streaming clásicas.- Se lee de una suscripción, no del topic. Con una suscripción, si el pipeline se detiene, los mensajes se acumulan y se recuperan al arrancar. Leer del topic directamente hace que Dataflow cree una suscripción efímera y se pierda todo lo publicado mientras el pipeline está parado.
- Plantillas de Dataflow: la vía práctica
Todo lo anterior es potente y también es trabajo. Para tareas comunes existe algo mucho más simple: las plantillas.
Una plantilla es un pipeline ya compilado y parametrizado que se lanza sin escribir ni compilar código. Google publica decenas.
| Plantilla | Qué hace |
|---|---|
| Pub/Sub Subscription to BigQuery | Lee JSON de una suscripción y lo inserta en una tabla |
| Cloud Storage Text to BigQuery | Carga ficheros con una función JavaScript de transformación |
| JDBC to BigQuery | Vuelca una base de datos relacional |
| BigQuery to Cloud Storage (Parquet) | Exporta |
| Datastream to BigQuery | Aplica CDC (lo veremos en 04-05) |
| Bulk Compress/Decompress | Utilidades sobre el bucket |
Lanzar la de Pub/Sub a BigQuery para AlpinaShop es una línea:
gcloud dataflow jobs run alpinashop-pedidos-a-bq \
--gcs-location gs://dataflow-templates-europe-west1/latest/PubSub_Subscription_to_BigQuery \
--region europe-west1 \
--service-account-email "[email protected]" \
--subnetwork regions/europe-west1/subnetworks/sn-datos-euw1 \
--disable-public-ips \
--max-workers 5 \
--parameters \
inputSubscription=projects/alpinashop-datos/subscriptions/sub-analitica,\
outputTableSpec=alpinashop-datos:alpinashop_analitica.pedidos_streaming,\
outputDeadletterTable=alpinashop-datos:alpinashop_analitica.pedidos_streaming_erroresFíjate en outputDeadletterTable: los mensajes que no encajen con el esquema van a una tabla de errores en lugar de bloquear el pipeline. Es el mismo patrón de cuarentena que programamos a mano, ya resuelto.
Hay dos sabores:
- Plantillas clásicas: el grafo se compila al crearla; los parámetros solo pueden ser valores en tiempo de ejecución (
ValueProvider). - Plantillas flexibles (Flex Templates): el pipeline se empaqueta como imagen de contenedor en Artifact Registry y el grafo se construye al lanzar. Son más flexibles y son la opción recomendada hoy para plantillas propias.
Crear una plantilla flexible del pipeline por lotes de AlpinaShop, para que Composer o Cloud Scheduler la invoquen en 04-06:
# 1) Construir la imagen y publicar la plantilla
gcloud dataflow flex-template build \
gs://alpinashop-dataflow/plantillas/pedidos-lote.json \
--image-gcr-path europe-west1-docker.pkg.dev/alpinashop-prod/alpinashop/dataflow-pedidos:1.0.0 \
--sdk-language PYTHON \
--flex-template-base-image PYTHON3 \
--py-path . \
--env FLEX_TEMPLATE_PYTHON_PY_FILE=pipeline_pedidos.py \
--env FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE=requirements.txt
# 2) Ejecutarla con parametros
gcloud dataflow flex-template run "pedidos-$(date +%Y%m%d-%H%M%S)" \
--template-file-gcs-location gs://alpinashop-dataflow/plantillas/pedidos-lote.json \
--region europe-west1 \
--service-account-email "[email protected]" \
--parameters fecha=2026-03-14La imagen va al mismo Artifact Registry que ya usa el catálogo (europe-west1-docker.pkg.dev/alpinashop-prod/alpinashop/), lo que mantiene un único inventario de artefactos, coherente con lo que veremos en el módulo 6.
El criterio de AlpinaShop: para "Pub/Sub a BigQuery" sin transformación, la plantilla de Google, sin escribir una línea. Para el pipeline por lotes con validación de negocio propia, código Beam empaquetado como plantilla flexible. Escribir código solo cuando aporta lógica que la plantilla no tiene.
- Escalado automático y Dataflow Prime
Dataflow ajusta el número de trabajadores durante la ejecución. En lote, mira el trabajo pendiente; en streaming, mira el retraso de la cola (backlog) y el uso de CPU.
--num_workers 2 # con cuantos empieza
--max_num_workers 20 # techo duro: el control de coste
--autoscaling_algorithm THROUGHPUT_BASED # por defecto en streamingPoner un max_num_workers bajo no siempre ahorra: un pipeline que tarda diez horas con dos trabajadores puede costar lo mismo que uno que tarda una hora con veinte, porque se factura por trabajador y hora. Lo que sí evita el techo es la sorpresa de un trabajo con un bucle patológico consumiendo cien máquinas toda la noche.
Dataflow Prime es la evolución del servicio con tres diferencias prácticas:
| Aspecto | Dataflow clásico | Dataflow Prime |
|---|---|---|
| Recursos | Eliges tipo de máquina para todo el pipeline | Ajuste vertical automático de memoria por paso |
| Escalado | Horizontal | Horizontal + vertical |
| Facturación | Por vCPU, memoria y disco por hora | Por Unidades de Cómputo de Datos (DCU) |
| Diagnóstico | Métricas | Recomendaciones automáticas de cuellos de botella |
| Cuándo usarlo | Pipelines estables y bien dimensionados | Pipelines con pasos de consumo muy desigual |
La ventaja real de Prime aparece cuando un paso del pipeline necesita mucha memoria y los demás no: en el modelo clásico dimensionas todas las máquinas para el paso más exigente y desperdicias recursos el resto del tiempo. Se activa con --dataflow_service_options=enable_prime.
Para AlpinaShop, con pipelines modestos y predecibles, el modelo clásico con n2-standard-2 es suficiente y más fácil de razonar en la factura. Prime queda anotado para cuando el volumen lo justifique.
- Monitorización, paralelismo y sesgo de datos
La consola de Dataflow muestra el grafo de ejecución con cada paso y, en cada uno, elementos procesados, tiempo de CPU y estado. Es la mejor interfaz de depuración de datos de la plataforma, y por eso insistimos en nombrar los pasos.
Métricas que hay que mirar:
| Métrica | Qué indica | Umbral de alarma |
|---|---|---|
| Retraso del sistema (system lag) | Segundos que lleva el elemento más antiguo sin procesar | En streaming, si crece de forma sostenida, el pipeline no da abasto |
| Frescura del dato (data freshness) | Antigüedad del dato más reciente ya emitido | Debe mantenerse estable |
| Elementos por segundo por paso | Dónde está el cuello de botella | El paso más lento manda |
| Uso de vCPU | Si el escalado sirve | Alto y con retraso creciente: hay un límite estructural |
| Trabajadores actuales | Comportamiento del autoescalado | Pegado al máximo: subir el techo o arreglar el pipeline |
El problema más frecuente y más difícil de diagnosticar es el sesgo de datos (data skew): una clave concentra una proporción enorme de los elementos. En AlpinaShop es muy fácil que ocurra: si agrupas visitas por sku y la mochila estrella acumula el 40 % del tráfico, un solo trabajador procesará el 40 % del trabajo mientras los otros diecinueve esperan. El síntoma es inconfundible: el escalado no mejora nada y hay un paso con un trabajador al 100 % y el resto ociosos.
Tres remedios, en orden de preferencia:
# 1) EL MEJOR: usar CombinePerKey en vez de GroupByKey.
# La agregacion parcial ocurre en cada trabajador antes del shuffle.
ventas = lineas | "Sumar" >> beam.CombinePerKey(sum)
# 2) Si necesitas GroupByKey de verdad: anadir sal a la clave
import random
def salar(elemento, n=20):
clave, valor = elemento
return (f"{clave}#{random.randint(0, n - 1)}", valor)
resultado = (
lineas
| "Salar" >> beam.Map(salar)
| "AgruparParcial" >> beam.CombinePerKey(sum)
| "QuitarSal" >> beam.Map(lambda kv: (kv[0].split("#")[0], kv[1]))
| "AgruparFinal" >> beam.CombinePerKey(sum)
)
# 3) Si un lado del JOIN es pequeno (el catalogo de productos):
# usar entrada lateral en lugar de CoGroupByKey, y evitar el shuffle
productos = p | "LeerProductos" >> beam.io.ReadFromBigQuery(query=SQL_PRODUCTOS)
catalogo = beam.pvalue.AsDict(productos | beam.Map(lambda r: (r["sku"], r)))
enriquecido = lineas | "Enriquecer" >> beam.Map(
lambda linea, cat: {**linea, "categoria": cat.get(linea["sku"], {}).get("categoria")},
cat=catalogo,
)La técnica 2, la sal, merece explicación: al añadir un sufijo aleatorio a la clave, la mochila estrella se convierte en veinte claves distintas que se reparten entre veinte trabajadores; luego se quita la sal y se hace una segunda agregación sobre veinte valores, que es trivial. Solo funciona con operaciones asociativas, pero cubre casi todos los casos de agregación.
La técnica 3, la entrada lateral (side input), es la que más usarás en AlpinaShop: el catálogo de productos son unos miles de filas y cabe en memoria de cada trabajador. Difundirlo evita por completo el shuffle del JOIN. Cuidado con el límite: si la entrada lateral no cabe en memoria, el pipeline se degrada muchísimo.
- Coste: qué se paga exactamente
Dataflow no factura por pipeline ni por dato procesado: factura los recursos consumidos por los trabajadores.
| Concepto | Cómo se mide | Orden de magnitud (verificar en la documentación oficial) |
|---|---|---|
| vCPU | Por vCPU y hora | ~0,05-0,07 $ (lote); algo más en streaming |
| Memoria | Por GB y hora | ~0,003-0,004 $ |
| Disco persistente | Por GB y hora | ~0,00005 $ estándar |
| Shuffle (lote) | Por GB procesados en el servicio de shuffle | ~0,011 $/GB |
| Streaming Engine | Por GB de datos en streaming procesados | ~0,018 $/GB |
| Dataflow Prime | Por DCU | Modelo unificado |
Un ejemplo realista para AlpinaShop: el pipeline por lotes nocturno, con 4 trabajadores n2-standard-2 (2 vCPU, 8 GB) durante 20 minutos, sale por céntimos. El pipeline de streaming, en cambio, funciona 24×7: 2 trabajadores permanentes son unas 1.440 horas de vCPU al mes, del orden de 80-100 € mensuales. Ese es el número que sorprende a todo el mundo, y es la razón de la advertencia del principio de la lección.
Consejos concretos para no pagar de más:
- Activa Streaming Engine y Shuffle Service (
--enable_streaming_engine,--experiments=shuffle_mode=service). Mueven el shuffle y el estado fuera de los trabajadores, permitiendo máquinas más pequeñas y discos mucho menores. Casi siempre sale a cuenta. - Reduce el disco. Con Streaming Engine,
--disk_size_gb=30es suficiente; el valor por defecto es mucho mayor y se paga por hora. - Pon
--max_num_workerssiempre. Es el freno de mano. - Usa VM Spot en lote tolerante a interrupciones:
--flexrs_goal=COST_OPTIMIZEDretrasa el arranque hasta 6 horas a cambio de un descuento notable. Perfecto para el volcado nocturno; inaceptable para streaming. - ¿De verdad necesitas streaming? Un micro-lote cada 15 minutos con la plantilla de Cloud Storage a BigQuery cuesta una fracción de un pipeline permanente. Si el negocio tolera 15 minutos de retraso, la respuesta es no.
- Vigila los pipelines de streaming olvidados. Un job de pruebas que nadie paró es la partida fantasma más común en la factura de datos. Lo detectarás con
gcloud dataflow jobs list --status=active.
- Cuándo Dataflow no es la respuesta
La honestidad sobre los límites es parte de saber usar una herramienta.
| Situación | Mejor opción | Por qué |
|---|---|---|
| Transformación expresable en SQL sobre datos que ya están en BigQuery | BigQuery (04-01) | No muevas datos para transformarlos; usa INSERT ... SELECT o una vista materializada |
| Ya tienes código Spark o el equipo sabe Spark, no Beam | Dataproc (04-03) | Reescribir a Beam es un coste sin retorno claro |
| Carga simple de ficheros a BigQuery, sin lógica | bq load |
Es gratis; Dataflow costaría dinero por hacer lo mismo |
| Integrar orígenes externos con un equipo sin perfil de programación | Data Fusion (04-05) | Interfaz visual, conectores listos |
| Decidir el orden en que se ejecutan varios procesos | Composer o Workflows (04-06) | Dataflow ejecuta un pipeline; no orquesta a los demás |
| Reaccionar a un evento puntual con poca lógica | Cloud Functions (06-03) | Un pipeline entero para procesar un fichero es desproporcionado |
| Copiar una base de datos con captura de cambios | Datastream (04-05) | CDC gestionado, sin código |
La confusión más habitual es la de la penúltima fila. Dataflow no es un orquestador. Puede ejecutar un pipeline complejísimo, pero no sabe "primero exporta Cloud SQL, luego carga en BigQuery, luego lanza este pipeline, y si algo falla avisa a Marta". Eso es exactamente lo que resuelve 04-06.
Errores Comunes y Consejos
Olvidar save_main_session=True. El pipeline funciona en local y falla en Dataflow con NameError: name 'PROYECTO' is not defined. Los trabajadores no reciben el ámbito global del módulo si no se lo dices.
No fijar la versión del SDK. Los trabajadores usan la versión con la que lanzaste el pipeline. Un pip install apache-beam[gcp] sin versión hace que el mismo código funcione hoy y falle mañana. Fija la versión en requirements.txt.
No poner nombres a los pasos. Sin "Nombre" >>, el grafo de la consola es ilegible y las métricas no dicen nada. Además, cambiar el nombre de un paso impide actualizar un pipeline de streaming en marcha (--update), porque Beam no puede mapear el estado del paso antiguo al nuevo.
Usar GroupByKey donde cabe CombinePerKey. Es la diferencia entre un pipeline que escala y uno que se cae con una clave caliente.
Leer de un topic en lugar de una suscripción. Con topic, Dataflow crea una suscripción temporal y se pierde todo lo publicado mientras el pipeline está parado. Con suscripción, los mensajes se acumulan y se recuperan.
Olvidar timestamp_attribute. Las ventanas se calculan sobre el momento de publicación en lugar del momento del evento, y los números no cuadran con los del proceso por lotes. Es sutil y grave.
Dejar un pipeline de streaming de pruebas encendido. Factura 24×7. Ponles etiquetas y revisa gcloud dataflow jobs list --status=active con periodicidad.
No usar la cuenta de servicio propia. Sin --service_account_email, los trabajadores usan la cuenta por defecto de Compute Engine, que en muchos proyectos es Editor. Un pipeline con permisos de Editor sobre producción es un riesgo innecesario.
Consejo: escribe siempre la rama de errores. Un pipeline sin salida de rechazos no es un pipeline de producción. Si no puedes ver qué se descartó y por qué, no puedes confiar en los números.
Consejo: --update para modificar un pipeline de streaming. Permite sustituir el código conservando el estado en curso, en lugar de drenarlo y arrancar de cero. Requiere compatibilidad del grafo, lo que es otra razón para no renombrar pasos alegremente.
Consejo: drena, no canceles. gcloud dataflow jobs drain deja de leer entradas nuevas y termina lo que tiene en vuelo. cancel mata el trabajo y puede perder datos en proceso.
Ejercicios
Ejercicio 1: pipeline de opiniones con cuarentena
Escribe un pipeline por lotes en Beam que lea gs://alpinashop-catalogo/exportaciones/2026/03/opiniones-*.csv con las columnas opinion_id,sku,fecha,puntuacion,texto,pais, y que:
- rechace las filas cuya
puntuacionno sea un entero entre 1 y 5, o cuyopaisno sea un código de dos letras; - normalice el
skua mayúsculas y recorte eltextoa 500 caracteres; - escriba las válidas en
alpinashop-datos:alpinashop_analitica.opiniones; - escriba las rechazadas, con su motivo, en
gs://alpinashop-catalogo/cuarentena/opiniones/; - lleve un contador de válidas y rechazadas visible en Dataflow.
Pruébalo con DirectRunner y un fichero local con al menos dos filas defectuosas.
Ejercicio 2: ventanas y disparadores para el panel de campaña
Dirección quiere un panel de la campaña de otoño con las ventas por hora y por país. Requisitos: agrupar por tiempo del evento, mostrar un avance parcial cada 30 segundos para que el panel se vea vivo, aceptar eventos con hasta 90 minutos de retraso emitiendo correcciones, y que cada emisión contenga el total acumulado de la hora. Escribe únicamente el fragmento de WindowInto y la agregación, y explica qué escribe exactamente el destino y por qué el modo de acumulación elegido obliga a un write_disposition concreto.
Ejercicio 3: diagnóstico de un pipeline que no escala
El pipeline de streaming de visitas lleva tres días funcionando. Desde ayer, el system lag ha pasado de 4 segundos a 22 minutos y sigue creciendo. El autoescalado ha llegado a los 20 trabajadores (su máximo), pero el uso medio de vCPU del conjunto es del 18 %. En el grafo, el paso AgruparPorSKU muestra un trabajador con 6 horas de tiempo de CPU y los demás con menos de 10 minutos. Ayer, marketing lanzó una campaña de la mochila MOCH-40L-AZ que ha multiplicado por doce sus visitas.
Diagnostica la causa, explica por qué subir --max_num_workers a 50 no arreglaría nada, y propón dos soluciones de código con su diferencia práctica.
Soluciones
Solución 1
import csv, io, logging, re
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions
PROYECTO, DATASET = "alpinashop-datos", "alpinashop_analitica"
BUCKET = "alpinashop-catalogo"
RE_PAIS = re.compile(r"^[A-Z]{2}$")
class ProcesarOpinion(beam.DoFn):
RECHAZOS = "rechazos"
def __init__(self):
self.ok = beam.metrics.Metrics.counter("opiniones", "validas")
self.ko = beam.metrics.Metrics.counter("opiniones", "rechazadas")
def process(self, linea):
try:
c = next(csv.reader(io.StringIO(linea)))
except Exception as exc:
self.ko.inc()
yield beam.pvalue.TaggedOutput(self.RECHAZOS,
{"linea": linea, "motivo": f"csv ilegible: {exc}"})
return
if len(c) != 6:
self.ko.inc()
yield beam.pvalue.TaggedOutput(self.RECHAZOS,
{"linea": linea, "motivo": f"esperadas 6 columnas, hay {len(c)}"})
return
opinion_id, sku, fecha, punt, texto, pais = [x.strip() for x in c]
motivos = []
try:
puntuacion = int(punt)
if not 1 <= puntuacion <= 5:
motivos.append("puntuacion fuera del rango 1-5")
except ValueError:
motivos.append("puntuacion no numerica")
puntuacion = None
pais = pais.upper()
if not RE_PAIS.match(pais):
motivos.append(f"pais invalido: {pais}")
if not sku:
motivos.append("sku vacio")
if motivos:
self.ko.inc()
yield beam.pvalue.TaggedOutput(self.RECHAZOS,
{"linea": linea, "motivo": "; ".join(motivos)})
return
self.ok.inc()
yield {
"opinion_id": opinion_id,
"sku": sku.upper(),
"fecha": fecha,
"puntuacion": puntuacion,
"texto": texto[:500],
"pais": pais,
}
def main():
opciones = PipelineOptions()
opciones.view_as(SetupOptions).save_main_session = True
with beam.Pipeline(options=opciones) as p:
salidas = (
p
| "Leer" >> beam.io.ReadFromText(
f"gs://{BUCKET}/exportaciones/2026/03/opiniones-*.csv",
skip_header_lines=1)
| "Procesar" >> beam.ParDo(ProcesarOpinion()).with_outputs(
ProcesarOpinion.RECHAZOS, main="validas")
)
(salidas.validas
| "EscribirBQ" >> beam.io.WriteToBigQuery(
table=f"{PROYECTO}:{DATASET}.opiniones",
schema="opinion_id:STRING,sku:STRING,fecha:DATE,"
"puntuacion:INTEGER,texto:STRING,pais:STRING",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
additional_bq_parameters={
"timePartitioning": {"type": "DAY", "field": "fecha"},
"clustering": {"fields": ["sku"]},
}))
(salidas[ProcesarOpinion.RECHAZOS]
| "Serializar" >> beam.Map(lambda r: f'{r["motivo"]}\t{r["linea"]}')
| "EscribirCuarentena" >> beam.io.WriteToText(
f"gs://{BUCKET}/cuarentena/opiniones/rechazos",
file_name_suffix=".tsv"))
if __name__ == "__main__":
logging.getLogger().setLevel(logging.INFO)
main()Fichero de prueba con defectos deliberados:
opinion_id,sku,fecha,puntuacion,texto,pais OPI-0001,moch-40l-az,2026-03-10,5,Muy comoda para travesias largas,es OPI-0002,FRON-300L,2026-03-10,9,Puntuacion imposible,FR OPI-0003,CRAM-12P,2026-03-11,cuatro,Puntuacion no numerica,ES OPI-0004,TIEN-2P,2026-03-11,4,Pais mal formado,ESPANA
La primera pasa (el sku se normaliza a mayúsculas y es a ES); las tres siguientes van a cuarentena con motivos distintos. Métricas esperadas: validas=1, rechazadas=3.
Solución 2
from apache_beam import window
from apache_beam.transforms.trigger import (
AfterWatermark, AfterProcessingTime, AccumulationMode
)
ventas_hora_pais = (
eventos
| "VentanaCampana" >> beam.WindowInto(
window.FixedWindows(3600), # 1 hora, tiempo del EVENTO
trigger=AfterWatermark(
early=AfterProcessingTime(30), # avance cada 30 s
late=AfterProcessingTime(300), # correcciones cada 5 min
),
allowed_lateness=90 * 60, # 90 minutos de retraso
accumulation_mode=AccumulationMode.ACCUMULATING,
)
| "ClavePais" >> beam.Map(lambda d: (d["pais"], float(d["total_pedido"])))
| "SumarPorPais" >> beam.CombinePerKey(sum)
| "Formatear" >> beam.Map(
lambda kv, w=beam.DoFn.WindowParam: {
"hora_inicio": w.start.to_utc_datetime().isoformat(),
"pais": kv[0],
"ventas_eur": round(kv[1], 2),
})
)Qué escribe exactamente el destino. Con ACCUMULATING, cada emisión de una ventana contiene el total acumulado desde el inicio de la ventana, no el incremento. Para la ventana 10:00-11:00 y el país ES, el destino recibirá una secuencia como: 120 € (a las 10:00:30), 345 € (10:01:00), … 4.210 € (al pasar la marca de agua), 4.235 € (a las 11:20, corrección por un evento tardío).
Por qué obliga a un write_disposition concreto. Si se hiciera WRITE_APPEND sobre la tabla final, esas seis emisiones se sumarían y el panel mostraría más de 9.000 € donde hay 4.235: el dato quedaría multiplicado. Las opciones correctas son:
- Escribir en una tabla de estados intermedios con
WRITE_APPENDy que el panel consulte solo la última emisión por ventana y clave, por ejemplo conQUALIFY ROW_NUMBER() OVER (PARTITION BY hora_inicio, pais ORDER BY momento_emision DESC) = 1. Es la opción recomendada en streaming, porque conserva el historial de correcciones y permite auditar. - O bien usar un destino que soporte sobrescritura por clave (
MERGEposterior sobrehora_inicio+pais).
Si en cambio se eligiera AccumulationMode.DISCARDING, cada emisión traería solo el incremento y entonces WRITE_APPEND con una suma posterior sería lo correcto. La combinación modo de acumulación + disposición de escritura debe decidirse conjuntamente; equivocarse produce paneles con cifras infladas que nadie detecta hasta que alguien cruza el número con contabilidad.
Solución 3
Diagnóstico: sesgo de datos (data skew) por clave caliente.
Las tres evidencias apuntan a lo mismo y se refuerzan entre sí:
- 20 trabajadores con 18 % de CPU media. Si el pipeline estuviera realmente saturado, la CPU estaría alta. Un uso bajo con retraso creciente significa que los trabajadores esperan, no que trabajen.
- Un trabajador con 6 horas de CPU y el resto con 10 minutos. Esa es la firma exacta del sesgo: el trabajo no se reparte.
- La campaña de
MOCH-40L-AZ. En unGroupByKeyporsku, todos los eventos de la mochila estrella se envían por el shuffle a un único trabajador, porque la clave determina el destino. Los demás no pueden ayudarle.
Por qué subir a 50 trabajadores no arregla nada. La unidad de paralelismo de una agrupación por clave es la clave, no el elemento. Los eventos de MOCH-40L-AZ seguirán yendo todos al mismo sitio, con 20 trabajadores o con 500. Los 30 nuevos estarían ociosos y facturando: el retraso seguiría creciendo y el coste se multiplicaría por 2,5. Escalar horizontalmente no resuelve un problema de distribución.
Solución A — sustituir GroupByKey por CombinePerKey (la buena, si la operación lo permite):
# ANTES: todos los eventos de la clave caliente viajan a un trabajador
visitas_por_sku = (
eventos
| "ClaveSKU" >> beam.Map(lambda e: (e["sku"], 1))
| "Agrupar" >> beam.GroupByKey()
| "Contar" >> beam.Map(lambda kv: (kv[0], len(list(kv[1]))))
)
# DESPUES: cada trabajador suma lo suyo antes del shuffle
visitas_por_sku = (
eventos
| "ClaveSKU" >> beam.Map(lambda e: (e["sku"], 1))
| "Contar" >> beam.CombinePerKey(sum)
)Por la red viajan sumas parciales, una por trabajador y clave, en lugar de millones de elementos individuales. Con 20 trabajadores, el trabajador de MOCH-40L-AZ recibe 20 números en lugar de doce millones de eventos. Es una línea de código y suele resolver el problema por completo.
Solución B — añadir sal a la clave (cuando hace falta GroupByKey de verdad, por ejemplo para conservar los elementos):
import random
N_SAL = 50
visitas_por_sku = (
eventos
| "SalarClave" >> beam.Map(
lambda e: (f'{e["sku"]}#{random.randint(0, N_SAL - 1)}', 1))
| "ParcialSalado" >> beam.CombinePerKey(sum)
| "QuitarSal" >> beam.Map(lambda kv: (kv[0].split("#")[0], kv[1]))
| "TotalFinal" >> beam.CombinePerKey(sum)
)Diferencia práctica entre A y B. La A es preferible siempre que se pueda: es más simple, más barata y no tiene parámetros que ajustar. La B añade un shuffle extra y un parámetro (N_SAL) que hay que dimensionar —demasiado bajo no reparte, demasiado alto crea sobrecarga—, pero es la única vía cuando la operación no es asociativa (por ejemplo, si necesitas la lista completa de eventos de la sesión para reconstruir un recorrido) o cuando el sesgo persiste incluso con agregación parcial.
Medida complementaria inmediata: antes de tocar código, bajar --max_num_workers a 5 para dejar de pagar 15 máquinas ociosas mientras se prepara el despliegue, y usar --update para sustituir el pipeline conservando el estado en curso, sin perder los datos en vuelo.
Conclusión
AlpinaShop ya tiene tuberías. En esta lección has visto por qué un script en una VM es una solución que funciona hasta que deja de funcionar, y qué aporta exactamente un servicio gestionado: paralelismo automático, reintentos, escalado durante la ejecución, semántica de escritura fiable y observabilidad de serie.
Has aprendido el modelo de Apache Beam —Pipeline, PCollection inmutable, PTransform, runner— y su idea central: un lote es un flujo acotado y un stream es un flujo no acotado, así que la misma lógica sirve para ambos. Conoces las transformaciones que cubren casi todos los casos y, sobre todo, sabes por qué CombinePerKey es casi siempre mejor que GroupByKey, y cuándo un ParDo con clase DoFn gana a un Map.
Has escrito el pipeline por lotes de AlpinaShop línea a línea: lee las exportaciones CSV de alpinashop-catalogo, parsea, valida con reglas de negocio reales —importes con coma, cantidades negativas, fechas rotas—, escribe el detalle limpio en lineas_pedido, calcula el agregado de ventas por SKU con agregación parcial, y manda todo lo defectuoso a una cuarentena en el bucket con su motivo, sin tumbar el proceso. Lo has probado en local con el DirectRunner, que es lento y estricto a propósito, con un fichero sucio a conciencia, y luego lo has lanzado en Dataflow con la cuenta de servicio sa-dataflow-pedidos, dentro de sn-datos-euw1, sin IP pública, con techo de trabajadores y etiquetas de facturación.
Has entrado en el territorio del tiempo, que es lo que de verdad distingue el procesamiento de datos serio: tiempo del evento frente a tiempo de proceso, con el cliente del metro como ejemplo de por qué agrupar por el segundo produce números falsos; ventanas fijas, deslizantes y de sesión; marcas de agua como estimación de "ya no espero más"; disparadores para emitir resultados parciales sin renunciar a la corrección posterior; y allowed_lateness como decisión de negocio explícita sobre cuánto se espera a los rezagados. Y has dejado escrito el pipeline de streaming que consumirá pedidos-nuevos en cuanto exista, con timestamp_attribute —la línea que hace que todo lo anterior sea cierto— y leyendo de una suscripción, no de un topic.
Has visto que para lo común hay atajo: las plantillas de Google resuelven "Pub/Sub a BigQuery" sin escribir código, y las plantillas flexibles empaquetan tu propio pipeline como imagen en Artifact Registry para que otro sistema lo invoque. Conoces el autoescalado y sus límites, Dataflow Prime y cuándo compensa, y sabes leer el grafo, el system lag y la frescura del dato para diagnosticar. Sabes reconocer el sesgo de datos —el escalado que no mejora nada— y corregirlo con agregación parcial, sal en la clave o entradas laterales. Y conoces la factura: por vCPU, memoria, disco y shuffle, con la advertencia grande de que un pipeline de streaming está encendido siempre y cuesta decenas de euros al mes aunque no pase nada.
Queda algo pendiente que has notado a lo largo de toda la lección. El pipeline por lotes funciona porque los datos ya están en el bucket, y el de streaming funciona porque alguien publica en pedidos-nuevos. Pero ese topic todavía no existe, y no hemos hablado de qué pasa con el resto de la casa: cuando entra un pedido, el almacén tiene que prepararlo, facturación tiene que emitir la factura, el cliente tiene que recibir su correo de confirmación y la analítica tiene que enterarse. Hoy la aplicación Flask tendría que llamar a los cuatro, uno detrás de otro, y quedarse colgada si el cuarto no responde.
Antes de resolverlo, sin embargo, hay una pieza del ecosistema de datos que conviene conocer, porque mucha gente llega a Google Cloud con ella ya puesta. En 04-03, Cloud Dataproc, veremos Spark y Hadoop gestionados: qué son, por qué siguen importando después de veinte años, y cuál es el patrón que convierte un clúster caro y permanente en uno efímero que vive cuatro minutos, hace su trabajo sobre los datos del bucket y se autodestruye. Crearemos alpinashop-spark, ejecutaremos un trabajo PySpark que calcula qué productos se compran juntos —la base del futuro recomendador del módulo 5— y compararemos con honestidad cuándo conviene Dataproc, cuándo Dataflow y cuándo no hace falta ninguno de los dos.
Curso de Google Cloud Platform (GCP)
Módulo 1: Introducción a Google Cloud Platform
- ¿Qué es Google Cloud Platform?
- Configuración de tu cuenta de GCP
- Descripción general de la consola de GCP
- Proyectos, jerarquía de recursos y facturación
- Regiones, zonas y modelo de responsabilidad compartida
- Cloud Shell y la CLI de gcloud
Módulo 2: Servicios principales de GCP
- Compute Engine: máquinas virtuales en Google Cloud
- Cloud Storage: almacenamiento de objetos
- Cloud SQL: bases de datos relacionales gestionadas
- App Engine: plataforma como servicio
- Google Kubernetes Engine (GKE)
- Bases de datos NoSQL: Firestore, Bigtable y Spanner
- Cómo elegir el servicio de cómputo adecuado
Módulo 3: Redes y seguridad
- Redes VPC
- Balanceo de carga en la nube
- Cloud CDN
- Gestión de identidad y acceso (IAM)
- Cloud Armor
- Secretos y cifrado: Secret Manager y Cloud KMS
- Cloud DNS, certificados TLS y publicación segura de servicios
Módulo 4: Datos y análisis
- BigQuery: el almacén de datos analítico
- Cloud Dataflow: procesamiento de datos por lotes y en streaming
- Cloud Dataproc: Spark y Hadoop gestionados
- Cloud Pub/Sub: mensajería asíncrona
- Cloud Data Fusion: integración de datos sin código
- Orquestación de pipelines con Cloud Composer y Workflows
- Gobierno del dato y cuadros de mando con Dataplex y Looker Studio
Módulo 5: Aprendizaje automático e IA
- Vertex AI: la plataforma de machine learning de GCP
- AutoML: modelos a medida sin escribir código
- TensorFlow en GCP: entrenamiento y servicio de modelos
- API de lenguaje natural
- API de visión
- IA generativa en Vertex AI: modelos Gemini y embeddings
- MLOps: del modelo al producto con Vertex AI Pipelines
Módulo 6: DevOps y monitoreo
- Cloud Build: integración continua en GCP
- Cloud Source Repositories y gestión del código fuente
- Cloud Functions: funciones sin servidor
- Cloud Monitoring (antes Stackdriver): métricas, paneles y alertas
- Cloud Deployment Manager e infraestructura como código nativa
- Cloud Logging y Cloud Trace: logs, trazas y diagnóstico
- Terraform en GCP: infraestructura como código en la práctica
Módulo 7: Temas avanzados de GCP
- Híbrido y multinube con Anthos
- Computación sin servidor con Cloud Run
- Redes avanzadas: VPC compartida, peering y conectividad híbrida
- Mejores prácticas de seguridad
- Gestión y optimización de costos
- Fiabilidad: SLO, alta disponibilidad y recuperación ante desastres
- Gobierno a escala: organización, políticas y auditoría
