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

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

- 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

- `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 即席查询"完整串起来的电商实时数仓教材工程**。代码最大的价值在于——每条链路的接入方式、每张表的建表选项、每个调优参数的动机，都以"一个类 + 注释"的形态平铺在你面前；后续三页会分别把"湖表机制""调优动作""业务口径"讲透，本页的地图到此就够用了。
