# Flink 电商实时数仓工程解读

> 一套 Maven 多模块的电商实时数仓教学工程，覆盖从埋点日志与业务库经 Kafka 到 Flink 分层加工（DIM/DWD/DWS）、落 ClickHouse 再由可视化工程查询的完整链路。值得研究的是它把数仓分层方法论落成了可运行的 Flink 代码，并在维表、去重、窗口聚合上给出了典型取舍。

## Context Links

- [Agent index](https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/llms.txt)
- [Human interactive wiki](https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69)
- [GitHub repository](https://github.com/aggaadfr/gmall-flink-3.0)

## Repository Metadata

- Repository: aggaadfr/gmall-flink-3.0

- Generated: 2026-09-08T05:31:55.835Z
- Updated: 2026-09-14T05:45:01.256Z
- Runtime: Claude Code · claude-fable-5
- Format: Custom
- Pages: 4

## Page Index

- 01. [导读：一套电商实时数仓的教学全景](https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/pages/01-page-1.md) - 从整体定位仓库：三个 Maven 工程各司其职，实时模块按数仓分层组织代码，两个发布工程负责把结果查询出来做可视化。作为开篇，先建立模块边界与技术栈印象，再给出从架构、数据流、计算存储、系统优化到业务实现的阅读路线，避免与后续各页重复。
- 02. [数仓分层与数据流转](https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/pages/02-page-2.md) - 讲解一条数据从哪来、往哪去：业务库与埋点日志先进 Kafka 主题，再按 DIM 维表层、DWD 明细层、DWS 汇总层的顺序被加工。以 DimApp、BaseLogApp 和交易预处理的代码为例，说明主题命名约定、脏数据侧输出流、维表配置的广播式加载，以及层与层之间如何靠主题衔接。
- 03. [计算模型与存储选型](https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/pages/03-page-3.md) - 围绕“算在哪、存哪里”讲清技术底座：DataStream 与 FlinkSQL 两种写法各自的使用场景，状态与 Checkpoint 在生产上的注意点，以及 Kafka、MySQL、Phoenix/HBase、ClickHouse、Redis 在链路中的分工。重点解释针对维度关联的 Redis 旁路缓存与异步查询、维表配置的广播状态等系统优化手段。
- 04. [业务域落地：指标口径与可复用经验](https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/pages/04-page-4.md) - 把前几页的通用机制放回业务场景，以流量、交易、用户三个域为线索还原典型指标在代码中的落地方式：新老用户判定与 UV、跳出行为、页面浏览窗口、交易预处理后的支付与下单汇总、注册回流统计。收尾给出从这套工程抽象出的实时数仓分层方法论与踩坑经验，作为全文总结。

## Source File Index

- `gmall-publisher-2022/pom.xml`
- `gmall-realtime/pom.xml`
- `gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradePayDetailSuc.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUserJumpDetail.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradeOrderWindow.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficVcChArIsNewPageViewWindow.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserRegisterWindow.java`
- `gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/JedisUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/MysqlUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/PhoenixUtil.java`
- `pom.xml`
- `README.md`

---

## 01. 导读：一套电商实时数仓的教学全景

> 从整体定位仓库：三个 Maven 工程各司其职，实时模块按数仓分层组织代码，两个发布工程负责把结果查询出来做可视化。作为开篇，先建立模块边界与技术栈印象，再给出从架构、数据流、计算存储、系统优化到业务实现的阅读路线，避免与后续各页重复。

- Page Markdown: https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/pages/01-page-1.md
- Generated: 2026-09-08T05:23:39.366Z

### Source Files

- `README.md`
- `pom.xml`
- `gmall-realtime/pom.xml`
- `gmall-publisher-2022/pom.xml`

<details>
<summary>相关源文件</summary>
以下文件用于生成此维基页面：
- [README.md](README.md)
- [pom.xml](pom.xml)
- [gmall-realtime/pom.xml](gmall-realtime/pom.xml)
- [gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java)
- [gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java](gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java](gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java)
- [gmall-publisher/pom.xml](gmall-publisher/pom.xml)
- [gmall-publisher-2022/pom.xml](gmall-publisher-2022/pom.xml)
- [gmall-publisher-2022/src/main/java/com/atguigu/gmall/publisher/controller/TrafficController.java](gmall-publisher-2022/src/main/java/com/atguigu/gmall/publisher/controller/TrafficController.java)
</details>

# 导读：一套电商实时数仓的教学全景

这个仓库是一份电商实时数仓的教学项目（Flink 版本 3.0），它的价值不在于业务复杂度，而在于把一套真实企业级数仓的**分层思想**压缩到了三个 Maven 工程里。读懂仓库之前，先看清三件事：谁在算、结果存到哪、谁把结果查出来。

本页只做“地图”工作：建立模块边界、技术栈与阅读顺序的印象，把「架构、数据流、计算存储、系统优化、业务实现」这五类话题指到正确的入口，避免与后续各页重复展开。

## 一个父工程，三个 Maven 模块

根 `pom.xml` 只是聚合壳，打包方式为 `pom`，声明了三个子模块与 Java 8 编译目标（[pom.xml:7-20](pom.xml)）。README 用一句话概括了分工（[README.md:24-31](README.md)）：

```
gmall-flink-3.0
├── gmall-publisher         测试数据可视化接口（实验版）
├── gmall-publisher-2022    完成数据可视化接口（完整版）
└── gmall-realtime          实时模块
```

三者定位差异很大，是阅读时的第一把钥匙：

| 模块 | 形态 | 职责 | 技术底座 |
| --- | --- | --- | --- |
| `gmall-realtime` | Maven 普通工程，无框架 | 数仓加工：消费 ODS、清洗分流、聚合落库 | Flink 1.13 + 各类连接器 |
| `gmall-publisher` | Spring Boot 2.4 | 可视化查询的实验版（GMV/UV 两个指标） | Web + MyBatis + ClickHouse JDBC |
| `gmall-publisher-2022` | Spring Boot 2.6.6 | 覆盖各业务域的完整可视化查询 | Web + MyBatis + ClickHouse JDBC |

注意一个易混淆点：README 称之为「测试数据可视化接口」的 `gmall-publisher` 与「完成数据可视化接口」的 `gmall-publisher-2022` 是**并列的两个发布工程**，不是新旧目录关系——教学上常拿前者做最小验证，后者是完整版。

## gmall-realtime：数仓的“加工车间”

这是仓库主体，59 个 Java 文件集中在 `com.atguigu` 下五个包，包名即分层（`gmall-realtime/src/main/java/com/atguigu/app`）。业务目录只有 `dim`、`dwd`、`dws`、`func` 四个——**没有 ODS 与 ADS 的包**，这两层以不同形态存在：ODS 就是 Kafka 里的原始主题，ADS/展示层则由两个发布工程承担。

Flink 版本统一收敛在 `gmall-realtime/pom.xml` 的属性里（[gmall-realtime/pom.xml:18-20](gmall-realtime/pom.xml)），依赖覆盖了这条教学主线：

- 计算与连接：`flink-java`、`flink-streaming-java`、`flink-connector-kafka`、`flink-json`、`flink-cep`、`flink-table-api` 与 `flink-table-planner-blink`（[gmall-realtime/pom.xml:25-59](gmall-realtime/pom.xml)、[gmall-realtime/pom.xml:142-150](gmall-realtime/pom.xml)）；
- 维表与配置同步：`flink-connector-mysql-cdc` 2.1.0、`mysql-connector-java`、`phoenix-spark` 5.0.0（HBase 2.0）（[gmall-realtime/pom.xml:97-137](gmall-realtime/pom.xml)）；
- 结果存储与关联：`clickhouse-jdbc`、`jedis`、分词器 `ikanalyzer` 等工具（[gmall-realtime/pom.xml:152-179](gmall-realtime/pom.xml)）。

### 包结构与分层对应

| 包/目录 | 对应数仓层 | 落点（Sink） | 典型代表 |
| --- | --- | --- | --- |
| `app/dim` | DIM 维度层 | Phoenix（HBase） | `DimApp` |
| `app/dwd/log` | DWD 流量日志明细 | Kafka 分流主题 | `BaseLogApp` |
| `app/dwd/db` | DWD 业务库明细（交易/工具/互动/用户） | Kafka | `DwdTradeOrderPreProcess` 等 12 个作业 |
| `app/dws` | DWS 汇总层 | ClickHouse | 11 个窗口聚合作业 |
| `app/func` | 通用函数 | — | `TableProcessFunction`、`DimSinkFunction`、异步/关联函数 |
| `bean` / `common` / `utils` | 支撑代码 | — | 数据模型、`GmallConfig`、`MyKafkaUtil` 等 |

**DIM 层**（[DimApp.java:35-80](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java)）体现了教学项目最值得看的设计：消费 Kafka 的 `topic_db`（装着全部业务表的 binlog），用 FlinkCDC 监听 MySQL 里的配置表 `table_process` 生成广播流，按配置动态决定哪些表进维度层；写 Phoenix 时用 `upsert` 动态拼表名与字段（[DimSinkFunction.java:27-52](gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java)），更新时先删 Redis 缓存保证一致性。

**DWD 层**以 `BaseLogApp` 为例（[BaseLogApp.java:35](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java)），核心动作是清洗、新老用户校验、再用侧输出流把页面/启动/曝光/动作/错误五类日志分到五个 Kafka 主题（[BaseLogApp.java:100](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java)、[BaseLogApp.java:172-176](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java)）。README 按 15 个主题细述了这层逻辑（[README.md:70-341](README.md)）。

**DWS 层**是窗口聚合，把结果写出到 ClickHouse。落库统一走泛型工具 `ClickHouseUtil.getJdbcSink`（[ClickHouseUtil.java:23-25](gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java)），用反射读字段、以 `TransientSink` 注解跳过不需要的字段。例如 `DwsTrafficPageViewWindow` 在头部注释里写明「将数据写出到 ClickHouse」（[DwsTrafficPageViewWindow.java:34-39](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java)），并在作业末尾调用该工具（[DwsTrafficPageViewWindow.java:202](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java)）。

存储地址集中在 `GmallConfig`（[GmallConfig.java:12-26](gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java)）：Phoenix 库 `GMALL211027_REALTIME`、ClickHouse 库 `gmall_211027`，主机是教学集群三节点。

## 两个 publisher：把结果“查”出来做可视化

DWS 写进 ClickHouse 后，剩下的问题就是怎么给大屏提供数据。两个发布工程都是标准的 Spring Boot 三层结构（`controller → service → mapper`），Mapper 层连 ClickHouse 查询预聚合结果；`gmall-publisher-2022` 在 `com.atguigu.gmall.publisher` 下按业务域拆出交易、流量、商品、用户、活动、优惠券六组统计模块（如 [TrafficController.java:17-24](gmall-publisher-2022/src/main/java/com/atguigu/gmall/publisher/controller/TrafficController.java)），并配了一批统计结果 Bean。两个工程都监听 8070 端口、连同一个 ClickHouse 库（[application.properties](gmall-publisher-2022/src/main/resources/application.properties)）。实验版 `gmall-publisher` 则以日活、GMV 两个指标做最小闭环（[SugarController.java:26-87](gmall-publisher/src/main/java/com/atguigu/gmallpublisher/controller/SugarController.java)），适合先跑通链路再看完整版。

## 全链路全景

```mermaid
flowchart TB
    subgraph 数据源头
        DB[("MySQL 业务库<br/>binlog")]
        LOG["埋点 / 模拟行为日志"]
    end

    subgraph RT["gmall-realtime（Flink 1.13 计算）"]
        ODS["Kafka ODS<br/>topic_db · topic_log"]
        CFG["MySQL 配置表<br/>table_process（FlinkCDC 广播）"]
        DIM["app/dim → DimApp<br/>维表落地 Phoenix"]
        DWD["app/dwd → log / db 十五个明细主题<br/>写回 Kafka"]
        DWS["app/dws → 窗口聚合<br/>写 ClickHouse"]
        DIMJOIN["Redis 维表缓存<br/>关联明细与维表"]
    end

    subgraph 结果发布
        P1["gmall-publisher（实验版）"]
        P2["gmall-publisher-2022（完整版）"]
        UI["数据大屏可视化"]
    end

    DB -->|FlinkCDC 实时抓取| ODS
    LOG -->|上报 Kafka| ODS
    ODS --> DIM
    CFG -.驱动维度表动态过滤.-> DIM
    ODS --> DWD
    DWD --> DWS
    DWS -.查询维表.-> DIMJOIN
    DWS --> P1
    DWS --> P2
    P1 --> UI
    P2 --> UI
```

维表读取 Redis 缓存的证据在 `DimUtil`：按 `DIM:表名:id` 拼 Key、提供删除与查询入口（[DimUtil.java:27-49](gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java)），配合 `app/func` 里的 `DimJoinFunction`、`DimAsyncFunction` 在 DWS 作业中做维表关联。

## 建议阅读路线

这是一份导读，后续页面请沿着下表推进；每页只深挖一个主题，本页不再展开：

| 顺序 | 主题 | 仓库里的主入口 | 后续页要回答的问题 |
| --- | --- | --- | --- |
| 1 | 架构定位 | 本文 + 三个 `pom.xml` | 模块边界为什么这样切，技术选型如何呼应教学 |
| 2 | 数据流 | README 分层逻辑 [README.md:35-70](README.md)、`BaseLogApp` | 一条日志/一个订单如何从源头流到各层 |
| 3 | 计算存储 | `DimApp`、`app/dws/*`、`GmallConfig` | 每层“加工什么、存到哪、为何这么存” |
| 4 | 系统优化 | README 调优部分 [README.md:651-1035](README.md) | 并行度、Checkpoint、反压、数据倾斜、FlinkSQL 优化的实战注解 |
| 5 | 业务实现 | 15 个 DWD 主题与 `gmall-publisher-2022` | 交易/流量/工具/互动/用户各域的指标如何一步步算出并展示 |

## 收尾

一句话记住这个仓库：**`gmall-realtime` 用 Flink 把「业务库 binlog + 埋点日志」按 DIM/DWD/DWS 加工成「维度（Phoenix）→ 明细（Kafka）→ 汇总（ClickHouse）」三层结果，两个 publisher 再把 ClickHouse 的汇总查出来喂给大屏**。阅读时始终带着这条主线，就不会在 59 个类、十几个 Kafka 主题里迷路；后续页面各自展开其中一段，这里的分工即为其索引。技术栈与模块定位同时见 [README.md:9-31](README.md) 与根 [pom.xml:7-15](pom.xml)。

---

## 02. 数仓分层与数据流转

> 讲解一条数据从哪来、往哪去：业务库与埋点日志先进 Kafka 主题，再按 DIM 维表层、DWD 明细层、DWS 汇总层的顺序被加工。以 DimApp、BaseLogApp 和交易预处理的代码为例，说明主题命名约定、脏数据侧输出流、维表配置的广播式加载，以及层与层之间如何靠主题衔接。

- Page Markdown: https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/pages/02-page-2.md
- Generated: 2026-09-08T05:23:40.275Z

### Source Files

- `gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java`

<details>
<summary>相关源文件</summary>
以下文件用于生成此维基页面：
- [gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java](gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/MysqlUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/MysqlUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/bean/TableProcess.java](gmall-realtime/src/main/java/com/atguigu/bean/TableProcess.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java)
- [gmall-realtime/src/main/java/com/atguigu/app/func/TableProcessFunction.java](gmall-realtime/src/main/java/com/atguigu/app/func/TableProcessFunction.java)
- [gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java](gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java)
- [README.md](README.md)
</details>

# 数仓分层与数据流转

这个电商实时数仓把「业务库里的订单、优惠券等数据」和「客户端上报的埋点日志」这两路数据，先汇入 Kafka，再按照 **DIM 维表层 → DWD 明细层 → DWS 汇总层** 的顺序逐层加工：每一层从 Kafka 的主题里读上游产物，处理完再写回 Kafka（或落库），供下一层继续消费。整条链路像一条流水线，而**主题就是层与层之间的传送带**。

本页用三个最有代表性的应用来讲解这条链路：`DimApp`（怎么把维表刷进宽表引擎）、`BaseLogApp`（怎么把一份原始日志拆成五类明细）、`DwdTradeOrderPreProcess`（怎么把多张业务表加工成一张订单宽表），并顺带说明主题命名约定、脏数据侧输出流和广播式维表配置加载。

## 一条数据从哪来：两个 ODS 入口主题

仓库并不包含采集端的代码，所有实时应用都以 **Kafka 里的两个 ODS 主题** 为起点（README 在各应用说明中也把这两个主题直接称为 ODS 输入）。

- **`topic_db`**：业务库 MySQL 的变更汇总。无论是订单还是字典表，只要业务库里发生了增删改，都会被 CDC 类工具抓成一条消息写进来。消息是统一结构，携带库名、表名、操作类型、变更后的新值 `data` 和变更前的旧值 `old`。工具类里有一段现成的 DDL 直接把这它声明成一张「表名就叫 topic_db」的流表，加工方用 SQL 的 `where` 就能按 `database`、`table`、`type` 筛选某一类业务变更。`DwdTradeOrderPreProcess` 里反复出现的 `where database='gmall-211027-flink' and table='order_info' and type in ('insert','update')` 就是这个写法的体现。

- **`topic_log`**：客户端上报的埋点日志汇总，一条日志同时可能携带 `common`（公共字段，含设备 mid、用户 uid、新老标记 is_new）、`page`（页面）、`start`（启动）、`displays`（曝光数组）、`actions`（动作数组）等多个部分。

Sources: [gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java:81-90](gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java:81-90) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:61-63](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:61-63) · [README.md:74](README.md)

