全文干货,涵盖 Flink CDC、Apache Hudi、Spark 自适应查询、ClickHouse 精确去重、DolphinScheduler 任务编排,拒绝理论堆砌。
在 2026 年的数据栈中,单一引擎无法解决所有问题。我们采用 “存储层统一(Hudi/Iceberg)+ 计算层多引擎(Flink + Spark)+ 服务层异构(ClickHouse + Redis)” 的混合架构。
层级 | 技术选型 | 核心职责 |
|---|---|---|
采集层 | Debezium (CDC) + Filebeat (日志) | 监听 MySQL Binlog、业务日志无侵入采集 |
传输层 | Apache Kafka (Confluent) | 削峰填谷,保留 7 天回溯能力 |
存储层(湖) | Apache Hudi (MOR) | 分钟级延迟 ODS,支持 Record-Level 更新 |
计算层(流) | Apache Flink 1.18 | 实时 ETL + 宽表拼接 + 动态规则告警 |
计算层(批) | Apache Spark 3.4 (AQE + Z-Order) | 天级/小时级离线聚合与回溯修复 |
加速层(仓) | ClickHouse (ReplicatedMergeTree) | 亚秒级多维度即席查询(ADS) |
调度层 | DolphinScheduler 3.2 | 复杂任务依赖管理 + 补数据脚本 |
治理层 | DataHub + Great Expectations | 元数据血缘 + 数据质量断言 |
为避免业务高峰期锁表,我们采用 Debezium + Snapshot 一致性快照,将 snapshot.mode 设为 initial(仅首次全量),后续增量基于 GTID。
Debezium MySQL Connector 核心配置 (JSON):
{
"name": "mysql-order-cdc",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "rm-xxx.mysql.rds.aliyuncs.com",
"database.port": "3306",
"database.user": "debezium_user",
"database.password": "****",
"database.server.id": "184054",
"database.server.name": "dbserver1",
"database.include.list": "trade_db",
"table.include.list": "trade_db.orders,trade_db.order_items",
"database.history.kafka.bootstrap.servers": "kafka-cluster:9092",
"database.history.kafka.topic": "schema-changes.trade",
"time.precision.mode": "connect",
"snapshot.locking.mode": "minimal", // 仅表级读锁,快速释放
"snapshot.isolation.mode": "repeatable_read",
"binary.handling.mode": "bytes",
"event.deserialization.failure.handling.mode": "ignore"
}
}坑点预警:
database.server.id必须全局唯一,且不能与 MySQL 从库 ID 冲突,否则 Binlog 断流。
启动 Flink CDC Source (基于 Flink SQL):
CREATE TABLE mysql_orders (
id BIGINT PRIMARY KEY NOT ENFORCED,
order_no STRING,
user_id BIGINT,
total_amount DECIMAL(12,2),
status STRING,
created_at TIMESTAMP(3),
updated_at TIMESTAMP(3),
-- 捕获 Debezium 元数据
`op_ts` TIMESTAMP_LTZ(3) METADATA FROM 'value.source.timestamp' VIRTUAL
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'rm-xxx.mysql.rds.aliyuncs.com',
'port' = '3306',
'username' = 'flink_user',
'password' = '***',
'database-name' = 'trade_db',
'table-name' = 'orders',
'scan.incremental.snapshot.enabled' = 'true', -- 分布式读分块
'scan.incremental.snapshot.chunk.size' = '8096'
);根据业务域隔离,设置不同分区数(订单主表 30 分区,明细表 60 分区)。关键参数调优:
# server.properties
num.network.threads=8
num.io.threads=16
socket.send.buffer.bytes=1048576
socket.receive.buffer.bytes=1048576
socket.request.max.bytes=104857600
# 日志保留策略(支持实时 + 离线回溯)
log.retention.hours=168 # 7天
log.retention.bytes=-1
log.segment.bytes=1073741824 # 1GB 滚动确保 Exactly-Once 语义的关键配置(Flink Kafka Producer):
-- Flink DDL 写入 Kafka
CREATE TABLE kafka_dwd_orders (
id BIGINT,
order_no STRING,
user_id BIGINT,
total_amount DECIMAL(12,2),
status STRING,
proc_time AS PROCTIME() -- 处理时间属性
) WITH (
'connector' = 'kafka',
'topic' = 'dwd_orders',
'properties.bootstrap.servers' = 'kafka-cluster:9092',
'properties.transaction.timeout.ms' = '600000',
'format' = 'debezium-json',
'debezium-json.encode.ignore' = 'false'
);放弃传统的 Hive 分区覆盖写入,采用 Hudi MOR (Merge-On-Read) 表,配合 Flink 流式 Upsert,实现订单状态更新(如配送状态变更)的实时入湖。
Flink SQL 创建 Hudi 目标表(与上游 Kafka 字段对齐):
CREATE TABLE hudi_ods_orders (
id BIGINT,
order_no STRING,
user_id BIGINT,
total_amount DECIMAL(12,2),
status STRING,
created_at TIMESTAMP(3),
updated_at TIMESTAMP(3),
`partition` STRING, -- 按天分区,例如 '2026-08-17'
PRIMARY KEY (id) NOT ENFORCED
) PARTITIONED BY (`partition`)
WITH (
'connector' = 'hudi',
'path' = 'hdfs://nameservice1/data/warehouse/ods/orders',
'table.type' = 'MERGE_ON_READ',
'compaction.schedule.enabled' = 'true',
'compaction.delta_commits' = '5', -- 每5个增量提交做一次压缩
'compaction.trigger.strategy' = 'num_commits',
'changelog.enabled' = 'true', -- 记录变更日志,供下游流表Join
'index.type' = 'BUCKET',
'bucket.index.num.buckets' = '64', -- 固定桶数,避免小文件爆炸
'write.tasks' = '4',
'hoodie.datasource.write.hive_style_partitioning' = 'true'
);实时流写入管道(Flink Insert Into):
INSERT INTO hudi_ods_orders
SELECT
id,
order_no,
user_id,
total_amount,
status,
created_at,
updated_at,
DATE_FORMAT(created_at, 'yyyy-MM-dd') AS `partition`
FROM mysql_orders;性能杀手锏:Hudi MOR 表的
compaction必须放在低峰期,否则 IO 飙升。配合 DolphinScheduler 定时触发hudi-cli离线压缩任务。
业务场景:实时计算“大额订单(>5000元)”并关联用户等级,写入 Redis 供前端大屏展示。
Flink SQL 双流 Join(带状态 TTL 防止状态无限膨胀):
-- 用户等级维表 (来自另一个CDC)
CREATE TABLE hudi_dim_user (
user_id BIGINT PRIMARY KEY NOT ENFORCED,
user_level INT,
phone STRING
) WITH (
'connector' = 'hudi',
'path' = 'hdfs://nameservice1/data/warehouse/dim/user',
'table.type' = 'MERGE_ON_READ',
'read.streaming.enabled' = 'true',
'read.streaming.start-commit' = 'latest'
);
-- 实时规则匹配
CREATE TABLE kafka_large_order_alert (
order_no STRING,
user_id BIGINT,
total_amount DECIMAL(12,2),
user_level INT,
alert_time TIMESTAMP(3)
) WITH ('connector' = 'kafka', ...);
-- 流 Join 并设置状态存活时间
SET table.exec.state.ttl = 3600000; -- 1小时过期
INSERT INTO kafka_large_order_alert
SELECT
o.order_no,
o.user_id,
o.total_amount,
u.user_level,
CURRENT_TIMESTAMP
FROM hudi_ods_orders o
LEFT JOIN hudi_dim_user u ON o.user_id = u.user_id
WHERE o.total_amount > 5000
AND o.status = 'PAID'
AND u.user_level >= 3;Flink Checkpoint 深度调优 (flink-conf.yaml):
state.backend: rocksdb
state.backend.incremental: true
state.backend.rocksdb.predefined-options: SPINNING_DISK_OPTIMIZED
state.backend.rocksdb.block.blocksize: 16kb
state.backend.rocksdb.writebuffer.size: 256mb
execution.checkpointing.interval: 180000 # 3分钟
execution.checkpointing.tolerable-failed-checkpoints: 3
execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION针对复杂的财务对账场景(T+1),利用 Spark 3.4 的 AQE (Adaptive Query Execution) 和 动态分区裁剪。
核心 Spark SQL(计算商家 GMV 及退款率):
-- 启用 AQE 和倾斜 JOIN 优化
SET spark.sql.adaptive.enabled=true;
SET spark.sql.adaptive.coalescePartitions.enabled=true;
SET spark.sql.adaptive.skewJoin.enabled=true;
SET spark.sql.adaptive.skewJoin.skewedPartitionFactor=5;
SET spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256MB;
-- 使用 Z-Order 优化 Hudi 表,加速过滤
CALL hudi.system.optimize(table => 'hudi_ods_orders', op => 'zorder', columns => 'order_no,user_id');
-- 离线聚合
CREATE TABLE hive_dws_merchant_gmv STORED AS PARQUET AS
WITH refund_orders AS (
SELECT order_id FROM hudi_ods_orders WHERE status = 'REFUNDED'
)
SELECT
DATE(created_at) AS stat_date,
merchant_id,
SUM(total_amount) AS gmv,
COUNT(DISTINCT user_id) AS paying_users,
COUNT(CASE WHEN status='REFUNDED' THEN 1 END) / COUNT(*) AS refund_rate
FROM hudi_ods_orders o
LEFT JOIN dim_merchant m ON o.merchant_id = m.id
WHERE o.status IN ('PAID', 'SHIPPED', 'REFUNDED')
GROUP BY DATE(created_at), merchant_id内存管理(避免 OOM):
spark-submit \
--class com.xxx.GmvJob \
--master yarn \
--deploy-mode cluster \
--driver-memory 4g \
--executor-memory 8g \
--executor-cores 4 \
--num-executors 20 \
--conf spark.memory.fraction=0.8 \
--conf spark.memory.storageFraction=0.3 \
--conf spark.sql.shuffle.partitions=400 \
--jars hudi-spark3.4-bundle.jar,iceberg-spark3.4.jar \
my-data-job.jar离线聚合结果写入 ClickHouse,支持运营秒级查询。重点利用 ReplicatedReplacingMergeTree 配合 ver 去重,解决离线任务重复写入的脏数据问题。
CK 建表 DDL(去重引擎):
CREATE TABLE ads_merchant_gmv ON CLUSTER ck_cluster
(
stat_date Date,
merchant_id UInt64,
gmv Decimal(12,2),
paying_users UInt32,
refund_rate Float32,
etl_time DateTime DEFAULT now()
) ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/ads_merchant_gmv', '{replica}')
PARTITION BY toYYYYMM(stat_date)
ORDER BY (stat_date, merchant_id)
SETTINGS index_granularity = 8192;
-- 创建物化视图,预聚合 Top10 商家(查询加速 10 倍)
CREATE MATERIALIZED VIEW mv_top_merchant
ENGINE = ReplicatedSummingMergeTree('/clickhouse/tables/{shard}/mv_top_merchant', '{replica}')
ORDER BY (stat_date)
AS SELECT stat_date, sum(gmv) AS total_gmv, count() AS cnt FROM ads_merchant_gmv GROUP BY stat_date;ClickHouse 查询优化(强制使用分区索引):
-- 查询时务必带上分区键,否则全表扫描
SELECT
merchant_id,
sum(gmv) AS total,
uniqExact(paying_users) AS uv -- 精确去重,使用 HLL 或 uniqExact
FROM ads_merchant_gmv
WHERE stat_date BETWEEN '2026-08-01' AND '2026-08-17'
GROUP BY merchant_id
ORDER BY total DESC
LIMIT 10
SETTINGS max_threads = 16, max_memory_usage = 10000000000;依赖关系:Flink实时任务(持续) -> Hudi 离线压缩(凌晨1点) -> Spark 离线聚合(凌晨2点) -> 数据质量校验(凌晨3点) -> CK 同步(凌晨3:30) -> 报表推送(凌晨4点)。
DolphinScheduler 参数传递(跨任务传参):
# Shell 任务 A (获取业务日期)
export biz_date=$(date -d "-1 day" +%Y%m%d)
echo "${biz_date}" > /tmp/date.txt
# Spark 任务 B (接收参数)
spark-submit --class GmvJob --conf biz_date=${biz_date} job.jarGreat Expectations 数据质量断言(集成至调度):
# data_quality_check.py
import great_expectations as gx
context = gx.get_context()
validator = context.sources.pandas_default.read_csv(
"hdfs://nameservice1/data/warehouse/dws/gmv.csv"
)
expectations = [
gx.expectations.ExpectColumnValuesToBeBetween(column="gmv", min_value=0, max_value=10000000),
gx.expectations.ExpectColumnValuesToNotBeNull(column="merchant_id"),
gx.expectations.ExpectColumnValuesToBeInSet(column="stat_date", value_set=[biz_date])
]
results = validator.validate(expectation_suite="suite")
if not results["success"]:
raise Exception("Data quality check failed!") # 调度任务置为失败,阻断下游taskmanager.memory.segment-size 为 64kb,增加 execution.buffer-timeout 到 100ms,同时将 table.exec.source.cdc-events-duplicate 设为 true 以兼容双写场景。write.insert.drop.duplicates,并设置 hoodie.parquet.small.file.limit 为 104857600 (100MB),自动复用未写满的 Parquet 文件组。-- 手动触发小文件合并 (Hudi CLI)
compaction schedule --tableName hudi_ods_orders --basePath hdfs://...
compaction run --tableName hudi_ods_orders --parallelism 10开启 spark.sql.adaptive.coalescePartitions.parallelismFirst=false,并设置 spark.sql.adaptive.advisoryPartitionSizeInBytes=256MB,使输出文件数稳定在合理范围。
利用 Flink Kubernetes Operator 管理 Flink 集群,实现自动重启和弹性伸缩。
FlinkDeployment 自定义资源 YAML:
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
name: realtime-etl
spec:
image: flink:1.18-scala_2.12
flinkVersion: v1_18
flinkConfiguration:
taskmanager.numberOfTaskSlots: "2"
state.checkpoints.dir: 's3a://flink-checkpoints/etl'
high-availability.type: kubernetes
high-availability.storageDir: 's3a://flink-ha'
serviceAccount: flink
jobManager:
replicas: 1
resource:
memory: "2048m"
cpu: 1
taskManager:
replicas: 3
resource:
memory: "4096m"
cpu: 2
podTemplate:
spec:
containers:
- name: flink-main-container
env:
- name: AWS_ACCESS_KEY_ID
valueFrom:
secretKeyRef:
name: s3-creds
key: accessKey
job:
jarURI: local:///opt/flink/usrlib/etl-job.jar
entryClass: com.xxx.RealtimeETL
args: ["--kafka.bootstrap", "kafka-svc:9092"]
parallelism: 4
upgradeMode: stateless维度 | 传统 Hadoop 数仓 | 本方案 (Lakehouse + Real-time) |
|---|---|---|
数据延迟 | T+1 批处理,延迟 > 12h | 流批一体,ODS 延迟 < 2min |
数据更新 | Hive 覆盖全表,成本极高 | Hudi MOR 支持行级 Upsert,代价降低 90% |
查询加速 | Presto 依赖资源扫描 | ClickHouse 预聚合 + 物化视图,P95 < 200ms |
精确一致性 | 依赖离线全量对账 | Flink Exactly-Once + Hudi 事务保证 |
资源成本 | 存算耦合,浪费严重 | 存储 HDFS,计算弹性 K8s,资源利用率提升 60% |
最佳实践清单:
primary key 尽量选择业务主键,不要用组合字段过长,否则索引膨胀。partition 按天或月,严禁按小时或按单条 ID。这篇文章从配置级细节到架构级决策,完整还原了一个日处理 10TB 级数据平台的建设路径。所有 YAML、SQL、Shell 脚本均可直接移植至企业生产环境。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。