Python 数据管道与 ETL 工程化:从 Airflow DAG 到增量处理的架构实践

一、脚本式 ETL 的维护困境:当数据管道变成"意大利面条"

数据管道的初始形态通常是一组 Python 脚本——从数据库导出、清洗转换、写入数据仓库。脚本之间通过 cron 定时调度,依赖关系靠执行顺序隐式保证。当管道数量增长到数十个,问题开始爆发:上游脚本失败但下游仍执行,产生脏数据;脚本的输入输出没有明确契约,格式变更导致级联失败;重跑历史数据需要手动修改时间参数,极易出错。

更严重的是缺乏可观测性——管道失败后没有告警,数据延迟数小时才被发现。排查问题时需要逐个检查脚本日志,耗时且低效。

二、ETL 工程化架构:从 DAG 调度到增量处理

工程化 ETL 的核心思路是将"脚本"转化为"可调度、可监控、可重跑"的数据管道。Apache Airflow 是最广泛使用的 DAG 调度器,配合增量处理策略和幂等写入,构建可靠的数据管道。

flowchart TD
    A[数据源] --> B[抽取层<br/>CDC / 全量快照]
    B --> C[暂存层<br/>Raw Zone]
    C --> D[转换层<br/>Airflow DAG]
    D --> E[清洗与标准化]
    E --> F[业务逻辑计算]
    F --> G[聚合与宽表]
    G --> H[服务层<br/>Serving Zone]
    H --> I[数据仓库 / 特征库]

    J[Airflow Scheduler] --> K[DAG 执行引擎]
    K --> L[任务状态追踪]
    L --> M[失败重试与告警]
    M --> N[Slack / PagerDuty]

增量处理是工程化 ETL 的关键策略——每次只处理新增或变更的数据,而非全量重算。增量处理依赖数据源提供变更标识(如更新时间戳、CDC 事件流),并在处理过程中维护水位线(Watermark)。

三、生产级代码实现:Airflow DAG、增量处理与幂等写入

3.1 Airflow DAG 定义

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from datetime import datetime, timedelta

default_args = {
    'owner': 'data-team',
    'retries': 3,
    'retry_delay': timedelta(minutes=5),
    'email_on_failure': True,
    'email': ['data-alerts@company.com'],
}

with DAG(
    dag_id='user_behavior_daily',
    default_args=default_args,
    description='用户行为数据日度 ETL 管道',
    schedule_interval='0 2 * * *',  # 每天凌晨 2 点
    start_date=datetime(2026, 1, 1),
    catchup=False,  # 不回填历史
    tags=['etl', 'user-behavior'],
) as dag:

    extract_task = PythonOperator(
        task_id='extract_user_events',
        python_callable=extract_user_events,
        op_kwargs={
            'date': '{{ ds }}',  # Airflow 模板变量,逻辑日期
        },
    )

    transform_task = PythonOperator(
        task_id='transform_and_clean',
        python_callable=transform_and_clean,
        op_kwargs={'date': '{{ ds }}'},
    )

    load_task = SQLExecuteQueryOperator(
        task_id='load_to_warehouse',
        conn_id='clickhouse_default',
        sql="""
            INSERT INTO user_behavior_daily
            SELECT * FROM staging.user_behavior_{{ ds_nodash }}
        """,
    )

    quality_check = PythonOperator(
        task_id='data_quality_check',
        python_callable=run_quality_check,
        op_kwargs={'date': '{{ ds }}'},
    )

    # 定义依赖关系
    extract_task >> transform_task >> load_task >> quality_check

3.2 增量抽取与水位线管理

import pandas as pd
from sqlalchemy import create_engine

class IncrementalExtractor:
    def __init__(self, source_db_url, watermark_store):
        self.engine = create_engine(source_db_url)
        self.watermark_store = watermark_store

    def extract(self, table_name, watermark_column, date):
        """基于水位线的增量抽取"""
        # 获取上次成功处理的水位值
        last_watermark = self.watermark_store.get_watermark(
            table_name, watermark_column
        )

        # 查询新增和变更的数据
        query = f"""
            SELECT * FROM {table_name}
            WHERE {watermark_column} > %s
              AND {watermark_column} <= %s
            ORDER BY {watermark_column}
        """
        df = pd.read_sql(
            query, self.engine,
            params=(last_watermark, date)
        )

        if df.empty:
            return df

        # 更新水位线(仅处理成功后才更新)
        new_watermark = df[watermark_column].max()
        return df, new_watermark