## 总体分层：谁加工成什么

从这两个 ODS 主题出发，各层的职责和去向可以概括为：**DIM 把「几乎不变的字典/主数据」从 ODS 里挑出来做成可联查的维表；DWD 把 ODS 里“数据库行 / 原始日志”重构成“面向业务过程的事件明细”；DWS 再按统计口径把明细聚合成指标。** 落库方式各不相同，这正是分层的目的——让每一层只依赖上一层、只暴露一种稳定形态。

```mermaid
flowchart TB
    subgraph SRC["数据源（在仓库之外采集）"]
        db["业务库 MySQL（gmall-211027-flink）"]
        lg["前端埋点日志"]
    end
    subgraph ODS["ODS：两个 Kafka 汇总主题"]
        tdb["topic_db（携带 库名/表名/type/data/old）"]
        tlg["topic_log（原始埋点日志）"]
    end
    subgraph DIM["DIM 维表层"]
        da["DimApp：按广播配置挑出维表数据"]
        ph["Phoenix/HBase：每张维表对应一张表"]
    end
    subgraph DWD["DWD 明细层"]
        bl["BaseLogApp：日志清洗分流"]
        pt["dwd_traffic_page_log 等 5 个流量明细主题"]
        pp["DwdTradeOrderPreProcess：订单宽表"]
        od["dwd_trade_order_detail（upsert-kafka 主题）"]
    end
    subgraph DWS["DWS 汇总层"]
        uu["DwsUserUserLoginWindow（读页面主题）"]
        ow["DwsTradeOrderWindow（读订单明细主题）"]
        ck["ClickHouse：dws_* 结果表"]
    end
    db -->|CDC 汇总| tdb
    lg -->|日志采集| tlg
    tdb --> da
    da --> ph
    tdb --> pp
    pp --> od
    tlg --> bl
    bl --> pt
    pt --> uu
    od --> ow
    uu --> ck
    ow --> ck
```

