tutoriales.com

Ingesta y Transformación de Datos Estructurados para Big Data: Un Enfoque Declarativo con Apache Airflow y dbt

Este tutorial te guiará a través de la creación de un pipeline de datos estructurados para Big Data, utilizando Apache Airflow para la orquestación de tareas y dbt (data build tool) para la transformación declarativa de datos. Aprenderás a integrar estas herramientas para gestionar tus flujos de trabajo de datos de manera eficiente y escalable.

Intermedio20 min de lectura19 views
Reportar error

🚀 Introducción al Procesamiento de Datos Estructurados en Big Data

En el mundo del Big Data, la ingesta y transformación de datos son procesos fundamentales. A menudo, nos enfrentamos a volúmenes masivos de información que necesitan ser cargados, limpiados, transformados y preparados para el análisis o para alimentar modelos de Machine Learning. Tradicionalmente, esto implica escribir una gran cantidad de código ETL (Extract, Transform, Load) complejo y mantenerlo, lo que puede ser un desafío.

Este tutorial se centra en un enfoque moderno y declarativo para construir pipelines de datos estructurados utilizando dos herramientas poderosas: Apache Airflow para la orquestación y dbt (data build tool) para la transformación de datos. Juntas, estas herramientas permiten crear flujos de trabajo de datos robustos, mantenibles y fáciles de escalar, aprovechando la potencia de SQL para las transformaciones.

¿Por qué Airflow y dbt Juntos? 🤔

  • Airflow: Es una plataforma para programar, ejecutar y monitorear flujos de trabajo de manera programática. Permite definir pipelines como DAGs (Directed Acyclic Graphs) en Python, lo que brinda flexibilidad y visibilidad sobre el estado de las tareas.
  • dbt: Se enfoca exclusivamente en la fase de 'Transformación' de ETL. Permite a los ingenieros de datos transformar datos en sus data warehouses (o data lakes con motores SQL como Spark SQL, Presto, etc.) utilizando solo SQL, pero con las mejores prácticas del desarrollo de software (control de versiones, pruebas, documentación, modularidad).

La combinación de ambos permite a Airflow ser el 'cerebro' que orquesta cuándo y cómo se ejecutan las transformaciones de dbt, mientras que dbt se encarga de la lógica de transformación en sí misma, haciendo que el proceso sea más eficiente y legible.

💡 Consejo: Piensa en Airflow como el director de orquesta que le dice a dbt cuándo tocar, y dbt como el músico experto que ejecuta la pieza SQL de forma impecable.

🛠️ Entorno de Trabajo: Configuración Inicial

Para seguir este tutorial, necesitaremos configurar un entorno básico que incluya Docker (para Airflow y una base de datos de ejemplo) y dbt. Asumiremos que ya tienes Docker y Docker Compose instalados en tu sistema.

1. Preparando la Base de Datos y Airflow con Docker Compose

Crearemos un archivo docker-compose.yml para levantar un entorno mínimo. Usaremos PostgreSQL como nuestra base de datos de ejemplo, que simulará nuestro data warehouse o data lake con capacidad SQL.

version: '3.8'
x-airflow-common:
  &airflow-common
  image: apache/airflow:2.8.1
  env_file:
    - .env
  volumes:
    - ./dags:/opt/airflow/dags
    - ./salsa:/opt/airflow/salsa
    - ./logs:/opt/airflow/logs
    - ./plugins:/opt/airflow/plugins
  ports:
    - "8080:8080"
  extra_hosts:
    - "host.docker.internal:host-gateway"
  # user: "${AIRFLOW_UID}:0"
  command: bash -c "airflow db migrate && airflow webserver"

services:
  postgres:
    image: postgres:13
    environment:
      POSTGRES_USER: ${POSTGRES_USER}
      POSTGRES_PASSWORD: ${POSTGRES_PASSWORD}
      POSTGRES_DB: ${POSTGRES_DB}
    ports:
      - "5432:5432"
    volumes:
      - pgdata:/var/lib/postgresql/data

  airflow-init:
    <<: *airflow-common
    container_name: airflow_init
    entrypoint: /bin/bash -c "mkdir -p /opt/airflow/logs && mkdir -p /opt/airflow/dags && mkdir -p /opt/airflow/salsa && airflow db migrate && airflow users create --username admin --firstname admin --lastname admin --role Admin --email admin@example.com --password admin && airflow webserver"

  airflow-scheduler:
    <<: *airflow-common
    container_name: airflow_scheduler
    command: airflow scheduler

  airflow-webserver:
    <<: *airflow-common
    container_name: airflow_webserver
    command: airflow webserver

