Inicio Experiencia Análisis Escríbeme View in English

Cómo crear un pipeline ELT robusto

Google Cloud Data Engineering Python BigQuery Serverless
Cómo crear un pipeline ELT robusto

Voy a empezar diciendo que esta será la primera de una serie de publicaciones relacionadas con la construcción de un pipeline de grado empresarial.

Anteriormente, tenía un pet project que vendrá muy bien para ilustrar los conceptos y las mejores prácticas. Y sí, mis queridos lectores, se trata de consumir información de Clash of Clans. Me interesa que vayamos ordenadamente, por lo que hoy solamente hablaremos del primer y segundo paso: extraer y cargar.

En las arquitecturas modernas en la nube, este enfoque se popularizó debido a la gran capacidad de cómputo que ofrecen los servicios de Data Warehouse en la nube y, en general, al potencial de escalado de estos. Antes de hablar del stack de GCP que utilizaremos (spoiler: BigQuery es parte de él), necesitamos hacer un diagnóstico inicial.

Situación actual

Clash of Clans tiene una API que esencialmente permite obtener datos (con el método GET). Para lograr obtener dicha información necesitamos utilizar una API KEY; no obstante, hay un detalle relevante: se requiere añadir la IP que hará la solicitud a una lista de permitidos (whitelist) por motivos de seguridad.

Clash of Clans Developer Portal

Objetivo

Crear un flujo de extracción resiliente que se ejecute diariamente.

Resiliencia

¿Qué implica realmente la resiliencia en este caso? En nuestro escenario particular implica:

  1. Entender que las APIs de terceros no dependen de nosotros, por lo que el cuerpo (body) de la respuesta puede cambiar.
  2. Si hay fallas de red, debemos ser capaces de monitorear y reprocesar oportunamente, sin que esto provoque duplicación de datos en capas posteriores.
  3. Si recibimos un archivo payload corrupto o anómalo, el proceso no debe romperse, pero debe ser rastreable para poder diagnosticar y actuar.

¿Qué otras cosas considerar para otros procesos?

  1. Gestión de rate limits: Clash of Clans no documenta cuáles son estos, pero como solo estamos solicitando información de un clan, no debería haber throttle. No obstante, en proyectos de grado empresarial, gestionar los rate limits es obligatorio.
  2. Caídas del servidor y errores 500-504: incorporar estrategias de reintento como exponential backoff. Hay que tomar en cuenta que ciertos procesos consumen información que cambia rápidamente. No es lo mismo información transaccional que puede cambiar por segundo a la información de la API de Clash of Clans, que casi no varía en todo el día. Si algo llega a fallar, podemos revisar los logs y ejecutar manualmente; recuerden, no siempre hay que matar moscas a cañonazos.

El Stack

Debido a las características de este caso, podemos implementar un flujo de trabajo robusto con una inversión muy baja aprovechando los servicios serverless que ofrece GCP. La arquitectura se ilustra en la siguiente imagen:

ELT Pipeline Egress Architecture

Algunos puntos a resaltar de la arquitectura:

  1. Utilizamos Cloud NAT porque necesitamos una IP pública estática para la whitelist; esta es una forma segura y robusta de garantizarlo.
  2. Workflows es nuestro orquestador serverless; aquí es donde vemos todo el flujo. Inicialmente la capa bronce, pero añadiremos silver y gold posteriormente.
  3. Un Cloud Run job es el que ejecuta la solicitud cada día e inserta los registros en BigQuery.
  4. Posteriormente, añadiremos los pipelines de transformación utilizando Dataform; por ahora podemos ignorar esto.
  5. Como dato de color: el Data Warehouse y el flujo de extracción están en proyectos diferentes; esto puede darse dependiendo del esquema de gobierno de datos.

El código

Puedes encontrar el código en: github.com/JoseChavezUriarte/coc_elt. Sin embargo, quiero resaltar algunas cosas:

  1. En models.py, no hay una validación estricta de atributos del API. Esto es porque queremos flexibilidad por si se añaden nuevos atributos, evitando así que se rompa el pipeline. La prioridad es garantizar la extracción.
from typing import Any, List
from pydantic import BaseModel, ConfigDict, model_validator

def normalize_envelope(data: Any) -> Any:
    if isinstance(data, dict):
        if "items" in data:
            return data
        raise ValueError(f"Expected pagination envelope with an 'items' array. Received keys: {list(data.keys())}")
    if isinstance(data, list):
        return {"items": data}
    raise ValueError(f"Root payload must be a dictionary envelope or a list. Received type: {type(data).__name__}")

class ClanRecord(BaseModel):
    model_config = ConfigDict(extra='allow')
    tag: str
    name: str

class MemberRecord(BaseModel):
    model_config = ConfigDict(extra='allow')
    tag: str
    name: str

class MemberListResponse(BaseModel):
    items: List[MemberRecord]

    @model_validator(mode='before')
    @classmethod
    def validate_envelope(cls, data: Any) -> Any:
        return normalize_envelope(data)

class WarRecord(BaseModel):
    model_config = ConfigDict(extra='allow')
    state: str

class CapitalRaidRecord(BaseModel):
    model_config = ConfigDict(extra='allow')
    state: str
    startTime: str

class CapitalRaidListResponse(BaseModel):
    items: List[CapitalRaidRecord]

    @model_validator(mode='before')
    @classmethod
    def validate_envelope(cls, data: Any) -> Any:
        return normalize_envelope(data)

class LeagueGroupRecord(BaseModel):
    model_config = ConfigDict(extra='allow')
    state: str
    season: str

class WarLeagueWarRecord(BaseModel):
    model_config = ConfigDict(extra='allow')
    state: str
  1. La infraestructura está desplegada enteramente usando Terraform; este es el estándar para IaC.
  2. Hice la migración apoyándose en IA, utilizando un arnés que me permita revisar los planes de implementación con la notación EARS.

¿Cómo se ven las tablas resultado?

Son tablas con dos columnas: timestamp y payload. Cada tabla está particionada por timestamp (DAY).

BigQuery Table Partitions BigQuery Table Preview

Comentarios finales

Hay algunas variables que vale la pena mencionar antes de catalogar este proyecto de “grado empresarial”:

  1. Cloud NAT puede llegar a costar unos 32 USD mensuales en promedio. Esto, dependiendo de la envergadura del proyecto, puede ser ineficiente; por ello, este diseño asume que utilizaremos Cloud NAT no solamente para el lote (batch) diario de este ELT, sino también para otros procesos.
  2. Por la naturaleza de la ingesta, preferimos un enfoque pragmático: hay duplicación de datos en la capa bronce si existen reprocesamientos. Es decir, estamos utilizando el patrón APPEND-ONLY en lugar de generar claves de idempotencia. Esto se justifica por el bajo volumen de datos, pues podemos obtener el registro más reciente en la fase de transformación sin incurrir en costos significativos ni complejizar la infraestructura.

Referencias