每一层的输入输出形态，决定了「衔接」靠的是什么：

| 层 | 输入 | 加工要点 | 输出 |
|---|---|---|---|
| DIM 维表层 | `topic_db` | 按表名挑出维表数据、只保留配置的列 | 写入 Phoenix/HBase（HBase 行键来自配置的主键），不再回 Kafka |
| DWD 明细层 | `topic_db` 或 `topic_log` | 清洗、去重/纠偏、展开、多表关联成事件明细 | 写回 Kafka 的 `dwd_*` 主题（宽表用 upsert-kafka） |
| DWS 汇总层 | 各类 `dwd_*` 主题 | 过滤口径、分组、开窗聚合 | 写入 ClickHouse 的 `dws_*` 结果表 |

所以更准确的说法是：Kafka 负责 **层与层之间** 的搬运（尤其是 DWD 的产出、DWS 的输入），而每层**最终面向查询的落库**在各自的存储里。DIM 的产物是维表库，DWS 的产物是指标库，它们不再作为 Kafka 主题被下一层搬运。

Sources: [gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java:13-26](gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java:13-26) · [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:174-175](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:174-175)

## 主题命名约定

Kafka 主题的命名把「属于哪一层、哪个域、什么内容」都写进了名字里，下游按名字就能找到上游产物：

- **ODS 汇总主题**：`topic_db`、`topic_log`，不分域，全量汇总。
- **DWD 明细主题**：`dwd_<域>_<业务内容>[_log]`。域前缀有 `traffic`（流量）、`trade`（交易）、`tool`（工具/优惠券）、`interaction`（互动）等；日志类明细习惯带 `_log` 后缀。同一个 ODS 主题会被**多个下游用不同消费组并行消费**——比如 `topic_db` 同时被 `DimApp`（消费组 `topic_db_DimApp`）、`DwdTradeOrderRefund`（消费组 `dwd_trade_order_refund`）等各自消费，互不影响进度。

一个值得留意的反例也在交易域：`DwdTradeOrderPreProcess` 输出时把主题写成了 `dwd_trade_order_detail`，而真正读「订单预处理宽表」的下游 `DwdTradeOrderDetail`、`DwdTradeCancelDetail` 却去读 `dwd_trade_order_pre_process`。也就是说**生产者写出的名字与消费者预期的预处理主题名并不一致**，阅读代码时会感到错位。这也说明这套命名约定是“约定”而非“强制”：主题名是否对得上，靠各应用手写的字符串保持一致，稍不留神就会漂移。

Sources: [gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java:46-47](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java:46-47) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java:172-176](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java:172-176) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:247](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:247) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java:80-81](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java:80-81)

## DIM 维表层：广播配置，把维表刷进 Phoenix

DIM 的典型问题是：业务库几十张表里，只有一部分是需要长期保留、供明细联查的维表（用户、SKU、品牌、字典等），而且**加一张新维表不该改代码重启任务**。`DimApp` 的解法是把“哪些表算维表、落到哪张表、留哪些列”做成**运行时可变的配置**。

主流程分四步：

1. **消费 `topic_db`** 拿到全部业务变更。
2. **校验 JSON**：能解析的进主流，解析失败的进“脏数据”侧输出流，不会让一条坏数据弄垮整个任务。
3. **用 Flink CDC 实时读 MySQL 里的配置表 `table_process`**，把配置做成广播流，连到主流上。
4. 主流数据按配置**过滤字段、补上目标表名**后，写入 Phoenix。

### 脏数据侧输出流

解析失败的数据没有简单丢弃，而是通过侧输出流单独摘出来打印，方便事后排查，主流程继续跑。

```java
OutputTag<String> dirtyDataTag = new OutputTag<String>("Dirty") { };
SingleOutputStreamOperator<JSONObject> jsonObjDS = kafkaDS.process(new ProcessFunction<String, JSONObject>() {
    public void processElement(String s, ...) throws Exception {
        try {
            JSONObject jsonObject = JSON.parseObject(s);
            collector.collect(jsonObject);          // 能解析 → 主流
        } catch (Exception e) {
            context.output(dirtyDataTag, s);        // 解析失败 → 脏数据侧输出流
        }
    }
});
DataStream<String> sideOutput = jsonObjDS.getSideOutput(dirtyDataTag);
sideOutput.print("dirtyDataTag >>>>>>>> ");
```

Sources: [gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java:54-71](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java:54-71)

### 维表配置的广播式加载

配置表 `table_process` 的每一行描述一张维表要怎么建、怎么写：

| 配置字段 | 含义 |
|---|---|
| `sourceTable` | 源业务表名（用来匹配 `topic_db` 里的 `table` 字段） |
| `sinkTable` | 落到 Phoenix 的表名 |
| `sinkColumns` | 要保留的列，逗号分隔 |
| `sinkPk` | 主键列，作为 HBase 行键 |
| `sinkExtend` | 建表时的扩展属性 |

配置的读取和下发在 `DimApp` 中完成：用 CDC 连接器指向 `gmall-211027-config.table_process`，初次启动就全量加载（`StartupOptions.initial()`），之后配置表的任何改动都会实时变成新的广播元素。广播状态用 `MapStateDescriptor` 声明，主流与广播流 `connect` 之后，每个并行子任务手里都有一份完整的“表名 → 建表/落列规则”。

```java
MapStateDescriptor<String, TableProcess> mapStateDescriptor =
        new MapStateDescriptor<>("map-state", String.class, TableProcess.class);
BroadcastStream<String> broadcastStream = mysqlSourceDS.broadcast(mapStateDescriptor);
BroadcastConnectedStream<JSONObject, String> connectedStream = jsonObjDS.connect(broadcastStream);
SingleOutputStreamOperator<JSONObject> hbaseDS =
        connectedStream.process(new TableProcessFunction(mapStateDescriptor));
```

广播流的处理逻辑做两件事：

- **收到一条新配置时自动建表**：按 `sinkColumns` 拼 `create table if not exists`，`sinkPk` 对应的列建成主键（即 HBase 的 rowkey），其余列一律 `varchar`，再拼上 `sinkExtend` 执行。这样新增一张维表不需要手工去 Phoenix 建表。
- **主流数据到来时按配置裁剪**：只保留 `sinkColumns` 里出现的列，并补一个 `sinkTable` 字段指明去向。

Sources: [gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java:74-94](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java:74-94) · [gmall-realtime/src/main/java/com/atguigu/bean/TableProcess.java:20-28](gmall-realtime/src/main/java/com/atguigu/bean/TableProcess.java:20-28) · [gmall-realtime/src/main/java/com/atguigu/app/func/TableProcessFunction.java:86-141](gmall-realtime/src/main/java/com/atguigu/app/func/TableProcessFunction.java:86-141)

### 落库与缓存一致性

最终写入用 Phoenix 的 `upsert` 语义（有则改、无则插），列与值来自裁剪后的 `data`。这里有一个容易被忽略的细节：因为 DWD 层联维表时通常会先查 Redis 缓存（维表缓存），所以当这条数据是 **update** 时，写入前要先删掉 Redis 里对应 id 的缓存，否则下游会一直读到旧值。

```java
String insertSql = "upsert into " + GmallConfig.HBASE_SCHEMA + "." + sinkTable + "(" + ...
if ("update".equals(value.getString("type"))) {
    DimUtil.delDimInfo(sinkTable.toUpperCase(), data.getString("id"));  // 更新维表时让缓存失效
}
```

Sources: [gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java:46-58](gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java:46-58)

## DWD 明细层：把 ODS 加工成面向业务过程的事件明细

DWD 内部又按数据来源分成两条线：**日志线**（从 `topic_log` 来）和**业务库线**（从 `topic_db` 来）。两条线做法差异很大：日志是多对多嵌套的 JSON，要先“拆开摊平”；业务库是多张关联表，要先“拼成宽表”。

### 日志线：BaseLogApp，一份日志拆成五类明细

`BaseLogApp` 是典型的“一进多出”，承担三层职责：

1. **清洗**：把 `topic_log` 的每条字符串解析成 JSON。
2. **新老访客纠偏**：埋点里的 `is_new` 标记不一定可靠（例如用户清了缓存后再访问，前端会误报成新用户）。代码按设备 `mid` 分组，用状态记录每个设备“最近一次访问的日期”：若上报 `is_new=1` 但状态里已经有其他日期的记录，就把它改写成 `0`；反过来，若上报 `is_new=0` 但状态为空，就回填一个“昨天”作为占位，避免首次启动就把老用户误判成新用户。
3. **分流**：一个处理函数同时判断日志里带有哪些部分，把它们分别送进不同的输出流——页面浏览走主输出流，启动、曝光、动作、错误走侧输出流。曝光 `displays` 和动作 `actions` 是数组，需要**逐条拆开**，并补上 `page_id`、`ts`、`common` 这些上下文，让每一条明细脱离父日志后依然自洽。

```mermaid
flowchart LR
    k["topic_log"] --> c["逐条解析 JSON + 按 mid 状态纠偏 is_new"]
    c --> s{"按日志携带的内容分流"}
    s -- "主输出流：含 page 的日志" --> p["dwd_traffic_page_log"]
    s -- "侧输出流：start 启动" --> st["dwd_traffic_start_log"]
    s -- "侧输出流：displays 逐条展开，补 page_id/ts/common" --> dp["dwd_traffic_display_log"]
    s -- "侧输出流：actions 逐条展开，补 page_id/ts/common" --> ac["dwd_traffic_action_log"]
    s -- "侧输出流：带 err 字段的日志" --> er["dwd_traffic_error_log"]
```

最后五个主题分别接上各自的生产者写出。原始日志在这里被“摊平并归类”，下游流量域 DWS 只需各取所需，例如算页面浏览量的读 `dwd_traffic_page_log`，算启动量的读 `dwd_traffic_start_log`。

Sources: [gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java:53-98](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java:53-98) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java:113-195](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java:113-195)

