Logo GCP con Eduardo

GCP con Eduardo

Descarga el código de la lección

De Cloud Storage a BigQuery con Dataflow y Python | ETL Serverless
✦ Guía Técnica & Arquitectura

De Cloud Storage a BigQuery con Dataflow y Python

Construir un pipeline ETL batch serverless de GCS a BigQuery parece sencillo hasta que aparecen los problemas reales: esquemas incompatibles, cargas duplicadas, tablas creadas accidentalmente, particionado, escalabilidad, costes y decisiones sobre el mecanismo de escritura. En esta guía se diseña una solución empresarial con Apache Beam, Dataflow y BigQueryIO, poniendo el foco en decisiones arquitectónicas y no únicamente en hacer que el primer job termine.

Lo que aprenderás

Diseñar un pipeline GCS → Dataflow → BigQuery

Separando ingestion, transformación y persistencia analítica.

Definir esquemas BigQuery explícitos

Controlando tipos, nulabilidad y evolución del contrato de datos.

Utilizar WRITE_APPEND y CREATE_IF_NEEDED correctamente

Entendiendo las consecuencias operativas de cada disposición.

Preparar el pipeline para producción

Con idempotencia, observabilidad, control de costes y tolerancia a fallos.

Python Apache Beam Dataflow Google Cloud Storage BigQuery BigQueryIO ETL Batch Serverless

Arquitectura del pipeline GCS → BigQuery

La arquitectura recomendada separa claramente el almacenamiento de objetos, el motor de procesamiento y el almacén analítico. GCS actúa como landing zone, Dataflow ejecuta el pipeline Apache Beam y BigQuery representa la capa de serving analítico.

Apache Beam proporciona el modelo de programación; Dataflow es el runner administrado que ejecuta el pipeline a escala. Esta separación permite probar localmente la lógica y posteriormente ejecutarla sobre infraestructura administrada.

                 ┌─────────────────────────┐
                 │   Google Cloud Storage   │
                 │        gs://landing/     │
                 │                           │
                 │ CSV / JSON / TXT / ...   │
                 └────────────┬──────────────┘
                              │
                              ▼
                 ┌─────────────────────────┐
                 │      Apache Beam        │
                 │       Pipeline          │
                 │                         │
                 │ Read → Parse → Validate │
                 │ → Transform → Normalize │
                 └────────────┬──────────────┘
                              │
                              ▼
                 ┌─────────────────────────┐
                 │       Dataflow          │
                 │ Managed Runner / Batch  │
                 │ Auto Scaling / Workers   │
                 └────────────┬──────────────┘
                              │
                              ▼
                 ┌─────────────────────────┐
                 │        BigQuery         │
                 │   dataset.fact_events   │
                 │                         │
                 │ Schema + Partitioning   │
                 │ + Clustering + IAM      │
                 └─────────────────────────┘

Matriz de decisión: ¿cómo escribir en BigQuery?

En un pipeline batch no conviene elegir el mecanismo de escritura únicamente por costumbre. El tamaño de la carga, el paralelismo, las cuotas, la necesidad de una dead-letter queue y la tolerancia a duplicados influyen directamente en la decisión.

Patrón Latencia Consistencia / semántica Coste / escala Uso recomendado
FILE_LOADS Batch Carga mediante ficheros intermedios Muy adecuado para grandes cargas batch Backfills y grandes volúmenes históricos
STORAGE_WRITE_API Baja API de escritura orientada a streams Excelente paralelismo; requiere controlar cuotas Batch/streaming con requisitos avanzados
STORAGE_WRITE_API_AT_LEAST_ONCE Baja Puede producir duplicados Puede resultar más económico Cuando la deduplicación posterior es aceptable
Streaming Inserts Muy baja Escritura orientada a eventos Más apropiado para casos de baja latencia Workloads donde la disponibilidad rápida importa
Decisión de arquitectura: para cargas batch grandes, la documentación actual de Google recomienda probar inicialmente FILE_LOADS. Storage Write API puede ser interesante para determinados pipelines batch de menor escala o cuando se necesitan capacidades concretas. No existe un único método óptimo para todos los workloads.

Referencias oficiales: Dataflow → BigQuery · Apache Beam BigQueryIO
Antipatrón crítico: WRITE_APPEND sin idempotencia

WRITE_APPEND no significa que tu proceso de negocio sea idempotente. Significa que las filas se añaden a la tabla existente. Si el mismo fichero de GCS se procesa dos veces, puedes terminar con dos conjuntos de registros.

Este problema aparece especialmente cuando un orquestador reintenta un job, cuando un fichero permanece en la landing zone después de una ejecución satisfactoria o cuando varios pipelines procesan simultáneamente la misma entrada.

Patrón recomendado: introduce una clave de negocio o un identificador de ingestión, registra el objeto de origen y su generación/versión cuando aplique, y diseña la capa posterior para poder detectar o eliminar duplicados. Para procesos críticos, considera una estrategia explícita de manifest files, control de estado y/o staging antes de publicar los datos definitivos.

