湖仓一体案例
案例目标:Flink CDC MySQL → Paimon 湖仓分层 → StarRocks 加速 BI
下面给一个电商场景的完整湖仓落地示例,正好可以用你 HDP 3.3.2 的现有组件:
┌──────────────┐ Flink CDC ┌───────────────────────────────────────────────────┐
│ MySQL 业务库 │ ──────────────────▶ │ Apache Paimon 湖仓 │
│ (ecommerce) │ Binlog 实时同步 │ ODS 贴源层 DWD 明细层 DWS 汇总层 │
└──────────────┘ │ order_info / dwd_order / dws_user_day / │
│ user_info / dwd_user / ads_order_day / │
│ pay_order dwd_pay dws_order_category │
└───────────────────────┬───────────────────────────┘
│
Flink / Spark 批流读 Paimon
│
▼
┌──────────────────────────────────┐ 报表
│ StarRocks / Doris 外表 / 导入 │ ───▶ BI 看板
│ (ADS 层加速 + 高并发查询) │ (Superset/Metabase)
└──────────────────────────────────┘
使用到的组件(跟 HDP 3.3.2 完全匹配):Flink CDC、Paimon、HDFS(HDP 自带)、Hive Metastore(HDP 自带)、StarRocks/Doris(你文档里有现成)。
1. 业务系统同步范围(MySQL 源表)
| MySQL 源表 | 业务含义 | 变化频率 | 对应 ODS 表 |
|---|---|---|---|
ecommerce.order_info | 订单主表 | 高(每秒数百条,Update 多) | ods.ods_order_info 主键表 |
ecommerce.order_item | 订单商品明细表 | 高(同订单) | ods.ods_order_item 主键表 |
ecommerce.user_info | 用户画像信息表 | 中(Update 少) | ods.ods_user_info 主键表(维表) |
ecommerce.pay_order | 支付流水表 | Append(少量回退状态) | ods.ods_pay_order 主键表 |
ecomm.goods_category | 商品类目维表 | 低(几乎不更新) | ods.ods_goods_category 主键表(小维表) |
2. ODS 层:Flink CDC → Paimon(贴源层,跟业务库 1:1)
2.1 在 Flink 里建 Paimon Catalog + MySQL CDC 源表(省略 CDC 建表,参考 Flink 集成实践第 3 节)
2.2 ODS 建表 SQL(每张业务表一张,全部主键表 + 动态分区)
USE CATALOG paimon;
CREATE DATABASE IF NOT EXISTS ods;
-- 订单主表 ODS(变化最多的一张,演示核心参数)
CREATE TABLE ods.ods_order_info (
id BIGINT,
user_id BIGINT,
goods_id BIGINT,
goods_num INT,
total_amount DECIMAL(18, 2),
pay_amount DECIMAL(18, 2),
pay_type TINYINT, -- 1 支付宝 2 微信 3 银行卡
order_status TINYINT, -- 1 待支付 2 已支付 3 已发货 4 已完成 5 已取消
receiver_name STRING,
receiver_phone STRING,
receiver_addr STRING,
create_time TIMESTAMP(3),
pay_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) PARTITIONED BY (dt STRING)
WITH (
'bucket' = '-1', -- 动态 bucket,自动分裂
'bucket-key' = 'id',
'changelog-producer' = 'input', -- 上游是 CDC,直接复用 before/after
'partition.fields.dt.source' = 'update_time', -- 分区字段 dt 由 update_time 推断
'partition.fields.dt.date-formatter' = 'yyyy-MM-dd',
'snapshot.time-retained' = '7 d', -- 快照保留 7 天
'snapshot.num-retained.min' = '10',
'target-file-size' = '128 MB'
);
-- 其他 ODS 表结构与上面一致,按需建:
-- ods.ods_order_item (order_item_id + order_id + goods_id + price + ...)
-- ods.ods_user_info (id + name + gender + phone + user_level + register_time ...)
-- ods.ods_pay_order (pay_id + order_id + pay_amount + pay_status + pay_time ...)
-- ods.ods_goods_category (id + cate_name + parent_id + level ...)
2.3 启动 5 条 Flink CDC 流作业入湖
SET 'execution.checkpointing.interval' = '3min';
SET 'pipeline.name' = 'cdc_to_paimon_ods_order';
INSERT INTO ods.ods_order_info SELECT * FROM default_catalog.default_database.cdc_order_info;
-- 另外 4 张表同理,每条作业单独一个 YARN Application,互相不影响稳定性
3. DWD 层:Flink 实时清洗 / 维度补全(明细层)
3.1 DWD 表设计(核心是「宽表 + 状态码翻译 + 脱敏 + 维度字段补全」)
CREATE DATABASE IF NOT EXISTS dwd;
CREATE TABLE dwd.dwd_order_detail (
order_id BIGINT,
user_id BIGINT,
user_level TINYINT, -- 从 ods_user_info 维表补
goods_id BIGINT,
cate_id BIGINT, -- 从商品类目表补
cate_name STRING,
goods_num INT,
order_amount DECIMAL(18,2),
pay_amount DECIMAL(18,2),
pay_type_name STRING, -- 状态码翻译:支付宝 / 微信 / 银行卡
order_status_name STRING, -- 状态码翻译:待支付 / 已支付 ...
province_name STRING, -- 从收货地址里抽省份(便于地市级统计)
city_name STRING,
receiver_name STRING,
receiver_phone STRING,
create_time TIMESTAMP(3),
pay_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (order_id, goods_id) NOT ENFORCED
) PARTITIONED BY (dt STRING)
WITH (
'bucket' = '-1',
'changelog-producer' = 'input',
'partition.fields.dt.source' = 'update_time'
);
3.2 Flink Streaming 清洗作业(维表 Join + 字段清洗)
SET 'execution.checkpointing.interval' = '3min';
SET 'table.exec.state.ttl' = '7 d'; -- Lookup Join 状态过期,防状态无限涨
SET 'pipeline.name' = 'ods_to_dwd_order_detail';
INSERT INTO dwd.dwd_order_detail
SELECT
o.id AS order_id,
o.user_id AS user_id,
CAST(COALESCE(u.user_level, 0) AS TINYINT) AS user_level,
i.goods_id AS goods_id,
c.id AS cate_id,
c.cate_name AS cate_name,
i.goods_num AS goods_num,
o.total_amount AS order_amount,
o.pay_amount AS pay_amount,
CASE o.pay_type
WHEN 1 THEN '支付宝'
WHEN 2 THEN '微信'
WHEN 3 THEN '银行卡'
ELSE '未知' END AS pay_type_name,
CASE o.order_status
WHEN 1 THEN '待支付'
WHEN 2 THEN '已支付'
WHEN 3 THEN '已发货'
WHEN 4 THEN '已完成'
WHEN 5 THEN '已取消'
ELSE '未知' END AS order_status_name,
split(o.receiver_addr, ' ')[1] AS province_name, -- 按实际地址格式正则抽
split(o.receiver_addr, ' ')[2] AS city_name,
o.receiver_name AS receiver_name,
regexp_replace(o.receiver_phone, '(\d{3})\d{4}(\d{4})','$1****$2') AS receiver_phone, -- 手机号脱敏
o.create_time AS create_time,
o.pay_time AS pay_time,
o.update_time AS update_time
FROM ods.ods_order_info o
LEFT JOIN ods.ods_order_item i ON o.id = i.order_id
LEFT JOIN ods.ods_goods_category c ON i.goods_id = c.id
LEFT JOIN ods.ods_user_info u FOR SYSTEM_TIME AS OF o.update_time
ON o.user_id = u.id;
4. DWS 层:Flink 分钟级窗口聚合 / Spark 小时级汇总(主题汇总层)
4.1 用户 GMV 汇总(Flink 滚动窗口 10 分钟,T+0)
CREATE DATABASE IF NOT EXISTS dws;
CREATE TABLE dws.dws_user_day_gmv (
user_id BIGINT,
user_level TINYINT,
dt STRING,
order_cnt BIGINT,
goods_cnt BIGINT,
gmv DECIMAL(18,2),
pay_amount DECIMAL(18,2),
PRIMARY KEY (user_id, dt) NOT ENFORCED
) PARTITIONED BY (dt)
WITH (
'bucket' = '4',
'changelog-producer' = 'input'
);
SET 'execution.checkpointing.interval' = '1min';
SET 'pipeline.name' = 'dwd_to_dws_user_day_gmv';
INSERT INTO dws.dws_user_day_gmv
SELECT
user_id,
MAX(user_level) AS user_level,
DATE_FORMAT(update_time, 'yyyy-MM-dd') AS dt,
COUNT(DISTINCT order_id) AS order_cnt,
SUM(goods_num) AS goods_cnt,
SUM(order_amount) AS gmv,
SUM(CASE WHEN order_status_name <> '已取消'
THEN pay_amount ELSE 0 END) AS pay_amount
FROM dwd.dwd_order_detail
GROUP BY user_id, DATE_FORMAT(update_time, 'yyyy-MM-dd');
4.2 类目天级汇总(Spark 每天凌晨 T+1 补跑一次全量,和 Flink 结果对账)
INSERT OVERWRITE dws.dws_category_day_gmv PARTITION (dt='2026-09-04')
SELECT
cate_id,
FIRST_VALUE(cate_name, true) AS cate_name,
COUNT(DISTINCT order_id) AS order_cnt,
COUNT(DISTINCT user_id) AS pay_user_cnt,
SUM(goods_num) AS goods_cnt,
SUM(pay_amount) AS pay_amount
FROM dwd.dwd_order_detail
WHERE dt = '2026-09-04' AND order_status_name <> '已取消'
GROUP BY cate_id;
5. 冷热分层加速:Paimon → StarRocks/Doris(ADS 层)
这里用 StarRocks 的 Paimon External Catalog(推荐,无需物理导入,0 拷贝就能查 Paimon 湖仓)
5.1 StarRocks 里建 Paimon Catalog
CREATE EXTERNAL CATALOG paimon_catalog
PROPERTIES (
"type" = "paimon",
"paimon.metastore"= "hive",
"hive.metastore.uris" = "thrift://master1:9083",
"paimon.warehouse" = "hdfs:///user/hive/warehouse/paimon"
);
SELECT * FROM paimon_catalog.dws.dws_user_day_gmv
WHERE dt = '2026-09-04'
ORDER BY gmv DESC
LIMIT 100;
如果 StarRocks 里要做高并发报表(1000 QPS 以上),就把 Paimon 天级汇总定时物理导入 StarRocks 本地表:
-- StarRocks 目标表(明细模型 or 主键模型,按主键幂等)
CREATE TABLE ads_user_day_gmv (
user_id BIGINT,
user_level TINYINT,
dt DATE,
order_cnt BIGINT,
gmv DECIMAL(18,2)
) UNIQUE KEY(user_id, dt)
DISTRIBUTED BY HASH(user_id) BUCKETS 8
PROPERTIES ("replication_num" = "3");
-- 每天凌晨定时 Insert Into(幂等,重跑也不会脏)
INSERT INTO ads_user_day_gmv
SELECT user_id, user_level, dt, order_cnt, gmv
FROM paimon_catalog.dws.dws_user_day_gmv
WHERE dt = date_sub(current_date(), 1);
6. 运维与治理(Paimon 湖仓落地必做的 5 件事)
6.1 每天凌晨 Tag 关键层(保留 7/30/365 天,可调度 DolphinScheduler)
CALL sys.create_tag('ods.ods_order_info', concat('daily_', date_sub(current_date(),1)),
TIMESTAMP concat(date_sub(current_date(),1), ' 23:59:59'));
CALL sys.create_tag('dwd.dwd_order_detail', concat('daily_', date_sub(current_date(),1)),
TIMESTAMP concat(date_sub(current_date(),1), ' 23:59:59'));
CALL sys.create_tag('dws.dws_user_day_gmv', concat('daily_', date_sub(current_date(),1)),
TIMESTAMP concat(date_sub(current_date(),1), ' 23:59:59'));
-- 30 天后删除日 Tag(只留月底 Tag)
CALL sys.delete_tag('ods.ods_order_info', 'daily_20260805');
6.2 大分区强制 Compaction + 孤儿文件清理
CALL sys.compact('dwd.dwd_order_detail', 'dt=2026-09-04');
CALL sys.compact('dws.dws_user_day_gmv');
CALL sys.remove_orphan_files('ods.ods_order_info', ' older_than 1 d ');
6.3 Ranger 权限控制
Paimon 复用 Hive Metastore,在 Ranger Hive Plugin 里直接建 Policy:
- 数仓团队:
dwd.*/dws.*ALL - BI 团队:
dws.*SELECT /ads_*SELECT - 业务分析师:只给
dws.dws_*_day_*SELECT,明细层不给
6.4 数据质量(Apache Griffin 或 自写 Flink SQL 校验)
- 非空校验:
dwd.dwd_order_detail里order_id / user_id / pay_amountIS NULL 数 = 0 - 一致性校验:
sum(ods.order_info.pay_amount) ≈ sum(dwd.order_detail.pay_amount)(误差 < 1%) - 唯一性校验:DWS 主键
(user_id, dt)重复数 = 0
6.5 监控告警(用 Prometheus + Grafana,别用 Ambari Metrics 踩坑了)
- Flink Checkpoint 失败 / 延迟
- Paimon DWD/DWS 分区延迟:当日写入量为 0 告警
- 小文件数:
SELECT * FROM dwd_order_detail$files里小文件 (> 64M) 比例 > 40% 告警,跑手动 compact
7. 与传统「Kafka + Hive (HDFS) + 调度脚本」架构对比
| 维度 | 老架构(Kafka → Hive Text/Parquet + 脚本合并) | 新架构(Flink CDC → Paimon 湖仓) |
|---|---|---|
| 入湖链路跳数 | Kafka 1 跳 → Hive 导入脚本第 2 跳 → 合并小文件第 3 跳 | 1 跳 Flink CDC 直接写 Paimon |
| ODS 层 Upsert/Delete | 不支持(只能 Overwrite 整分区) | 主键 Upsert / Delete 原生支持 |
| 小文件治理 | 凌晨单独跑定时脚本合并(经常漏/挂) | Checkpoint 后自动 Compaction,无需外部调度 |
| 下游实时加工能力 | 要重复读 Kafka(一份在 Kafka,一份在 Hive,链路冗余) | 湖本身可以流读,下游 DWD/DWS Flink 直接流读 Paimon changelog |
| 历史回溯能力 | 靠 Hive 分区(分区丢了就没了) | Snapshot + Tag 长期保留,按时间戳/Tag 精确回溯 |
| 运维复杂度 | Kafka + Hive + Flume/SeaTunnel + 调度脚本 × N | 只有 Paimon + HMS + Flink,技术栈简化一半 |