### 业务库线：DwdTradeOrderPreProcess，四张源表拼一张订单宽表

交易域的事实是“订单明细”，但它在业务库里分散在订单表、订单明细表、明细活动关联表、明细优惠券关联表里，还带着许多需要解释的编码。`DwdTradeOrderPreProcess` 用 Flink SQL 一次性把这些问题解决掉，过程清晰得像在写查询：

1. 用现成 DDL 把 `topic_db` 声明成流表，消费组就叫 `dwd_trade_order_preprocess`。
2. 分别用四条 `select` 从流表里筛出 `order_detail`（仅 insert）、`order_info`（insert/update，并保留 `type`、`old` 供后续判断）、`order_detail_activity`（仅 insert）、`order_detail_coupon`（仅 insert），各自建成临时视图。
3. 引入一张 MySQL 字典维表 `base_dic`——注意这次不是广播流，而是 **Lookup 维表**：查询时实时去 MySQL 按 `dic_code` 取 `dic_name`，并带有一小时、100 条的本地缓存。
4. 把上述五路 `join` 起来：明细与订单做内连接，活动与优惠券做左连接（明细可能没有），字典表用“处理时间关联”查出 `source_type` 对应的 `source_type_name`（来源类型名，如“用户下单/小程序下单”），得到一张**已经退化好维度的宽表**。
5. 以 `order_detail_id` 为主键，用 upsert-kafka 写回 Kafka 主题，供后续交易域子任务继续消费。

一个关键点：**宽表是按主键在 Kafka 里以 upsert 形式维护的**，所以订单表后续的 update（比如状态流转）也能以同一主键刷新明细行——这是普通 append 主题做不到的，也是为什么这里必须用 upsert-kafka。

Sources: [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:40-64](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:40-64) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:71-100](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:71-100) · [gmall-realtime/src/main/java/com/atguigu/utils/MysqlUtil.java:20-48](gmall-realtime/src/main/java/com/atguigu/utils/MysqlUtil.java:20-48) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:141-195](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:141-195) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:201-251](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:201-251)

### DWD 内部还会再切一刀

订单宽表本身还没到“一个主题只表达一个业务过程”的程度。下游 `DwdTradeOrderDetail`、`DwdTradeCancelDetail` 会**从预处理宽表主题再各取所需**：一个筛出下单成功的数据，一个筛出取消的数据，各自拼装成更聚焦的明细主题。这说明 DWD 内部也不是一步到位，而是“宽表打底、子任务分主题”的多级加工，每一级都继续使用 Kafka 主题做衔接。

Sources: [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java:50-81](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java:50-81) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java:120-137](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java:120-137)

## DWS 汇总层：从 DWD 主题到 ClickHouse 指标

DWS 是“出口转指标”的一层：读 DWD 主题 → 按统计口径过滤 → 分组 + 开窗聚合 → 写 ClickHouse。以用户登录汇总 `DwsUserUserLoginWindow` 为例，它消费的就是上面日志线产出的 `dwd_traffic_page_log`，完整演示了“DWD 主题如何被 DWS 拾取”。

处理链路上的几个判断很典型：

1. **口径过滤**：只保留“登录后进入的首页浏览”这类事件——条件是 `common.uid` 不为空且 `page.last_page_id` 为空。有 uid 说明用户已登录，`last_page_id` 为空说明这是本次会话的第一页，二者合起来代表一次“登录访问”。
2. **事件时间与水位线**：日志里的 `ts` 才是事件发生时间，按它生成水位线，并容忍 2 秒乱序。
3. **分组算回流与独立**：按 `uid` 分组后，用状态记录每个用户最近一次访问日期。状态为空说明当天首次见到该用户，计 1 个独立访客；状态里是更早的日期且距今天数 **≥ 8 天**，则除了独立访客再计 1 个回流用户。
4. **开窗聚合**：按 10 秒滚动事件时间窗口做全量聚合，在窗口收尾时补上窗口起止时间 `stt`/`edt`，形成一条可落库的指标记录。
5. **落库**：通过 ClickHouse 的 JDBC sink 写入 `dws_user_user_login_window` 表，每条记录对应窗口内的独立访客数与回流用户数。

Sources: [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:66-92](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:66-92) · [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:95-147](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:95-147) · [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:150-175](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:150-175)

## 小结

把这条链路串起来看：**业务库与埋点日志先汇入 `topic_db`、`topic_log` 两个 ODS 主题；`DimApp` 用广播的维表配置从 `topic_db` 中把字典、主数据刷进 Phoenix/HBase 维表库；`BaseLogApp` 把 `topic_log` 拆成五类 `dwd_traffic_*_log` 明细主题；`DwdTradeOrderPreProcess` 等交易任务把 `topic_db` 里的多张业务表拼成订单宽表主题再细分；DWS 任务读取各类 `dwd_*` 主题，过滤口径、开窗聚合后把指标写进 ClickHouse。** Kafka 主题贯穿始终，是层与层之间最核心的衔接机制；DIM 与 DWS 各自的存储（Phoenix、ClickHouse）则是面向查询的终点。理解这张图，就等于理解了整个实时数仓里每一条数据从哪来、被谁加工、又往哪去。

Sources: [gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java:23-25](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java:23-25) · [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:13-17](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:13-17) · [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:29-34](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:29-34)

---

## 03. 计算模型与存储选型

> 围绕“算在哪、存哪里”讲清技术底座：DataStream 与 FlinkSQL 两种写法各自的使用场景，状态与 Checkpoint 在生产上的注意点，以及 Kafka、MySQL、Phoenix/HBase、ClickHouse、Redis 在链路中的分工。重点解释针对维度关联的 Redis 旁路缓存与异步查询、维表配置的广播状态等系统优化手段。

- Page Markdown: https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/pages/03-page-3.md
- Generated: 2026-09-08T05:31:55.833Z

### Source Files

- `gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/PhoenixUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/MysqlUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/utils/JedisUtil.java`
- `gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java`

<details>
<summary>相关源文件</summary>
以下文件用于生成此维基页面：
- [gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/PhoenixUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/PhoenixUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/MysqlUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/MysqlUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/JedisUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/JedisUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/utils/ThreadPoolUtil.java](gmall-realtime/src/main/java/com/atguigu/utils/ThreadPoolUtil.java)
- [gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java](gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java)
- [gmall-realtime/src/main/java/com/atguigu/app/func/TableProcessFunction.java](gmall-realtime/src/main/java/com/atguigu/app/func/TableProcessFunction.java)
- [gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java](gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java)
- [gmall-realtime/src/main/java/com/atguigu/app/func/DimAsyncFunction.java](gmall-realtime/src/main/java/com/atguigu/app/func/DimAsyncFunction.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdInteractionComment.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdInteractionComment.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradeProvinceOrderWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradeProvinceOrderWindow.java)
</details>

# 计算模型与存储选型

这一页讲两件互相咬合的事：**业务逻辑用哪种方式算**，以及**算出来的数据存在哪里**。项目里同一个 Flink 运行时里混着两套计算写法——过程式的 DataStream API 和声明式的 FlinkSQL，它们各有顺手的分层；而 Kafka、MySQL、HBase/Phoenix、ClickHouse、Redis 各管一段存储，谁存什么由"读写模式 + 查询时延"决定。把这层关系看懂，再去看每个作业的代码就会清楚"为什么这个算子放这里、为什么这张表放这个库"。

同时，维度关联是实时数仓里最容易拖垮吞吐的环节，本页会重点拆解仓库里针对它的三套手段：**Redis 旁路缓存、异步 IO、广播状态维表配置**，以及状态与 Checkpoint 的生产注意点。Sources: [DimUtil.java:47-109](gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java)

## 一个先入为主的整体画面

下面的图不是"数据流水线逐层倒一遍"（那是另一页的话题），而是强调**每一层由哪套计算模型驱动、结果落到哪个存储**——这是本页真正要教的空间结构：

```mermaid
flowchart TB
    subgraph SRC["数据源头"]
        BIZ[("MySQL 业务库<br/>订单/用户/商品")]
        CFG[("MySQL 配置与字典<br/>table_process / base_dic")]
        LOG["日志埋点"]
    end

    KAFKA[("Kafka 消息总线<br/>topic_db → dwd_* → dws_*")]

    subgraph FLINK["Flink 实时计算"]
        DIM["DIM 层 DimApp<br/>DataStream + 广播状态<br/>维表入 Phoenix"]
        DWD["DWD 明细层<br/>FlinkSQL 声明式关联"]
        DWS["DWS 汇总层<br/>DataStream 状态 + 窗口聚合"]
    end

    subgraph STORE["维表与结果存储"]
        PHOENIX[("HBase/Phoenix<br/>维表层")]
        CH[("ClickHouse<br/>OLAP 结果")]
        REDIS[("Redis<br/>维表旁路缓存")]
    end

    PUB["Publisher 查询服务"]

    BIZ --> KAFKA
    LOG --> KAFKA
    KAFKA --> DIM
    KAFKA --> DWD
    KAFKA --> DWS
    DWD --> KAFKA
    DIM -->|"写入/按配置建表"| PHOENIX
    DIM -. "配置流(CDC)" .-> CFG
    DWD -. "Lookup 字典" .-> CFG
    DWS -. "异步读维表" .-> REDIS
    REDIS -. "未命中回落" .-> PHOENIX
    DIM -. "更新后删缓存" .-> REDIS
    DWS --> CH
    CH --> PUB
```

## 两种计算模型的分工

Flink 应用层只有两种写法，本项目都用了，且分工稳定。

**DataStream API（过程式，细粒度掌控）**：一切需要"我精确知道每一步在干嘛"的逻辑都走这里。典型场景是 DIM 层维表处理（广播配置、按表过滤字段）、DWS 层窗口聚合（事件时间水位线、`windowAll`/`keyedWindow`、定时器状态去重）。它把状态、窗口、定时器、异步函数、广播状态都显式暴露给开发者。Sources: [DwsTrafficPageViewWindow.java:51-66](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java)

