为什么需要批量处理

ETL 的本质是数据的搬移和转换。当数据量达到一定规模后,逐条处理就不再可行。不管是每天几百万条订单记录,还是每小时数 GB 的日志文件,批量处理都是数据工程的基础操作方式。

批量处理的核心思路很简单:把数据分成一个个批次,每个批次作为一个整体来处理。这样做的好处很明显:

  • 吞吐量高:一次处理一批数据,而不是一条一条来
  • 资源利用好:可以集中分配计算和存储资源
  • 可恢复性强:一个批次失败不会影响其他批次
  • 成本可控:按批处理,容易预估时间和开销

举一个实际的例子。假设你每天凌晨需要把前一天的所有订单从业务数据库同步到数据仓库。如果逐条同步一百万条订单,每条都在 10 毫秒,总耗时超过 2.5 小时;但如果用批量方式,每次插入一万条,网络往返次数从一百万次减少到一百次,整个流程可能十几分钟就完成了。

调度策略的三种模式

调度就是决定"什么时候跑哪个任务"。在 ETL 场景中,调度策略主要分三种。

时间驱动调度

最简单也最常用的方式。规定好时间,到点就执行。常见的形式有:

  • 固定间隔:每 5 分钟跑一次增量抽取
  • 定时执行:每天凌晨 3 点跑全量刷新
  • 日历触发:每月 1 号跑月度汇总

时间驱动的优点是实现简单,依赖少。缺点是不够灵活——你无法感知上游数据是否准备好,如果上游系统故障导致数据延迟到达,定时任务可能跑了个空。

事件驱动调度

不是看时间,而是看条件。当某个事件发生时触发执行。事件可以是:

  • 文件到达:监控某个 FTP 目录,新文件出现就触发处理
  • 消息通知:上游系统发送一条 Kafka 消息,通知数据已就绪
  • API 回调:第三方系统通过 Webhook 通知数据变更
  • 数据库变更:CDC 工具捕获到变更后触发后续流程

事件驱动的优势是反应迅速,不会浪费计算资源跑空任务。缺点是引入事件中间件会增加系统复杂度,而且事件丢失后需要有补偿机制。

依赖驱动调度

任务之间有关系。B 必须在 A 完成后才能跑,C 必须等 A 和 B 都完成。这种"先做这个,再做那个"的编排方式就是依赖驱动调度。

依赖驱动的典型实现是 DAG(有向无环图)。每个任务是一个节点,任务之间的依赖关系是边。调度引擎按照 DAG 的拓扑顺序执行任务,保证每个任务只在其所有前置依赖完成后才启动。

Cron 表达式与定时任务

Cron 是 Linux 世界最经典的定时任务工具。几乎所有调度框架都兼容 cron 表达式。来看一个基本的 cron 格式:

┌─── 分钟 (0-59)
│ ┌─── 小时 (0-23)
│ │ ┌─── 日期 (1-31)
│ │ │ ┌─── 月份 (1-12)
│ │ │ │ ┌─── 星期 (0-7, 0和7都表示周日)
│ │ │ │ │
* * * * * command

常见的 cron 表达式示例:

# 每天凌晨 3 点执行
0 3 * * * /opt/scripts/daily_etl.sh

# 每 15 分钟执行一次
*/15 * * * * /opt/scripts/poll_files.sh

# 工作日早上 9 点执行
0 9 * * 1-5 /opt/scripts/business_daily.sh

# 每月 1 号和 15 号凌晨 2 点执行
0 2 1,15 * * /opt/scripts/monthly_report.sh

# 每周日凌晨 4 点执行全量刷新
0 4 * * 0 /opt/scripts/full_refresh.sh

Cron 的局限

Cron 虽然简单,但在 ETL 场景中有几个明显问题:

  • 没有依赖管理:B 必须在 A 之后跑?cron 做不到,只能用 shell 脚本串联
  • 失败处理粗糙:任务失败了最多发个邮件,没有重试机制
  • 没有可视化:看不到任务执行历史、耗时趋势
  • 不能动态调整:如果上一个任务跑了 2 小时导致下一个任务延迟,cron 不会自动做任何调整

