Python 数据管道与 ETL 工程化:从 Airflow DAG 到增量处理的架构实践
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,并建立数据质量告警机制。
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)