**FlinkSQL（声明式，交给引擎优化）**：主要用于 DWD 层的事实明细关联。把 Kafka 主题声明成"表"，剩下用一条条 `select` 表达过滤、多流关联、维度退化，Join 顺序与状态保留由引擎处理。SQL 写法让"关联订单、明细、活动、优惠券再加字典表"这类五表联动的逻辑可读性远高于手写 state。Sources: [DwdTradeOrderPreProcess.java:27-40](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java)

| 对比维度 | DataStream API | FlinkSQL |
|---|---|---|
| 编写心智 | 过程式：算子 + 回调 | 声明式：建表 DDL + SQL |
| 在库代表 | DIM 层、DWS 层全部作业、日志清洗 | DWD 明细关联类作业 |
| 状态控制 | 显式声明 ValueState/广播状态 | 引擎隐式管理，靠 TTL 约束 |
| 维表关联 | 手写异步函数 + 缓存 | 声明 Lookup Join，连接器自带缓存 |
| 适合问题 | 去重、窗口聚合、配置驱动分流 | 多流 Join、按主键 Upsert 的宽表加工 |

两种模型还能混在同一作业里：`StreamTableEnvironment` 桥接 DataStream 与 Table，SQL 算完可以再转回流做过程式收尾。Sources: [DwdInteractionComment.java:29-48](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdInteractionComment.java)

## 存储分工：谁存什么、为什么

| 存储 | 负责的层 | 读写模式 | 承担的角色 |
|---|---|---|---|
| Kafka | 层间传输总线 | 流式读写、追加 | 各层之间唯一的"话传话"通道，主题天然带业务主键语义 |
| MySQL | 业务源 + 元数据 | 联机事务 | 业务库经 CDC 进 Kafka；`table_process` 配置库和 `base_dic` 字典表直接被实时任务当"活配置"用 |
| HBase/Phoenix | DIM 维表层 | 主键点查 + 批量 Upsert | 维表唯一事实源，SQL 语法 + 行键模型，天然适合"一维一条记录" |
| ClickHouse | DWS 结果层 | 批量写、OLAP 聚合读 | 汇总结果落库，供查询服务按 `sum(order_amount)` 这类大范围聚合秒回 |
| Redis | 维表加速缓存 | 点查 + TTL | 挡住对 Phoenix 的重复主键查询 |
| HDFS | 状态后端 | 分布式文件 | Checkpoint 落盘位置，恢复语义的根基 |

几个值得注意的映射逻辑：维表写入 Phoenix 时用小写 schema `GMALL211027_REALTIME` 前缀拼库表，Phoenix 本身用 `upsert into`（有则改无则插），主键语义靠连接串指向 ZooKeeper；ClickHouse 则没有"改"，只有整块 `insert`，所以它只承接**只追加的窗口聚合结果**，不适合放会被反复修改的业务明细。Sources: [GmallConfig.java:13-26](gmall-realtime/src/main/java/com/atguigu/common/GmallConfig.java)、[PhoenixUtil.java:52-96](gmall-realtime/src/main/java/com/atguigu/utils/PhoenixUtil.java)、[ClickHouseUtil.java:24-60](gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java)

另一条容易被忽略的规则：**写到 Kafka 的下游主题如果是主键 Upsert 语义，就要用 `upsert-kafka` 连接器而不是普通 `kafka` 连接器**。否则 `left join` 产生的空右表（`null`）无法表达"删掉旧值"的意图，项目里为此约定把 null 序列化成空串、下游再判重兜底。Sources: [MyKafkaUtil.java:108-124](gmall-realtime/src/main/java/com/atguigu/utils/MyKafkaUtil.java)

### ClickHouse 批量写的两个默认值

DWS 结果落 ClickHouse 走的是统一的 JDBC Sink，靠反射把 Bean 字段按序填进 `PreparedStatement` 占位符，用 `TransientSink` 注解跳过不需要入库的字段。写批有两个硬编码节流参数：**每 5 秒刷一次、每批 5 条**——这是为了教学环境控制频率，生产上这两个值都偏小，吞吐受限于此，需要按业务量调大。Sources: [ClickHouseUtil.java:31-55](gmall-realtime/src/main/java/com/atguigu/utils/ClickHouseUtil.java)

## 系统优化一：维表关联的 Redis 旁路缓存

DWS 关联省份、品牌、SPU 这类维度时，若每条事实都直连 Phoenix，主键查询会变成热点。仓库的通用解法是**缓存旁路（Cache-Aside）**：先查 Redis，命中直接返回；未命中才查 Phoenix，查到后回填 Redis 并设一天 TTL；查询结果永远以 Phoenix 为准，缓存只负责加速。

```java
// gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java
StringBuilder redisKey = new StringBuilder("dim:" + tableName.toLowerCase() + ":");
...
dimJsonStr = jedis.get(redisKey.toString());          // 1. 先查缓存
if (dimJsonStr != null && dimJsonStr.length() > 0) {
    dimJsonObj = JSON.parseObject(dimJsonStr);        // 命中即返回
} else {
    List<JSONObject> dimList = PhoenixUtil.queryList(...); // 2. 未命中查 Phoenix
    jedis.setex(redisKey.toString(), 3600 * 24, dimJsonObj.toJSONString()); // 3. 回填
}
```

配套的**失效**动作在维表更新侧：`DimSinkFunction` 收到 `update` 类型记录时，先删对应 Redis key 再写 Phoenix，保证下次查询能拿到新值——这就是"写更新时主动清缓存、读未命中时被动重建"的标准 Cache-Aside 闭环。Redis 连接用单例连接池（最大 100、阻塞等待上限 2 秒），避免每条记录都建连。Sources: [DimUtil.java:47-109](gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java)、[DimSinkFunction.java:50-53](gmall-realtime/src/main/java/com/atguigu/app/func/DimSinkFunction.java)、[JedisUtil.java:19-37](gmall-realtime/src/main/java/com/atguigu/utils/JedisUtil.java)

用时序图看这条闭环最清楚：

```mermaid
sequenceDiagram
    participant Fact as DWS 事实流
    participant R as Redis
    participant P as Phoenix 维表
    participant D as DimSink 更新流

    Fact->>R: GET dim:表名:主键
    alt 命中
        R-->>Fact: 返回缓存 JSON
    else 未命中
        Fact->>P: 主键查询维表
        P-->>Fact: 维度数据
        Fact->>R: SETEX（TTL = 1 天）
    end
    D->>R: 收到 update 类型 → 删除缓存 key
    D->>P: upsert 新维表数据
```

**一个生产注意点**：仓库里同一套"删缓存"逻辑写了两份，key 的拼法不一致——`delDimInfo` 用 `"DIM:" + 表名(大写)` 前缀，而查询与 `deleteCached` 用的是 `"dim:" + 表名(小写)`。若直接照搬现状，更新维表时主动删除大概率打不中真正被查询的 key，缓存只能等一天 TTL 自然过期，导致维表变更后一段时间内读到旧值。上线前需要统一 key 的大小写规则。Sources: [DimUtil.java:27-32](gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java)、[DimUtil.java:145-147](gmall-realtime/src/main/java/com/atguigu/utils/DimUtil.java)

## 系统优化二：异步 IO 关联维表

即便有 Redis，网络往返仍是每条记录串行等待。仓库把维表查询包进 **`RichAsyncFunction` + 线程池**：每个元素到来后不是阻塞当前算子线程，而是把"取主键 → 查缓存/库 → 结果回填"丢给共享线程池，主线程立即去处理下一条。调用侧用 `AsyncDataStream.unorderedWait` 挂接，并给整体等待设了 60 秒超时兜底。Sources: [DimAsyncFunction.java:33-65](gmall-realtime/src/main/java/com/atguigu/app/func/DimAsyncFunction.java)、[DwsTradeProvinceOrderWindow.java:131-146](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradeProvinceOrderWindow.java)

线程池是懒加载单例，核心 4 线程、上限 20，队列无界——意味着查询慢时请求会排队而不是丢数据，但也提示生产要按维度查询时延评估池子大小与队列上限，避免 OOM。Sources: [ThreadPoolUtil.java:18-32](gmall-realtime/src/main/java/com/atguigu/utils/ThreadPoolUtil.java)

值得注意的是：这套"异步 + 旁路缓存"的读维表姿势，只在 DWS 的 DataStream 作业里出现；**DWD 的 SQL 作业不写这套代码**，因为 FlinkSQL 的 Lookup Join 把缓存下推给了连接器本身。

## 系统优化三：广播状态管理维表配置（配置驱动建表）

DIM 层最聪明的一手是"维表该不该建、建哪些列"不是写死在代码里的，而是存在 MySQL 配置库的一张 `table_process` 表里，运行时用 FlinkCDC 持续监听这张表，变化广播到所有并行子任务。

```java
// gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java
MapStateDescriptor<String, TableProcess> mapStateDescriptor =
        new MapStateDescriptor<>("map-state", String.class, TableProcess.class);
BroadcastStream<String> broadcastStream = mysqlSourceDS.broadcast(mapStateDescriptor);
BroadcastConnectedStream<JSONObject, String> connectedStream =
        jsonObjDS.connect(broadcastStream);   // 主流(业务CDC) + 广播流(配置)
```

广播流侧收到一条新配置时，`TableProcessFunction` 会**当场拼 `create table if not exists` 把 Phoenix 维表建出来**（主键列由 `sinkPk` 指定），再把 `sourceTable → 目标表/列清单` 写进广播状态；主流侧每来一条业务数据，就按 `sinkColumns` 白名单过滤字段、补上 `sinkTable` 名，然后交给 Sink 写 Phoenix。Sources: [DimApp.java:73-99](gmall-realtime/src/main/java/com/atguigu/app/dim/DimApp.java)、[TableProcessFunction.java:49-72](gmall-realtime/src/main/java/com/atguigu/app/func/TableProcessFunction.java)、[TableProcessFunction.java:78-142](gmall-realtime/src/main/java/com/atguigu/app/func/TableProcessFunction.java)