Implementación práctica con Apache Beam y Python

El siguiente ejemplo implementa un pipeline batch que descubre ficheros CSV en GCS, transforma cada línea en un diccionario compatible con BigQuery y utiliza WriteToBigQuery con un esquema explícito, CREATE_IF_NEEDED y WRITE_APPEND.

pipeline.py Apache Beam / Python
import apache_beam as beam

from apache_beam.options.pipeline_options import (
    PipelineOptions,
    StandardOptions,
)


PROJECT_ID = "my-gcp-project"
REGION = "europe-west1"

INPUT_PATTERN = "gs://my-landing-bucket/events/*.csv"
BQ_TABLE = f"{PROJECT_ID}:analytics.events"

TABLE_SCHEMA = {
    "fields": [
        {
            "name": "event_id",
            "type": "STRING",
            "mode": "REQUIRED",
        },
        {
            "name": "event_timestamp",
            "type": "TIMESTAMP",
            "mode": "REQUIRED",
        },
        {
            "name": "customer_id",
            "type": "STRING",
            "mode": "NULLABLE",
        },
        {
            "name": "event_type",
            "type": "STRING",
            "mode": "NULLABLE",
        },
        {
            "name": "amount",
            "type": "NUMERIC",
            "mode": "NULLABLE",
        },
    ]
}


def parse_csv_line(line: str) -> dict:
    """
    Convierte una línea CSV en una fila compatible
    con el esquema de BigQuery.

    En producción, utilizar un parser CSV real cuando
    existan campos entrecomillados o delimitadores complejos.
    """
    parts = line.split(",")

    return {
        "event_id": parts[0],
        "event_timestamp": parts[1],
        "customer_id": parts[2] or None,
        "event_type": parts[3] or None,
        "amount": float(parts[4]) if parts[4] else None,
    }


def run(argv=None):
    pipeline_options = PipelineOptions(argv)

    # Indicamos explícitamente que el pipeline será batch.
    pipeline_options.view_as(StandardOptions).streaming = False

    with beam.Pipeline(options=pipeline_options) as pipeline:

        rows = (
            pipeline
            | "ReadFromGCS"
            >> beam.io.ReadFromText(
                INPUT_PATTERN,
                skip_header_lines=1,
            )
            | "ParseCSV"
            >> beam.Map(parse_csv_line)
        )

        (
            rows
            | "WriteToBigQuery"
            >> beam.io.WriteToBigQuery(
                table=BQ_TABLE,
                schema=TABLE_SCHEMA,

                # Si la tabla no existe, BigQuery puede crearla
                # utilizando el esquema anterior.
                create_disposition=(
                    beam.io.BigQueryDisposition.CREATE_IF_NEEDED
                ),

                # Añade las filas a los datos existentes.
                write_disposition=(
                    beam.io.BigQueryDisposition.WRITE_APPEND
                ),
            )
        )


if __name__ == "__main__":
    run()

La API de Beam documenta WriteToBigQuery para escribir una PCollection de diccionarios y permite especificar schema, create_disposition y write_disposition. Cuando se utiliza CREATE_IF_NEEDED, debe existir un esquema disponible si la tabla de destino necesita ser creada.

CREATE_IF_NEEDED vs CREATE_NEVER

CREATE_IF_NEEDED

El pipeline puede crear la tabla si todavía no existe. Es cómodo para datasets gestionados completamente por el pipeline.

  • Requiere un esquema si crea la tabla.
  • Reduce pasos manuales de provisioning.
  • Útil para entornos controlados.

CREATE_NEVER

La tabla debe existir previamente. Si no existe, la escritura falla. Es una opción interesante cuando BigQuery está gobernado como contrato.

  • Schema gestionado fuera del pipeline.
  • Evita creación accidental de tablas.
  • Muy apropiado para producción regulada.

Decisión Staff

En una plataforma empresarial, el provisioning de datasets y tablas suele pertenecer a IaC o a una capa de gobierno de datos.

  • Terraform para infraestructura.
  • Data contracts para esquemas.
  • Pipeline centrado en procesamiento.

WRITE_APPEND: cuándo utilizarlo correctamente

WRITE_APPEND es una opción natural para cargas incrementales: cada ejecución añade nuevas filas a la tabla existente. El problema aparece cuando se confunde la operación física de append con la semántica de negocio.

Si una ejecución se reinicia después de haber persistido parcialmente datos, el diseño debe poder determinar qué registros ya fueron procesados. Una estrategia robusta suele incorporar una combinación de:

1. Identidad

Cada evento debe disponer de una clave estable, por ejemplo event_id.

2. Trazabilidad

Conserva metadatos como fichero de origen, fecha de ingestión y versión del objeto cuando el caso de uso lo requiera.

3. Publicación

Separa staging y serving cuando una carga parcialmente procesada no deba aparecer directamente en la tabla analítica final.