class WatermarkStore:
    """水位线持久化存储"""
    def __init__(self, redis_client):
        self.redis = redis_client

    def get_watermark(self, table, column):
        key = f"watermark:{table}:{column}"
        value = self.redis.get(key)
        return value.decode() if value else '1970-01-01'

    def update_watermark(self, table, column, value):
        key = f"watermark:{table}:{column}"
        self.redis.set(key, value)

3.3 幂等写入与数据质量检查

def load_to_warehouse(df, table_name, date, engine):
    """幂等写入:先删后插,保证重跑安全"""
    with engine.begin() as conn:
        # 删除目标日期的旧数据
        conn.execute(
            f"DELETE FROM {table_name} WHERE dt = %s", (date,)
        )
        # 写入新数据
        df.to_sql(
            table_name, conn,
            if_exists='append', index=False,
            method='multi', chunksize=10000
        )


def run_quality_check(table_name, date, engine):
    """数据质量检查:行数、空值率、唯一性"""
    with engine.connect() as conn:
        # 检查1:行数不低于历史均值的 50%
        row_count = conn.execute(
            f"SELECT COUNT(*) FROM {table_name} WHERE dt = %s",
            (date,)
        ).scalar()

        avg_count = conn.execute("""
            SELECT AVG(daily_count) FROM (
                SELECT COUNT(*) as daily_count
                FROM {table}
                WHERE dt >= DATE_SUB(CURDATE(), INTERVAL 30 DAY)
                GROUP BY dt
            ) t
        """.format(table=table_name)).scalar()

        if row_count < avg_count * 0.5:
            raise ValueError(
                f"数据量异常: {row_count} 行,"
                f"低于30日均值 {avg_count:.0f} 的50%")

        # 检查2:关键字段空值率不超过 5%
        null_rate = conn.execute(f"""
            SELECT AVG(CASE WHEN user_id IS NULL THEN 1 ELSE 0 END)
            FROM {table_name} WHERE dt = %s
        """, (date,)).scalar()

        if null_rate > 0.05:
            raise ValueError(
                f"user_id 空值率 {null_rate:.2%} 超过阈值 5%")

四、ETL 工程化的隐性成本与架构权衡

Airflow 的调度延迟:Airflow Scheduler 的最小调度粒度为分钟级,且任务从调度到实际执行存在秒级延迟。对于需要秒级延迟的实时管道,Airflow 不是合适的选择,应考虑 Flink 或 Kafka Streams。

增量处理的一致性窗口:基于时间戳的增量抽取,在数据源写入存在延迟时可能遗漏数据。例如,事务 T1 在 23:59:58 提交,但数据库的 updated_at 字段在 00:00:02 才更新。日度 ETL 在 00:00 执行时,T1 不会被包含在当日数据中。缓解方案是在水位线后增加一个安全窗口(如 1 小时),但会增加数据延迟。

幂等写入的存储开销:"先删后插"的幂等策略在重跑时会短暂删除数据,查询该时间窗口的用户可能看到空结果。更安全的方案是写入临时表,验证通过后原子替换,但增加了存储和计算开销。

DAG 的复杂度膨胀:当 DAG 数量超过 100 个时,Airflow 的 UI 和调度性能开始下降。DAG 之间的依赖关系也需要显式管理(ExternalTaskSensor),否则可能出现循环等待或死锁。

五、总结

ETL 工程化的本质是将"脚本式数据处理"转化为"可调度、可监控、可重跑"的管道系统。本文方案的核心链路为:Airflow DAG 调度 → 增量抽取(水位线管理)→ 清洗转换 → 幂等写入 → 数据质量检查。落地时需重点关注三个参数:重试次数(建议 3 次)、水位线安全窗口(建议 1 小时)、数据质量阈值(建议行数不低于均值 50%,空值率不超过 5%)。建议从核心业务管道开始工程化改造,逐步替换脚本式 ETL,并建立数据质量告警机制。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