volumes:
  pgdata:

Crea un archivo .env en el mismo directorio con las siguientes variables de entorno:

AIRFLOW_UID=50000 # O tu UID local, ej. id -u
POSTGRES_USER=airflow
POSTGRES_PASSWORD=airflow
POSTGRES_DB=airflow
POSTGRES_HOST=postgres
POSTGRES_PORT=5432

Levanta el entorno con docker-compose up -d.

2. Instalación de dbt y sus Adaptadores

Instalaremos dbt-core y el adaptador para PostgreSQL. Si estás usando otra base de datos (Snowflake, BigQuery, etc.), necesitarás el adaptador correspondiente.

# Crea un entorno virtual (opcional pero recomendado)
python -m venv dbt_env
source dbt_env/bin/activate

pip install dbt-core dbt-postgres

3. Configuración de dbt (profiles.yml) ⚙️

dbt necesita saber cómo conectarse a tu base de datos. Crea un directorio ~/.dbt/ (si no existe) y dentro de él, el archivo profiles.yml:

# ~/.dbt/profiles.yml
my_project:
  target: dev
  outputs:
    dev:
      type: postgres
      host: host.docker.internal # O la IP de tu VM si no usas Docker Desktop
      port: 5432
      user: airflow
      password: airflow
      dbname: airflow
      schema: public
      threads: 1

Importante: host.docker.internal permite que dbt (ejecutado en tu host local) se conecte al servicio postgres dentro de Docker. Si estás en Linux, puede que necesites usar la IP de tu docker0 bridge (ej. 172.17.0.1).


📝 Definiendo Nuestro Proyecto dbt

Un proyecto dbt es un directorio con un archivo dbt_project.yml y subdirectorios para modelos, tests, seeds, etc.

1. Creando el Proyecto dbt

En la raíz de tu proyecto (donde está docker-compose.yml), crea un nuevo directorio llamado salsa (como referencia a la carpeta que montamos en Airflow) y dentro, ejecuta dbt init my_first_dbt_project.

mkdir salsa
cd salsa
dbt init my_first_dbt_project

Esto creará una estructura como esta:

salsa/
├── my_first_dbt_project/
│   ├── models/
│   │   ├── example/
│   │   │   ├── my_first_dbt_model.sql
│   │   │   └── my_second_dbt_model.sql
│   ├── dbt_project.yml
│   ├── seeds/
│   ├── macros/
│   ├── tests/
│   └── ...

Modifica my_first_dbt_project/dbt_project.yml para asegurarte de que el nombre del perfil coincide con el de profiles.yml:

# salsa/my_first_dbt_project/dbt_project.yml

name: 'my_first_dbt_project'
version: '1.0.0'
config-version: 2

profile: 'my_project' # ¡Este debe coincidir!

model-paths: ["models"]
analysis-paths: ["analyses"]
seed-paths: ["seeds"]
test-paths: ["tests"]
macro-paths: ["macros"]
snapshot-paths: ["snapshots"]

target-path: "target"  # directory which will store compiled SQL files
clean-targets:
  - "target"
  - "dbt_packages"
  - "dbt_modules"

# Configure models
models:
  my_first_dbt_project:
    # Config all models in this project to be materialized as views.
    # See https://docs.getdbt.com/docs/build/materializations
    materialized: view

    example:
      materialized: table # Sobrescribimos para las tablas de ejemplo

2. Ingesta de Datos de Ejemplo (Seeds) 🌱

dbt puede cargar archivos CSV estáticos a tu base de datos como tablas. Esto es útil para datos de referencia o pequeños datasets de prueba. Crearemos un archivo users.csv en salsa/my_first_dbt_project/seeds/:

id,name,email,registration_date,status
1,Alice,alice@example.com,2023-01-01,active
2,Bob,bob@example.com,2023-01-05,inactive
3,Charlie,charlie@example.com,2023-01-10,active
4,David,david@example.com,2023-01-12,active
5,Eve,eve@example.com,2023-01-15,inactive

Ahora, carga los datos a la base de datos ejecutando el comando dbt desde el directorio del proyecto dbt:

cd salsa/my_first_dbt_project
dbt seed

Verifica en tu base de datos (usando psql o DBeaver) que la tabla users ha sido creada en el esquema public.

3. Modelado de Datos con dbt (SQL) 🏗️