收益很实在：**新增一张维表或调整字段，只改配置库一行，任务热更新配置即可，不必改代码重启作业**。这是"配置驱动数据管道"的典型落地。

## 状态与 Checkpoint 的生产注意点

仓库里状态和 Checkpoint 的用法分散，把它们汇总成一份可操作的清单：

- **状态是用来做"记忆"的**。典型如订单明细去重：按明细 ID 分组，把第一条数据存进 `ValueState` 并注册一个 2 秒后的处理时间定时器，期间到达的重复数据用时间戳比较取新弃旧，定时器触发时只输出一次——先到先得 + 定时兜底，天然抗 Kafka 重放。Sources: [OrderDetailFilterFunction.java:46-88](gmall-realtime/src/main/java/com/atguigu/app/func/OrderDetailFilterFunction.java)
- **SQL 作业要主动设状态保留时长**。SQL 关联产生的隐式状态默认永不过期，仓库在订单预处理作业里显式设了 3 天空闲保留，避免长尾 key 撑爆内存。Sources: [DwdTradeOrderPreProcess.java:37](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java)
- **Checkpoint 的标准模板**：`EXACTLY_ONCE`、3 秒一次、超时 60 秒、两次 Checkpoint 最小间隔 3 秒；取消作业时**保留**外部 Checkpoint 以便从最近一次恢复；失败重启用"一天内最多 3 次、每次间隔 1 分钟"的失败率策略。状态后端是 `HashMapStateBackend`，落盘位置是 HDFS。Sources: [DwsTrafficPageViewWindow.java:51-66](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java)、[DwdInteractionComment.java:35-48](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdInteractionComment.java)
- **教学代码常把 Checkpoint 注释掉**。DIM 作业里状态后端与 Checkpoint 整段是注释状态，部分 DWS 作业亦然——这适合本地演示，但**生产必须打开**，否则任何故障都是全量重算。此外多个作业硬编码 `setParallelism(1)`，同样只适合单机教学。
- **SQL 时区要显式设**。Kafka 主题里的时间戳是 `timestamp_ltz(3)`，仓库把表环境时区固定为 `GMT+8`，避免窗口按 UTC 切导致"8 小时错位"。Sources: [DwdTradeOrderDetail.java:44-48](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java)

## 小结

一句话收束这张图景：**DataStream 管"要精细控制"的 DIM/DWS，FlinkSQL 管"声明式关联"的 DWD；Kafka 当层间总线，HBase/Phoenix 当维表事实源，ClickHouse 当结果库，Redis 当维表缓存，MySQL 既是业务源头又是"活配置"的存放处。** 维表关联这个最大瓶颈，仓库用三件套应对——Redis 旁路缓存挡重复查询、异步 IO 挡串行等待、广播状态让维表结构与建表动作配置化；而它们能否在生产站住脚，取决于 Checkpoint 是否真正打开、缓存 key 是否一致、批参数是否按吞吐调过。把这几条对上了，实时链路才算"能上生产"而不是"能跑通"。

---

## 04. 业务域落地：指标口径与可复用经验

> 把前几页的通用机制放回业务场景，以流量、交易、用户三个域为线索还原典型指标在代码中的落地方式：新老用户判定与 UV、跳出行为、页面浏览窗口、交易预处理后的支付与下单汇总、注册回流统计。收尾给出从这套工程抽象出的实时数仓分层方法论与踩坑经验，作为全文总结。

- Page Markdown: https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/pages/04-page-4.md
- Generated: 2026-09-08T05:24:55.785Z

### Source Files

- `gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUserJumpDetail.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficVcChArIsNewPageViewWindow.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradePayDetailSuc.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradeOrderWindow.java`
- `gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserRegisterWindow.java`

<details>
<summary>相关源文件</summary>
以下文件用于生成此维基页面：

- [gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUniqueVisitorDetail.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUniqueVisitorDetail.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUserJumpDetail.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUserJumpDetail.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficVcChArIsNewPageViewWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficVcChArIsNewPageViewWindow.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradePayDetailSuc.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradePayDetailSuc.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradeOrderWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradeOrderWindow.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradePaymentSucWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradePaymentSucWindow.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdUserRegister.java](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdUserRegister.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserRegisterWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserRegisterWindow.java)
- [gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java)
- [gmall-realtime/src/main/java/com/atguigu/bean/TrafficPageViewBean.java](gmall-realtime/src/main/java/com/atguigu/bean/TrafficPageViewBean.java)
- [gmall-realtime/src/main/java/com/atguigu/bean/UserLoginBean.java](gmall-realtime/src/main/java/com/atguigu/bean/UserLoginBean.java)

</details>

# 业务域落地：指标口径与可复用经验

前几页介绍了实时数仓的分层骨架与通用算子，这一页把这些机制放回三个典型业务域——流量、交易、用户，回答一个更实际的问题：**同一个指标在真实代码里到底是怎么算出来的**。你会发现"新用户"在流量域、交易域、用户域是三套不同的口径；"跳出"在注释里描述的算法和实际生效的 CEP 模式并不完全一致；注册数是一回事，回流统计又是另一回事。理清这些口径差异，比记住某个 API 更接近这套工程的核心经验。

阅读方式是"顺着一个事件看它如何被层层加工"：ODS 的埋点与业务变更流先进 DWD 变成**语义明确的细粒度事实**（一次跳出、一条 UV、一条下单、一条支付），再由 DWS 按窗口与维度**预聚合**成可直接查询的汇总。本节把三个域串起来，收尾提炼出可复用的分层方法与踩坑清单。

## 一、三个域的落地全景

三个域的数据源头、加工链路与最终出口不同，但骨架一致。先看全局对照：

| 域 | ODS 入口 | DWD 关键加工 | DWS 汇总（窗口粒度） | 出口存储 |
|---|---|---|---|---|
| 流量 | `topic_log`（埋点日志） | 新老访客修正、分流、UV 去重、跳出识别 | 访客类别页面浏览、首页/详情 UV 窗口 | ClickHouse |
| 交易 | `topic_db`（业务库 CDC） | 订单预处理宽表、下单/支付拆分 | 下单汇总、支付成功汇总 | ClickHouse |
| 用户 | `topic_db`（用户表 insert） | 注册事实 | 注册窗口、登录/回流窗口 | ClickHouse |

其中流量域的分工最典型，值得用一张图说明"DWD 越做越窄、DWS 越做越宽"：

```mermaid
flowchart LR
    subgraph ODS
        L["topic_log<br/>全量埋点日志"]
    end
    subgraph DWD_流量域["DWD 流量域（每条任务只回答一个问题）"]
        B["BaseLogApp<br/>新老访客修正 + 五路分流"]
        U["DwdTrafficUniqueVisitorDetail<br/>取会话首跳 + 按 mid 去重"]
        J["DwdTrafficUserJumpDetail<br/>CEP 识别跳出（10s 窗口）"]
    end
    subgraph DWS_流量域["DWS 流量域（多主题合并回宽）"]
        W["DwsTrafficVcChArIsNewPageViewWindow<br/>vc/ch/ar/is_new × 10s 窗口"]
        P["DwsTrafficPageViewWindow<br/>首页 / 商品详情 UV"]
    end
    L -->|topic_log| B
    B -->|dwd_traffic_page_log| U
    B -->|dwd_traffic_page_log| J
    B -->|dwd_traffic_page_log| W
    B -->|dwd_traffic_page_log| P
    U -->|dwd_traffic_unique_visitor_detail| W
    J -->|dwd_traffic_user_jump_detail| W
    W --> CH1["ClickHouse: dws_traffic_vc_ch_ar_is_new_page_view_window"]
    P --> CH2["ClickHouse: dws_traffic_page_view_window"]
```

> 这张图的关键是 DWD 的三个任务各自从同一条页面日志里切出"互斥且完整"的一类事件，而 DWS 的窗口任务再把它们重新 union 起来、按业务键聚合。**一个指标只在一个地方算第一次**（在 DWD 产出 0/1 增量），DWS 只负责累加——这是本仓库最值得复用的设计。

Sources: [DwsTrafficVcChArIsNewPageViewWindow.java:66-81](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficVcChArIsNewPageViewWindow.java:66-81)

## 二、流量域：新老访客、UV 与跳出

### 1. 新老访客判定：一次"防作弊"修正

埋点 SDK 会随每条日志上报 `common.is_new`（1=新访客）。但这个标记是前端在**设备首次访问**时打的，一旦用户清缓存或换入口，前端可能重复上报"新"。`BaseLogApp` 在分流前先按 `mid`（设备/访客标识）分组，用状态记住**最近一次访问日期**做二次校验：

```java
// BaseLogApp.java: 以 mid 分组后的状态校验
if ("1".equals(isNew)) {                     // 前端声称是新访客
    String curDt = DateFormatUtil.toDate(ts);
    if (lastVisitDt == null) {
        lastVisitDtState.update(curDt);       // 第一次见，信任它
    } else if (!lastVisitDt.equals(curDt)) {
        value.getJSONObject("common").put("is_new", "0");  // 之前来过 -> 改判老访客
    }
} else if (lastVisitDt == null) {
    // 状态为空却声称老访客：回填昨天，防止"当天先见老客、后见新客"被误判
    String yesterday = DateFormatUtil.toDate(ts - 24 * 60 * 60 * 1000L);
    lastVisitDtState.update(yesterday);
}
```

