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.
🚀 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.
🛠️ 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'
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
🔄 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
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.
📈 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.
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 🔄
🔮 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, ProfileConfigwith 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.
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
- Optimización de Consultas SQL en Entornos Big Data: El Poder de Apache Calciteadvanced20 min
- Ingeniería de Características en Big Data: Potenciando Modelos con Feature Engineering Distribuidointermediate20 min
- Data Governance en Microservicios: Un Enfoque Descentralizado para Big Dataadvanced18 min
- Optimización de Consultas en Data Lakes: Estrategias con Apache Parquet y Presto/Trinointermediate10 min
- Gobernanza de Datos en Big Data: Clave para la Confianza y la Eficienciaintermediate15 min
Comentarios (0)
Aún no hay comentarios. ¡Sé el primero!