# 电商实时数据湖仓技术 Wiki（Flink + Paimon + StarRocks）

> 一个从 0 到 1 构建电商实时数据湖仓的教学型工程：用 Spring 生成器模拟日志与业务数据，用 Flink SQL 与 Paimon 完成 ODS/DIM/DWD/DWS/ADS 五层湖仓搭建，再由 StarRocks 直读 Paimon 实现可视化分析。值得研究的是它把「数据生成、双链路入湖、湖仓分层、系统调优、业务指标实现」串成了一条可完整跑通的最小闭环。

## Context Links

- [Agent index](https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/llms.txt)
- [Human interactive wiki](https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533)
- [GitHub repository](https://github.com/fengyu-eng/paimon-datalake)

## Repository Metadata

- Repository: fengyu-eng/paimon-datalake

- Generated: 2026-09-08T05:47:07.310Z
- Updated: 2026-09-14T05:23:42.536Z
- Runtime: Claude Code · claude-fable-5
- Format: Custom
- Pages: 4

## Page Index

- 01. [开篇导览：项目定位与阅读地图](https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/pages/01-page-1.md) - 按自定义格式的五个视角（架构、数据流、计算存储、系统优化、业务实现）介绍本仓库是什么、由哪些组件拼成，并给出后三页各自的阅读入口，帮助读者先建立整体心智模型。
- 02. [数据源头：业务数据与用户日志双生成器](https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/pages/02-page-2.md) - 讲解仿真数据是怎么造出来的：一个 Spring 应用把商品、订单、用户等全链路业务数据写入 MySQL，另一个 Spring 应用持续构造带嵌套结构的用户行为日志并发往 Kafka，为湖仓提供真实感输入。
- 03. [湖仓五层结构与实时数据流](https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/pages/03-page-3.md) - 沿着数据流讲清湖仓分层：日志经 Kafka Connector、业务数据经 MySQL CDC 两条链路进入 ODS，再经 DIM 维度加工、DWD 明细与枚举解码、DWS 近 N 日聚合、ADS 指标产出，落成可按天分区、可主键更新的 Paimon 表。
- 04. [计算存储调优、业务实现与经验收尾](https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/pages/04-page-4.md) - 收尾页：从统一配置入口出发解读 Flink 状态与检查点、微批、Paimon 文件策略等系统优化点，说明元数据在 Hive Metastore、数据落 HDFS 的存储形态，介绍 StarRocks 外部分区目录直读与可视化，最后把交易、流量、用户等业务域如何落到各层做一次收束，并沉淀可复用经验。

## Source File Index

- `dependency-reduced-pom.xml`
- `pom.xml`
- `README.md`
- `src/main/datagenerator/business_code/BusinessApplication.java`
- `src/main/datagenerator/business_code/generator/BaseDataGenerator.java`
- `src/main/datagenerator/business_code/util/DbUtil.java`
- `src/main/datagenerator/business_code/util/RandomUtil.java`
- `src/main/datagenerator/userlog_code/generator/UserLogGenerator.java`
- `src/main/datagenerator/userlog_code/UserLogApplication.java`
- `src/main/java/com/fy/warehouse/ads/ads_coupon_stats_full.java`
- `src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java`
- `src/main/java/com/fy/warehouse/DIM/dim_sku_full.java`
- `src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java`
- `src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java`
- `src/main/java/com/fy/warehouse/dwd/dwd_traffic_action_full.java`
- `src/main/java/com/fy/warehouse/dws/dws_trade_user_sku_order_nd_full.java`
- `src/main/java/com/fy/warehouse/ODS/ods_log_inc.java`
- `src/main/resources/core-site.xml`
- `src/main/resources/hdfs-site.xml`
- `StarRocks.sql`

---

## 01. 开篇导览：项目定位与阅读地图

> 按自定义格式的五个视角（架构、数据流、计算存储、系统优化、业务实现）介绍本仓库是什么、由哪些组件拼成，并给出后三页各自的阅读入口，帮助读者先建立整体心智模型。

- Page Markdown: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/pages/01-page-1.md
- Generated: 2026-09-08T05:46:09.178Z

### Source Files

- `README.md`
- `pom.xml`
- `StarRocks.sql`
- `src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java`

<details>
<summary>相关源文件</summary>
以下文件用于生成此维基页面：
- [README.md](README.md)
- [pom.xml](pom.xml)
- [StarRocks.sql](StarRocks.sql)
- [src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)
- [src/main/java/com/fy/warehouse/ODS/ods_log_inc.java](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)
- [src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java](src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java)
- [src/main/java/com/fy/warehouse/DIM/dim_sku_full.java](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java)
- [src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java](src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java)
- [src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java](src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java)
- [src/main/java/com/fy/warehouse/ads/ads_coupon_stats_full.java](src/main/java/com/fy/warehouse/ads/ads_coupon_stats_full.java)
- [src/main/java/com/fy/warehouse/Test.java](src/main/java/com/fy/warehouse/Test.java)
- [src/main/datagenerator/userlog_code/README.md](src/main/datagenerator/userlog_code/README.md)
</details>

# 开篇导览：项目定位与阅读地图

本页是整个 Wiki 的入口页。它先回答三个问题：**这个仓库是什么**、**由哪些组件拼成**、**按什么视角去读它**。看完本页，你应该能在大脑中建立起一张"数据从哪来、经过几层加工、落在哪里、被谁查询"的整体心智模型，然后带着方向进入后续的深度页面。

本仓库（fengyu-eng/paimon-datalake）是一套以电商为业务背景、用 **Flink + Paimon** 从 0 到 1 搭建的**实时数据湖仓**参考实现：既有完整的建仓代码，也自带模拟数据源，可以按 README 的步骤复现整条链路。官方定位与技术选型见 [README.md:1-17](README.md)、[README.md:19-30](README.md)（版本表）。

## 一、架构视角：组件与拼装方式

抛开细节，这个仓库由四层组件拼成，每一层在代码里都有清晰的归属目录：

| 层次 | 组件 | 代码/配置位置 | 一句话职责 |
| --- | --- | --- | --- |
| ① 数据源模拟 | Spring 数据生成器 | `src/main/datagenerator/` | 日志写入 Kafka，业务数据写入 MySQL |
| ② 计算与接入 | 一整套 Flink SQL 作业 | `src/main/java/com/fy/warehouse/` | 完成"接数 + 分层加工" |
| ③ 湖存储 | Paimon（Hive Metastore 管元数据、HDFS 管数据） | 各作业内的 `CREATE CATALOG`/`CREATE TABLE` | 落数据、按主键更新、按天分区 |
| ④ 查询与可视化 | StarRocks + DataRT | [StarRocks.sql](StarRocks.sql) | 用 MPP 引擎直接读湖，再可视化 |

数据源模拟器分为两套：`userlog_code` 生成用户行为日志（Kafka 主题为 `ODS_BASE_LOG`），`business_code` 生成订单、用户、优惠券等业务数据写入 MySQL 库 `gmall`；旁边还预置了对应的 `userlog_jar`/`business_jar` 可运行产物与 `gmall.sql` 建库脚本。参见 [src/main/datagenerator/userlog_code/README.md:1-39](src/main/datagenerator/userlog_code/README.md) 中对数据结构与 Kafka 发送逻辑的说明。

计算层是仓库的主体：`com.fy.warehouse` 下按数仓分层建包，从目录看作业类分布为 **ODS 28 个、DIM 6 个、dwd 18 个、dws 12 个、ads 2 个**，外加 `udf`（1 个自定义函数）、`config`（统一配置类）。每个作业是一个带 `main()` 的类，类名即目标表名。README 明确数仓分层为 ODS→DIM→DWD→DWS→ADS（[README.md:9-17](README.md)）。

```mermaid
flowchart LR
    subgraph SRC["① 数据源模拟（Spring）"]
        ULOG["userlog_code 用户行为日志"]
        BIZ["business_code 订单/用户/优惠券"]
    end

    subgraph FLINK["② Flink SQL 作业（com.fy.warehouse 包）"]
        ODSJ["ODS 层 28 个作业：Kafka/CDC 接入"]
        DIMJ["DIM 层 6 个作业：维度加工"]
        DWDJ["dwd 层 18 个作业：明细加工"]
        DWSJ["dws 层 12 个作业：汇总聚合"]
        ADSJ["ads 层 2 个作业：应用指标"]
    end

    subgraph LAKE["③ Paimon 湖（paimon_hive Catalog）"]
        ODSDB[("ods 库")]
        DIMDB[("dim 库")]
        DWDDB[("dwd 库")]
        DWSDB[("dws 库")]
        ADSDB[("ads 库")]
        META["Hive Metastore 元数据 + HDFS 数据文件"]
    end

    subgraph OLAP["④ 查询与可视化"]
        SR["StarRocks paimon_catalog"]
        DRT["DataRT"]
    end

    ULOG -->|"Kafka: ODS_BASE_LOG"| ODSJ
    BIZ -->|"MySQL: gmall"| ODSJ
    ODSJ --> ODSDB
    ODSDB --> DIMJ --> DIMDB
    ODSDB --> DWDJ --> DWDDB
    DWDDB --> DWSJ
    DIMDB --> DWSJ
    DWSJ --> DWSDB --> ADSJ --> ADSDB
    ADSDB --> SR --> DRT
    ODSDB -. 元数据/文件 .-> META
```

需要强调一个关键设计：**计算与存储是解耦的**。Flink 只负责"算并写"，Paimon 的元数据放在 Hive Metastore、数据文件落在 HDFS（README 原文"元数据存储在 Hive 中，实际数据存储在 HDFS"）；查询阶段则由 StarRocks 通过 `paimon_catalog` 外部目录**直接读湖上的 ads 表**做可视化，不再把数据导回一份副本。`StarRocks.sql` 里建的就是这个外部目录，随后 `use ads; select * from ads_coupon_stats_full;`（[StarRocks.sql:1-13](StarRocks.sql)）。

## 二、数据流视角：两条入口，五层加工

数据进湖有**两条互不相同的入口**，这是理解本仓库数据流的钥匙：

1. **日志数据（增量、无界）**：生成器 → Kafka `ODS_BASE_LOG` → Flink Kafka Connector → 写入 ODS。代表作业 `ods_log_inc`：注册 Kafka 源表后建 Paimon 目标表，把嵌套的 `common`/`page`/`start` ROW 拍平成 `common_*`/`page_*` 前缀列写入（[src/main/java/com/fy/warehouse/ODS/ods_log_inc.java:24-47](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)、[:87-122](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)）。
2. **业务数据（有界快照 + CDC）**：生成器 → MySQL `gmall` → Flink MySQL CDC（`scan.startup.mode = 'initial'`，先全量后增量）→ 写入 ODS。代表作业 `ods_sku_info_full`（[src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java:22-45](src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java)）。

两条入口进入同一套分层管道后，作业之间的读取依赖是我逐层核对过的：DIM 作业从 `ods.*` 读、DWD 作业从 `ods.*` 读、DWS 作业从 `dwd.*` 读（少量还从 `dim.*` 取维度补齐维度名）、ADS 作业只从 `dws.*` 读。也就是说，代码里严格保持了"只读下层、只写本层"的依赖方向：

```text
数据入口                              数仓分层加工                          出口
                        ┌─► DIM(dim库) ────────────────┐
Kafka ─► ODS ──────────┼─────────────────────────────► DWS(dws库) ─► ADS(ads库) ─► StarRocks/DataRT
MySQL ─► (ods库)       └─► DWD(dwd库) ─────────────────┘
```

这条链路里还有一个贯穿始终的统一约定：**每张加工表都用 `k1` 做天级分区**，`k1` 由事件/创建时间格式化而来（例如 `DATE_FORMAT(create_time, 'yyyy-MM-dd')`，见 [src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java:93-103](src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java)）。因此后续层做"最近 N 日"统计时，本质是**按分区取数再聚合**，而非依赖真正的流式时间窗口——例如 DWS 的 30 日订单汇总用 `timestampdiff(day, k1, '2024-05-31') between 0 and 30` 圈定分区范围（[src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java:60-82](src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java)）。表名后缀 `_full`/`_inc` 也是同一套约定：`_inc` 表示持续追加的增量表（主要用于日志），`_full` 表示按天分区重算的全量表。

## 三、计算存储视角：Flink 计算 + Paimon 湖表

**计算侧**统一是"Flink Table API + SQL"形态：所有作业先取 `FlinkConfigUtil.getFlinkConfig()` 得到 `Configuration`，再据此创建 `TableEnvironment`，然后执行三段 SQL——建 Catalog、建表、`INSERT INTO ... SELECT`（[src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java:12-29](src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java)）。没有 DataStream API、没有手写算子，整条数仓就是"DDL + 一条 SQL"的集合。

**存储侧**，每张 Paimon 表的建表语句高度模板化，重复出现如下特征（以 [src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java:32-63](src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java) 为代表）：

- `PRIMARY KEY (..., k1) NOT ENFORCED`：主键表，靠主键做 upsert；
- `PARTITIONED BY (k1)` + `'metastore.partitioned-table' = 'true'`：天分区；
- `'file.format' = 'parquet'`：数据文件格式；
- 写缓冲与分区过期：`write-buffer-size 512mb`、`write-buffer-spillable true`、分区过期 1 天、每小时检查一次。

DIM 层的建表与此同构，只是加工语义不同：用一串 `LEFT JOIN` 把多张 ODS 表组装成宽维度（商品 SKU 维关联 spu、三级分类、品牌，并用 `collect()` 把属性聚成列表），见 [src/main/java/com/fy/warehouse/DIM/dim_sku_full.java:33-65](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java) 的建表与 [:70-159](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java) 的加工 SQL。用户维度还保留了拉链结构的痕迹（`start_date`/`end_date` 字段），见 [src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java:33-59](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java)。

对外部系统的依赖全部集中在 `src/main/resources/` 的 Hadoop 配置（NameNode 地址、临时目录、root 代理等，[core-site.xml:19-50](src/main/resources/core-site.xml)、[hdfs-site.xml:19-41](src/main/resources/hdfs-site.xml)）以及 `Test.java` 这类连通性验证入口（直接以 root 身份访问 HDFS 上某张 Paimon 表的目录与文件，[Test.java:11-32](src/main/java/com/fy/warehouse/Test.java)）。

## 四、系统优化视角：一个配置类，集中了全部调优

仓库把运行层面的调优几乎全部收敛在 **`FlinkConfigUtil`** 这一个类里，是理解"系统如何被调稳"的最佳入口。代码注释清楚交代了这些配置是为了解决哪些真实问题：

| 关注点 | 配置 | 注释里的动机 |
| --- | --- | --- |
| 状态持久化 | 开启 Checkpoint；状态后端 `filesystem`，路径指向 HDFS | "原配置为 false，导致状态无法持久化"、"解决状态超限" |
| 精确一次 | checkpoint 间隔 60s、超时 300s、容忍 3 次失败、`EXACTLY_ONCE` | 精确一次语义 |
| 状态体积 | 状态 TTL 24 小时 | 控制状态增长 |
| 吞吐延迟 | 微批 500ms / 1000 条 | 小批量聚合 |
| Paimon 写入 | `paimon.sink.batch-size=100`、`buffer-time=0s`、`file.size=64mb` | "小批次快速写入"、"小文件即可，避免合并" |
| HDFS 稳定性 | 经 `env.java.opts` 调大 RPC socket 与连接超时 | "解决 RPC 中断" |
| 时区与约束 | `Asia/Shanghai`、非空约束 `DROP` | 统一时区 |

以上证据见 [src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java:24-73](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)。另有两点值得读者留意：其一，默认并行度被设为 1（[:23](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)），叠加 DWS 里硬编码的统计日期，说明这套工程的目标是**链路可复现、易讲解**，而非追求生产级吞吐——把它当"调优动作大全"来读比当"生产参数"来抄更合适；其二，每个作业文件里残留的注释（"关键修改 1/2/3"）本身就是一份排障日志，记录了一个新手从状态丢失、RPC 中断一路调稳的过程。

## 五、业务实现视角：电商域的五个主题

把表按业务语义归类，可以看见这个仓库覆盖了电商数仓的典型主题，而不是零散的表：

| 业务主题 | 典型作业/表 | 落点 |
| --- | --- | --- |
| 交易域 | 下单明细、支付成功、退款、购物车 | `dwd_trade_order_detail_full`、`dwd_trade_pay_detail_suc_full`、`dws_trade_user_order_nd_full` 等 |
| 流量域 | 启动、页面浏览、动作、曝光、报错 | `ods_log_inc` 拆分出的 `dwd_traffic_*` 系列 |
| 工具域（营销） | 优惠券领取、下单、支付 | `dwd_tool_coupon_*`、`dws_trade_coupon_order_nd_full` |
| 互动域 | 评论、收藏 | `dwd_interaction_comment_full`、`dwd_interaction_favor_add_full` |
| 基础维度 | 用户、商品、地区、日期、字典 | `dim_user_zip_full`、`dim_sku_full`、`dim_date_full` |

"业务含义→宽表→指标"的写法很典型：DWD 层用字典表把 `source_type` 代码翻译成名称（[src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java:69-141](src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java)），ADS 层再从 DWS 汇总表计算业务比率——优惠券统计里的 `reduce_rate`（补贴率）就是用 `coupon_reduce_amount_30d / original_amount_30d` 相除得到的（[src/main/java/com/fy/warehouse/ads/ads_coupon_stats_full.java:30-58](src/main/java/com/fy/warehouse/ads/ads_coupon_stats_full.java)）。

## 六、阅读地图：先看什么，后三页从哪里入口

**建议的首读顺序**：先读本页（建立整体模型）→ 挑一条最完整的链路逐类跟读一遍 → 再按自己关心的视角深入。最省力的一条"样板链路"是：`ods_sku_info_full`（CDC 进湖）→ `dim_sku_full`（维度加工）→ `dwd_trade_order_detail_full`（明细宽表）→ `dws_trade_user_order_nd_full`（30 日汇总）→ `ads_coupon_stats_full`（应用指标）。读通这五个类，就等于读通了仓库 90% 的写法。

后续三个深度页面与入口文件对应如下：

| 后续页面（视角） | 本文已铺垫的 | 建议入口文件 |
| --- | --- | --- |
| 计算与存储 | Catalog/主键/分区/格式的统一模板 | 任意建表作业的 `CREATE TABLE` 段 + `Test.java` |
| 系统优化 | 配置类解决的问题清单 | `FlinkConfigUtil.java` 全文 + ODS 作业里的"关键修改"注释 |
| 业务实现 | 五类业务主题与代表表 | `src/main/datagenerator/` 各类生成器 + 各层代表类 |

一句话总结：**这是一个把"模拟数据源 + Flink SQL 分层加工 + Paimon 湖存储 + StarRocks 即席查询"完整串起来的电商实时数仓教材工程**。代码最大的价值在于——每条链路的接入方式、每张表的建表选项、每个调优参数的动机，都以"一个类 + 注释"的形态平铺在你面前；后续三页会分别把"湖表机制""调优动作""业务口径"讲透，本页的地图到此就够用了。

---

## 02. 数据源头：业务数据与用户日志双生成器

> 讲解仿真数据是怎么造出来的：一个 Spring 应用把商品、订单、用户等全链路业务数据写入 MySQL，另一个 Spring 应用持续构造带嵌套结构的用户行为日志并发往 Kafka，为湖仓提供真实感输入。

- Page Markdown: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/pages/02-page-2.md
- Generated: 2026-09-08T05:45:29.800Z

### Source Files

- `src/main/datagenerator/business_code/BusinessApplication.java`
- `src/main/datagenerator/business_code/generator/BaseDataGenerator.java`
- `src/main/datagenerator/userlog_code/UserLogApplication.java`
- `src/main/datagenerator/userlog_code/generator/UserLogGenerator.java`
- `src/main/datagenerator/business_code/util/DbUtil.java`
- `src/main/datagenerator/business_code/util/RandomUtil.java`

<details>
<summary>相关源文件</summary>
以下文件用于生成此维基页面：
- [src/main/datagenerator/business_code/BusinessApplication.java](src/main/datagenerator/business_code/BusinessApplication.java)
- [src/main/datagenerator/business_code/generator/BaseDataGenerator.java](src/main/datagenerator/business_code/generator/BaseDataGenerator.java)
- [src/main/datagenerator/business_code/generator/UserDataGenerator.java](src/main/datagenerator/business_code/generator/UserDataGenerator.java)
- [src/main/datagenerator/business_code/generator/OrderDataGenerator.java](src/main/datagenerator/business_code/generator/OrderDataGenerator.java)
- [src/main/datagenerator/business_code/generator/UserBehaviorGenerator.java](src/main/datagenerator/business_code/generator/UserBehaviorGenerator.java)
- [src/main/datagenerator/business_code/util/DbUtil.java](src/main/datagenerator/business_code/util/DbUtil.java)
- [src/main/datagenerator/business_code/util/RandomUtil.java](src/main/datagenerator/business_code/util/RandomUtil.java)
- [src/main/datagenerator/userlog_code/UserLogApplication.java](src/main/datagenerator/userlog_code/UserLogApplication.java)
- [src/main/datagenerator/userlog_code/generator/UserLogGenerator.java](src/main/datagenerator/userlog_code/generator/UserLogGenerator.java)
- [src/main/datagenerator/userlog_code/model/UserLog.java](src/main/datagenerator/userlog_code/model/UserLog.java)
- [src/main/datagenerator/userlog_code/model/Common.java](src/main/datagenerator/userlog_code/model/Common.java)
</details>

# 数据源头：业务数据与用户日志双生成器

## 页面概述

这个仓库做的是"从 0 到 1 构建电商实时数据湖仓"，整个链路需要有源源不断的输入数据。本页讲解的就是这套输入数据是怎么"造"出来的：仓库里放了两套互相独立的 Spring 程序——**业务数据生成器**负责把商品、订单、用户等全链路电商业务数据批量写进 MySQL，**用户日志生成器**负责按固定节奏构造带嵌套结构的用户行为日志并发往 Kafka。

之所以拆成两条线，是因为湖仓对两类数据的采集方式天然不同：业务表走 Flink CDC（读 MySQL 的变更），用户日志走 Kafka 流。两个生成器正好把两类"有真实感"的输入喂到下游，是整条 ODS 层数据的第一来源。

Sources: [README.md:6-10]()；[src/main/datagenerator/userlog_code/README.md:1-4]()；[src/main/datagenerator/business_code/README.md:1-4]()

## 两个程序在整个湖仓里的位置

可以把数据源头理解为一条"双叉路口"：左边是结构化业务表，右边是半结构化日志流，汇合点在下游的 ODS 层。

```mermaid
flowchart LR
    subgraph SRC["数据源头（src/main/datagenerator 下的 Spring 进程）"]
        BA["BusinessApplication<br/>业务数据生成器"]
        UA["UserLogApplication<br/>用户日志生成器"]
    end

    subgraph CARRIER["数据载体"]
        DB[("MySQL<br/>gmall 库")]
        KAFKA["Kafka 主题<br/>ODS_BASE_LOG"]
    end

    subgraph LAKE["湖仓采集（Flink SQL）"]
        CDC["Flink CDC<br/>mysql-cdc 连接器"]
        KC["Flink Kafka 连接器"]
        ODS["Paimon ODS 层"]
    end

    BA -- "HikariCP 批量 INSERT" --> DB
    UA -- "逐条 JSON 消息" --> KAFKA
    DB --> CDC --> ODS
    KAFKA --> KC --> ODS
```

Sources: [src/main/datagenerator/business_code/BusinessApplication.java:13-78]()；[src/main/datagenerator/userlog_code/UserLogApplication.java:22-100]()；[src/main/java/com/fy/warehouse/ODS/ods_log_inc.java:24-45]()；[src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java:21-47]()

一个容易混淆的点先说明：`business_jar/` 与 `userlog_jar/` 目录里放的是另一套预构建 mock 工具的部署脚本与配置（`log.sh` 启动 `datageneration2024-mock-log-2024-05-01.jar`，配置键以 `mock.*` 开头），仓库内并没有这个 jar 的本体。本页描述的是仓库内带源码的 `*_code` 两套实现。

Sources: [src/main/datagenerator/business_jar/log.sh:1-4]()；[src/main/datagenerator/userlog_jar/application.yml:1-12]()；[src/main/datagenerator/business_jar/application.properties:1-9]()

## 两个生成器的整体对照

| 对比维度 | 业务数据生成器 | 用户日志生成器 |
| --- | --- | --- |
| 入口类 | `BusinessApplication`（business_code） | `UserLogApplication`（userlog_code） |
| 产物形态 | 一张张关系表里的行 | 一条条嵌套 JSON 日志 |
| 数据出口 | MySQL `gmall` 库（JDBC 批量插入） | Kafka 主题（默认 `ODS_BASE_LOG`） |
| 生成节奏 | 主循环每 `interval` 毫秒跑一轮，每轮写一批（默认 1000） | 每 `interval` 毫秒生成并同步发送一条 |
| 关键配置 | `spring.datasource.*`、`generator.batch-size`、`generator.interval` | `kafka.bootstrap-servers`、`kafka.topic`、`generator.interval` |
| 连不上目标时的降级 | 把 SQL 打印到控制台 | 把日志 JSON 打印到控制台 |
| 下游消费方式 | Flink CDC 读 MySQL 变更 | Flink Kafka Connector 读主题 |

Sources: [src/main/datagenerator/business_code/BusinessApplication.java:17-21]()；[src/main/datagenerator/userlog_code/UserLogApplication.java:27-51]()；[src/main/datagenerator/userlog_code/UserLogApplication.java:61-89]()

## 业务数据生成器：把全链路业务表写进 MySQL

### 启动与主循环

`BusinessApplication` 是标准的 Spring Boot 应用（`@SpringBootApplication`），启动后先做一次**基础数据**生成，然后进入无限循环，每一轮依次调用各类生成器，最后 `Thread.sleep(interval)` 控制节奏（默认 5 秒一轮）：

```java
@Override
public void run(String... args) throws Exception {
    baseDataGenerator.generateBaseData(batchSize);      // 一次性：分类/品牌/属性/地区/字典

    while (true) {
        userDataGenerator.generateUserData(batchSize / 10);
        productDataGenerator.generateProductData(batchSize / 10, batchSize / 20);
        activityDataGenerator.generateActivityData(batchSize / 10, batchSize / 20);
        couponDataGenerator.generateCouponData(batchSize / 10, batchSize / 20);
        orderDataGenerator.generateOrderData(batchSize);       // 每轮量最大的交易数据
        userBehaviorGenerator.generateUserBehaviorData(batchSize);
        warehouseDataGenerator.generateWarehouseData(batchSize / 5);
        cmsDataGenerator.generateCMSData(batchSize / 20, batchSize / 40, batchSize / 100);
        Thread.sleep(interval);
    }
}
```

Sources: [src/main/datagenerator/business_code/BusinessApplication.java:57-78]()

### 每轮覆盖哪些数据

在默认 `batch-size=1000` 下，每个循环周期各类生成器会新增大约下面这些行（各子表数量基本随主参数等比缩放）：

| 生成器 | 覆盖的数据/表 | 每轮新增量（默认值） |
| --- | --- | --- |
| `BaseDataGenerator` | 三级分类、品牌、商品属性、销售属性、省份地区、字典、前端参数 | 启动时一次性生成 |
| `UserDataGenerator` | `user_info`、`user_address` | 各约 100 行 |
| `ProductDataGenerator` | `spu_info`、`sku_info` 及其图片/属性/销售属性 | 约 100 个 SPU + 50 个 SKU |
| `ActivityDataGenerator` | `activity_info`、活动规则、活动商品、秒杀 | 约 100 个活动 |
| `CouponDataGenerator` | `coupon_info`、券范围、领券记录 | 约 100 张券 |
| `OrderDataGenerator` | 订单主表 + 明细 + 支付 + 状态流水 + 优惠分摊 + 退款 | 每个子表约 1000 行 |
| `UserBehaviorGenerator` | `cart_info` 购物车、`comment_info` 评价、`favor_info` 收藏 | 每个子表约 1000 行 |
| `WarehouseDataGenerator` / `CMSDataGenerator` | 仓储库存 / Banner、主题、评论 | 相对少量 |

其中 `OrderDataGenerator` 内部又串了订单主表、订单明细、支付信息、订单状态流水、明细-活动、明细-优惠券、退款单、退款流水 8 张表；金额上会用"原价 − 活动优惠 − 券优惠 + 运费"拼出 `total_amount`，让数值之间自洽。

Sources: [src/main/datagenerator/business_code/generator/OrderDataGenerator.java:22-31]()、[OrderDataGenerator.java:49-56]()；[src/main/datagenerator/business_code/generator/UserBehaviorGenerator.java:24-28]()；[src/main/datagenerator/business_code/README.md:6-45]()、[README.md:91-105]()

### 主键续接与"随机外键"

每个生成方法开头都会先查这张表的当前最大主键，再接着往下排：

```java
String maxIdSql = "SELECT COALESCE(MAX(id), 0) FROM base_category1";
int startId = dbUtil.queryForInt(maxIdSql) + 1;
```

Sources: [src/main/datagenerator/business_code/generator/BaseDataGenerator.java:38-41]()；[src/main/datagenerator/business_code/generator/UserDataGenerator.java:27-30]()；[src/main/datagenerator/business_code/util/DbUtil.java:161-178]()

这样做的效果是：**主键单调递增、进程重启不冲突**，配合 CDC 增量同步时能持续产出"新行"。但要特别说明：表之间的关联（订单→用户、明细→SKU、地址→省份）不是真正去查已存在的主数据，而是用随机数落进一个区间，例如 `user_id = RandomUtil.generateNumber(1, 1000)`、省份 ID 取 1–34。所以这是一种"概率意义上的弱外键"——数据形态看起来真实，但同一轮生成的行之间并不保证严格可 join。仓库 README 中"SKU 关联 SPU、订单关联用户"等说法，本质是靠这个随机区间近似实现的。

Sources: [src/main/datagenerator/business_code/generator/OrderDataGenerator.java:63]()、[OrderDataGenerator.java:108-109]()、[OrderDataGenerator.java:139]()；[src/main/datagenerator/business_code/generator/UserDataGenerator.java:73-76]()；[src/main/datagenerator/business_code/README.md:100-105]()

### DbUtil：连接池、批量写与"无库可写"降级

`DbUtil` 用 HikariCP 维护一个连接池，连接串、账号、池大小等全部来自 `spring.datasource.*` 配置。批量写通过 `PreparedStatement.addBatch()` + `executeBatch()` 一次提交：

Sources: [src/main/datagenerator/business_code/util/DbUtil.java:60-71]()、[DbUtil.java:88-113]()

最有意思的是它的**降级设计**：如果启动时连不上数据库，或某次批量插入抛 SQL 异常，它不会让进程崩溃，而是把 SQL 和这批数据以可读文本打印到控制台（`printToConsole`），生成流程照常推进：

```java
public void batchInsert(String sql, List<Object[]> params) {
    if (!isConnected) {
        printToConsole(sql, params);   // 数据库不可用 → 把"将要执行的 SQL"打出来
        return;
    }
    ...
}
```

Sources: [src/main/datagenerator/business_code/util/DbUtil.java:88-113]()、[DbUtil.java:115-136]()

对教学/演示环境来说，这意味着即使还没有 MySQL，也能先看到程序"会写什么"；对正式跑数仓来说，则要求配置好连接，否则数据只落在控制台而进不了库。

### RandomUtil：仿真感的来源

"像不像真实数据"主要靠 `RandomUtil` 里的词库与取数规则：中文姓名按"张/王/李… × 伟/芳/娜…"拼接，昵称用"快乐/阳光… × 小天使/小星星…"，手机号按 1 开头加 10 位数字，价格保留两位小数，日期往过去随机偏移形成历史分布。分类、品牌、地区也都有固定的候选数组做轮询或随机选取。

Sources: [src/main/datagenerator/business_code/util/RandomUtil.java:15-24]()、[RandomUtil.java:26-78]()、[RandomUtil.java:100-106]()

## 用户日志生成器：把嵌套 JSON 发往 Kafka

### 一条日志的结构

`UserLogGenerator.generateLog()` 每次组装一条完整的日志对象，模型用 `UserLog` 承载，包含 `common`（公共信息：地区/品牌/渠道/设备 ID/用户 ID/OS/版本）、`start`（启动信息：入口/加载时长/广告）、`page`（页面信息：停留时长/页面 ID/来源类型）、`actions`（行为数组）、`displays`（曝光数组）、`err`（错误，多数为空）、`ts`（毫秒时间戳）：

```java
public static UserLog generateLog() {
    UserLog log = new UserLog();
    log.setCommon(generateCommon());
    log.setStart(generateStart());
    log.setPage(generatePage());
    log.setActions(generateActions());
    log.setDisplays(generateDisplays());
    log.setErr(generateError());
    log.setTs(Instant.now().toEpochMilli());
    return log;
}
```

Sources: [src/main/datagenerator/userlog_code/model/UserLog.java:9-15]()；[src/main/datagenerator/userlog_code/generator/UserLogGenerator.java:21-31]()；[src/main/datagenerator/userlog_code/model/Common.java:6-15]()

数组字段的数量也做成随机的：`actions` 是 0–4 条随机行为，`displays` 是 0–9 条曝光，每条含 `item_id`、`item_type`、曝光位等。`err` 只在约 5% 的日志里出现（`random.nextInt(100) < 95` 时为空），模拟偶发错误。

Sources: [src/main/datagenerator/userlog_code/generator/UserLogGenerator.java:68-84]()、[UserLogGenerator.java:86-103]()、[UserLogGenerator.java:105-113]()

### 嵌套数组怎么变成"真 JSON"

这里有个值得一提的实现细节：模型里 `actions`、`displays` 字段声明成 `String`，Jackson 序列化时会把它们当普通字符串，结果是"JSON 里套着被转义的一长串文本"。所以 `UserLogApplication` 发送前用 `processNestedJson` 补救：先把对象变成 JSON 节点，再把 `actions`/`displays` 的字符串内容重新解析成真正的数组节点后写回，最后整体序列化输出：

```java
// 先将对象转换为 JSON 节点
JsonNode rootNode = objectMapper.valueToTree(log);
// 把 actions 字符串解析回 JSON 数组，覆盖原字段
JsonNode actionsNode = objectMapper.readTree(actionsStr);
((ObjectNode) rootNode).set("actions", actionsNode);
```

Sources: [src/main/datagenerator/userlog_code/UserLogApplication.java:105-131]()；[src/main/datagenerator/userlog_code/generator/UserLogGenerator.java:79-83]()、[UserLogGenerator.java:98-102]()

经过这一步，下游拿到的就是 `actions`、`displays` 为原生数组的 JSON，而不是带反斜杠转义的字符串。

### 发送节奏与降级

`UserLogApplication` 启动时先按配置建 `KafkaProducer`（`StringSerializer`、`acks=all`、`retries=3`），并调用 `partitionsFor(topic)` 探活；连不上 Kafka 就标记为本地模式。主循环每 `interval` 毫秒生成一条日志：Kafka 可用则 `send(...).get()` **同步等待发送结果**（保证消息落盘），失败则把该条日志打印到控制台；Kafka 不可用则一直走本地打印。

Sources: [src/main/datagenerator/userlog_code/UserLogApplication.java:45-64]()、[UserLogApplication.java:66-92]()

也就是说，默认配置下它的吞吐并不高（一条一条来、每条约隔 1 秒，同步确认），`userlog_code/README.md` 里"批量发送、异步提高吞吐"的描述和当前代码的实际行为并不一致——代码里是逐条同步发送，以本页代码为准。演示环境可以靠调小 `generator.interval` 来提高数据密度。

Sources: [src/main/datagenerator/userlog_code/UserLogApplication.java:71-91]()；[src/main/datagenerator/userlog_code/README.md:108-122]()

## 小结

整套仿真数据的源头是 `src/main/datagenerator` 下两个独立的 Spring Boot 进程：`BusinessApplication` 用"先打基础数据、再按固定间隔无限循环"的方式，把电商交易链路上的分类、商品、用户、营销、订单、购物车、仓储等几十张表的数据批量写进 MySQL 的 `gmall` 库——主键通过 `MAX(id)` 续接、跨表靠随机数区间近似关联、数据库不可用时降级为控制台打印 SQL；`UserLogApplication` 则把 `common/start/page/actions/displays/err/ts` 组成的嵌套 JSON 日志逐条同步发往 Kafka 主题，并在发送前把嵌套数组还原成真正的 JSON 数组，Kafka 不可用时同样降级为本地打印。MySQL 里的业务表随后由 Flink CDC 同步进 ODS，Kafka 里的日志流由 Flink Kafka Connector 接入 ODS，两条数据线在此汇合，共同支撑后续 DIM/DWD/DWS/ADS 的数仓分层。理解这两个生成器，就能明白湖仓每一条"看起来像真实电商"的明细数据从何而来，也清楚它们的节奏、关联强度与降级边界。

Sources: [README.md:1-11]()；[src/main/datagenerator/business_code/BusinessApplication.java:57-78]()；[src/main/datagenerator/userlog_code/UserLogApplication.java:44-92]()；[src/main/java/com/fy/warehouse/ODS/ods_log_inc.java:24-45]()；[src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java:21-47]()

---

## 03. 湖仓五层结构与实时数据流

> 沿着数据流讲清湖仓分层：日志经 Kafka Connector、业务数据经 MySQL CDC 两条链路进入 ODS，再经 DIM 维度加工、DWD 明细与枚举解码、DWS 近 N 日聚合、ADS 指标产出，落成可按天分区、可主键更新的 Paimon 表。

- Page Markdown: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/pages/03-page-3.md
- Generated: 2026-09-08T05:47:07.309Z

### Source Files

- `src/main/java/com/fy/warehouse/ODS/ods_log_inc.java`
- `src/main/java/com/fy/warehouse/DIM/dim_sku_full.java`
- `src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java`
- `src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java`
- `src/main/java/com/fy/warehouse/dwd/dwd_traffic_action_full.java`
- `src/main/java/com/fy/warehouse/dws/dws_trade_user_sku_order_nd_full.java`
- `src/main/java/com/fy/warehouse/ads/ads_coupon_stats_full.java`

<details>
<summary>相关源文件</summary>
以下文件用于生成此维基页面：
- [ODS/ods_log_inc.java](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)
- [ODS/ods_order_info_full.java](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java)
- [DIM/dim_sku_full.java](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java)
- [DIM/dim_user_zip_full.java](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java)
- [dwd/dwd_trade_order_detail_full.java](src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java)
- [dwd/dwd_traffic_action_full.java](src/main/java/com/fy/warehouse/dwd/dwd_traffic_action_full.java)
- [dws/dws_trade_user_sku_order_nd_full.java](src/main/java/com/fy/warehouse/dws/dws_trade_user_sku_order_nd_full.java)
- [dws/dws_trade_coupon_order_nd_full.java](src/main/java/com/fy/warehouse/dws/dws_trade_coupon_order_nd_full.java)
- [ads/ads_coupon_stats_full.java](src/main/java/com/fy/warehouse/ads/ads_coupon_stats_full.java)
- [ads/ads_activity_stats_full.java](src/main/java/com/fy/warehouse/ads/ads_activity_stats_full.java)
- [config/FlinkConfigUtil.java](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)
- [udf/JsonActionsArrayParser.java](src/main/java/com/fy/warehouse/udf/JsonActionsArrayParser.java)
</details>

# 湖仓五层结构与实时数据流

本文沿着一条数据的完整旅程，讲解这个基于 Flink + Paimon 的电商实时湖仓如何分层：用户行为日志与 MySQL 业务数据从两条入口进入 ODS 原始层，经过 DIM 维度加工、DWD 明细落地与枚举解码、DWS 近 N 日聚合，最终在 ADS 产出可直接用于分析的指标。全链路作业都是 Java 里拼装的 Flink SQL，每一个目标表都是「按天分区、可主键更新」的 Paimon 表。

理解这条链路的意义在于：分层不是目录摆设，而是每一层解决一个确定的问题——ODS 保真、DIM 补语义、DWD 建模到业务过程、DWS 预聚合、ADS 出指标。五层共用同一套 Paimon 建表范式，数据在层间靠 SQL 流转，链路长但规律一致。

## 数据入口与湖仓分层总览

仓库把作业按包名划分层次：`ODS`、`DIM`、`dwd`、`dws`、`ads`，另有一个 `config` 提供统一 Flink 配置、一个 `udf` 提供自定义函数。数据流自下而上，越往上越贴近指标口径。

```mermaid
flowchart LR
    subgraph 源端
        K[Kafka<br/>topic ODS_BASE_LOG]
        M[(MySQL<br/>gmall 业务库)]
    end
    subgraph ODS[ODS 原始数据层]
        L[ods_log_inc 用户行为日志]
        F1[ods_order_info_full 等业务全量表]
    end
    subgraph DIM[DIM 维度层]
        D1[dim_sku_full 商品维度]
        D2[dim_user_zip_full 用户拉链]
    end
    subgraph DWD[DWD 明细层]
        T[dwd_trade_order_detail_full 交易明细]
        A[dwd_traffic_action_full 行为动作明细]
    end
    subgraph DWS[DWS 汇总层]
        S1[dws_trade_user_sku_order_nd_full]
        S2[dws_trade_coupon_order_nd_full]
    end
    subgraph ADS[ADS 应用层]
        AD[ads_coupon_stats_full 补贴率]
    end
    K -->|Kafka Connector| L
    M -->|Flink CDC| F1
    F1 --> D1
    F1 --> D2
    F1 --> T
    L --> A
    D1 -->|补充维度属性| T
    D1 -->|补充维度属性| S1
    T --> S1
    T --> S2
    D1 -->|优惠券维度| S2
    S2 --> AD
```

## ODS 原始数据层：两条入口、一个落点

ODS 的作用是「原样接入」，把不同来源的数据先落地成 Paimon 表。链路有两条，对应仓库里 `*_inc` 与 `*_full` 两类表。

### 日志链路：Kafka Connector

模拟 App 端产生的用户行为日志先写入 Kafka 主题 `ODS_BASE_LOG`。作业里先注册一个 Kafka 源表，用 `METADATA` 语法把分区、offset、时间戳带出来，`common`/`start`/`page` 是 ROW 嵌套结构，`actions`/`displays` 以字符串原样保留（为后续 UDF 解析留口子），开启 `earliest-offset` 与容错解析。 Sources: [ODS/ods_log_inc.java:24-47](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)

随后执行 `INSERT ... SELECT` 把嵌套字段「拍平」成一张 ODS 宽表：`k1` 列由事件时间戳 `ts` 转成 `yyyy-MM-dd` 字符串作为天级业务日期，`id` 由分区+offset+时间戳拼接而成。日志表的物理分区未开启，但主键与 `k1` 仍保留，业务上仍可按天管理。 Sources: [ODS/ods_log_inc.java:87-122](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)、[ODS/ods_log_inc.java:129-162](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)

### 业务数据链路：MySQL CDC

电商业务数据存放在 MySQL 的 `gmall` 库中。以订单为例，源表用 `mysql-cdc` 连接器声明，`scan.startup.mode = 'initial'` 表示先做一次全量快照再持续读取增量 binlog，这解释了为什么同步出来的业务表多叫 `*_full`——它们是「当前全量 + 后续变更」的语义。 Sources: [ODS/ods_order_info_full.java:22-59](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java)

写入 Paimon 时，`create_time` 被格式化成 `k1` 作为分区值，主键为 `(id, k1)`，即同一张 MySQL 表的同一主键在当天内更新、跨天则自然分开。ODS 层的三十余张表覆盖订单、订单明细、商品(SKU/SPU)、分类、品牌、优惠券、优惠券使用、支付、退款、收藏、评价、省份地区、字典等，几乎与 MySQL 表一一对应。 Sources: [ODS/ods_order_info_full.java:78-116](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java)、[ODS/ods_order_info_full.java:121-139](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java)

两条链路对比如下：

| 维度 | 日志链路 | 业务数据链路 |
| --- | --- | --- |
| 源端 | Kafka 主题 `ODS_BASE_LOG` | MySQL 的 `gmall` 库多张业务表 |
| 接入方式 | Kafka Connector，JSON 格式 | MySQL CDC，initial 模式先全量后增量 |
| 主键来源 | Kafka 分区/offset/时间戳拼接 | MySQL 表原主键 + `k1` |
| 表命名 | `ods_log_inc`（增量） | `ods_*_full`（全量快照） |
| 数据形态 | 嵌套日志拍平成宽表 | 与源表字段近似对齐 |

## DIM 维度层：给事实补语义

DIM 层把散落在多张 ODS 业务表里的维度信息加工成可直接 join 的维表，并为下游 DWD/DWS 提供「一行一实体」的查询视图。

### 商品维度：一张宽表收纳分类与品牌

`dim_sku_full` 把 SKU 主信息与 SPU、一/二/三级分类、品牌、平台属性、销售属性全部合并进一张表。SQL 用多个 `LEFT JOIN` 把各自主键表带进来：`ods_sku_info_full` 为事实主体，`ods_spu_info_full` 提供 spu 名，三级分类表逐级向上带出二级、一级分类名称，品牌表补 `tm_name`；平台属性与销售属性则先各自 `GROUP BY sku_id` 后用 `collect(id)` 聚成数组再转字符串挂到行上。这属于典型的维度退化——把多级维度提前摊平，让下游明细与汇总层免去多次关联。 Sources: [DIM/dim_sku_full.java:70-159](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java)

### 用户维度：拉链与脱敏

`dim_user_zip_full` 是拉链表：字段带 `start_date`/`end_date` 两个 CHAR(10) 边界，初始导入把 `start_date` 设为 `operate_time` 当天、`end_date` 设为 `9999-12-31` 表示当前有效。取数上它先从 `ods_user_info_full` 按用户用 `row_number()` 按操作时间倒序取最新一行（`rk = 1`），避免同一用户多条变更记录撑爆维度；对姓名、手机号、邮箱这些敏感字段则用 `md5()` 脱敏后再入库。 Sources: [DIM/dim_user_zip_full.java:33-59](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java)、[DIM/dim_user_zip_full.java:65-86](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java)

## DWD 明细层：按业务过程建模并解码

DWD 把 ODS 的「贴源」数据重组成面向业务过程的事实明细，同时完成两类加工：关联补齐上下文、把编码翻译成人话。仓库里 DWD 至少覆盖交易、流量、互动、工具(优惠券)四个域。

### 交易域：明细宽表与金额口径

`dwd_trade_order_detail_full` 是交易域的核心明细。它从 `ods_order_detail_full` 出发，先当场计算分摊前金额 `sku_num * order_price` 作为 `split_original_amount`，再分别关联订单主表取 `user_id`/`province_id`、关联 `ods_order_detail_activity_full` 取活动、关联 `ods_order_detail_coupon_full` 取优惠券，最终形成一条含原始金额、活动优惠、优惠券优惠、最终金额四段金额的订单明细。 Sources: [dwd/dwd_trade_order_detail_full.java:69-141](src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java)

「枚举解码」在这张表里体现得很直接：订单来源只存了 `source_type` 编码，SQL 用 `ods_base_dic_full` 字典表（限定 `parent_code='24'` 的来源类型分组）把编码翻译成 `dic_name`，落到 `source_type_name` 字段，让下游直接用中文名称而无需再查字典。 Sources: [dwd/dwd_trade_order_detail_full.java:134-141](src/main/java/com/fy/warehouse/dwd/dwd_trade_order_detail_full.java)

### 流量域：UDF 解析嵌套行为

日志里每个用户会有多个动作，ODS 层把 `actions` 整体存成了 JSON 字符串。`dwd_traffic_action_full` 先注册自定义函数 `json_actions_array_parser`，再对每条日志取数组首个元素，拆出 `action_id`、`item`、`item_type`、`ts` 四个字段，`ts` 同时被格式化成 `date_id`（天）与 `action_time`（秒级）两个时间列。解析后再按 `common_ar` 区域码关联省份维表 `ods_base_province_full` 换出 `province_id`。 Sources: [dwd/dwd_traffic_action_full.java:72-156](src/main/java/com/fy/warehouse/dwd/dwd_traffic_action_full.java)

自定义函数本身是一个 Flink `ScalarFunction`：把输入按 JSON 数组解析并取下标 0 的元素，逐字段取值；遇到空串、非数组或解析异常时不抛错，而是返回一组默认值（空串/0），保证流作业不会因单条脏数据中断。 Sources: [udf/JsonActionsArrayParser.java:16-52](src/main/java/com/fy/warehouse/udf/JsonActionsArrayParser.java)

## DWS 汇总层：近 N 日聚合

DWS 按分析粒度把 DWD 明细预聚合，仓库中表名带 `_nd` 表示近 N 日（本仓库统一用 30 日）。以 `dws_trade_user_sku_order_nd_full` 为例，它把 DWD 交易明细限制在 30 天窗口内（用 `k1` 与一个基准日做 `timestampdiff` 过滤），再以 `(user_id, sku_id)` 为粒度 join 商品维表 `dim_sku_full` 补上 sku 名称、三级分类、品牌等描述属性，最后 `group by` 用户+商品+属性后一次性算出：下单次数 `count(order_id)`、下单件数 `sum(sku_num)`、下单原始金额、活动优惠、优惠券优惠、下单最终金额共六类 30 日指标。 Sources: [dws/dws_trade_user_sku_order_nd_full.java:64-124](src/main/java/com/fy/warehouse/dws/dws_trade_user_sku_order_nd_full.java)

同类写法也出现在优惠券维度上：`dws_trade_coupon_order_nd_full` 用优惠券维表 `dim_coupon_full` 关联 DWD 明细中 `coupon_id` 不为空的行，聚出某张券近 30 日带来的原始金额与优惠金额。 Sources: [dws/dws_trade_coupon_order_nd_full.java:56-90](src/main/java/com/fy/warehouse/dws/dws_trade_coupon_order_nd_full.java)

## ADS 应用层：从汇总出指标

ADS 直接消费 DWS 汇总结果，产出面向报表的窄表。两个 ADS 作业结构几乎一致：`ads_coupon_stats_full` 把优惠券 30 日汇总转成补贴率口径，`reduce_rate = coupon_reduce_amount_30d / original_amount_30d`，即平台为这张券实际补贴的比例；同时带上券名、发布日期、优惠规则等解释性字段，形成一张可直接看数的表。 Sources: [ads/ads_coupon_stats_full.java:30-58](src/main/java/com/fy/warehouse/ads/ads_coupon_stats_full.java)

活动侧对应 `ads_activity_stats_full`，从 `dws_trade_activity_order_nd_full` 算活动补贴率 `activity_reduce_amount_30d / original_amount_30d`，指标口径与优惠券完全对称。值得一提的是 ADS 表不再按天分区，只保留 `dt` 统计日期列，因为这一层是「结果」而非「过程」。 Sources: [ads/ads_activity_stats_full.java:31-49](src/main/java/com/fy/warehouse/ads/ads_activity_stats_full.java)

## 贯穿五层的工程约定

五层代码读下来，能抽出一套完全一致的工程范式，这也是本项目最值得复用的部分。

| 约定 | 实现 | 出处 |
| --- | --- | --- |
| 统一 Catalog | `paimon_hive`：metastore 走 Hive thrift，warehouse 指向 HDFS | 几乎每个作业的开头 |
| 天级分区 | 业务表 `PARTITIONED BY (k1)`，`k1` 为 `yyyy-MM-dd`；主键带 `k1` | [DIM/dim_sku_full.java:54-64](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java) |
| 主键更新 | 每张表声明 `PRIMARY KEY (...) NOT ENFORCED`，靠 Paimon 主键语义完成更新 | [ODS/ods_order_info_full.java:105-106](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java) |
| 分区过期 | `partition.expiration-time = 1 d`，按 `k1` 每天清理过期分区 | [dwd/dwd_traffic_action_full.java:56-66](src/main/java/com/fy/warehouse/dwd/dwd_traffic_action_full.java) |
| 写入参数 | parquet 格式、512MB 可溢写写缓冲，配合全局 Paimon sink 参数 | [dws/dws_trade_user_sku_order_nd_full.java:49-59](src/main/java/com/fy/warehouse/dws/dws_trade_user_sku_order_nd_full.java) |
| 表命名 | `*_inc` 增量日志、`*_full` 全量快照、`*_nd` 近 N 日、`_zip` 拉链 | 见各层源文件 |

运行环境相关的调优集中在 `FlinkConfigUtil`：流式执行并默认并行度 1；开启 60 秒间隔的 `EXACTLY_ONCE` Checkpoint、状态后端落 HDFS；状态 TTL 24 小时；开启 500ms/1000 条微批；时区固定 `Asia/Shanghai`。这些配置对「少量作业、教学/演示集群、以天分区回刷为主」的场景是自洽的：并行度低便于观察，Checkpoint 兜底故障恢复，微批与 `paimon.sink.batch-size` 控制小文件写入节奏。 Sources: [config/FlinkConfigUtil.java:21-58](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)

## 小结

这个湖仓的全部逻辑可以概括为一句话：Flink SQL 把 Kafka 日志与 MySQL 业务数据分别接入 ODS，随后五层之间靠「查上一层、写下一层」的 INSERT...SELECT 完成维度退化、明细建模、枚举解码、近 30 日聚合与指标产出，所有中间与结果表统一落在按 `k1` 天级分区、可主键更新的 Paimon 表上。ODS 保真、DIM 补语义、DWD 建模到业务过程、DWS 预聚合、ADS 出指标——每一层职责单一，层与层只通过表契约交互，这正是这套代码可以横向扩展（新增域、新增指标）而无需推翻结构的原因。

---

## 04. 计算存储调优、业务实现与经验收尾

> 收尾页：从统一配置入口出发解读 Flink 状态与检查点、微批、Paimon 文件策略等系统优化点，说明元数据在 Hive Metastore、数据落 HDFS 的存储形态，介绍 StarRocks 外部分区目录直读与可视化，最后把交易、流量、用户等业务域如何落到各层做一次收束，并沉淀可复用经验。

- Page Markdown: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/pages/04-page-4.md
- Generated: 2026-09-08T05:45:45.250Z

### Source Files

- `src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java`
- `src/main/resources/core-site.xml`
- `src/main/resources/hdfs-site.xml`
- `pom.xml`
- `dependency-reduced-pom.xml`
- `StarRocks.sql`

<details>
<summary>相关源文件</summary>

以下文件用于生成此维基页面:

- [src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)
- [src/main/resources/core-site.xml](src/main/resources/core-site.xml)
- [src/main/resources/hdfs-site.xml](src/main/resources/hdfs-site.xml)
- [pom.xml](pom.xml)
- [StarRocks.sql](StarRocks.sql)
- [src/main/java/com/fy/warehouse/ODS/ods_log_inc.java](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)
- [src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java)
- [src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java)
- [src/main/java/com/fy/warehouse/DIM/dim_sku_full.java](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java)
- [src/main/java/com/fy/warehouse/dwd/dwd_trade_pay_detail_suc_full.java](src/main/java/com/fy/warehouse/dwd/dwd_trade_pay_detail_suc_full.java)
- [src/main/java/com/fy/warehouse/dwd/dwd_user_login_full.java](src/main/java/com/fy/warehouse/dwd/dwd_user_login_full.java)
- [src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java](src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java)
- [src/main/java/com/fy/warehouse/dws/dws_traffic_page_visitor_page_view_nd_full.java](src/main/java/com/fy/warehouse/dws/dws_traffic_page_visitor_page_view_nd_full.java)
- [src/main/java/com/fy/warehouse/ads/ads_activity_stats_full.java](src/main/java/com/fy/warehouse/ads/ads_activity_stats_full.java)
- [src/main/java/com/fy/warehouse/udf/JsonActionsArrayParser.java](src/main/java/com/fy/warehouse/udf/JsonActionsArrayParser.java)
- [src/main/datagenerator/business_code/README.md](src/main/datagenerator/business_code/README.md)
- [src/main/datagenerator/userlog_code/README.md](src/main/datagenerator/userlog_code/README.md)

</details>

# 计算存储调优、业务实现与经验收尾

这是电商实时数据湖仓系列的收尾页。前面的页面分别讲清了架构、各数据分层的数据流与明细实现,本页把这些线索合到一处:先从代码里唯一的统一配置入口出发,解释 Flink 状态与检查点、微批、Paimon 文件策略这些"计算与存储调优点"到底解决了什么问题;再说明"元数据进 Hive Metastore、数据落 HDFS"的双层存储形态,以及 StarRocks 如何以外部分区表的方式直读并交给 DataRT 可视化;最后把交易、流量、用户、互动、工具等业务域在 ODS→DIM→DWD→DWS→ADS 每层落成哪些表做一次收束,沉淀一批可直接复用的经验。读完本页,应该能回答"这套湖仓为什么这样配、每个业务域的指标是怎么一层层算出来的"。

仓库的全部 Flink 作业共用同一份配置,因此看懂了它,就同时看懂了所有作业的运行与存储方式。

## 统一配置入口:状态、检查点、微批与文件策略

所有作业的入口逻辑一致:`main` 方法先取 `FlinkConfigUtil.getFlinkConfig()` 得到一份带满默认项的 `Configuration`,再以 `.withConfiguration(config)` 创建流式 `TableEnvironment`,之后所有 DDL/DML 都跑在这个环境上([FlinkConfigUtil.java:9-23](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java))。这份集中配置是按"边踩坑边修"的方式沉淀出来的,注释里保留了大量"关键修改"标记,本身就是最有价值的排障记录。

### Flink 状态与检查点

配置里最核心的一组开关围绕"状态能否持久化、能否恢复":

| 配置项 | 取值 | 作用与背景 |
| --- | --- | --- |
| `execution.checkpointing.enabled` | `true` | 注释明确:原为 `false`,导致状态无法持久化,这是第一个被修复的点 |
| `state.backend` | `filesystem` | 状态后端选文件系统;注释提示大状态可换 RocksDB(并预留了 `rocksdb` 与对应 checkpoint 目录两行) |
| `execution.checkpointing.interval` | `60s` | 检查点间隔 |
| `execution.checkpointing.timeout` | `300s` | 5 分钟超时,避免检查点长时间卡住 |
| `execution.checkpointing.tolerable-failure-number` | `3` | 容忍 3 次连续失败,网络抖动不至于直接让作业失败 |
| `execution.checkpointing.mode` | `EXACTLY_ONCE` | 精确一次语义 |
| `table.exec.state.ttl` | `86400000` | 状态 TTL 24 小时,控制状态体积 |
| `taskmanager.memory.process.size` / `task.heap.size` | `1024m` / `512m` | 明确 TaskManager 进程与堆内存,配合文件系统状态后端控制状态占用 |

关键结论:检查点目录直接指到 HDFS(`hdfs://192.168.10.102:8020/flink/checkpoints`),状态与检查点都落在分布式文件系统上;这正是"解决了状态超限"的手段——状态不再只存在于 TaskManager 本地内存,而是周期性快照到 HDFS,作业重启后可恢复。

### 微批(Mini-Batch)配置

```java
config.setString("table.exec.mini-batch.enabled", "true");
config.setString("table.exec.mini-batch.allow-latency", "500ms");
config.setString("table.exec.mini-batch.size", "1000");
```

微批把同一 key 上 500ms 内到达、或攒够 1000 条的数据攒批处理,让聚合、去重等高成本算子减少逐条触发,以亚秒级延迟换取吞吐。由于运行时是 `STREAMING`(见 21 行),微批是在流上叠加的批化优化,不是改成批处理。

### Paimon 侧的文件写入策略

```java
config.setString("paimon.sink.batch-size", "100");   // 小批次快速写入
config.setString("paimon.sink.buffer-time", "0s");   // 关闭缓冲等待
config.setString("paimon.file.size", "64mb");        // 小文件即可,避免合并
```

这三行体现了教学场景的文件策略取舍:批量 100、缓冲 0s 是为了"来一条写一条、尽快可见",`file.size = 64mb` 的注释直说"小文件即可,避免合并",即不做 aggressive 的 small-file 合并,换取写入即时性和实现简单。这与生产库追求大文件、低文件数的方向不同,是本项目定位决定的。

### 与 Hive / HDFS 协同的必要配置

同一份配置里还埋了几个容易踩的环境类配置:Hive 版本必须与集群一致(`hive.version=3.1.3`);开启 Hive 兼容读取器;关闭分区剪枝以免 Flink 扫描分区卡住;把 `HADOOP_USER_NAME` 与 `hadoop.user.name` 设为 `root` 并对齐 HDFS 地址(`fs.defaultFS`)、Hive Metastore(`thrift://192.168.10.102:9083`);HDFS RPC 超时通过 `env.java.opts` 注入(`dfs.client.socket-timeout`、`ipc.client.connect.timeout`、`ipc.client.connect.max.retries=5`),注释记录这是为"解决 RPC 中断"加的。也就是说,一份配置同时决定了**运行时语义(流/检查点/微批)、落盘文件行为(Paimon)、以及外围依赖(Hive/HDFS)怎么连**。Sources:[FlinkConfigUtil.java:24-73](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)

## 存储形态:元数据在 Hive Metastore,数据落 HDFS

湖仓的表不把数据存在 Hive 表里,而是"借壳":建一个 Paimon Catalog,`metastore = hive` 指向 Hive Metastore,`warehouse` 指向 HDFS 路径——表结构(库、表、字段、分区元数据)注册进 Hive Metastore,真正的列式数据文件写在 HDFS 的 warehouse 目录下,由 Paimon 管理:

```sql
CREATE CATALOG paimon_hive WITH (
    'type' = 'paimon',
    'metastore' = 'hive',
    'uri' = 'thrift://192.168.10.102:9083',
    'warehouse' = 'hdfs://192.168.10.102/user/hive/warehouse'
);
```

这段 Catalog DDL 在几乎每个作业里原样出现(如 [ods_log_inc.java:72-79](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)、[dim_user_zip_full.java:19-24](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java))。落到 Paimon 表的数据文件统一用 `parquet` 格式,并普遍开启两类存储参数:写缓冲 `write-buffer-size = 512mb` + `write-buffer-spillable = true`(缓冲写不下就溢写到磁盘);以及**分区生命周期管理**:

```java
'partition.expiration-time' = '1 d',
'partition.expiration-check-interval' = '1 h',
'partition.timestamp-formatter' = 'yyyy-MM-dd',
'partition.timestamp-pattern' = '$k1'
```

分区字段是每个表都有的 `k1`(业务日期字符串,如 `2024-05-01`),Paimon 按 `$k1` 从分区名解析时间戳,1 小时检查一次、自动清理超过 1 天的分区,避免 HDFS 上堆积过多历史分区。ODS 层业务表建表时即声明 `PRIMARY KEY (id, k1) NOT ENFORCED` 与 `PARTITIONED BY (k1)`,形成"主键 + 日期分区"的通用骨架([ods_order_info_full.java:78-116](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java))。

配套的 Hadoop 环境配置在 XML 中:`core-site.xml` 声明默认文件系统为 `hdfs://hadoop102:8020`、本地数据目录,并把 `root` 配成全局可代理用户(`hadoop.proxyuser.root.*` 均为 `*`,对应代码里统一用 root 访问 HDFS);`hdfs-site.xml` 配置 NameNode/SecondaryNameNode 的 Web 地址与节点白名单、黑名单。换言之:**代码里的 IP/端口与 XML 里的集群地址是一体两面**,迁移集群时两处要同步改。Sources:[core-site.xml:20-49](src/main/resources/core-site.xml)、[hdfs-site.xml:20-39](src/main/resources/hdfs-site.xml)

从依赖侧看,工程同时引入 Paimon 的 Flink 连接器与 Hive 3.1 连接器(`paimon-flink-1.17` + `paimon-hive-connector-3.1`,见 [pom.xml:91-100](pom.xml)),并统一了 Flink 1.17.2、Hadoop 3.1.3、Hive 3.1.3、MySQL CDC 2.4.2 等版本;打包用 shade 插件合并 `META-INF/services`,正是为了让 Hive Catalog 等工厂类在打成 fat jar 后仍能被 SPI 找到([pom.xml:205-245](pom.xml))。

## StarRocks 外部分区表直读与可视化

Hive 元数据 + HDFS 数据的形态带来一个直接收益:**下游分析引擎可以不改数据、以外部目录方式直读**。仓库的 `StarRocks.sql` 演示了在 StarRocks 侧建一个 Paimon External Catalog:

```sql
CREATE EXTERNAL CATALOG paimon_catalog
PROPERTIES (
    'type' = 'paimon',
    'paimon.catalog.type' = 'hive',
    'warehouse' = 'hdfs://192.168.10.102:8020/user/hive/warehouse',
    'hive.metastore.uris' = 'thrift://192.168.10.102:9083',
    'hadoop.username' = 'root'
);
use ads;
show tables;
select * from ads_coupon_stats_full;
```

要点是 `paimon.catalog.type = hive` + 与 Flink 侧完全相同的 warehouse 与 metastore 地址:StarRocks 通过 Hive Metastore 拿到各层分区表的元数据,再按 HDFS 路径直接读取 Paimon 写的 parquet 分区数据,整个链路**不复制、不导入、不经过 HiveServer2**,做到了"建一个 Catalog 即可查数仓各层"。示例里直接 `use ads` 查 ADS 应用层表,正是面向可视化取数的场景——README 中说明可视化层由 StarRocks(MPP 查询)+ DataRT(可视化工具)构成,Flink 只负责把各层数据写进 Paimon,查询压力由 StarRocks 承接。Sources:[StarRocks.sql:1-13](StarRocks.sql)

## 业务域落地收束:从数据生成到各层表

要讲清"业务如何落到各层",先明确最上游的数据从哪来:仓库自带两套 Spring 模拟生成器。业务数据生成器面向 MySQL 的 `gmall` 库,按基础数据、商品、营销、交易、用户、仓储等模块持续写入([business_code/README.md](src/main/datagenerator/business_code/README.md));用户行为日志生成器输出 `common / start / page / actions / displays / err / ts` 结构的 JSON,发往 Kafka 主题 `ODS_BASE_LOG`([userlog_code/README.md](src/main/datagenerator/userlog_code/README.md))。

一张图概括全链路与各层归属:

```mermaid
flowchart LR
    subgraph SRC[数据源]
        M[(MySQL gmall 业务库)]
        K[(Kafka ODS_BASE_LOG)]
    end
    subgraph FLINK[Flink 流式作业]
        CDC[MySQL-CDC 同步]
        KV[Kafka Connector 读取]
    end
    subgraph LAKE[Paimon 湖仓 · Hive Metastore 元数据 + HDFS 数据]
        ODS[ODS 原始层]
        DIM[DIM 维度层]
        DWD[DWD 明细层]
        DWS[DWS 汇总层]
        ADS[ADS 应用层]
    end
    subgraph BI[可视化]
        SR[StarRocks External Catalog 直读]
        RT[DataRT 报表]
    end
    M --> CDC --> ODS
    K --> KV --> ODS
    ODS --> DIM
    ODS --> DWD
    DIM --> DWD
    DWD --> DWS --> ADS
    ADS --> SR --> RT
```

### ODS:两条采集路径

ODS 层是两类源的分岔口。日志类(`ods_log_inc`)用 Kafka Connector 订阅 `ODS_BASE_LOG`,source 侧把 `common/start/page` 声明成嵌套 `ROW`、`actions/displays` 暂存为字符串;落到 Paimon 时通过 `SELECT` 把嵌套列全部展开成 `common_ar、common_mid、page_during_time……` 这样的平铺大宽列,主键用 `kafka_partition+offset+timestamp` 拼接保证唯一([ods_log_inc.java:24-47](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)、[ods_log_inc.java:87-162](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java))。业务类(如 `ods_order_info_full`)则用 MySQL CDC(`scan.startup.mode = initial`,先全量后增量)把 `gmall.order_info` 等表搬到 Paimon,并按 `create_time` 生成 `k1` 天分区([ods_order_info_full.java:50-59](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java)、[ods_order_info_full.java:121-139](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java))。

### DIM:维表做宽表与拉链

DIM 层把分散在多个 ODS 业务表的属性收敛成可直接 join 的宽维表。`dim_sku_full` 是典型:**sku 主体 left join spu、三级/二级/一级分类、品牌、属性聚合**,一次拼出带全链路分类名的商品维表([dim_sku_full.java:64-140](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java) 附近,insert 段见文件尾部多表 left join);`dim_user_zip_full` 则是**拉链表**:对 `ods_user_info_full` 按 `id` 开窗取最新一条,`start_date` 取 operate_time 当天、`end_date` 置为 `9999-12-31`,并以 `md5()` 处理姓名、手机、邮箱等敏感字段([dim_user_zip_full.java:33-86](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java));`dim_date_full` 这类日期维表由 MySQL 侧 `dim_date` 表 CDC 同步而来([dim_date_full.java:20-39](src/main/java/com/fy/warehouse/DIM/dim_date_full.java))。

### DWD:按业务域建模的事实明细

DWD 是"业务域"归属最明显的一层。每个 Java 类开头的注释即声明域归属,汇总如下:

| 业务域 | 事实表 | 计算要点 |
| --- | --- | --- |
| 交易域 | 下单 `dwd_trade_order_detail_full`、支付成功 `dwd_trade_pay_detail_suc_full`、退单 `dwd_trade_order_refund_full`、退款支付成功 `dwd_trade_refund_pay_suc_full`、取消 `dwd_trade_cancel_detail_full`、加购 `dwd_trade_cart_add_full`、购物车快照 `dwd_trade_cart_full` | 以订单明细为主体,多路 left join 支付、订单、活动、优惠券、字典表,算各类**金额分摊**(如 `split_original_amount = sku_num * order_price`),并把支付类型、来源类型等编码关联字典表翻译成名称 |
| 流量域 | 页面浏览 `dwd_traffic_page_view_full`、启动 `dwd_traffic_start_full`、动作 `dwd_traffic_action_full`、曝光 `dwd_traffic_display_full`、错误 `dwd_traffic_error_full` | 由 `ods_log_inc` 展开而来,`actions/displays` 这类数组用自定义 UDF `JsonActionsArrayParser`(Jackson 解析数组首元素成 `action_id/item/item_type/ts` 行)配合维表补地区、品牌 |
| 用户域 | 登录 `dwd_user_login_full` | 从日志里按"设备 + 会话起点"拼 `session_id`、开窗取会话首个事件界定登录,再 join 省份维表补省份;并对单表 `ALTER TABLE ... SET ('sink.parallelism'='10')` 单独调大写入并行度 |
| 互动域 | 收藏 `dwd_interaction_favor_add_full`、评论 `dwd_interaction_comment_full` | 收藏/评论业务事实 |
| 工具域 | 优惠券领取 `dwd_tool_coupon_get_full`、下单用券 `dwd_tool_coupon_order_full`、支付用券 `dwd_tool_coupon_pay_full` | 券的生命周期分阶段建模 |

以支付成功事实为例,其 SQL 结构是"明细子查询 × 多路维表/字典 left join + 字典过滤(parent_code='11' 支付方式、'24' 来源类型)"([dwd_trade_pay_detail_suc_full.java:94-185](src/main/java/com/fy/warehouse/dwd/dwd_trade_pay_detail_suc_full.java))。注意这里 DWD 仍然在 Flink SQL 里用 `left join` 直查 ODS Paimon 表——维表关联走的是**批量小表 join 明细流**的简化路线,这也是为什么全局并行度默认为 1、多数表按分区追加。

### DWS:按粒度的 N 日/累计汇总

DWS 把 DWD 明细按"分析粒度 + 时间跨度"聚合,命名规律是 `域_粒度_指标_nd/td`:用户粒度最近 30 日(`dws_trade_user_order_nd_full`)、用户粒度累计(`dws_trade_user_order_td_full`,带首末次下单日期 `order_date_first/last` 与 `_td` 累计口径)、用户 × SKU 粒度 30 日、省份粒度 30 日(`dws_trade_province_order_nd_full`)、活动/优惠券粒度订单 30 日,以及流量侧的访客 × 页面粒度 30 日。统计口径统一为"从明细里取 `k1` 与基准日(代码中硬编码如 `2024-05-31`)相差 0~30 天的记录再 `group by` 聚合",例如下单次数 `count(order_id)`、金额 `sum(split_*)`([dws_trade_user_order_nd_full.java:60-82](src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java));流量侧同理,按 `mid_id/brand/model/operate_system/page_id` 分组算 30 日浏览时长与访问次数([dws_traffic_page_visitor_page_view_nd_full.java:55-84](src/main/java/com/fy/warehouse/dws/dws_traffic_page_visitor_page_view_nd_full.java))。

### ADS:面向指标的极薄应用层

ADS 直接消费 DWS 汇总,做最轻的二次加工,产出可直接画图的口径。例如活动补贴率表(`ads_activity_stats_full`)就是把 DWS 的活动 30 日数据转成"补贴率 = 活动优惠金额 / 原始金额"的比率指标,字段只有 `dt / activity_id / activity_name / start_date / reduce_rate`([ads_activity_stats_full.java:42-49](src/main/java/com/fy/warehouse/ads/ads_activity_stats_full.java));优惠券统计同理落 `ads_coupon_stats_full`。由于 ADS 建表不带 `k1` 分区,它更像"轻量结果集",也正是 StarRocks `use ads; select * from ads_coupon_stats_full;` 直接可视化的对象。

## 经验沉淀

把前面所有细节收拢,可复用的经验大致有四条:

1. **集中配置、注释留痕排障过程。** 全仓库只有一个配置入口,每处调优都带着"原值是什么、导致什么问题、改成什么"的注释(如开启 checkpoint、换文件系统状态后端、加 HDFS RPC 超时)。新环境复刻时,逐条对注释检查即可,不必重新踩坑。
2. **"读时展开、写时平铺"的日志处理套路。** Kafka 日志源以嵌套 ROW + JSON 字符串进 ODS,落地时一次 SELECT 展开成平铺列;数组类字段(actions/displays)留给 DWD 用 UDF 拆。既保住了原始数据,又让下游好查。
3. **统一的"主键 + k1 天分区 + 分区过期"表骨架。** 每层表都以 `(业务键, k1)` 为主键、按天分区,配 `partition.expiration-time=1 d` 自动清历史;Flink 侧文件走 `file.size=64mb` 小文件快速写。这套骨架被 ODS/DIM/DWD/DWS 数百行建表 SQL 原样复用,是一致性的来源。
4. **多引擎共享一份湖上数据。** 写入用 Flink + Paimon(Hive Catalog),查询用 StarRocks External Catalog 指向同一 metastore/warehouse 直读,可视化在 DataRT。切引擎只需改连接声明,数据文件不复制——这是"湖"相对传统数仓最直接的收益。

一点提醒(同样来自代码):DWS 的 30 日窗口基准日(`2024-05-31`、`2024-05-01`)与 `application.properties` 的 `mock.date=2024-05-01` 是**硬编码的教学日期**,业务日期推进时必须同步改,否则聚合窗口会空;这是把本仓库迁移到真实日期口径时最容易漏的改动点。Sources:[dws_trade_user_order_nd_full.java:80](src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java)、[src/main/datagenerator/business_jar/application.properties:16](src/main/datagenerator/business_jar/application.properties)

综上:这套实时数据湖仓以"Flink 流式作业 + Paimon 湖表 + Hive 元数据/HDFS 存储"为底座,用一个统一配置入口管住状态与检查点、微批、文件策略等全部计算存储调优点;业务数据与日志数据分别在 ODS 汇聚,经 DIM 宽表与拉链、DWD 分域事实、DWS 分粒度汇总,最后在 ADS 形成指标,交给 StarRocks 外部分区直读并可视化。各业务域(交易、流量、用户、互动、工具)在每一层都有清晰、命名规律统一的落点,配合集中配置与通用表骨架,整条链路从数据生成到报表呈现保持了一致的可复现性——这套"配置入口 + 统一表骨架 + 分层分域"的组织方式,是比任何单张表都更值得带走的东西。

---