这条规则解决了一个很隐蔽的时序问题：如果一天内某设备第一条日志是老访客标记、当天后来才来一条新访客标记，若不回填昨天，后者会误判为"今天的新用户"。修正后的 `is_new` 会写回 `common`，随主流一起进入下游，DWS 直接把它当维度切分（如 `is_new` 进 group key）即可。值得注意的是，这里"新老"以 **mid（设备）** 为单位，与后文用户域以 `uid` 为单位的"新用户"完全是两个口径。

Sources: [BaseLogApp.java:55-98](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java:55-98)

### 2. UV 独立访客：会话起点 + 按日去重

`DwdTrafficUniqueVisitorDetail` 处理"日活 UV"：先过滤出 `page.last_page_id == null` 的事件——`last_page_id` 为空代表这是一段会话的**第一个页面**，只有它才可能代表一次"新的访问"；然后按 `mid` 分组，用状态存"最近一次见到的日期"，同一天内重复出现的 mid 直接丢弃：

```java
// DwdTrafficUniqueVisitorDetail.java: 状态带 1 天 TTL
ValueStateDescriptor<String> d = new ValueStateDescriptor<>("visit-dt", String.class);
d.enableTimeToLive(new StateTtlConfig.Builder(Time.days(1))
        .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite).build());
// filter 中：dt == null || !dt.equals(curDt) 才放行并更新状态
```

也就是说，每个 mid 每天最多向 `dwd_traffic_unique_visitor_detail` 输出一次。下游 DWS 统计 UV 时不需要再去重，**只需把收到的事件计为 1**。一个细节：UV 的主题名定义在变量 `targetTopic` 中，但写出的 sink 却写死了字符串 `"targetTopic"`，这是示例代码里的一处笔误，判断真实连通主题名时应以下游 DWS 消费的 `dwd_traffic_unique_visitor_detail` 为准（详见"踩坑"小节）。

Sources: [DwdTrafficUniqueVisitorDetail.java:56-101](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUniqueVisitorDetail.java:56-101)

### 3. 跳出行为：把"用户只看了这一页就走了"翻译成模式

跳出在业务上定义为"一次会话只浏览了一个页面就离开"。代码把它操作化为可计算的 CEP 模式：仍按 `mid` 分组、基于事件时间（允许 2 秒乱序），匹配"**两条紧邻的会话首跳事件**（`last_page_id == null`）在 **10 秒内**先后出现"——第一条首跳代表会话 A 的进入，第二条首跳代表会话 A 期间没有任何站内跳转、直接开启新会话，因此会话 A 是一次跳出：

```java
// DwdTrafficUserJumpDetail.java:105-116 实际生效的模式
Pattern.<JSONObject>begin("first").where(v -> v.getJSONObject("page").getString("last_page_id") == null)
        .next("second").where(v -> v.getJSONObject("page").getString("last_page_id") == null)
        .within(Time.seconds(10));
```

应用模式后分别从 `select` 主结果与 `timeout` 侧输出流中取"第一条首跳事件"，再 union 到一起写回 Kafka。文件注释与 README 还描述了一套"状态 + 10 秒定时器"的等价算法（来首跳存状态、注册定时器；再来一条首跳先清空再输出；定时器到点也输出），这是**同一个口径的两种教学实现**；类里另有一段 `times(2).consecutive()` 的严格近邻模式 `p2` 定义了但从未挂到 `CEP.pattern(...)` 上，属于死代码。看到这类双实现时要记住：**生效的口径以真正 `.pattern(...)` 应用的那个为准**。

Sources: [DwdTrafficUserJumpDetail.java:67-152](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUserJumpDetail.java:67-152)

### 4. 页面浏览窗口：把三类事件 union 回同一张表

`DwsTrafficVcChArIsNewPageViewWindow` 是这套设计的高潮：它同时消费**页面日志、跳出明细、UV 明细**三条主题，各自转成同一张 `TrafficPageViewBean`，用字段位表达"这条事件贡献了什么"：

- 页面日志每条计 `pvCt=1`；若 `last_page_id == null` 再计 `svCt=1`（新会话）；`during_time` 累加到 `durSum`。
- 跳出明细计 `ujCt=1`，且时间戳**顺延 10 秒**——跳出是在超时/下一首跳发生时才知道的，把事件归到它真正结束的窗口。
- UV 明细计 `uvCt=1`，时间戳沿用首跳原时间。

```java
// 合并三类流后：水位线延迟 13s（= 最慢的跳出流延迟 10s + 余量）
DataStream<TrafficPageViewBean> unionDS = pageViewBeanDS.assignTimestampsAndWatermarks(
    WatermarkStrategy.<TrafficPageViewBean>forBoundedOutOfOrderness(Duration.ofSeconds(13)) ...);

// 按 vc / ch / ar / is_new 四维 keyBy，10s 滚动事件时间窗口，allowedLateness 10s
windowedStream.reduce(累加五个度量, new ProcessWindowFunction<...>() { /* 补 stt/edt 窗口起止 */ });
```

合并后按 `vc(版本) / ch(渠道) / ar(地区) / is_new(新老访客)` 做 keyed 窗口——**只有这个窗口任务是按键开窗**，其它多为 `windowAll` 全局窗口。窗口处理是"增量 Reduce + 末尾 ProcessWindowFunction 补窗口起止时间"的组合：增量累加省内存，窗口函数只做一次性的时间标注。最终写 ClickHouse 的 `dws_traffic_vc_ch_ar_is_new_page_view_window`。

Sources: [DwsTrafficVcChArIsNewPageViewWindow.java:86-278](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficVcChArIsNewPageViewWindow.java:86-278)

### 5. 跨天不重复的"首页/详情 UV"

`DwsTrafficPageViewWindow` 回答的是更细的问题：**首页 UV 有多少、商品详情页 UV 有多少**。做法与全站 UV 不同——它不过滤首跳，而是只保留 `page_id ∈ {home, good_detail}` 的事件，按 `mid` 分组后为两类页面**各维护一个"最近访问日期"状态**，同一天重复访问不计：

```java
// DwsTrafficPageViewWindow.java:106-163 核心：页面级跨天去重
if (pageId.equals("home") && (homeLastDt == null || !homeLastDt.equals(visitDt))) {
    homeUvCt = 1L;  homeLastVisitDt.update(visitDt);
}
if (pageId.equals("good_detail") && (detailLastDt == null || !detailLastDt.equals(visitDt))) {
    detailUvCt = 1L;  detailLastVisitDt.update(visitDt);
}
```

注意它把"当天是否首次访问"的判断收敛在 **DWS 的状态算子**里（不是 DWD），输出的是 0/1 增量；随后 `windowAll` 10 秒滚动窗口只做纯累加。同一个"UV 去重"问题，全站口径放在 DWD、页面口径放在 DWS，说明**去重位置由下游聚合键决定**——键细、只有少数消费者时，就地判重更简单。

Sources: [DwsTrafficPageViewWindow.java:79-205](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTrafficPageViewWindow.java:79-205)

## 三、交易域：预处理宽表之上的下单与支付

### 1. 订单预处理：一张"可回放"的宽表

交易域的复杂度在于**一个订单明细会随订单状态反复变更**。`DwdTradeOrderPreProcess` 用 Flink SQL 把五路数据打成一张宽主题 `dwd_trade_order_detail`：

- `order_detail`（明细，仅 insert）与 `order_info`（订单主表，insert **和 update**，且保留 `type`、`old` 字段用于下游追踪变更）；
- 两张关联表 `order_detail_activity`、`order_detail_coupon`（均仅 insert）；
- `base_dic` 字典表以 **`FOR SYSTEM_TIME AS OF` Lookup Join** 补充来源类型名称。

```sql
-- DwdTradeOrderPreProcess.java:186-194（语义摘要）
from order_detail od
join order_info oi on od.order_id = oi.id
left join order_activity oa on od.order_detail_id = oa.order_detail_id
left join order_coupon  oc on od.order_detail_id = oc.order_detail_id
join base_dic FOR SYSTEM_TIME AS OF od.pt as dic on od.source_type = dic.dic_code
```

写出的目标是 **Upsert-Kafka**，主键 `order_detail_id`：同一明细的多次变更在主题里按主键更新，而不是追加多行。它还顺带完成了度量预计算：`split_original_amount = sku_num * order_price`，并把 `current_row_timestamp()` 记作 `row_op_ts`（行操作时间），为下游"同一条明细以哪个版本为准"提供依据。

Sources: [DwdTradeOrderPreProcess.java:44-251](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderPreProcess.java:44-251)

### 2. 一张宽表，各自取数：下单与取消

预处理表只解决"关联"，不回答"这条明细现在算什么"。每个业务由独立的 DWD 消费者按 `type`/`old` 切片：

- `DwdTradeOrderDetail`：只取 `type='insert'` 的行——**新增订单明细才算下单**，输出下单明细主题并带上 `date_id`、以 `od_ts` 为事件时间、透传 `row_op_ts`。
- 取消订单等场景（README 同节描述的同类消费者）则反过来监听 `order_status` 的变更。

这是 DWD 层的又一条可复用经验：**宽表只做 join 与退化维度，用操作类型与变更前后值驱动下游切流**，让"下单""取消""支付"各自成为独立的明细流，后续各 DWS 不必互相干扰。

Sources: [DwdTradeOrderDetail.java:83-139](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradeOrderDetail.java:83-139)、[README.md:157-180](README.md)

### 3. 支付成功：三表再关联出"支付宽表"

`DwdTradePayDetailSuc` 把支付事实与订单明细第二次拉宽：从 `topic_db` 筛出 `payment_info`（目前**只按表名过滤**，代码中 `type='update'`、`payment_status='1602'` 两个条件被注释掉了），按 `order_id` 与下单明细 join（一次支付对应多个 SKU 行），再 Lookup `base_dic` 把 `payment_type` 变成中文名。结果同样写 Upsert-Kafka、主键 `order_detail_id`。