这些局限正是专业调度框架要解决的问题。

Python 实现简单调度器

在引入 Airflow 这类重型调度框架之前,先用 Python 自己实现一个调度器,有助于理解调度系统的核心原理。

import time
import schedule
from datetime import datetime
from typing import Callable


def etl_extract():
    print(f"[{datetime.now()}] 开始抽取数据...")
    time.sleep(2)
    print(f"[{datetime.now()}] 抽取完成")


def etl_transform():
    print(f"[{datetime.now()}] 开始转换数据...")
    time.sleep(3)
    print(f"[{datetime.now()}] 转换完成")


def etl_load():
    print(f"[{datetime.now()}] 开始加载数据...")
    time.sleep(1)
    print(f"[{datetime.now()}] 加载完成")


def full_pipeline():
    """完整 ETL 流水线"""
    print("=" * 50)
    print(f"ETL 流水线开始: {datetime.now()}")
    etl_extract()
    etl_transform()
    etl_load()
    print(f"ETL 流水线结束: {datetime.now()}")
    print("=" * 50)


# 配置调度
schedule.every().day.at("03:00").do(full_pipeline)
schedule.every(10).minutes.do(etl_extract)

# 调度循环
if __name__ == "__main__":
    print("调度器已启动...")
    while True:
        schedule.run_pending()
        time.sleep(1)

这个例子用的 schedule 库实现了一个极简调度器。在实际生产中,你还需要处理失败重试、并发控制、日志记录等问题。让我们加一个带重试机制的版本:

import time
import functools
from datetime import datetime


def retry(max_attempts=3, delay=5):
    """重试装饰器"""
    def decorator(func):
        @functools.wraps(func)
        def wrapper(*args, **kwargs):
            last_exception = None
            for attempt in range(1, max_attempts + 1):
                try:
                    return func(*args, **kwargs)
                except Exception as e:
                    last_exception = e
                    print(f"[{datetime.now()}] 第 {attempt} 次执行失败: {e}")
                    if attempt < max_attempts:
                        print(f"等待 {delay} 秒后重试...")
                        time.sleep(delay)
            raise last_exception
        return wrapper
    return decorator


@retry(max_attempts=3, delay=5)
def fetch_from_api(url):
    """模拟从 API 抽取数据"""
    print(f"正在请求: {url}")
    # 模拟网络不稳定
    import random
    if random.random() < 0.5:
        raise ConnectionError("网络连接失败")
    return {"status": "ok", "data": [1, 2, 3]}


if __name__ == "__main__":
    result = fetch_from_api("https://api.example.com/orders")
    print(f"结果: {result}")

工作流编排与 DAG

当 ETL 任务增加到几十甚至上百个时,手动管理依赖关系就变得不可行。工作流编排引擎应运而生,其中最流行的当属 Apache Airflow。

DAG 的基本概念

DAG(有向无环图)是工作流编排的核心抽象。它由三部分组成:

  • 任务(Task):一个具体的执行单元,比如"抽取订单数据"
  • 依赖(Dependency):任务之间的关系,比如"转换必须在抽取之后"
  • 无环(Acyclic):不能有循环依赖,否则任务永远跑不完

一个典型的 ETL DAG 如下所示:

         ┌──────────┐
         │ 抽取订单  │
         └────┬─────┘
              │
         ┌────▼─────┐
         │ 抽取用户  │
         └────┬─────┘
              │
         ┌────▼─────┐
         │ 转换订单  │
         └────┬─────┘
              │
         ┌────▼─────┐
         │ 合并数据  │
         └────┬─────┘
              │
         ┌────▼─────┐
         │ 加载入库  │
         └──────────┘

Airflow DAG 示例

来看一个实际的 Airflow DAG 定义:

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from airflow.operators.dummy_operator import DummyOperator

default_args = {
    "owner": "data_team",
    "depends_on_past": False,
    "email_on_failure": True,
    "email": ["alert@example.com"],
    "retries": 3,
    "retry_delay": timedelta(minutes=5),
    "execution_timeout": timedelta(hours=2),
}

