数据管道设计模式:ETL vs ELT、增量同步与流批一体

深入剖析数据工程中的核心管道设计模式,涵盖 ETL 与 ELT 架构选型、批量与增量同步策略、CDC 捕获、流批一体实现以及数据质量保障体系。

数据管道设计是数据工程的核心。本文梳理六大关键模式:ETL/ELT 对比、批量处理、增量同步、CDC、流批一体与数据质量治理。


一、ETL 与 ELT:架构哲学的根本分野

ETL 在中间层完成清洗与转换,适合一致性要求高、目标算力有限的场景。ELT 将原始数据加载到 Snowflake、BigQuery 等弹性算力平台,以 SQL/dbt 在仓内转换,适合敏捷分析。

1.1 ETL vs ELT 对比表

维度ETLELT
转换位置中间 ETL 服务器/集群目标数据仓库内部
硬件要求需要独立计算资源依赖目标端弹性算力
灵活性低,变更需重跑全量高,支持按需回溯与调整
原始数据保留不保留完整保留,支持审计
适用数据量GB 级至 TB 级TB 级至 PB 级
技术栈Informatica、Talend、Airflow + Pythondbt、SQLMesh、BigQuery/Snowflake
延迟分钟级至小时级分钟级至小时级
运维复杂度

1.2 ETL 管道示例

import pandas as pd
from sqlalchemy import create_engine
import hashlib

src = create_engine("mysql+pymysql://user:pass@source:3306/shop")
df = pd.read_sql(
    "SELECT order_id, user_id, amount, status, created_at "
    "FROM orders WHERE created_at >= CURDATE() - INTERVAL 7 DAY", con=src)

df["user_hash"] = df["user_id"].apply(lambda u: hashlib.sha256(u.encode()).hexdigest()[:16])
df["amount_cny"] = df["amount"].astype(float).round(2)
df["is_vip"] = df["amount_cny"] >= 1000
df = df.dropna(subset=["order_id", "amount_cny"])

dst = create_engine("postgresql://user:pass@warehouse:5432/dwh")
df.to_sql("fct_orders", con=dst, if_exists="append", index=False, chunksize=1000)

1.3 ELT 管道示例(dbt)

-- models/staging/stg_orders.sql
select
    order_id::bigint, user_id::text,
    amount::numeric(18,2), status::text, created_at::timestamp
from {{ source('shop', 'orders') }}

-- models/marts/fct_orders.sql
with s as (select * from {{ ref('stg_orders') }})
select order_id, sha256(user_id::bytea)::text as user_hash,
       amount as amount_cny, amount >= 1000 as is_vip,
       status, created_at, current_timestamp as dwh_loaded_at
from s
models:
  my_dwh:
    staging: {+materialized: view}
    marts: {+materialized: table, +sort: created_at, +dist: order_id}

二、批量处理:可控性与可观测性的基石

批量处理在资源利用率、失败回滚与依赖调度方面不可替代。幂等性是其生命线,手段包括覆盖写、分区交换与带主键的 upsert。

2.1 Airflow DAG 批量管道

from airflow import DAG
from airflow.providers.postgres.operators.postgres import PostgresOperator
from datetime import datetime, timedelta

default_args = {
    "owner": "data-platform", "depends_on_past": False,
    "email_on_failure": True, "retries": 2,
    "retry_delay": timedelta(minutes=5),
}

with DAG(
    dag_id="daily_sales_pipeline",
    default_args=default_args,
    start_date=datetime(2026, 1, 1),
    schedule_interval="0 3 * * *",
    catchup=False, tags=["batch", "sales"], max_active_runs=1,
) as dag:

    extract = PostgresOperator(
        task_id="extract", postgres_conn_id="src_shop_db",
        sql="""
            COPY (SELECT * FROM orders WHERE created_at >= '{{ ds }}'
                  AND created_at < '{{ next_ds }}')
            TO '/tmp/orders_{{ ds }}.csv' WITH CSV HEADER;
        """
    )
    load_staging = PostgresOperator(
        task_id="load_staging", postgres_conn_id="dwh_conn",
        sql="""
            DELETE FROM staging.orders WHERE dt = '{{ ds }}';
            COPY staging.orders FROM '/tmp/orders_{{ ds }}.csv' CSV HEADER;
        """
    )
    transform_mart = PostgresOperator(
        task_id="transform_mart", postgres_conn_id="dwh_conn",
        sql="""
            DELETE FROM marts.daily_sales WHERE sales_date = '{{ ds }}';
            INSERT INTO marts.daily_sales
            SELECT user_id, SUM(amount), COUNT(*), '{{ ds }}'::date
            FROM staging.orders WHERE dt = '{{ ds }}' GROUP BY user_id;
        """
    )
    extract >> load_staging >> transform_mart