```sql
-- DwdTradePayDetailSuc.java:88-135（语义摘要）
select od.*, pi.payment_type payment_type_code,
       dic.dic_name payment_type_name, pi.callback_time,
       od.split_total_amount split_payment_amount
from payment_info pi
join dwd_trade_order_detail od on pi.order_id = od.order_id
left join `base_dic` for system_time as of pi.proc_time as dic
  on pi.payment_type = dic.dic_code
```

"支付金额 = 订单明细拆分金额（split_payment_amount）"这一口径在这里被固化——下游 DWS 不再关心支付与 SKU 金额的换算。

Sources: [DwdTradePayDetailSuc.java:50-167](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradePayDetailSuc.java:50-167)

### 4. 下单与支付的窗口汇总：两套很相似的"当日/新增"口径

两个 DWS 窗口结构几乎一致，但"判新"依据的时间字段不同：

| 维度 | `DwsTradeOrderWindow`（下单） | `DwsTradePaymentSucWindow`（支付成功） |
|---|---|---|
| 输入主题 | `dwd_trade_order_detail` | `dwd_trade_pay_detail_suc` |
| 第一道去重 | 按 `order_detail_id`，用 `row_op_ts` 比较保留新值 + 5s 定时器兜底 | 按 `order_detail_id`，TTL 5s 状态"保留第一条" |
| "当日独立用户"依据 | 状态日期 vs `date_id`（**业务下单日**） | 状态日期 vs `callback_time`（**支付回调日**） |
| "新用户"判定 | 状态为空即首次下单 | 状态为空即首次支付 |
| 金额度量 | 原始金额、活动减免、券减免 | （支付窗口以人数为主，承接 DWD 的金额口径） |

`DwsTradeOrderWindow` 值得展开的是两道状态：第一道按 `order_detail_id` 分组，用 `ValueState` 暂存当前最优行，比较 `row_op_ts` 谁更新、再注册 5 秒处理时间定时器输出，**解决 Upsert 主题里同一明细的旧版本/迟到更新**；第二道按 `user_id` 分组，`lastOrderDtState` 为 null 时"新增下单用户 + 当日下单独立用户"都计 1，否则仅跨天计独立用户 1，三个金额累加后走 `windowAll` 10 秒窗口聚合。

```java
// DwsTradeOrderWindow.java:190-198 下单侧"新用户"即"状态首次见"
if (lastOrderDt == null) {
    orderNewUserCount = 1L;
    orderUniqueUserCount = 1L;
} else if (!lastOrderDt.equals(orderDt)) {
    orderUniqueUserCount = 1L;
}
lastOrderDtState.update(orderDt);
```

一个必须点破的口径陷阱：这里"下单新用户"不是"注册后首次下单"，而是**"任务启动后状态里第一次见到的用户"**。由于示例普遍注释掉了 checkpoint（见下文踩坑），任务重启会清空状态，"新用户数"会被重复累计。这类度量的正确性依赖状态长期保留。

Sources: [DwsTradeOrderWindow.java:98-298](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradeOrderWindow.java:98-298)、[DwsTradePaymentSucWindow.java:88-197](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradePaymentSucWindow.java:88-197)

## 四、用户域：注册是事实，回访是状态

### 1. 注册：最朴素的"insert 过滤 + 计数"

注册事实很直接：业务库里 `user_info` 表的 insert 就是新增注册。`DwdUserRegister` 用 SQL 从 `topic_db` 过滤 `table='user_info' and type='insert'`，把 `create_time` 格式化成 `date_id` 写入 Upsert-Kafka 的 `dwd_user_register`（主键 `user_id`）。

`DwsUserUserRegisterWindow` 则是全套工程里最"标准"的窗口任务：读主题 → 转 `UserRegisterBean`（每条先置 1）→ 事件时间 + 2s 水位线 → `windowAll` 10 秒滚动窗口 reduce 求和 → 补窗口起止 → 写 ClickHouse。它没有任何状态去重，因为**去重责任已经由 DWD 的主键 Upsert 与"insert 才输出"的过滤承担**——这正是"度量尽量在 DWD 完成"原则的最小范例。

Sources: [DwdUserRegister.java:40-80](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdUserRegister.java)、[DwsUserUserRegisterWindow.java:60-108](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserRegisterWindow.java:60-108)

### 2. 回流：藏在"登录窗口"里的状态机

"回流统计"并不在注册窗口，而在用户登录窗口 `DwsUserUserLoginWindow` 里，且口径完全不同：它关心的是**老用户的回访**，而不是新注册。任务先对页面日志过滤出 `uid != null && last_page_id == null`——即有登录态的会话第一条事件，然后按 `uid` 分组，用状态记录上次访问日期：

```java
// DwsUserUserLoginWindow.java:118-135 回流判定
if (lastDt == null) {                 // 状态为空 -> 当天首个独立登录用户
    uuCt = 1L;  lastVisitDt.update(curDt);
} else if (!lastDt.equals(curDt)) {   // 跨天再来
    uuCt = 1L;  lastVisitDt.update(curDt);
    if ((ts - DateFormatUtil.toTs(lastDt)) / (1000L*60*60*24) >= 8L) {
        backCt = 1;                   // 距上次 ≥8 天 -> 回流用户
    }
}
```

于是同一段代码产出两个互补指标：`uuCt`（当日独立登录用户）与 `backCt`（回流用户，间隔 ≥8 天的回访）。对照来看，"注册"是业务事实（拉新效果），"回流"是基于行为的回访判定（召回效果），前者消费 `dwd_user_register`，后者消费登录会话事件——**两个指标放在不同链路里，因为它们回答的问题不同**。

Sources: [DwsUserUserLoginWindow.java:73-175](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsUserUserLoginWindow.java:73-175)

## 五、从这套工程抽象出的方法与踩坑

### 可复用的分层方法论

1. **度量在 DWD 算第一次，DWS 只累加。** UV、跳出、新增注册的 0/1 判定都在 DWD 完成，窗口任务普遍只有 `reduce` 求和。好处是同一指标不会在多个 DWS 里用不同的方式重复实现。
2. **一个宽表 + 按操作类型切多路明细。** 交易域先打一张带主键的预处理宽表，下游各自按 `type`/`old`/`row_op_ts` 取数，避免为每个业务各做一遍五表 join。
3. **两套时间要分清。** "当日"判定用业务时间（`date_id`/`callback_time`），窗口归属用事件时间（日志 `ts`、`row_op_ts`）。混用两套时间时，要为"度量流"与"窗口流"各自选择正确的字段。
4. **多流合并先对齐再聚合。** 跳出明细延迟 10 秒才产生，合并任务就把它的 `ts` 顺延 10 秒、并放宽水位线到 13 秒，让各流的度量落在同一批窗口里。
5. **两种技术栈按场景取舍。** 日志域用 DataStream API + 状态/CEP（轻量、算子自由），业务库强关联场景用 Flink SQL + Lookup Join + Upsert-Kafka（声明式、免去手动 join）。

### 踩坑清单

- **`is_new` 必须二次修正。** 前端标记可被重复上报，`BaseLogApp` 用按 mid 的最近访问日期兜底，并回填"昨天"防止首日排序导致的误判；下游用 `is_new` 当维度前要确认它已过这道修正。Sources: [BaseLogApp.java:73-94](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/BaseLogApp.java:73-94)
- **示例代码存在未生效/笔误处，别照抄主题名。** 跳出类的 `p2` 模式从未挂到 `pattern()` 上；UV 明细的 sink 写死字符串 `"targetTopic"`，与变量声明的 `dwd_traffic_unique_visitor_detail` 不一致。排查连通性时应以消费端实际订阅的主题为准。Sources: [DwdTrafficUserJumpDetail.java:118-130](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUserJumpDetail.java:118-130)、[DwdTrafficUniqueVisitorDetail.java:99-101](gmall-realtime/src/main/java/com/atguigu/app/dwd/log/DwdTrafficUniqueVisitorDetail.java:99-101)
- **支付 DWD 的状态过滤被注释，去重被推给下游。** `payment_status='1602'` 等条件未启用，任何 `payment_info` 行都会进宽表；下游支付窗口的按 `order_detail_id` 去重状态只有 5 秒 TTL，**超过 5 秒的重复或迟到事件仍可能重复计数**。Sources: [DwdTradePayDetailSuc.java:96-100](gmall-realtime/src/main/java/com/atguigu/app/dwd/db/DwdTradePayDetailSuc.java:96-100)、[DwsTradePaymentSucWindow.java:99-116](gmall-realtime/src/main/java/com/atguigu/app/dws/DwsTradePaymentSucWindow.java:99-116)
- **"新用户"类指标都是"状态首次见"语义，且多处 checkpoint 被注释。** `DwsTradeOrderWindow`、`DwsUserUserRegisterWindow` 等把 checkpoint/状态后端配置整体注释掉，教学可跑；生产必须开启 Exactly-Once + 持久化状态存储，否则重启即清空判定状态，新用户/回流等累积型指标会失真。
- **全局窗口与按键窗口的选择。** 多数汇总用 `windowAll`（数据量小、简单），只有多维度下钻的访客类别窗口用 `keyBy(vc,ch,ar,isNew)` 分担负载；决定前先估算单窗口吞吐，避免 `windowAll` 成为单点瓶颈。

### 小结

把三域放在一起看，这套工程的核心不是某个框架特性，而是一条纪律：**在 DWD 把口径固化成细粒度事实，在 DWS 只做带窗口的累加，跨业务复用"主键 Upsert + 状态判重 + 事件时间对齐"三种武器**。流量域教你怎么把"跳出""UV"这类模糊行为翻译成可计算的事件序列；交易域教你用一张可回放的宽表支撑下单、支付、取消等多个口径；用户域提醒你注册与回流分属事实与状态两条链路。理解了这些口径从哪来、在哪一步被算定，再去看任何一张 ClickHouse 汇总表的数字，都能沿着 Kafka 主题一路追回最初的埋点或 CDC 变更——这正是实时数仓最值得沉淀的经验。

---
