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
Separando ingestion, transformación y persistencia analítica.
Controlando tipos, nulabilidad y evolución del contrato de datos.
Entendiendo las consecuencias operativas de cada disposición.
Con idempotencia, observabilidad, control de costes y tolerancia a fallos.
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 |
Referencias oficiales: Dataflow → BigQuery · Apache Beam BigQueryIO
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.
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.
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.