dag = DAG(
    "daily_order_etl",
    default_args=default_args,
    description="每日订单 ETL 处理",
    schedule_interval="0 3 * * *",
    start_date=datetime(2026, 1, 1),
    catchup=False,
    tags=["etl", "订单"],
)


def extract_orders(**context):
    """从业务库抽取订单"""
    # 实际的抽取逻辑
    execution_date = context["execution_date"]
    print(f"抽取 {execution_date} 的订单数据")
    return {"records_count": 100000}


def transform_orders(**context):
    """清洗和转换订单数据"""
    ti = context["ti"]
    extract_result = ti.xcom_pull(task_ids="extract_orders")
    print(f"转换 {extract_result['records_count']} 条订单记录")
    return {"status": "success"}


def load_orders(**context):
    """加载到数据仓库"""
    print("写入数据仓库")
    return {"loaded": True}


start = DummyOperator(task_id="start", dag=dag)

extract = PythonOperator(
    task_id="extract_orders",
    python_callable=extract_orders,
    provide_context=True,
    dag=dag,
)

transform = PythonOperator(
    task_id="transform_orders",
    python_callable=transform_orders,
    provide_context=True,
    dag=dag,
)

load = PythonOperator(
    task_id="load_orders",
    python_callable=load_orders,
    provide_context=True,
    dag=dag,
)

end = DummyOperator(task_id="end", dag=dag)

# 定义依赖关系
start >> extract >> transform >> load >> end

Airflow 的核心优势在于它提供了一整套基础设施:调度器(Scheduler)负责按时触发 DAG、执行器(Executor)负责分配任务给 worker、元数据库(Metadata DB)记录所有执行历史。相比简单的 cron 脚本,Airflow 让 ETL 调度变得可观测、可管理、可扩展。

重试逻辑与错误处理

批量任务不可避免地会失败。网络超时、数据库连不上、磁盘空间不足、数据格式异常——失败的原因五花八门。好的调度系统需要在这几个层面处理失败:

任务级重试

对每个任务设置重试次数和重试间隔。大部分临时性失败(网络抖动、服务重启)通过 1-3 次重试就能解决。

def etl_with_retry(task_func, max_retries=3, backoff_factor=2):
    """带指数退避的重试"""
    for attempt in range(max_retries):
        try:
            task_func()
            return
        except Exception as e:
            wait = backoff_factor ** attempt * 10
            print(f"第 {attempt + 1} 次失败,{wait} 秒后重试: {e}")
            if attempt < max_retries - 1:
                time.sleep(wait)
            else:
                print("已达最大重试次数,任务失败")


# 使用示例
def risky_extract():
    import random
    if random.random() < 0.7:
        raise TimeoutError("连接超时")
    print("抽取成功")


etl_with_retry(risky_extract)

死信队列

对于多次重试仍然失败的数据,不要丢弃,而是放入死信队列(Dead Letter Queue)供人工处理。

报警通知

任务超过重试次数后触发报警。常见的报警渠道有邮件、短信、Slack、钉钉等。

def send_alert(task_name, error_message):
    """发送任务失败告警"""
    alert_message = f"""
    【ETL 任务失败】
    任务名称: {task_name}
    失败时间: {datetime.now()}
    错误信息: {error_message}
    影响: 下游报表将停滞
    """
    # 这里可以集成邮件、Slack、钉钉等通知渠道
    print(f"发送告警: {alert_message}")
    # requests.post(slack_webhook_url, json={"text": alert_message})

幂等性设计

批处理任务最重要的设计原则之一:同一个任务重复执行多次,结果必须一致。这就是幂等性。

实现幂等性的常见做法:

  • 全量覆盖:每次写入前先清空目标分区
  • 使用 UPSERT:存在则更新,不存在则插入
  • 基于时间戳去重:通过唯一键和更新时间戳判断是否处理过
-- 使用 MERGE (UPSERT) 实现幂等加载
MERGE INTO dw.orders AS target
USING (
    SELECT order_id, amount, status, updated_at
    FROM staging.orders
    WHERE dt = CURRENT_DATE
) AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN
    UPDATE SET amount = source.amount, status = source.status,
              updated_at = source.updated_at
