# 数仓分层与数据流转

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

- 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/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)