Ahora, crearemos un modelo dbt para transformar nuestros datos. El objetivo es crear una tabla de usuarios activos y otra que contenga métricas agregadas.

Modelo 1: active_users.sql

Crea salsa/my_first_dbt_project/models/active_users.sql:

-- models/active_users.sql

SELECT
  id,
  name,
  email,
  registration_date
FROM {{ source('my_first_dbt_project', 'users') }}
WHERE status = 'active'
📌 Nota: Usamos `{{ source(...) }}` para referenciar la tabla 'users' que cargamos como seed. dbt generará el SQL de forma segura. También, puedes definir `sources.yml` para tipar tus fuentes de datos.

Para que dbt sepa dónde encontrar la fuente users, crea salsa/my_first_dbt_project/models/sources.yml:

version: 2

sources:
  - name: my_first_dbt_project
    schema: public # El esquema donde se encuentra tu tabla 'users'
    tables:
      - name: users

Modelo 2: user_registration_summary.sql

Crea salsa/my_first_dbt_project/models/user_registration_summary.sql:

-- models/user_registration_summary.sql

SELECT
  registration_date,
  COUNT(id) AS total_registrations,
  COUNT(CASE WHEN status = 'active' THEN id END) AS active_registrations
FROM {{ source('my_first_dbt_project', 'users') }}
GROUP BY registration_date
ORDER BY registration_date

Ahora, ejecuta los modelos con dbt:

cd salsa/my_first_dbt_project
dbt run

Esto creará las vistas o tablas (dependiendo de tu configuración materialized) active_users y user_registration_summary en tu base de datos.

4. Prueba de Modelos con dbt (Tests) ✅

dbt permite añadir pruebas a tus modelos para asegurar la calidad de los datos. Crearemos algunas pruebas simples.

Crea salsa/my_first_dbt_project/models/schema.yml (o añade a un schema.yml existente):

version: 2

models:
  - name: active_users
    description: "Usuarios registrados y activos en la plataforma."
    columns:
      - name: id
        description: "Identificador único del usuario."
        tests:
          - unique
          - not_null
      - name: email
        tests:
          - unique

  - name: user_registration_summary
    description: "Resumen de registros de usuarios por fecha."
    columns:
      - name: registration_date
        tests:
          - unique
          - not_null

Ejecuta las pruebas:

cd salsa/my_first_dbt_project
dbt test
🔥 Importante: Las pruebas de dbt son cruciales para la fiabilidad de tu pipeline. No subestimes su importancia en un entorno de Big Data.

🔄 Orquestación con Apache Airflow

Ahora que tenemos nuestro proyecto dbt configurado, es el momento de orquestarlo con Airflow. Crearemos un DAG que ejecute las tareas de dbt en el orden correcto.

1. Preparando la Integración Airflow-dbt

Para que Airflow pueda ejecutar dbt, necesitamos asegurarnos de que dbt esté disponible en el entorno de Airflow. Una forma sencilla para este tutorial es instalar dbt-postgres directamente en la imagen de Airflow (en producción, preferirías crear una imagen Docker personalizada).

Para este ejercicio, modificaremos el docker-compose.yml para instalar dbt al iniciar el contenedor webserver y scheduler, o lo haremos manualmente en el contenedor en ejecución.

Opción 1: Manual (para pruebas rápidas)

Accede al contenedor airflow-webserver y airflow-scheduler y ejecuta:

docker exec -it airflow_webserver bash
pip install dbt-core dbt-postgres
# Salir del contenedor
exit

docker exec -it airflow_scheduler bash
pip install dbt-core dbt-postgres
# Salir del contenedor
exit

Opción 2: Añadir al Dockerfile (recomendado para producción)

Crearías un Dockerfile personalizado basado en la imagen de Airflow y añadirías la instalación de dbt. Por simplicidad en este tutorial, usaremos la opción 1 o la opción de instalarlo a través del DAG si usas un operador BashOperator con pip install.

2. Creando el DAG de Airflow para dbt

Crearemos un archivo DAG llamado dbt_pipeline_dag.py en el directorio dags/ que montamos con Docker. Asegúrate de que este directorio existe: mkdir dags.

# dags/dbt_pipeline_dag.py

from __future__ import annotations

import pendulum

from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator

# Ruta al directorio de nuestro proyecto dbt dentro del contenedor Airflow
# Esto asume que el volumen se monta en /opt/airflow/salsa
DBT_PROJECT_DIR = '/opt/airflow/salsa/my_first_dbt_project'
DBT_PROFILES_DIR = '/opt/airflow/salsa/my_first_dbt_project' # dbt buscará profiles.yml aquí por defecto