FinOps: no confundas serverless con coste cero

Dataflow elimina la administración directa de clusters, pero no elimina el coste de procesamiento. Una arquitectura serverless sigue necesitando controles de volumen, paralelismo y duración del job.

  • Filtra lo antes posible: no envíes a BigQuery columnas o registros que nunca serán utilizados.
  • Evita reprocesar: los reintentos funcionales no deben convertirse en reprocesamientos innecesarios de todo el histórico.
  • Controla el paralelismo: más workers no implica automáticamente menor coste total.
  • Usa particionamiento: una tabla BigQuery bien particionada permite reducir el volumen leído por consultas posteriores.
  • Mide: duración, workers, bytes procesados, errores, throughput y coste por ejecución deben formar parte de la observabilidad.

Patrones de diseño para producción

Contrato de datos

Define el schema de BigQuery explícitamente y versiona los cambios. Evita que la inferencia automática se convierta en el mecanismo principal de gobierno.

Landing → Staging → Serving

GCS conserva el raw input; Dataflow normaliza y valida; BigQuery expone la capa analítica. Esta separación facilita auditoría y replay.

Dead Letter

Los registros inválidos no deberían provocar necesariamente la pérdida de toda la carga. Diseña una salida de errores con suficiente contexto para poder corregir y reprocesar.

Idempotencia

El pipeline debe poder ejecutarse de nuevo sin generar efectos de negocio inesperados.

Observabilidad

Instrumenta métricas de registros leídos, transformados, rechazados y escritos, además de duración y errores.

Seguridad

Aplica mínimo privilegio a la identidad de ejecución de Dataflow, separando permisos de lectura de GCS y escritura sobre BigQuery.

Checklist de implementación empresarial

Define el contrato de entrada

Especifica formato, delimitador, encoding, columnas obligatorias, reglas de calidad y comportamiento ante registros inválidos.

Diseña la landing zone

Organiza los objetos de GCS por dominio, fecha, fuente o ejecución para facilitar descubrimiento y replay.

Implementa Parse → Validate → Transform

Separa la lectura de la lógica de negocio. Los errores de parsing y los errores de validación deben poder diagnosticarse.

Define el schema de BigQuery

Establece tipos, modos NULLABLE/REQUIRED y, cuando corresponda, particionamiento y clustering.

Decide la disposición de escritura

Utiliza WRITE_APPEND para cargas incrementales cuando la estrategia de idempotencia esté resuelta. Evalúa WRITE_TRUNCATE para reemplazos controlados y CREATE_NEVER cuando el schema sea gestionado externamente.

Prueba localmente

Ejecuta Apache Beam con un runner local y un dataset pequeño antes de desplegar en Dataflow.

Despliega en Dataflow

Configura proyecto, región, staging, temporary location, identidad de ejecución y límites operativos apropiados.

Opera y observa

Monitoriza errores, throughput, workers, duración, registros rechazados y coste. Define alertas antes de considerar el pipeline terminado.

Preguntas frecuentes sobre Dataflow GCS → BigQuery

¿Qué diferencia hay entre WRITE_APPEND y WRITE_TRUNCATE?

WRITE_APPEND añade nuevas filas a una tabla existente. WRITE_TRUNCATE reemplaza los datos existentes de la tabla. Para cargas incrementales, WRITE_APPEND suele ser apropiado, pero debe acompañarse de un diseño de idempotencia si el mismo input puede procesarse más de una vez.

¿Qué hace CREATE_IF_NEEDED en WriteToBigQuery?

Permite crear la tabla de destino cuando no existe. Si el pipeline puede crear la tabla, debe proporcionar un schema compatible con BigQuery. En arquitecturas donde las tablas están gobernadas externamente, CREATE_NEVER puede ser una opción más segura.

¿Por qué utilizar Dataflow en lugar de una carga directa de GCS a BigQuery?

Dataflow aporta valor cuando existe lógica de procesamiento: parsing, limpieza, validación, enriquecimiento, joins, deduplicación, transformaciones o necesidades de escalado. Para un simple movimiento de datos sin transformación, una solución más sencilla puede ser suficiente.

Sobre el Autor: Eduardo Martínez Agrelo

AI & Data Architect

Eduardo Martínez Agrelo es AI & Data Architect, especializado en arquitectura de datos, ingeniería de datos, plataformas cloud, inteligencia artificial y diseño de sistemas analíticos escalables. Su enfoque combina fundamentos técnicos, decisiones arquitectónicas, automatización y criterios de operación en producción.

En sus contenidos técnicos aborda tecnologías y patrones de arquitectura orientados a profesionales que necesitan pasar de ejemplos de laboratorio a soluciones robustas y mantenibles en entornos empresariales.

© 2026 Eduardo Martínez Agrelo · AI & Data Architect

Arquitectura de datos · Apache Beam · Dataflow · BigQuery · Cloud Storage