2.2 Spark 分布式批处理

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, to_date, sum as spark_sum, count

spark = SparkSession.builder \
    .appName("DailySalesBatch") \
    .config("spark.sql.adaptive.enabled", "true") \
    .getOrCreate()

df = spark.read.parquet("s3://datalake/landing/orders/")
daily = df.filter(
    (col("created_at") >= "2026-08-31") & (col("created_at") < "2026-09-01")
).withColumn("sales_date", to_date(col("created_at")))

daily.groupBy("user_id", "sales_date").agg(
    spark_sum("amount").alias("total_amount"),
    count("*").alias("order_count")
).write.mode("overwrite").partitionBy("sales_date") \
    .parquet("s3://datalake/marts/daily_sales/")
spark.stop()

三、增量同步:从全量覆盖到精准投递

全量抽取对大型表不可持续。增量同步仅捕获变更数据,大幅降低抽取窗口与网络负载。时间戳最简单但无法感知删除;自增 ID 仅适用于追加表;哈希比对开销高,仅适用于小表。

3.1 增量合并(Merge)

MERGE INTO dwh.fct_orders AS target
USING staging.orders_delta AS source
ON target.order_id = source.order_id
WHEN MATCHED AND source._op = 'D' THEN DELETE
WHEN MATCHED AND source._op IN ('U', 'I') THEN
    UPDATE SET
        target.user_id = source.user_id,
        target.amount_cny = source.amount,
        target.status = source.status,
        target.updated_at = source.updated_at,
        target.etl_loaded_at = current_timestamp
WHEN NOT MATCHED AND source._op IN ('I', 'U') THEN
    INSERT (order_id, user_id, amount_cny, status, created_at, updated_at, etl_loaded_at)
    VALUES (source.order_id, source.user_id, source.amount, source.status,
            source.created_at, source.updated_at, current_timestamp);

3.2 基于水印的增量抽取框架

from dataclasses import dataclass
from typing import Optional
import pandas as pd
from sqlalchemy import text

@dataclass
class IncrementalConfig:
    table: str; watermark_column: str; watermark_type: str; primary_key: str

class IncrementalSync:
    def __init__(self, src_engine, dst_engine):
        self.src = src_engine; self.dst = dst_engine

    def get_watermark(self, table: str) -> Optional[str]:
        with self.dst.connect() as conn:
            row = conn.execute(text(
                "SELECT max_watermark FROM etl.watermarks WHERE table_name = :t"
            ), {"t": table}).fetchone()
            return str(row[0]) if row else None

    def extract_delta(self, cfg: IncrementalConfig) -> pd.DataFrame:
        wm = self.get_watermark(cfg.table) or (
            "1970-01-01" if cfg.watermark_type == "timestamp" else "0")
        q = text(f"""
            SELECT * FROM {cfg.table} WHERE {cfg.watermark_column} > :wm
            ORDER BY {cfg.watermark_column} ASC LIMIT 500000
        """)
        return pd.read_sql(q, self.src, params={"wm": wm})

    def sync(self, cfg: IncrementalConfig):
        df = self.extract_delta(cfg)
        if df.empty: print(f"[{cfg.table}] No new data."); return
        df.to_sql("_tmp_delta", self.dst, if_exists="replace", index=False)
        merge = text(f"""
            MERGE INTO {cfg.table} AS t
            USING _tmp_delta AS s ON t.{cfg.primary_key} = s.{cfg.primary_key}
            WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT *
        """)
        with self.dst.begin() as conn:
            conn.execute(merge)
            new_wm = str(df[cfg.watermark_column].max())
            conn.execute(text("""
                INSERT INTO etl.watermarks (table_name, max_watermark) VALUES (:t, :wm)
                ON CONFLICT (table_name) DO UPDATE SET max_watermark = EXCLUDED.max_watermark
            """), {"t": cfg.table, "wm": new_wm})
        print(f"[{cfg.table}] Synced {len(df)} rows -> {new_wm}")