with DAG(
    dag_id="dbt_etl_pipeline",
    start_date=pendulum.datetime(2023, 1, 1, tz="UTC"),
    schedule=None,
    catchup=False,
    tags=["dbt", "etl", "data_transformation"],
    description="Un DAG para ejecutar un pipeline ETL usando dbt para transformaciones.",
) as dag:

    # Tarea 1: Verificar la conectividad de dbt
    dbt_debug = BashOperator(
        task_id="dbt_debug",
        bash_command=f"dbt debug --project-dir {DBT_PROJECT_DIR} --profiles-dir {DBT_PROFILES_DIR}",
    )

    # Tarea 2: Cargar datos 'seeds'
    dbt_seed = BashOperator(
        task_id="dbt_seed_data",
        bash_command=f"dbt seed --project-dir {DBT_PROJECT_DIR} --profiles-dir {DBT_PROFILES_DIR}",
    )

    # Tarea 3: Ejecutar los modelos dbt (transformaciones)
    dbt_run = BashOperator(
        task_id="dbt_run_models",
        bash_command=f"dbt run --project-dir {DBT_PROJECT_DIR} --profiles-dir {DBT_PROFILES_DIR}",
    )

    # Tarea 4: Ejecutar las pruebas dbt
    dbt_test = BashOperator(
        task_id="dbt_test_models",
        bash_command=f"dbt test --project-dir {DBT_PROJECT_DIR} --profiles-dir {DBT_PROFILES_DIR}",
    )

    # Definir el orden de las tareas
    dbt_debug >> dbt_seed >> dbt_run >> dbt_test

⚠️ Advertencia: Para la instalación de dbt dentro del contenedor, el `BashOperator` asume que `dbt` ya está en el PATH. Si no es así, deberías especificar la ruta completa o asegurarte de que tu imagen de Airflow personalizada lo incluye. En entornos de producción, se recomienda usar `dbt-airflow-plugin` o `Astronomer Cosmos` para una integración más robusta y segura.

3. Visualizando el DAG en la UI de Airflow

Una vez que el archivo dbt_pipeline_dag.py esté en el directorio dags/, Airflow debería detectarlo automáticamente. Abre tu navegador y ve a http://localhost:8080 (usuario: admin, contraseña: admin).

Deberías ver el DAG dbt_etl_pipeline en la lista. Actívalo y ejecútalo manualmente para observar el flujo de trabajo. Puedes monitorear el progreso y los logs de cada tarea en la interfaz de usuario de Airflow.

dbt_debug dbt_seed_data dbt_run_models dbt_test_models

📈 Beneficios y Mejores Prácticas

La combinación de Airflow y dbt ofrece varios beneficios clave para los ingenieros de datos:

1. Orquestación y Observabilidad Centralizadas 🔭

Airflow proporciona una vista centralizada de todos tus pipelines de datos. Puedes ver el estado de cada tarea, los logs, los tiempos de ejecución y configurar alertas en caso de fallos. Esto es crucial para la gestión de errores y el monitoreo de SLAs (Service Level Agreements) en entornos de Big Data.

💡 Consejo: Utiliza la funcionalidad de retries de Airflow y configura notificaciones para fallos críticos.

2. Desarrollo Declarativo y Reutilizable con SQL ✨

dbt permite a los equipos de datos aplicar principios de ingeniería de software a sus transformaciones SQL:

  • Control de Versiones: Los modelos dbt son archivos SQL planos que pueden ser versionados en Git.
  • Modularidad: Puedes construir modelos sobre otros modelos, creando un gráfico de dependencias claro.
  • Pruebas: Define pruebas de calidad de datos directamente en tu código.
  • Documentación: Documenta tus modelos, columnas y métricas junto al código.
  • Materializaciones: Controla cómo se persisten tus modelos (vistas, tablas, incrementales).
-- Ejemplo de modelo incremental en dbt
{{ config(
    materialized='incremental',
    unique_key='id'
) }}

SELECT *
FROM {{ source('my_first_dbt_project', 'events') }}

{% if is_incremental() %}
  -- Esto se ejecuta solo en cargas incrementales
  WHERE event_timestamp > (SELECT MAX(event_timestamp) FROM {{ this }})
{% endif %}

3. Escalabilidad y Mantenibilidad a Largo Plazo 🚀

