# 计算模型与存储选型

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

- Repository: aggaadfr/gmall-flink-3.0
- GitHub: https://github.com/aggaadfr/gmall-flink-3.0
- Human wiki: https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69
- Complete Markdown: https://grok-wiki.com/public/wiki/aggaadfr-gmall-flink-3-0-446e0d5adb69/llms-full.txt

## 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 是否一致、批参数是否按吞吐调过。把这几条对上了，实时链路才算"能上生产"而不是"能跑通"。