四、CDC 与 Debezium:捕捉每一次数据脉搏

CDC 从数据库事务日志解析事件,无需侵入业务表结构即可捕捉完整变更流。相比时间戳轮询,CDC 无需审计字段、能捕获物理删除,且对源库几乎无负载。

4.1 Debezium 连接器配置

{
  "name": "shop_mysql_connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    "database.hostname": "mysql.shop.internal",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "${secrets:dbz-password}",
    "database.server.id": "184054",
    "database.server.name": "shop_db",
    "database.include.list": "shop",
    "table.include.list": "shop.orders,shop.users,shop.inventory",
    "snapshot.mode": "when_needed",
    "tombstones.on.delete": "true",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "database.history.kafka.topic": "schema-changes.shop",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false",
    "transforms.unwrap.delete.handling.mode": "rewrite",
    "transforms.unwrap.add.fields": "op,source.ts_ms"
  }
}

when_needed 仅在初始时触发快照;ExtractNewRecordState 扁平化 envelope 并保留 op 字段;tombstones.on.delete=true 保留删除墓碑。

4.2 Spark 消费 CDC 事件写入 Delta Lake

from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, when
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType

spark = SparkSession.builder \
    .appName("CDC_Streaming") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog",
            "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

schema = StructType([
    StructField("order_id", StringType()), StructField("user_id", StringType()),
    StructField("amount", DoubleType()), StructField("status", StringType()),
    StructField("created_at", TimestampType()),
    StructField("__op", StringType()), StructField("__source_ts_ms", StringType())
])

raw = spark.readStream.format("kafka") \
    .option("kafka.bootstrap.servers", "kafka:9092") \
    .option("subscribe", "shop_db.shop.orders") \
    .option("startingOffsets", "latest").load()

df = raw.select(from_json(col("value").cast("string"), schema).alias("v")).select("v.*")
df = df.withColumn("_is_deleted", when(col("__op") == "d", True).otherwise(False))

df.writeStream.format("delta").outputMode("append") \
    .option("checkpointLocation", "/checkpoints/orders_cdc") \
    .table("delta_cdc.orders").awaitTermination()

4.3 Kafka Connect JDBC Sink 连接数仓

{
  "name": "postgres_sink_connector",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "tasks.max": "2",
    "topics": "shop_db.shop.orders",
    "connection.url": "jdbc:postgresql://dwh:5432/analytics",
    "connection.user": "loader",
    "connection.password": "${secrets:pg-password}",
    "insert.mode": "upsert",
    "pk.mode": "record_key",
    "pk.fields": "order_id",
    "delete.enabled": "true",
    "auto.create": "false",
    "auto.evolve": "true"
  }
}

五、流批一体:统一语义与统一引擎

流批一体以同一套 API 处理有限数据集(批)与无限数据流(流)。Flink DataStream 将批视为有界流;Spark Structured Streaming 将流微批化为增量批处理。核心在于时间语义与状态管理:事件时间与水印保证乱序聚合正确;可重放状态保障 Exactly-Once。

CREATE TABLE user_events (
    user_id STRING, event_type STRING, amount DECIMAL(18,2),
    event_time TIMESTAMP(3),
    WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka', 'topic' = 'user_events',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json', 'scan.startup.mode' = 'latest-offset'
);

CREATE TABLE user_profile (
    user_id STRING, user_name STRING, register_date DATE,
    PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
    'connector' = 'jdbc', 'url' = 'jdbc:mysql://mysql:3306/shop',
    'table-name' = 'users', 'username' = 'flink',
    'password' = '${secrets:flink-mysql}'
);