Al separar la orquestación de la lógica de transformación, tu pipeline se vuelve más escalable. Airflow puede manejar miles de tareas, mientras que dbt puede procesar petabytes de datos aprovechando el poder de tu data warehouse/lake. La mantenibilidad mejora significativamente porque la lógica de negocio está encapsulada en SQL, fácil de entender y probar.

4. Flujo de Trabajo Típico 🔄

Paso 1: Ingesta Bruta - Datos cargados en el data lake/warehouse (ej. con Kafka, Flink, Nifi, Airbyte, Fivetran).
Paso 2: Staging con dbt - Modelos dbt para limpiar y estandarizar datos brutos en tablas de staging.
Paso 3: Transformación Central con dbt - Modelos dbt para construir tablas intermedias (ej. dimensiones, hechos).
Paso 4: Consumo con dbt - Modelos dbt finales para tablas de reporting o para alimentar ML.
Paso 5: Orquestación Airflow - Airflow programa y ejecuta todos los comandos `dbt run`, `dbt test`, etc., en secuencia.

🔮 Más allá de lo Básico: Consideraciones Avanzadas

Una vez que domines lo básico, hay varias áreas para explorar y mejorar tu pipeline.

1. Integración con Operadores dbt Avanzados

Para una integración más estrecha, considera el uso de:

  • dbt-airflow-factory: Una librería que genera dinámicamente tareas de Airflow para cada modelo, test, etc., de tu proyecto dbt. Esto crea un DAG que refleja el grafo de dependencias de dbt.
  • Astronomer Cosmos: Un operador de Airflow que te permite ejecutar proyectos dbt directamente desde Airflow, generando automáticamente un DAG con las dependencias de dbt. Simplifica mucho la orquestación.
Ejemplo de uso de Cosmos ```python # Ejemplo conceptual con Cosmos from airflow.models.dag import DAG from airflow.utils.dates import days_ago from cosmos import DbtDag, ProjectConfig, ProfileConfig

with DbtDag( project_config=ProjectConfig( dbt_project_path="/opt/airflow/salsa/my_first_dbt_project", ), profile_config=ProfileConfig( profile_name="my_project", target_name="dev", conn_id="postgres_default" # Necesitarías configurar una conexión DB en Airflow UI ), schedule_interval="@daily", start_date=days_ago(1), catchup=False, dag_id="dbt_cosmos_pipeline", tags=["dbt", "cosmos"], ) as dbt_dag: pass

</details>

### 2. Uso de Variables y Entornos dbt 🏷️

dbt permite usar variables para hacer tus modelos más dinámicos, por ejemplo, para filtrar datos por fecha o para configurar la materialización de forma condicional. También puedes gestionar diferentes entornos (desarrollo, staging, producción) con perfiles dbt.

<div class="progress-bar"><div class="progress-fill" style="width: 80%; background: #28A745;">Control de Entornos: 80%</div></div>

### 3. Documentación y Generación de Linaje 📖

dbt genera automáticamente documentación HTML para tu proyecto, incluyendo el linaje (graph) de tus modelos. Airflow también puede visualizar el DAG. Juntas, estas herramientas ofrecen una comprensión profunda de cómo fluyen los datos en tu sistema.

users FUENTE active_users MODELO user_registration_summary MODELO Linaje de Datos

4. Control de Acceso y Seguridad

En entornos de producción, es crucial gestionar el acceso a Airflow y a la base de datos subyacente. Utiliza las características de RBAC (Role-Based Access Control) de Airflow y las mejores prácticas de seguridad de tu data warehouse/lake. Las credenciales de dbt no deben estar en profiles.yml en producción, sino gestionadas por un sistema de secretos (ej. HashiCorp Vault, AWS Secrets Manager) y accedidas por Airflow.

Avanzado


🔚 Conclusión

La combinación de Apache Airflow y dbt ofrece una solución potente y moderna para la ingesta y transformación de datos estructurados en el ámbito del Big Data. Airflow proporciona la orquestación robusta y la observabilidad necesarias para gestionar pipelines complejos, mientras que dbt permite a los equipos de datos construir transformaciones declarativas, testeables y mantenibles utilizando la ubicuidad de SQL.

Al adoptar estas herramientas y seguir las mejores prácticas, puedes construir un sistema de procesamiento de datos que no solo sea eficiente y escalable, sino también fácil de entender, auditar y evolucionar a medida que tus necesidades de datos crecen.

¡Esperamos que este tutorial te haya proporcionado una base sólida para comenzar tu viaje con Airflow y dbt!

Tutoriales relacionados

Comentarios (0)

Aún no hay comentarios. ¡Sé el primero!