**作者:** 数据工程团队
**关键词:** 分布式ETL、Airflow、Celery、任务调度、性能优化、监控
一、前言
ETL(Extract-Transform-Load)作为数据仓库建设的核心环节,承载着企业数据从业务系统到分析平台的流转使命。当我们从单机脚本迈入分布式时代,ETL的架构设计、调度策略与性能优化便成为数据工程师必须攻克的技术高地。
本文将围绕分布式ETL调度监控与性能优化这一主题,从架构设计到生产实践,层层深入,与各位读者分享我们在千万级数据量场景下的实战经验。
二、分布式ETL架构设计
2.1 主从模式(Master-Slave)
主从模式是最经典的分布式ETL架构。一个Master节点负责任务解析、DAG编排、状态管理和Worker调度;多个Slave节点实际执行数据处理任务。
┌─────────────────────────────────────────┐
│ Master Node │
│ ┌──────────┐ ┌──────────┐ │
│ │ Scheduler │ │ Executor │ │
│ └─────┬────┘ └────┬─────┘ │
│ │ │ │
│ ┌─────▼─────────────▼──────┐ │
│ │ Metadata DB │ │
│ └──────────────────────────┘ │
└────────────────┬────────────────────────┘
│ 任务分发
┌────────────┼────────────┐
▼ ▼ ▼
┌─────────┐ ┌─────────┐ ┌─────────┐
│ Worker1 │ │ Worker2 │ │ Worker3 │
│ (Slave) │ │ (Slave) │ │ (Slave) │
└─────────┘ └─────────┘ └─────────┘
优点: 架构清晰,管理集中,调度逻辑与执行逻辑分离。
缺点: Master存在单点故障风险,水平扩展受限于Master处理能力。
适用场景: 中小规模集群(10台以内),调度频率适中(分钟级以上)。
2.2 对等模式(Peer-to-Peer)
对等模式中每个节点角色相同,通过一致性协议(如Raft)选举Leader协调任务分配,节点宕机后自动重新分发任务。
# 伪代码:对等模式任务分配示意
class PeerNode:
def __init__(self, node_id, peers):
self.node_id = node_id
self.peers = peers # 其他节点列表
self.leader = self.elect_leader()
def elect_leader(self):
# Raft协议实现Leader选举
pass
def assign_task(self, task):
# 一致性哈希分配任务
target = consistent_hash(task.task_id, self.peers)
if target.node_id == self.node_id:
self.execute(task)
else:
self.forward(target, task)
优点: 高可用,无单点瓶颈,弹性伸缩能力强。
缺点: 实现复杂,元数据一致性维护成本高。
2.3 微服务化架构
将ETL的每个环节拆解为独立微服务,通过消息队列串联,各服务独立部署、独立扩缩容。
数据源 → [Extract Service] → Kafka/TubeMQ → [Transform Service] →
Kafka → [Load Service] → 目标存储
↑ ↑ ↑
[Schema Registry] [UDF Engine] [DLQ Manager]
微服务化架构的关键考量:
- **服务发现:** 使用 Consul / Nacos 实现动态注册
- **流量控制:** 基于消息队列的分区机制实现背压(Backpressure)
- **数据一致性:** 借助 Exactly-Once 语义和幂等写入保证端到端一致性
三、高并发数据处理方案对比
3.1 多线程(Multi-threading)
适合I/O密集型任务——网络请求、文件读写、数据库操作。
from concurrent.futures import ThreadPoolExecutor, as_completed
def extract_table(table_name):
"""从源库抽取一张表数据(I/O密集)"""
return query_source(f"SELECT * FROM {table_name}")
def multi_thread_extract(tables: list, max_workers=8):
with ThreadPoolExecutor(max_workers=max_workers) as executor:
futures = {executor.submit(extract_table, t): t for t in tables}
results = {}
for future in as_completed(futures):
table = futures[future]
try:
results[table] = future.result()
except Exception as e:
log.error(f"Extract failed for {table}: {e}")
raise
return results
注意: Python GIL(Global Interpreter Lock)限制了纯计算场景的线程并行度,**多线程不适合CPU密集型变换**。
3.2 多进程(Multi-processing)
适合CPU密集型任务——数据清洗、复杂转换、加密解密。
import multiprocessing as mp
from functools import partial
def transform_chunk(chunk, rules: dict):
"""对数据分片执行转换规则(CPU密集)"""
return chunk.apply(rules)
def parallel_transform(df, rules: dict, num_workers=None):
num_workers = num_workers or mp.cpu_count()
chunks = np.array_split(df, num_workers)
with mp.Pool(num_workers) as pool:
results = pool.map(partial(transform_chunk, rules=rules), chunks)
return pd.concat(results)
注意: 进程间通信(IPC)和序列化开销不可忽视,建议在进程内批量处理数据后再传递结果。
3.3 异步IO(Async IO)
适合高并发网络请求场景——调用外部API、微服务编排。
import asyncio
import aiohttp
async def fetch_api(session, endpoint):
async with session.get(endpoint) as resp:
return await resp.json()
async def async_extract(endpoints: list):
async with aiohttp.ClientSession() as session:
tasks = [fetch_api(session, ep) for ep in endpoints]
return await asyncio.gather(*tasks, return_exceptions=True)
3.4 方案选型总览
| 维度 | 多线程 | 多进程 | 异步IO |
|------|--------|--------|--------|
| CPU密集 | ❌ | ✅ | ❌ |
| I/O密集 | ✅ | 可行(浪费) | ✅✅ |
| 网络IO高并发 | 可行 | 可行 | ✅✅✅ |
| 实现复杂度 | 低 | 中 | 中高 |
| 资源开销 | 低 | 高 | 极低 |
最佳实践: 在实际ETL中往往是混合使用——多线程/异步IO做数据抽取和加载,多进程池做核心变换逻辑。
四、Apache Airflow调度框架深入
4.1 DAG设计原则
Airflow的核心是DAG(有向无环图),好的DAG设计是调度系统的基石。
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.dummy import DummyOperator
from airflow.utils.task_group import TaskGroup
default_args = {
'owner': 'data_team',
'depends_on_past': False,
'email_on_failure': True,
'email': ['oncall@company.com'],
'retries': 3,
'retry_delay': timedelta(minutes=5),
'execution_timeout': timedelta(hours=2),
}
with DAG(
dag_id='etl_order_pipeline',
start_date=datetime(2024, 1, 1),
schedule_interval='0 2 * * *', # 每天凌晨2点
catchup=False,
max_active_runs=1,
default_args=default_args,
tags=['etl', 'order', 'production'],
) as dag:
start = DummyOperator(task_id='start')
with TaskGroup(group_id='extract_layer') as extract:
extract_orders = PythonOperator(
task_id='extract_orders',
python_callable=extract_table,
op_kwargs={'table': 'orders'},
pool='extract_pool',
)
extract_users = PythonOperator(
task_id='extract_users',
python_callable=extract_table,
op_kwargs={'table': 'users'},
pool='extract_pool',
)
with TaskGroup(group_id='transform_layer') as transform:
clean_orders = PythonOperator(
task_id='clean_orders',
python_callable=clean_data,
sla=timedelta(minutes=30),
)
join_data = PythonOperator(task_id='join_data', python_callable=join_orders_users)
with TaskGroup(group_id='load_layer') as load:
load_dwd = PythonOperator(task_id='load_dwd', python_callable=load_to_dwd)
load_dws = PythonOperator(task_id='load_dws', python_callable=load_to_dws)
start >> extract >> transform >> load
4.2 Task编排与依赖管理
Airflow支持多种依赖模式:
# 1. 链式依赖
task1 >> task2 >> task3
# 2. 扇入(Fan-in)——多个前置任务完成后执行
[task_a, task_b, task_c] >> task_join
# 3. 扇出(Fan-out)——一个任务触发多个下游
task_source >> [task_x, task_y, task_z]
# 4. 条件分支
from airflow.operators.python import BranchPythonOperator
def decide_branch(**context):
execution_date = context['execution_date']
if execution_date.isoweekday() <= 5:
return 'weekday_process'
return 'weekend_process'
branch = BranchPythonOperator(
task_id='branch_decision',
python_callable=decide_branch,
)
branch >> [weekday_process, weekend_process]
# 5. 跨DAG依赖——使用ExternalTaskSensor
from airflow.sensors.external_task import ExternalTaskSensor
wait_for_upstream = ExternalTaskSensor(
task_id='wait_for_dw_layer',
external_dag_id='dw_daily_load',
external_task_id='dwd_complete',
timeout=3600,
poke_interval=60,
)
4.3 DAG设计避坑指南
| 反模式 | 说明 | 正确做法 |
|--------|------|----------|
| 超长DAG | 单DAG包含100+任务,难以维护 | 按业务域拆分,每个DAG 20-30个任务 |
| 过度依赖 | 使用depends_on_past=True且backfill时死锁 | 仅在必要时使用,配合`wait_for_downstream` |
| 硬编码连接信息 | Connection信息硬编码在DAG文件 | 使用Airflow Connection/Variables管理 |
| 忽略任务幂等性 | 重复执行产生重复数据 | 确保每个任务支持幂等重跑(upsert/overwrite) |
五、分布式任务调度方案
5.1 Celery + Redis / RabbitMQ
Celery是Python生态最成熟的分布式任务队列,与Airflow结合可实现高吞吐调度。
# celery_app.py
from celery import Celery
from celery.signals import task_failure, task_success
celery_app = Celery(
'etl_worker',
broker='redis://redis-cluster:6379/0',
backend='redis://redis-cluster:6379/1',
)
celery_app.conf.update(
task_serializer='json',
accept_content=['json'],
result_serializer='json',
timezone='Asia/Shanghai',
enable_utc=True,
task_track_started=True,
task_acks_late=True, # 任务完成后才ACK
worker_prefetch_multiplier=1, # 每次只取一个任务,防止倾斜
task_reject_on_worker_lost=True,
)
# tasks.py
@celery_app.task(bind=True, max_retries=5, default_retry_delay=60)
def etl_sync_task(self, source_table, target_table):
"""分布式ETL任务"""
try:
data = extract_data(source_table)
transformed = transform_data(data)
load_data(transformed, target_table)
return {'status': 'success', 'rows': len(transformed)}
except ConnectionError as exc:
# 网络异常——可重试
raise self.retry(exc=exc, countdown=2 ** self.request.retries * 60)
except Exception as exc:
# 业务异常——记录死信
send_to_dlq(source_table, target_table, exc)
raise
生产配置要点:
- `task_acks_late=True` + `worker_prefetch_multiplier=1`:避免Worke积压
- `task_reject_on_worker_lost=True`:Worker宕机任务重新分发
- 结合Flower监控Worker和任务状态
5.2 Airflow Executor选型
| Executor类型 | 适用规模 | 依赖组件 | 任务分发方式 |
|-------------|---------|----------|-------------|
| SequentialExecutor | 开发调试 | 无 | 串行执行 |
| LocalExecutor | 单机/小规模 | 无 | 多进程并行 |
| CeleryExecutor | 中等规模 | Redis/RabbitMQ + Celery | 分布式Worker |
| KubernetesExecutor | 大规模/云原生 | Kubernetes集群 | 每个任务Pod |
| CeleryKubernetesExecutor | 混合场景 | Redis + K8s | 混合调度 |
# airflow.cfg 关键配置(CeleryExecutor模式)
[core]
executor = CeleryExecutor
parallelism = 32
dag_concurrency = 16
[celery]
worker_concurrency = 8
broker_url = redis://redis-cluster:6379/0
result_backend = redis://redis-cluster:6379/1
flower_url_prefix = /flower
[celery_kubernetes]
# 如果使用混合模式,可配置K8S相关参数
namespace = airflow
worker_container_image = airflow-worker:2.8.0
5.3 DolphinScheduler——国产分布式调度
Apache DolphinScheduler(海豚调度)是专为大数据场景设计的分布式调度系统,支持可视化DAG编辑、Spark/Flink任务原生支持。
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Master-1 │ │ Master-2 │ │ Master-3 │
│ (Active) │◄──►│ (Standby) │◄──►│ (Standby) │
└──────┬──────┘ └─────────────┘ └─────────────┘
│ RPC
┌──────▼────────────────────────────────────┐
│ Zookeeper (选主+服务发现) │
└──────┬────────────────────────────────────┘
│
┌──────▼──────┐ ┌──────▼──────┐ ┌──────▼──────┐
│ Worker-1 │ │ Worker-2 │ │ Worker-3 │
│ (Queue-A) │ │ (Queue-B) │ │ (Queue-A) │
└─────────────┘ └─────────────┘ └─────────────┘
DolphinScheduler任务提交示例:
# 通过DolphinScheduler REST API提交ETL任务
import requests
import json
def submit_ds_process(process_name, tasks, tenant='etl_user'):
"""向DolphinScheduler提交工作流"""
api_url = "http://ds-master:12345/dolphinscheduler/projects/1/process-definition"
# 构建工作流定义
process_def = {
"name": process_name,
"description": "分布式ETL任务",
"tenantCode": tenant,
"globalParams": [],
"tasks": tasks,
"timeout": 7200,
}
resp = requests.post(api_url, json=process_def)
return resp.json()
# SQL任务示例——无需编写Python代码
sql_task = {
"type": "SQL",
"id": "task_extract_orders",
"name": "抽取订单数据",
"params": {
"type": "MYSQL",
"datasource": 1, # 数据源ID
"sql": "SELECT * FROM orders WHERE dt = '${dt}'",
"sqlType": "0", # 查询
},
"runFlag": "NORMAL",
"maxRetryTimes": 3,
"retryInterval": 5,
}
shell_task = {
"type": "SHELL",
"id": "task_transform",
"name": "数据转换",
"params": {
"resourceList": [],
"localParams": [],
"rawScript": "spark-submit --class TransformJob transform.jar --date ${dt}"
},
}
与Airflow的核心差异:
| 维度 | Airflow | DolphinScheduler |
|------|---------|-----------------|
| 部署复杂度 | 中(需管理DB+Executor) | 低(一键集群部署) |
| 任务类型 | Python原生化 | 原生支持Shell/SQL/Spark/Flink |
| 调度精度 | 分钟级 | 秒级 |
| 可视化 | 需独立部署,运维监控在外部 | 内置完备运维界面 |
| 社区活跃度 | 极高 | 高(Apache项目) |
六、性能优化策略
6.1 分区策略(Partitioning)
合理分区是ETL性能优化的第一要务。
def build_incremental_query(table, last_max_id, batch_size=100000):
"""基于主键的分页增量抽取"""
return f"""
SELECT * FROM {table}
WHERE id > {last_max_id}
ORDER BY id
LIMIT {batch_size}
"""
def partition_by_date(base_date, num_partitions=8):
"""按日期范围分片"""
from datetime import timedelta
chunk_size = timedelta(hours=24 // num_partitions)
for i in range(num_partitions):
start = base_date + i * chunk_size
end = start + chunk_size
yield start, end
数据分片建议:
- **时间分区:** 按天/小时分区,适合增量ETL
- **哈希分区:** 按主键哈希分散热点,适合大规模JOIN
- **列表分区:** 按业务域(地域/渠道)隔离,适合多租户场景
6.2 并行处理(Parallelism)
# 动态调整并行度
def calc_optimal_parallelism(total_records, row_size_bytes, worker_memory_mb):
"""根据数据量和内存计算最佳并行度"""
per_worker_limit = worker_memory_mb * 0.6 * 1024 * 1024 # 60%可用内存
records_per_worker = per_worker_limit // row_size_bytes
parallelism = max(1, total_records // records_per_worker)
# 限制最大并行度
return min(parallelism, 32)
6.3 内存管理
ETL任务中的内存泄漏是生产环境最常见的问题之一。
import gc
import psutil
import os
def memory_safe_transform(iterable, batch_size=10000):
"""流式处理——每次只保留一个batch在内存"""
batch = []
for i, row in enumerate(iterable, 1):
batch.append(row)
if i % batch_size == 0:
yield process_batch(batch)
batch.clear()
if i % (batch_size * 10) == 0:
# 每处理10万行强制GC
gc.collect()
# 检查内存水位
mem = psutil.Process(os.getpid()).memory_percent()
if mem > 80:
log.warning(f"Memory usage high: {mem:.1f}%")
if batch:
yield process_batch(batch)
6.4 数据倾斜处理
数据倾斜是分布式ETL中最棘手的性能问题之一,直接导致部分Worker过载而其他Worker空闲。
# 倾斜检测与自动处理
class SkewDetector:
def __init__(self, threshold_ratio=3.0):
self.threshold_ratio = threshold_ratio # 最大/最小比值阈值
def detect(self, partition_sizes: dict):
"""检测是否存在数据倾斜"""
sizes = list(partition_sizes.values())
if len(sizes) < 2:
return False
max_size = max(sizes)
min_size = min(sizes)
return max_size / max(min_size, 1) > self.threshold_ratio
def resolve(self, key, skewed_partition, num_splits=8):
"""对倾斜分区进行打散处理——加盐(Salting)"""
def add_salt(row):
salt = random.randint(0, num_splits - 1)
return f"{row[key]}_{salt}"
return add_salt
# 使用示例:在Spark中处理倾斜JOIN
"""
-- 大表加盐,小表膨胀,解决倾斜Join
WITH salted_orders AS (
SELECT *,
CONCAT(order_id, '_', FLOOR(RAND() * 10)) AS salted_key
FROM orders
),
expanded_users AS (
SELECT *,
EXPLODE(ARRAY(0,1,2,3,4,5,6,7,8,9)) AS salt
FROM users
)
SELECT /*+ BROADCAST(expanded_users) */ *
FROM salted_orders o
JOIN expanded_users u
ON o.salted_key = CONCAT(u.user_id, '_', u.salt)
"""
6.5 I/O优化策略
# 批量写入vs逐条写入
def batch_load(rows, target_table, batch_size=5000):
"""批量写入——减少IO次数"""
conn = get_target_conn()
cursor = conn.cursor()
batch = []
for row in rows:
batch.append(row)
if len(batch) >= batch_size:
cursor.executemany(
f"INSERT INTO {target_table} VALUES (%s, %s, %s)",
batch
)
conn.commit()
batch.clear()
if batch:
cursor.executemany(...)
conn.commit()
I/O优化要点汇总:
| 策略 | 说明 | 预期提升 |
|------|------|---------|
| 列式存储格式 | Parquet/ORC替代CSV | 读取速度提升5-10倍 |
| 压缩传输 | Snappy/LZ4压缩中间数据 | 网络IO降低60-80% |
| 连接池复用 | 避免反复创建数据库连接 | 减少连接开销80%+ |
| 文件合并 | 小文件合并为大文件 | HDFS场景提升显著 |
| SSD + RAID | 本地临时数据使用SSD | 随机读写提升100倍 |
七、调度监控系统搭建
7.1 任务状态追踪
完备的任务状态机是监控的基础:
┌─────────┐
│ queued │
└────┬────┘
↓
┌─────────┐
┌────────►│ running │◄────────┐
│ └────┬────┘ │
│ │ │
│ ┌────▼────┐ │
│ ┌────│ success │────┐ │
│ │ └─────────┘ │ │
│ │ │ │
▼ ▼ ▼ ▼
┌──────┐ ┌──────┐ ┌──────────┐
│retry │ │skipped│ │ failed │
└──────┘ └──────┘ └────┬─────┘
│
┌─────▼─────┐
│ dead_letter│
└───────────┘
# 自定义任务事件监听
from airflow.models import TaskInstance
from airflow.utils.state import State
def task_state_listener(context):
"""任务状态变更回调——发送到监控系统"""
ti: TaskInstance = context['task_instance']
dag_id = ti.dag_id
task_id = ti.task_id
state = ti.state
duration = (ti.end_date - ti.start_date).total_seconds() if ti.end_date else 0
metrics = {
'dag_id': dag_id,
'task_id': task_id,
'state': state,
'duration_seconds': duration,
'execution_date': str(ti.execution_date),
'hostname': ti.hostname,
'try_number': ti.try_number,
}
# 推送到Prometheus / InfluxDB
push_to_monitoring(metrics)
# 失败告警
if state == State.FAILED:
send_alert(
title=f"[ETL Failed] {dag_id}.{task_id}",
message=f"执行次数: {ti.try_number}\n耗时: {duration}s",
severity='critical',
)
7.2 失败重试机制
# 智能重试策略
class SmartRetryPolicy:
"""根据异常类型决定重试行为"""
RETRYABLE_EXCEPTIONS = (
ConnectionError,
TimeoutError,
DatabaseError,
)
FATAL_EXCEPTIONS = (
DataValidationError,
SchemaMismatchError,
PermissionError,
)
@classmethod
def should_retry(cls, exception, try_number, max_retries=3):
if isinstance(exception, cls.FATAL_EXCEPTIONS):
return False, "不可恢复异常,停止重试"
if try_number >= max_retries:
return False, f"已达最大重试次数({max_retries})"
return True, f"第{try_number}次重试"
@classmethod
def calc_backoff(cls, try_number):
"""指数退避:1min, 2min, 4min, 8min..."""
return min(2 ** (try_number - 1) * 60, 3600) # 上限1小时
7.3 告警通知体系
告警层级:
┌──────────────────────────────────────────────┐
│ P0(Critical):链路中断、数据丢失 │
│ 通知方式:电话 + 短信 + 即时消息 + 邮件 │
│ 响应要求:15分钟 │
├──────────────────────────────────────────────┤
│ P1(Error):任务失败、数据延迟 │
│ 通知方式:即时消息 + 邮件 │
│ 响应要求:1小时 │
├──────────────────────────────────────────────┤
│ P2(Warning):数据量波动、执行时间异常 │
│ 通知方式:日报/周报汇总 │
│ 响应要求:24小时 │
└──────────────────────────────────────────────┘
7.4 告警通道集成
# 多通道告警分发实现
class AlertDispatcher:
"""告警分发器——按级别路由到不同通道"""
def __init__(self):
self.channels = {
'p0': [SmsChannel(), PhoneChannel(), IMChannel(), EmailChannel()],
'p1': [IMChannel(), EmailChannel()],
'p2': [EmailChannel()],
}
def dispatch(self, alert: Alert):
channels = self.channels.get(alert.severity, self.channels['p2'])
results = []
for channel in channels:
try:
channel.send(alert)
results.append(True)
except Exception as e:
log.error(f"Alert channel {channel.name} failed: {e}")
results.append(False)
# 如果所有通道都失败,升级告警
if not any(results):
self.escalate(alert)
def escalate(self, alert: Alert):
"""告警升级——通知值班经理"""
escalation_channel = PhoneChannel()
escalation_alert = Alert(
title=f"[ESCALATION] {alert.title}",
message=f"原始告警通道全部失败,需人工介入\n{alert.message}",
severity='p0',
)
escalation_channel.send(escalation_alert)
# 集成企业微信/钉钉机器人
class IMChannel(AlertChannel):
def send(self, alert: Alert):
webhook_url = self.get_webhook_url(alert.severity)
payload = {
"msgtype": "markdown",
"markdown": {
"content": (
f"## ⚠ ETL告警通知\n"
f"**级别:** {alert.severity}\n"
f"**任务:** {alert.dag_id}.{alert.task_id}\n"
f"**时间:** {alert.timestamp}\n"
f"**详情:** {alert.message}\n"
)
}
}
requests.post(webhook_url, json=payload, timeout=5)
7.5 监控指标与Dashboard
推荐监控关键指标(SRE黄金信号在ETL场景的应用):
1. 延迟(Latency): DAG整体执行耗时、各Task执行耗时 P50/P95/P99
2. 流量(Traffic): 数据处理行数、数据量大小(Bytes)
3. 错误(Errors): 任务失败率、重试次数、死信队列积压
4. 饱和度(Saturation): Worker CPU/内存、数据库连接池使用率、消息队列积压
# Prometheus指标暴露示例
from prometheus_client import Histogram, Counter, Gauge, start_http_server
etl_duration = Histogram(
'etl_task_duration_seconds',
'ETL任务执行耗时分布',
['dag_id', 'task_id'],
buckets=(60, 300, 600, 1800, 3600, 7200),
)
etl_rows_processed = Counter(
'etl_rows_processed_total',
'处理数据行数',
['dag_id', 'task_id', 'table_name'],
)
etl_queue_backlog = Gauge(
'etl_queue_backlog',
'任务队列积压数',
['queue_name'],
)
八、完整示例:基于Airflow + Celery的分布式ETL调度
下面给出一个端到端的生产级示例,涵盖从数据源抽取到目标加载的完整链路。
8.1 项目结构
airflow_etl/
├── dags/
│ ├── etl_order_pipeline.py # 主DAG定义
│ └── utils/
│ ├── extractors.py # 抽取模块
│ ├── transformers.py # 转换模块
│ ├── loaders.py # 加载模块
│ └── monitoring.py # 监控埋点
├── plugins/
│ └── operators/
│ └── etl_operators.py # 自定义Operator
├── config/
│ ├── airflow.cfg # Airflow配置
│ └── connections.yaml # 连接信息
└── requirements.txt
8.2 自定义ETL Operator
# plugins/operators/etl_operators.py
from airflow.models import BaseOperator
from airflow.utils.decorators import apply_defaults
from typing import Callable, Dict, Any
import time
class SmartETLOperator(BaseOperator):
"""
智能ETL Operator —— 内置监控埋点、重试、性能统计
"""
@apply_defaults
def __init__(
self,
extract_fn: Callable,
transform_fn: Callable,
load_fn: Callable,
source_config: Dict[str, Any],
target_config: Dict[str, Any],
batch_size: int = 10000,
enable_metric: bool = True,
*args, **kwargs
):
super().__init__(*args, **kwargs)
self.extract_fn = extract_fn
self.transform_fn = transform_fn
self.load_fn = load_fn
self.source_config = source_config
self.target_config = target_config
self.batch_size = batch_size
self.enable_metric = enable_metric
def execute(self, context):
dag_id = context['dag'].dag_id
task_id = self.task_id
log.info(f"Starting ETL: {dag_id}.{task_id}")
start_time = time.time()
total_rows = 0
try:
# 1. Extract——流式读取
for batch_idx, raw_batch in enumerate(self.extract_fn(
self.source_config, self.batch_size
)):
# 2. Transform
transformed = self.transform_fn(raw_batch)
# 3. Load
self.load_fn(transformed, self.target_config)
total_rows += len(transformed)
if batch_idx % 10 == 0:
log.info(f"Processed {total_rows} rows...")
# 上报指标
duration = time.time() - start_time
if self.enable_metric:
report_etl_metrics(dag_id, task_id, total_rows, duration)
log.info(f"ETL completed: {total_rows} rows in {duration:.2f}s")
return {'rows': total_rows, 'duration': duration}
except Exception as e:
log.error(f"ETL failed: {e}")
raise
8.3 Celery Worker配置
# 启动Celery Worker(生产建议使用Supervisor或K8S管理)
celery -A etl_worker.celery_app worker \
--loglevel=INFO \
--concurrency=8 \
--queues=etl_high,etl_normal \
--hostname=worker1@%h \
--max-tasks-per-child=1000 \
--time-limit=7200 \
--soft-time-limit=6600
8.4 生产部署架构
┌──────────────────┐
│ Nginx反向代理 │
└────────┬─────────┘
│
┌───────────────┼───────────────┐
▼ ▼ ▼
┌────────────┐ ┌────────────┐ ┌────────────┐
│ Airflow │ │ Airflow │ │ Airflow │
│ Webserver │ │ Scheduler │ │ Flower │
└────────────┘ └─────┬──────┘ └────────────┘
│
┌─────────────┼─────────────┐
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Celery │ │ Celery │ │ Celery │
│ Worker │ │ Worker │ │ Worker │ ← 动态扩缩容
└──────────┘ └──────────┘ └──────────┘
│ │ │
└─────────────┼─────────────┘
▼
┌──────────────┐
│ Redis │
│ (Broker+ │
│ Backend) │
└──────┬──────┘
▼
┌──────────────┐
│ PostgreSQL │
│ (MetadataDB) │
└──────────────┘
数据源: MySQL ──► Kafka ──► HDFS/Hive
九、生产环境最佳实践
9.1 容灾设计
- **多活架构:** 关键任务在两地三中心部署,通过DAG级别的Active-Standby切换
- **数据校验:** ETL完成后执行行数校验(Count Check)和校验和校验(Checksum)
- **回滚预案:** 保留最近N个分区的快照,支持快速回滚
- **限流保护:** 在Extract端和Load端都配置限流,防止压垮源/目标系统
# 数据校验示例
def validate_data_integrity(source_conn, target_conn, table):
"""基于行数和校验和的端到端校验"""
source_count = source_conn.fetch_val(f"SELECT COUNT(*) FROM {table}")
target_count = target_conn.fetch_val(f"SELECT COUNT(*) FROM {table}")
source_checksum = source_conn.fetch_val(
f"SELECT BIT_XOR(CAST(CRC32(CONCAT_WS('|', *)) AS UNSIGNED)) FROM {table}"
)
target_checksum = target_conn.fetch_val(
f"SELECT BIT_XOR(CAST(CRC32(CONCAT_WS('|', *)) AS UNSIGNED)) FROM {table}"
)
assert source_count == target_count, f"行数不匹配: {source_count} vs {target_count}"
assert source_checksum == target_checksum, f"校验和不匹配"
return True
9.2 数据血缘(Data Lineage)
数据血缘是元数据管理的核心能力,记录数据从源到目标的完整链路。
# 血缘元数据Schema(可存入Neo4j或Atlas)
lineage_record = {
'source': {
'type': 'mysql',
'database': 'source_db',
'table': 'orders',
'columns': ['order_id', 'user_id', 'amount', 'created_at'],
},
'transform': [
{'type': 'filter', 'description': '过滤取消订单'},
{'type': 'join', 'with': 'users', 'keys': ['user_id']},
{'type': 'aggregation', 'description': '按天聚合'},
],
'target': {
'type': 'hive',
'database': 'dwd',
'table': 'dwd_order_daily',
'columns': ['dt', 'total_orders', 'total_amount'],
},
'execution': {
'dag_id': 'etl_order_pipeline',
'task_id': 'transform_and_load',
'execution_time': '2024-06-15T02:30:00Z',
'rows_written': 125000,
},
}
建设路径:
1. 自动采集: 通过Operator Hook自动记录输入输出
2. 统一存储: Apache Atlas / DataHub / 自建Neo4j图谱
3. 血缘回溯: 支持从目标表追溯到原始字段,定位数据质量问题根因
9.3 数据质量监控
数据质量是ETL的生命线,在分布式调度中必须将质量监控嵌入每个环节。
# 数据质量规则引擎示例
class DataQualityEngine:
"""内置规则引擎,在ETL过程中实时校验数据质量"""
RULES = {
'not_null': lambda col, df: df[col].isnull().sum() == 0,
'unique': lambda col, df: df[col].is_unique,
'range': lambda col, df, mn, mx: df[col].between(mn, mx).all(),
'regex': lambda col, df, pattern: df[col].str.match(pattern).all(),
'referential': lambda col, df, ref_tbl, ref_col: ...,
}
def validate(self, df, checks: list):
for check in checks:
rule_type = check['type']
rule_fn = self.RULES[rule_type]
result = rule_fn(check['column'], df, **check.get('params', {}))
if not result:
raise DataQualityError(
f"质量规则未通过: {check}"
)
# 可选:将异常数据写入DLQ
质量监控三层体系:
| 层次 | 校验时机 | 校验内容 | 失败处理 |
|------|---------|---------|---------|
| 事前 | 任务执行前 | 上游数据可用性、Schema兼容性 | 阻塞等待或跳过 |
| 事中 | 数据变换中 | 空值、重复值、格式合法性 | 异常隔离+DLQ |
| 事后 | 加载完成后 | 行数比对、汇总校验、业务逻辑验证 | 告警+回滚 |
9.4 元数据管理最佳实践
# 元数据配置示例 (YAML)
tables:
- name: dwd_order_daily
owner: data_team
sla: 08:00 # 每日8点前必须完成
upstream:
- source_db.orders
- source_db.users
downstream:
- ads.order_report
retention: 30 # 保留30天分区
partition:
column: dt
granularity: DAY
quality_checks:
- type: not_null
columns: [order_id, dt]
- type: unique
columns: [order_id]
- type: row_count
min: 1000
max: 10000000
十、写在最后
分布式ETL调度监控与性能优化是一条持续演进之路。从单体脚本到主从架构,再到云原生的KubernetesExecutor模式,基础设施的演进不断推动着ETL工程能力的升级。关键在于:
1. 架构先行: 根据数据规模、时效性要求和技术栈选择合适的分布式架构,切忌过度设计
2. 可观测性: 没有完善的监控与链路追踪,分布式系统无异于黑盒运行,故障定位将极为困难
3. 持续优化: 性能优化不是一次性工作,而是持续的循环——测量→分析→优化→验证,每次迭代追求20%的提升
4. 规范驱动: 标准化DAG设计、任务命名、日志规范、元数据管理,规范的团队产出规范的ETL
5. 拥抱云原生: 容器化部署、动态扩缩容、KubernetesExecutor等云原生技术正在重塑ETL工程的交付方式
回顾本文的核心脉络:
- 我们从**架构设计**出发,对比了主从模式、对等模式和微服务化三种方案的优劣
- 深入**高并发处理**的内部机制,理清多线程/多进程/异步IO的适用边界
- 剖析了**Airflow调度框架**的DAG设计哲学与任务编排技巧
- 对比了**主流调度方案**的技术选型,并给出了完整的代码示例
- 提供了从分区策略到内存管理的**性能优化工具箱**
- 最后分享了生产环境中的**容灾、血缘和质量监控**最佳实践
ETL工程化是一个永无止境的精进过程。正如我们在生产环境中反复验证的那样——没有最好的架构,只有最适合业务的架构。推荐文章阅读的下一站:Apache Kafka在实时ETL中的应用、Flink CDC实时数据同步实战、数据湖湖仓一体架构演进。
欢迎在评论区分享你在生产环境中遇到的ETL挑战与独家解决方案。
*本文首发于公司技术博客,转载请注明出处。*
评论