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

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

- Repository: fengyu-eng/paimon-datalake
- GitHub: https://github.com/fengyu-eng/paimon-datalake
- Human wiki: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533
- Complete Markdown: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/llms-full.txt

## 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 出指标——每一层职责单一，层与层只通过表契约交互，这正是这套代码可以横向扩展（新增域、新增指标）而无需推翻结构的原因。
