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

  1. Qué es un pipeline de datos y por qué un script no basta
  2. Apache Beam: el modelo unificado
  3. Los cuatro conceptos: Pipeline, PCollection, PTransform, runner
  4. Las transformaciones que usarás el 90 % del tiempo
  5. Primer pipeline por lotes, explicado línea a línea
  6. Ejecución local con DirectRunner
  7. Ejecución gestionada con DataflowRunner
  8. El tiempo en streaming: evento frente a proceso
  9. Marcas de agua, ventanas, disparadores y datos tardíos
  10. El pipeline de streaming de pedidos-nuevos
  11. Plantillas de Dataflow: la vía práctica
  12. Escalado automático y Dataflow Prime
  13. Monitorización, paralelismo y sesgo de datos
  14. Coste: qué se paga exactamente
  15. Cuándo Dataflow no es la respuesta

  1. 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.

  1. 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.

  1. Los cuatro conceptos: Pipeline, PCollection, PTransform, runner

Pipeline 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:

resultado = entrada | "Nombre descriptivo del paso" >> beam.Map(funcion)

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.

  1. 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.

  1. 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:

  • yield en lugar de return. Un DoFn es un generador: puede emitir cero, uno o muchos elementos por entrada. Un return con valor no funciona como esperas.
  • TaggedOutput marca 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 exporta 89,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:

  1. with beam.Pipeline(...) as p: al salir del bloque with, Beam llama a run() y espera. Sin el with, hay que llamar a p.run().wait_until_finish() explícitamente.
  2. WRITE_APPEND frente a WRITE_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.
  3. additional_bq_parameters crea 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.
  4. beam.Flatten() une varias PCollection del mismo tipo en una. Es la unión de ramas del grafo, el equivalente a un UNION ALL.
  5. save_main_session=True serializa el ámbito global del módulo para los trabajadores. Sin esto, un pipeline que funciona en local falla en Dataflow con NameError sobre las constantes o los imports. Es el error de novato número uno.

  1. Ejecución local con DirectRunner

Antes 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 DirectRunner

Ese 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,90 mete 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.

  1. Ejecución gestionada con DataflowRunner

Con 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=analitica

Opció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.

  1. 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.

  1. 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.

  1. El pipeline de streaming de pedidos-nuevos

En 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.

  1. 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_errores

Fí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-14

La 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.

  1. 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 streaming

Poner 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.

  1. 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.

  1. 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=30 es suficiente; el valor por defecto es mucho mayor y se paga por hora.
  • Pon --max_num_workers siempre. Es el freno de mano.
  • Usa VM Spot en lote tolerante a interrupciones: --flexrs_goal=COST_OPTIMIZED retrasa 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.

  1. 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 puntuacion no sea un entero entre 1 y 5, o cuyo pais no sea un código de dos letras;
  • normalice el sku a mayúsculas y recorte el texto a 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_APPEND y que el panel consulte solo la última emisión por ventana y clave, por ejemplo con QUALIFY 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 (MERGE posterior sobre hora_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í:

  1. 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.
  2. 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.
  3. La campaña de MOCH-40L-AZ. En un GroupByKey por sku, 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

Módulo 2: Servicios principales de GCP

Módulo 3: Redes y seguridad

Módulo 4: Datos y análisis

Módulo 5: Aprendizaje automático e IA

Módulo 6: DevOps y monitoreo

Módulo 7: Temas avanzados de GCP

Módulo 8: Proyecto final

© Copyright 2026. Todos los derechos reservados