Comprobaciones de calidad de datos

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Por qué importan los quality gates

 

  • Un pipeline sin checks es un pipeline en el que no puedes confiar
  • Un cambio upstream introduce nulos o valores negativos
  • Todos los paneles downstream muestran números erróneos
  • Los quality gates detectan problemas en el origen

Datos erróneos río abajo

Creación de canalizaciones de datos con Airflow

SQLColumnCheckOperator

from airflow.providers.common.sql.operators.sql import
    (SQLExecuteQueryOperator, SQLColumnCheckOperator)

check_columns = SQLColumnCheckOperator( task_id="check_columns", conn_id="duckdb_analytics", table="daily_summary", column_mapping={ "total_revenue": {"min": {"greater_than": 0}}, "total_orders": {"null_check": {"equal_to": 0}}, }, )
Creación de canalizaciones de datos con Airflow

Tipos de comprobaciones de columna

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: valida el valor mínimo de una columna
  • null_check: cuenta valores nulos (usa equal_to: 0 para no tener nulos)
  • distinct_check: cuenta valores distintos
  • Checks disponibles: null_check, unique_check, distinct_check, min, max
  • Comparadores disponibles: equal_to, greater_than, geq_to, less_than, leq_to
Creación de canalizaciones de datos con Airflow

SQLTableCheckOperator

from airflow.providers.common.sql.operators.sql import SQLTableCheckOperator

check_table = SQLTableCheckOperator( task_id="check_table", conn_id="duckdb_analytics", table="daily_summary", checks={ "row_count_check": { "check_statement": "COUNT(*) > 0", }, }, )
  • Valida propiedades de la tabla completa
  • check_statement: cualquier expresión SQL que evalúe a true o false
Creación de canalizaciones de datos con Airflow

Encadenar comprobaciones de calidad

@dag(schedule="@daily", template_searchpath="/path/to/sql")
def sales_pipeline():

    load = SQLExecuteQueryOperator(
        task_id="load_daily_sales",
        conn_id="duckdb_analytics",
        sql="aggregate_sales.sql",
    )

    check_columns = SQLColumnCheckOperator(...)
    check_table = SQLTableCheckOperator(...)

    load >> check_columns >> check_table
  • Carga primero y luego valida
  • Si falla algún check, el pipeline se detiene
  • Quienes consumen río abajo nunca ven datos erróneos
Creación de canalizaciones de datos con Airflow

Comprobaciones de columna vs. de tabla

SQLColumnCheckOperator

check_columns = SQLColumnCheckOperator(
    task_id="check_columns",
    conn_id="duckdb_analytics",
    table="daily_summary",
    column_mapping={
        "total_revenue": {"min": {"greater_than": 0}},
        "total_orders": {"null_check": {"equal_to": 0}},
    },
)
  • ¿Hay nulos?
  • ¿Los valores están dentro del rango?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • ¿La tabla tiene filas?
  • ¿El total está dentro de los límites esperados?
Creación de canalizaciones de datos con Airflow

Supervisar la calidad a escala con Astro

$$

  • Astro Observe de Astronomer ayuda a monitorizar la salud del pipeline
  • Seguimiento de SLA con alertas antes de los plazos
  • Lineage del pipeline entre DAGs y tablas
  • Supervisión de calidad de datos a escala
  • Análisis de causa raíz con IA

Lineage en Astro Observe

Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...