WHEN NOT MATCHED THEN
    INSERT (order_id, amount, status, updated_at)
    VALUES (source.order_id, source.amount, source.status, source.updated_at);

资源管理

批量任务往往消耗大量计算和存储资源。如果没有资源管理机制,多个大任务同时运行可能导致资源耗尽、系统崩溃。

并发控制

限制同时运行的任务数量,避免资源争抢。Airflow 中有多个层级的并发控制:

  • 全局池:限制整个系统的并发任务数
  • 任务池:按优先级或资源类型分配槽位
  • 资源槽:类似 Apache YARN 的资源管理方式

资源隔离

将不同类型的任务分配到不同的执行环境:

  • CPU 密集型任务(如数据转换)分配到计算资源充足的节点
  • IO 密集型任务(如数据抽取)分配到网络带宽充裕的节点
  • 重资源任务(如大规模排序)使用独立集群

增量批处理

不是所有数据都需要全量刷新。增量批处理只处理上次运行以来发生变化的数据,大幅减少处理量和执行时间。

def incremental_etl(db_conn, watermark_table, source_table, target_table):
    """基于水印表的增量 ETL"""
    
    # 1. 读取上次的水印(处理到哪里的标记)
    cursor = db_conn.cursor()
    cursor.execute(f"SELECT last_processed_id FROM {watermark_table}")
    last_id = cursor.fetchone()[0]
    
    # 2. 抽取增量数据
    cursor.execute(f"""
        SELECT * FROM {source_table}
        WHERE id > {last_id}
        ORDER BY id
        LIMIT 10000
    """)
    batch = cursor.fetchall()
    
    while batch:
        # 3. 转换并加载这批数据
        load_to_target(batch, target_table)
        
        # 4. 更新水印
        max_id = batch[-1]["id"]
        cursor.execute(f"""
            UPDATE {watermark_table}
            SET last_processed_id = {max_id}, 
                updated_at = NOW()
        """)
        db_conn.commit()
        
        # 5. 处理下一批
        cursor.execute(f"""
            SELECT * FROM {source_table}
            WHERE id > {max_id}
            ORDER BY id
            LIMIT 10000
        """)
        batch = cursor.fetchall()
    
    print("增量 ETL 完成")

生产环境调度架构

最后来看一个真实生产环境的调度架构图。这不是一个具体的工具选型,而是一种架构模式:

┌──────────────────────────────────────────┐
│              调度控制层                    │
│  ┌──────────┐  ┌──────────┐  ┌────────┐  │
│  │ Airflow  │  │ 事件总线  │  │ 手动触发│  │
│  └────┬─────┘  └────┬─────┘  └───┬────┘  │
└───────┼─────────────┼──────────────┼───────┘
        │             │              │
┌───────▼─────────────▼──────────────▼───────┐
│              任务执行层                     │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐ │
│  │ Spark 作业│  │ Python  │  │ SQL 脚本  │ │
│  └──────────┘  └──────────┘  └──────────┘ │
└─────────────────────────────────────────────┘
        │             │              │
┌───────▼─────────────▼──────────────▼───────┐
│              数据存储层                     │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐ │
│  │ 业务数据库│  │ 数据仓库  │  │ 对象存储  │ │
│  └──────────┘  └──────────┘  └──────────┘ │
└─────────────────────────────────────────────┘

在这个架构中:

  • 控制层负责任务定义、触发和依赖管理
  • 执行层负责具体的计算逻辑
  • 存储层负责数据持久化

三层分离的好处是每层可以独立扩展和优化。控制层挂了不影响已提交的任务,执行层可以动态扩缩容。

小结

批量处理是 ETL 的基石。本文从最简单的 cron 定时任务出发,逐步深入到 DAG 工作流编排、重试策略和资源管理。核心要点是:选择调度策略时要考虑数据量和业务需求,重试要配合幂等设计保证数据一致性,资源管理要防止任务间相互干扰。增量批处理是降低处理量的关键手段,也是下一篇文章"增量抽取与 CDC"的基础。

Summary: 批量处理调度策略与工作流编排的核心概念和实践。从 cron 到 Airflow DAG,涵盖重试、幂等和资源管理。