CREATE TABLE event_summary (
    window_start TIMESTAMP(3), window_end TIMESTAMP(3),
    user_name STRING, total_amount DECIMAL(18,2), event_count BIGINT,
    PRIMARY KEY (window_start, user_name) NOT ENFORCED
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:postgresql://dwh:5432/analytics',
    'table-name' = 'event_summary', 'username' = 'flink',
    'password' = '${secrets:flink-pg}'
);

INSERT INTO event_summary
SELECT TUMBLE_START(e.event_time, INTERVAL '1' MINUTE),
       TUMBLE_END(e.event_time, INTERVAL '1' MINUTE),
       p.user_name, SUM(e.amount), COUNT(*)
FROM user_events e
LEFT JOIN user_profile FOR SYSTEM_TIME AS OF e.event_time AS p
    ON e.user_id = p.user_id
GROUP BY TUMBLE(e.event_time, INTERVAL '1' MINUTE), p.user_name;

切换模式即可重跑历史:

./bin/sql-client.sh -Dexecution.runtime-mode=streaming
./bin/sql-client.sh -Dexecution.runtime-mode=batch

5.2 Spark Structured Streaming 微批

from pyspark.sql import SparkSession
from pyspark.sql.functions import window, col, sum as spark_sum, count

spark = SparkSession.builder.appName("StreamingBatchUnified").getOrCreate()

stream_df = spark.readStream.format("kafka") \
    .option("kafka.bootstrap.servers", "kafka:9092") \
    .option("subscribe", "transactions") \
    .option("startingOffsets", "earliest") \
    .load() \
    .selectExpr("CAST(value AS STRING) as json") \
    .selectExpr("from_json(json, 't_id STRING, u_id STRING, amt DOUBLE, ts TIMESTAMP') as d") \
    .select("d.*")

agg = stream_df.withWatermark("ts", "10 minutes") \
    .groupBy(window(col("ts"), "5 minutes"), col("u_id")) \
    .agg(spark_sum("amt").alias("total"), count("*").alias("cnt"))

agg.writeStream.format("delta").outputMode("append") \
    .option("checkpointLocation", "/checkpoints/txn_stream") \
    .table("delta_stream.txn_summary")

六、数据质量检查:管道的免疫系统

质量检查必须在关键节点建立关卡。企业通常从六维度度量:完整性、唯一性、时效性、有效性、一致性、准确性。

6.1 Great Expectations 质量套件

import great_expectations as gx

context = gx.get_context()
datasource = context.sources.add_pandas("pandas_dwh")
data_asset = datasource.add_csv_asset(
    asset_name="orders_csv", filepath_or_buffer="/data/landing/orders.csv"
)

context.add_or_update_expectation_suite(expectation_suite_name="orders_suite")
validator = context.get_validator(
    batch_request=data_asset.build_batch_request(),
    expectation_suite_name="orders_suite"
)

validator.expect_column_values_to_not_be_null(column="order_id")
validator.expect_column_values_to_be_unique(column="order_id")
validator.expect_column_values_to_be_in_set(
    column="status",
    value_set=["created", "paid", "shipped", "completed", "cancelled"]
)
validator.expect_column_values_to_be_between(
    column="amount", min_value=0.01, max_value=1000000.00, mostly=0.995
)
validator.expect_column_values_to_be_between(
    column="created_at", min_value="2020-01-01", max_value="2026-12-31"
)
validator.expect_compound_columns_to_be_unique(column_list=["order_id", "user_id"])
validator.save_expectation_suite(discard_failed_expectations=False)

checkpoint = context.add_or_update_checkpoint(
    name="orders_checkpoint",
    expectation_suite_name="orders_suite",
    action_list=[
        {"name": "store_result", "action": {"class_name": "StoreValidationResultAction"}},
        {"name": "update_docs", "action": {"class_name": "UpdateDataDocsAction"}},
        {"name": "slack_alert", "action": {
            "class_name": "SlackNotificationAction",
            "slack_webhook": "${SLACK_WEBHOOK_URL}",
            "notify_on": "failure"
        }}
    ]
)
result = checkpoint.run()
if not result.success:
    raise ValueError("数据质量校验未通过,终止下游任务。")

