数据管道设计是数据工程的核心。本文梳理六大关键模式:ETL/ELT 对比、批量处理、增量同步、CDC、流批一体与数据质量治理。
一、ETL 与 ELT:架构哲学的根本分野
ETL 在中间层完成清洗与转换,适合一致性要求高、目标算力有限的场景。ELT 将原始数据加载到 Snowflake、BigQuery 等弹性算力平台,以 SQL/dbt 在仓内转换,适合敏捷分析。
1.1 ETL vs ELT 对比表
| 维度 | ETL | ELT |
|---|---|---|
| 转换位置 | 中间 ETL 服务器/集群 | 目标数据仓库内部 |
| 硬件要求 | 需要独立计算资源 | 依赖目标端弹性算力 |
| 灵活性 | 低,变更需重跑全量 | 高,支持按需回溯与调整 |
| 原始数据保留 | 不保留 | 完整保留,支持审计 |
| 适用数据量 | GB 级至 TB 级 | TB 级至 PB 级 |
| 技术栈 | Informatica、Talend、Airflow + Python | dbt、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。
5.1 Flink SQL 流批一体
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 与流批一体不断拓宽实时分析的边界。优秀的团队基于数据规模、延迟要求与源系统约束,组合运用六大模式,并在关键节点嵌入数据质量卡口。最终目标只有一个:让干净、及时、可信的数据,以最低成本流向需要它的人。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。