首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >多层次构建企业级大数据平台:从CDC采集到湖仓一体实时计算全栈实战

多层次构建企业级大数据平台:从CDC采集到湖仓一体实时计算全栈实战

原创
作者头像
用户12678265
发布2026-08-17 14:17:55
发布2026-08-17 14:17:55
210
举报

多层次构建企业级大数据平台:从CDC采集到湖仓一体实时计算全栈实战

全文干货,涵盖 Flink CDC、Apache Hudi、Spark 自适应查询、ClickHouse 精确去重、DolphinScheduler 任务编排,拒绝理论堆砌。


一、架构总览:不止于 Lambda 的“流批一体”全景图

在 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 的 MySQL CDC 无锁配置

为避免业务高峰期锁表,我们采用 Debezium + Snapshot 一致性快照,将 snapshot.mode 设为 initial(仅首次全量),后续增量基于 GTID。

Debezium MySQL Connector 核心配置 (JSON):

代码语言:javascript
复制
{
  "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):

代码语言:javascript
复制
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'
);

三、第二层:Kafka 主题设计与分区策略

根据业务域隔离,设置不同分区数(订单主表 30 分区,明细表 60 分区)。关键参数调优:

代码语言:javascript
复制
# 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):

代码语言:javascript
复制
-- 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'
);

四、第三层:Apache Hudi 湖存储 —— 分钟级 ODS 构建

放弃传统的 Hive 分区覆盖写入,采用 Hudi MOR (Merge-On-Read) 表,配合 Flink 流式 Upsert,实现订单状态更新(如配送状态变更)的实时入湖。

Flink SQL 创建 Hudi 目标表(与上游 Kafka 字段对齐):

代码语言:javascript
复制
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):

代码语言:javascript
复制
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 离线压缩任务。


五、第四层:Flink 实时计算 —— 动态规则双流 Join

业务场景:实时计算“大额订单(>5000元)”并关联用户等级,写入 Redis 供前端大屏展示。

Flink SQL 双流 Join(带状态 TTL 防止状态无限膨胀):

代码语言:javascript
复制
-- 用户等级维表 (来自另一个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):

代码语言:javascript
复制
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

六、第五层:Spark 离线批处理(自适应查询 + Z-Order 优化)

针对复杂的财务对账场景(T+1),利用 Spark 3.4 的 AQE (Adaptive Query Execution)动态分区裁剪

核心 Spark SQL(计算商家 GMV 及退款率):

代码语言:javascript
复制
-- 启用 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):

代码语言:javascript
复制
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 加速层(物化视图 + 精确去重)

离线聚合结果写入 ClickHouse,支持运营秒级查询。重点利用 ReplicatedReplacingMergeTree 配合 ver 去重,解决离线任务重复写入的脏数据问题。

CK 建表 DDL(去重引擎):

代码语言:javascript
复制
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 查询优化(强制使用分区索引):

代码语言:javascript
复制
-- 查询时务必带上分区键,否则全表扫描
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;

八、第七层:DolphinScheduler 任务编排与数据质量闭环

依赖关系:Flink实时任务(持续) -> Hudi 离线压缩(凌晨1点) -> Spark 离线聚合(凌晨2点) -> 数据质量校验(凌晨3点) -> CK 同步(凌晨3:30) -> 报表推送(凌晨4点)

DolphinScheduler 参数传递(跨任务传参):

代码语言:javascript
复制
# 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.jar

Great Expectations 数据质量断言(集成至调度):

代码语言:javascript
复制
# 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!")  # 调度任务置为失败,阻断下游

九、性能调优实战:解决数据倾斜与反压

9.1 Flink 反压处理(Watermark 对齐问题)

  • 现象:Kafka Source 消费 lag 持续增加,但下游算子繁忙。
  • 解法:调整 taskmanager.memory.segment-size 为 64kb,增加 execution.buffer-timeout 到 100ms,同时将 table.exec.source.cdc-events-duplicate 设为 true 以兼容双写场景。

9.2 Hudi 小文件合并策略

  • 现象:ODS 表小文件超过 10 万个,NameNode 内存告警。
  • 解法:在 Flink 写入时开启 write.insert.drop.duplicates,并设置 hoodie.parquet.small.file.limit 为 104857600 (100MB),自动复用未写满的 Parquet 文件组。

代码语言:javascript
复制
-- 手动触发小文件合并 (Hudi CLI)
compaction schedule --tableName hudi_ods_orders --basePath hdfs://... 
compaction run --tableName hudi_ods_orders --parallelism 10

9.3 Spark 动态分区优化(避免产生 2000+ 个小分区)

开启 spark.sql.adaptive.coalescePartitions.parallelismFirst=false,并设置 spark.sql.adaptive.advisoryPartitionSizeInBytes=256MB,使输出文件数稳定在合理范围。


十、高可用部署:K8s + Operator 云原生实战

利用 Flink Kubernetes Operator 管理 Flink 集群,实现自动重启和弹性伸缩。

FlinkDeployment 自定义资源 YAML:

代码语言:javascript
复制
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%

最佳实践清单

  1. 绝对不要将 Debezium 和 Flink 部署在同一台机器,避免网卡打满。
  2. Hudi 表的 primary key 尽量选择业务主键,不要用组合字段过长,否则索引膨胀。
  3. ClickHouse 的 partition 按天或月,严禁按小时或按单条 ID。
  4. 数据质量校验(Great Expectations)必须阻断下游,否则脏数据蔓延后修复成本指数级增长。
  5. 离线任务和实时任务必须读写分离(实时写 Hudi,离线读 Hive 外表),避免文件锁冲突。

这篇文章从配置级细节到架构级决策,完整还原了一个日处理 10TB 级数据平台的建设路径。所有 YAML、SQL、Shell 脚本均可直接移植至企业生产环境。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 多层次构建企业级大数据平台:从CDC采集到湖仓一体实时计算全栈实战
    • 一、架构总览:不止于 Lambda 的“流批一体”全景图
    • 二、第一层:基于 Debezium 的 MySQL CDC 无锁配置
    • 三、第二层:Kafka 主题设计与分区策略
    • 四、第三层:Apache Hudi 湖存储 —— 分钟级 ODS 构建
    • 五、第四层:Flink 实时计算 —— 动态规则双流 Join
    • 六、第五层:Spark 离线批处理(自适应查询 + Z-Order 优化)
    • 七、第六层:ClickHouse 加速层(物化视图 + 精确去重)
    • 八、第七层:DolphinScheduler 任务编排与数据质量闭环
    • 九、性能调优实战:解决数据倾斜与反压
      • 9.1 Flink 反压处理(Watermark 对齐问题)
      • 9.2 Hudi 小文件合并策略
      • 9.3 Spark 动态分区优化(避免产生 2000+ 个小分区)
    • 十、高可用部署:K8s + Operator 云原生实战
    • 十一、总结:三层架构的核心差异化
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档