6.2 仓库内 SQL 质量监控

-- dbt 原生测试(schema.yml)
version: 2
models:
  - name: fct_orders
    columns:
      - name: order_id
        tests: [not_null, unique]
      - name: status
        tests:
          - accepted_values:
              values: ['created', 'paid', 'shipped', 'completed', 'cancelled']
      - name: amount_cny
        tests:
          - dbt_utils.expression_is_true:
              expression: "> 0"
    tests:
      - dbt_utils.unique_combination_of_columns:
          combination_of_columns: [order_id, user_id]
CREATE OR REPLACE VIEW dwh.qa_daily_monitor AS
WITH checks AS (
    SELECT 'null_orders' as check_name,
           COUNT(*) FILTER (WHERE order_id IS NULL) as vio, COUNT(*) as total
    FROM dwh.fct_orders WHERE dt = CURRENT_DATE - 1
    UNION ALL
    SELECT 'dup_orders',
           COUNT(*) - COUNT(DISTINCT order_id), COUNT(*)
    FROM dwh.fct_orders WHERE dt = CURRENT_DATE - 1
    UNION ALL
    SELECT 'negative_amount',
           COUNT(*) FILTER (WHERE amount_cny <= 0), COUNT(*)
    FROM dwh.fct_orders WHERE dt = CURRENT_DATE - 1
    UNION ALL
    SELECT 'late_arrival',
           COUNT(*) FILTER (WHERE etl_loaded_at > created_at + INTERVAL '24 hours'),
           COUNT(*)
    FROM dwh.fct_orders WHERE dt = CURRENT_DATE - 1
)
SELECT CURRENT_DATE as check_date, check_name, vio, total,
       ROUND(100.0 * vio / NULLIF(total, 0), 4) as vio_pct,
       CASE
           WHEN check_name IN ('null_orders','dup_orders') AND vio > 0 THEN 'CRITICAL'
           WHEN check_name = 'negative_amount' AND vio_pct > 0.1 THEN 'WARNING'
           WHEN check_name = 'late_arrival' AND vio_pct > 5.0 THEN 'WARNING'
           ELSE 'PASS'
       END as severity
FROM checks;

FAQ

Q1: ETL 和 ELT 如何选择?是否有混合模式?

选型取决于源系统复杂度与目标算力。实际生产中存在混合架构:对敏感字段多的源系统先轻量 ETL(脱敏、过滤),加载至 staging 后通过 ELT 完成关联聚合。

Q2: CDC 引入后,如何处理 Schema 变更?

Debezium 的 schema_change topic 记录所有 DDL 事件。下游应结合 Schema Registry 或 Iceberg/Delta 的自动 schema merge 处理兼容变更,破坏性变更需显式校验并触发人工审批。

Q3: 流批一体是否意味着放弃批处理?

并非如此。流批一体解决开发效率与语义一致性问题,但纯流处理在资源与调试成本方面并非最优。推荐 Lambda 进化形态:流处理支撑实时看板,同一套代码的批模式重跑支撑深度分析。

Q4: 数据质量检查应嵌入哪个环节?频率如何设置?

质检覆盖三个关键点:入仓后验证完整性;转换后验证正确性;对外表增加出口关卡。批处理每次全量校验;流处理采用样本滑动窗口(每 5 分钟抽样 1 万条),结合异常检测监控延迟与体量。


总结

数据管道设计没有银弹。ETL 与 ELT 各司其职,批量与增量互为补充,CDC 与流批一体不断拓宽实时分析的边界。优秀的团队基于数据规模、延迟要求与源系统约束,组合运用六大模式,并在关键节点嵌入数据质量卡口。最终目标只有一个:让干净、及时、可信的数据,以最低成本流向需要它的人。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「Data Engineering」更多文章

  1. 数据治理与质量管理:数据血缘、质量监控与合规设计