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

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

- 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/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 变更——这正是实时数仓最值得沉淀的经验。
