跳到主要内容

湖仓一体案例

下面给一个电商场景的完整湖仓落地示例,正好可以用你 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.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 ...)
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.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'
);
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;

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');
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,明细层不给
  • 非空校验:dwd.dwd_order_detailorder_id / user_id / pay_amount IS 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,技术栈简化一半