Orquestación de ETL con dbt: Una guía para directores técnicos (CTO)
En la empresa moderna, el almacén de datos (DW) se ha convertido en el sistema nervioso central, el cerebro estratégico que impulsa el análisis, la inteligencia artificial y la toma de decisiones. Sin embargo, para muchas organizaciones, este activo crucial está limitado por procesos de transformación frágiles, opacos y manuales. Los procesos ETL tradicionales son lentos, difíciles de mantener y carecen fundamentalmente de las mejores prácticas de ingeniería de software que exigimos de nuestro código de aplicación.
Aquí es donde dbt (la herramienta de construcción de datos) cambia fundamentalmente el paradigma. Al pasar a un modelo ELT (Extracción, Carga, Transformación), dbt introduce los principios de la ingeniería analítica: modularidad, control de versiones, pruebas y CI/CD, en la capa de transformación (T) que se encuentra dentro de tu almacén de datos.
Este artículo no es una descripción general de alto nivel. Es un plan técnico y práctico para directores técnicos (CTOs) e ingenieros senior sobre cómo implementar un flujo de trabajo dbt totalmente automatizado y para producción. Nos centraremos en construir una robusta infraestructura CI/CD y una estrategia de orquestación consciente del estado, para llevar su práctica de datos desde scripts improvisados hasta una fábrica de datos de alta velocidad y confiable.
Servicios de Ingeniería de Productos
Trabaje con nuestros gerentes de proyecto, ingenieros de software y testers de calidad internos para desarrollar su nuevo producto de software personalizado o para apoyar su flujo de trabajo actual, siguiendo metodologías Agile, DevOps y Lean.
El argumento de '¿Por qué dbt?', para el liderazgo técnico
Antes de sumergirnos en la automatización, debemos tener claro qué estamos automatizando. dbt es un compilador y un marco de trabajo para la lógica de datos. No extrae ni carga datos; gestiona magistralmente la transformación de los datos brutos que ya están en su almacén (por ejemplo, Snowflake, BigQuery, Redshift, Databricks).
Para un líder técnico, el valor de dbt es evidente:
- Es simplemente SQL (y Jinja): Aprovecha las habilidades existentes de SQL de tu equipo de análisis, envolviéndolas en un potente motor de plantillas (Jinja) para la modularidad y los principios DRY (Don't Repeat Yourself).
- Flujo de trabajo nativo de Git: Cada cambio en un modelo de datos es un commit. Cada nueva característica es una rama. Cada lanzamiento a producción es una fusión a
main. Esto proporciona auditabilidad, colaboración y capacidades de reversión. - Pruebas integradas: dbt proporciona un marco de primer nivel para las afirmaciones de calidad de datos. Puedes definir pruebas (por ejemplo,
not_null,unique,accepted_values, basadas en SQL personalizadas) en simples archivos YAML, que dbt ejecuta contra tu almacén. - Gestión de dependencias y trazabilidad: La función
ref()de dbt es su característica más crítica. En lugar de codificar nombres de tablas, haces referencia a otros modelos. dbt utiliza esto para inferir automáticamente un Gráfico Acíclico Dirigido (DAG) de dependencias. Sabe exactamente qué orden seguir para construir los modelos y puede generar un gráfico completo de trazabilidad de datos para todo tu almacén (documentación de dbt).
En resumen, dbt trata la lógica de transformación de tus datos como un producto de software, lo cual es una condición previa para la automatización significativa.
Implementar un flujo de trabajo CI/CD para dbt
La regla fundamental de una configuración de producción de dbt es:Ningún desarrollador, analista o ingeniero debe ejecutardbt run directamente contra el entorno de producción desde su portátil. Todos los cambios deben ser revisados y validados a través de una pipeline de CI (Integración Continua) automatizada.
El objetivo de CI es responder a una pregunta para cada solicitud de extracción (PR): "¿Este cambio propuesto causa algún problema?".
La estrategia más eficaz para esto es"Slim CI", que aprovecha las capacidades de estado de dbt. En lugar de reconstruir todo el almacén de datos en cada PR (lo cual es computacional y financieramente inviable), Slim CI identifica y ejecuta solo los modelos que han sido modificados y sus dependencias.
Proceso de CI paso a paso (Ejemplo de GitHub Actions)
Esta tarea se ejecuta en cada solicitud de extracción contra la rama principal . Crea los modelos modificados en un esquema temporal e independiente dentro de tu almacén de datos.
Requisito: Esta estrategia requiere que almacene su archivo manifest.json en una ubicación accesible (como S3, GCS o Azure Blob Storage) después de cada ejecución exitosa del entorno de producción. Este manifiesto es un mapa de su entorno de producción, que dbt utiliza para calcular la diferencia "de estado".
Aquí tienes un archivo detallado .github/workflows/dbt-ci.yml:
name: dbt CI (Slim)
on:
pull_request:
branches:
- main
env:
DBT_PROFILES_DIR: . # Tell dbt where to find profiles.yml
# Store all warehouse credentials as GitHub Secrets
DBT_USER: ${{ secrets.DBT_USER_CI }}
DBT_PASSWORD: ${{ secrets.DBT_PASSWORD_CI }}
DBT_ACCOUNT: ${{ secrets.DBT_ACCOUNT }}
DBT_ROLE: ${{ secrets.DBT_ROLE_CI }}
DBT_WAREHOUSE: ${{ secrets.DBT_WAREHOUSE_CI }}
DBT_DATABASE: ${{ secrets.DBT_DATABASE_PROD }} # Use prod database
jobs:
run_dbt_slim_ci:
runs-on: ubuntu-latest
steps:
- name: Checkout code
uses: actions/checkout@v3
- name: Download production manifest
# This step downloads the manifest from your artifact storage.
# This example uses a GCS bucket.
uses: 'google-github-actions/auth@v1'
with:
credentials_json: '${{ secrets.GCS_SA_KEY }}'
env:
GCP_BUCKET: "your-dbt-artifacts-bucket"
- name: 'Download Manifest'
uses: 'google-github-actions/storage-transfer@v1'
with:
source: 'gs://${{ env.GCP_BUCKET }}/prod/manifest.json'
destination: './prod-manifest'
# Continue if no manifest exists (e.g., first run)
continue-on-error: true
- name: Set up Python
uses: actions/setup-python@v4
with:
python-version: '3.10'
- name: Install dbt and adapter
run: |
pip install dbt-snowflake==1.8.0 # Use your specific adapter and version
dbt --version
- name: Configure dbt profile (for CI)
# Dynamically create profiles.yml from secrets
# This builds into a unique, temporary schema for this PR
run: |
echo "default:" > profiles.yml
echo " target: ci" >> profiles.yml
echo " outputs:" >> profiles.yml
echo " ci:" >> profiles.yml
echo " type: snowflake"
echo " account: $DBT_ACCOUNT"
echo " user: $DBT_USER"
echo " password: $DBT_PASSWORD"
echo " role: $DBT_ROLE"
echo " warehouse: $DBT_WAREHOUSE"
echo " database: $DBT_DATABASE"
# CRITICAL: Isolate the build
echo " schema: dbt_ci_pr_${{ github.event.pull_request.number }}"
echo " threads: 4"
echo " client_session_keep_alive: False"
- name: Install dbt dependencies
run: dbt deps
- name: Run and test only modified models (Slim CI)
# 'state:modified+' selector finds all modified models AND
# all downstream models that depend on them.
# --state flag points to the downloaded production manifest.
run: |
dbt build --select state:modified+ --state ./prod-manifest
# If no manifest was found, fall back to a full build (slower)
continue-on-error: false
# TODO: Add a fallback step here if the first run fails
# e.g., if: failure()
# run: dbt build --select state:modified+
- name: Clean up temporary schema
# This runs whether the build succeeds or fails
if: always()
run: |
dbt run-operation drop_pr_schema --args '{schema_name: "dbt_ci_pr_${{ github.event.pull_request.number }}"}'
Al fusionar con main
Cuando una PR se fusiona, la rama principal representa tu artefacto desplegable. Una tarea en la rama principal no debería desplegarse en producción. En cambio, debería:
- Realice una última
construcción dbten un entorno de Staging dedicado (una copia sin copia del entorno de producción es ideal). - Si tiene éxito, genere y guarde el
manifest.json. Este artefacto ahora representa el estado deproducción principal, y se utilizará tanto por su orquestador de producción como por futuras ejecuciones de CI.
Servicios de Ingeniería de Productos
Trabaje con nuestros gestores de proyectos, ingenieros de software y probadores de calidad internos para desarrollar su nuevo producto de software personalizado o para apoyar su flujo de trabajo actual, siguiendo metodologías Agile, DevOps y Lean.
Orquestación de la Producción (La parte de 'CD')
CI construye y prueba código. Un Orchestrador construye y prueba datos según un horario. Su herramienta de orquestación (por ejemplo, Airflow, Dagster, Prefect, dbt Cloud) es responsable de ejecutar su proyecto dbt en producción.
Principios clave:
- Separación de preocupaciones: El DAG (Grafo de Dependencias) del orquestador debe ser simple. Obtiene la última
rama principal, obtiene el últimomanifest.json de producción, y ejecuta comandos dbt. - Ejecuciones conscientes del estado: Las ejecuciones de producción deben también utilizar el
manifest.jsonelemento. Esto permite que dbt pueda retomar el proceso de forma inteligente en caso de un fallo parcial.dbt run --state ./ --manifest ./manifest.json), ahorrando un tiempo de cálculo significativo. - Puertas de control de calidad de datos: El pipeline debe tener puertas de control de calidad de datos.
Ejemplo: Orquestación de la producción con Airflow
Este ejemplo asume que estás utilizando Airflow con el BashOperator o DockerOperator para ejecutar comandos de dbt.
from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.utils.dates import days_ago
# Define default arguments
default_args = {
'owner': 'data_engineering',
'depends_on_past': False,
'start_date': days_ago(1),
'retries': 1,
}
with DAG(
dag_id="dbt_production_pipeline",
default_args=default_args,
schedule_interval="0 6 * * *", # 6 AM daily
catchup=False,
tags=["dbt", "production"],
) as dag:
# 1. Pull the latest dbt project code
git_pull = BashOperator(
task_id="git_pull",
bash_command="cd /path/to/dbt-project && git pull origin main",
)
# 2. Download the latest production manifest from storage
# (This step is analogous to the CI step)
download_manifest = BashOperator(
task_id="download_manifest",
bash_command="gcloud storage cp gs://your-dbt-artifacts-bucket/prod/manifest.json /path/to/dbt-project/prod-artifacts/",
)
# 3. Check that raw source data has arrived
check_source_freshness = BashOperator(
task_id="check_source_freshness",
bash_command="cd /path/to/dbt-project && dbt source freshness",
# If this fails, the whole pipeline fails and alerts
)
# 4. Run dbt models (state-aware)
run_dbt_models = BashOperator(
task_id="run_dbt_models",
bash_command=(
"cd /path/to/dbt-project && "
"dbt run --state ./prod-artifacts --manifest ./prod-artifacts/manifest.json"
),
)
# 5. Test dbt models
test_dbt_models = BashOperator(
task_id="test_dbt_models",
bash_command="cd /path/to/dbt-project && dbt test",
)
# 6. Generate docs and the NEW manifest
generate_artifacts = BashOperator(
task_id="generate_artifacts",
bash_command="cd /path/to/dbt-project && dbt docs generate",
)
# 7. Upload the new manifest for the next CI/prod run
upload_new_manifest = BashOperator(
task_id="upload_new_manifest",
bash_command=(
"cd /path/to/dbt-project && "
"gcloud storage cp ./target/manifest.json gs://your-dbt-artifacts-bucket/prod/manifest.json"
),
)
# Define the DAG's execution order
(
git_pull
>> download_manifest
>> check_source_freshness
>> run_dbt_models
>> test_dbt_models
>> generate_artifacts
>> upload_new_manifest
)
Nota: Para una configuración realmente robusta, utilizarías el KubernetesPodOperator o el DockerOperator para ejecutar cada paso en un contenedor aislado que contenga el código y las credenciales de tu proyecto dbt.
Estrategia avanzada: Implementación sin interrupciones
Para datos de misión crítica, no puede permitirse que un proceso dbt falle en medio, dejando sus tablas de producción en un estado defectuoso. La solución es una estrategia de implementación "azul/verde".
- Tu orquestador (por ejemplo, Airflow) no escribe en la
PRODbase de datos. - En cambio, utiliza
dbt clone(o la copia "ZERO-COPY CLONE" de Snowflake) para crear una copia perfecta, que solo contiene metadatos:CLONING SIN COPIA) para crear una copia perfecta, que solo contenga metadatos:PROD_CLONE. - Toda la
ejecución de dbty laspruebas de dbtse ejecutan contra estaPROD_CLONEbase de datos. - Si todos los pasos tienen éxito, el "desplazamiento" es un único comando atómico que intercambia
Intercambiocomando que intercambiaPRODyPROD_CLONE. - Tus herramientas de BI y tus usuarios, que antes apuntaban a
PRODUCTOS, ahora apuntan sin problemas a los nuevos datos, que han sido probados completamente. La antigua<s30>PRODbase de datos (ahora llamadaPROD_CLONE) puede ser eliminada.
Este enfoque garantiza que tus consumidores de datos nunca vean una tabla incompleta y que todas las comprobaciones de calidad de los datos hayan sido completadas antes de que se sirvan los datos.anteslos datos están disponibles.
Servicios de Ingeniería de Productos
Trabaje con nuestros gestores de proyectos, ingenieros de software y probadores de calidad internos para desarrollar su nuevo producto de software personalizado o para apoyar su flujo de trabajo actual, siguiendo metodologías Agile, DevOps y Lean.
Finalmente
Automatizar tu almacén de datos con dbt no se trata solo de programar una ejecución de dbt, sino de adoptar una metodología completa de ingeniería analítica que considera las transformaciones de tus datos como un servicio de software de producción.
Al implementar una línea de desarrollo CI eficiente,
Al crear unaplataforma de orquestación con conocimiento del estado,se crea una fábrica de datos resiliente, eficiente y monitorizable que puede recuperarse de fallos y garantizar la calidad de los datos.
Esta capa de automatización es la infraestructura fundamental que libera todo el potencial de tu equipo de datos. Permite que se centren en tareas de alto valor, como el análisis y la creación de productos de datos, en lugar de tener que lidiar con problemas y realizar validaciones manuales.
Preguntas frecuentes
¿Qué es dbt y por qué se utiliza para la transformación de datos?
dbt (data build tool) es un marco de trabajo para la transformación que opera sobre los datos ya cargados en un almacén de datos. Se utiliza para aplicar las mejores prácticas de ingeniería de software a los códigos analíticos. Su valor proviene de permitir el control de versiones (flujo de trabajo nativo de Git), modularidad (utilizando SQL y Jinja), pruebas automatizadas para la calidad de los datos, y gestión automática de dependencias, lo que permite construir modelos de datos en el orden correcto.
¿Cómo funciona una "pipeline Slim CI" para dbt?
Una "pipeline Slim CI" (Integración Continua) automatiza las pruebas de los cambios en el código de dbt en una solicitud de extracción. En lugar de reconstruir todo el almacén de datos, utiliza un enfoque "con estado". Al comparar los cambios de código con un archivo de manifiesto de producción, identifica e ejecuta pruebas de forma inteligente solo en los modelos que han sido modificados y cualquier modelo dependiente. Esto hace que el proceso de validación sea significativamente más rápido y rentable.
¿Qué es la orquestación consciente del estado en un entorno de producción dbt?
La orquestación consciente del estado se refiere a ejecutar trabajos de dbt en un entorno de producción (por ejemplo, según una programación diaria) utilizando herramientas como Airflow o dbt Cloud. Este proceso también utiliza el archivo de manifiesto de producción para comprender el estado actual del almacén de datos. Si una ejecución en producción falla a mitad de camino, esta conciencia del estado permite que dbt continúe de forma inteligente desde el punto de fallo, en lugar de comenzar la ejecución completa desde cero, lo que ahorra un tiempo y recursos informáticos significativos.