|
@@ -73,18 +73,19 @@ DS 项目级 globalParams(poyee-data-warehouse 项目,所有工作流继承
|
|
|
|
|
|
|
|
## 4. raw 层 dt 语义
|
|
## 4. raw 层 dt 语义
|
|
|
|
|
|
|
|
-- **抓窗**:`COALESCE(update_time, create_time) ∈ [${dt}, ${tdt}) = [T-1, T+1)`,48h 宽窗(raw ini reader.where)
|
|
|
|
|
|
|
+- **抓窗(OR 双锚)**:`create_time ∈ [${dt}, ${tdt}) OR update_time ∈ [${dt}, ${tdt})`,即 `[T-1, T+1)` 48h 宽窗(raw ini reader.where)。**不用 COALESCE**——`COALESCE(update, create)` 取 update,会把 `update < create` 的脏单甩出窗漏掉;OR 双锚只要 create 或 update 任一落窗即抓(业务会改 create_time,二者皆活动性锚点,见 §5)
|
|
|
- **写入分区**:dt = `${dt}` = T-1
|
|
- **写入分区**:dt = `${dt}` = T-1
|
|
|
-- **48h 宽窗设计目的**:覆盖跨日 update_time 漂移(详见 ADR-03 §零点漂移决策)
|
|
|
|
|
|
|
+- **48h 宽窗设计目的**:覆盖跨日漂移(详见 ADR-03 §零点漂移决策)
|
|
|
- **重跑幂等**:DataX hdfs writer 写 `dt=T-1` 单分区,INSERT OVERWRITE 单分区语义 → 同 sched 重跑结果一致
|
|
- **重跑幂等**:DataX hdfs writer 写 `dt=T-1` 单分区,INSERT OVERWRITE 单分区语义 → 同 sched 重跑结果一致
|
|
|
|
|
|
|
|
## 5. ods 层 dt 语义
|
|
## 5. ods 层 dt 语义
|
|
|
|
|
|
|
|
|
|
+- **归位锚 = 活动日 = `GREATEST(create_time, update_time)`**:业务会改 create_time,create / update 皆"活动性"时间戳、无可靠"创建日"概念,取较晚者 = 记录最后一次被动过的日子。Spark 侧写 `GREATEST(create_time, NULLIF(update_time, ''))`(空串当 null 跳过,避免 greatest 遇空返错)
|
|
|
- **来源**:raw 双源 union — `WHERE dt IN ('${dt}', '${pdt}')` 即 raw dt=T-1 + raw dt=T-2 两个分区
|
|
- **来源**:raw 双源 union — `WHERE dt IN ('${dt}', '${pdt}')` 即 raw dt=T-1 + raw dt=T-2 两个分区
|
|
|
-- **过滤**:`DATE_FORMAT(COALESCE(update_time, create_time), 'yyyyMMdd') = '${dt}'`,只保留 COALESCE(update_time, create_time) 落在 T-1 那天的记录
|
|
|
|
|
-- **写入分区**:动态分区 `PARTITION (dt)`,行级 dt = `DATE_FORMAT(COALESCE(update_time, create_time), 'yyyyMMdd')`,配合过滤实际只写到 dt=T-1 一个分区
|
|
|
|
|
-- **跨日漂移修正**:raw dt=T-2 因 48h 宽窗抓到的"漂到 T-1"的部分,被 union 进 ods dt=T-1(详见 ADR-03)
|
|
|
|
|
-- **dedupe**:`ROW_NUMBER() OVER (PARTITION BY id, DATE_FORMAT(COALESCE(update_time, create_time), 'yyyyMMdd') ORDER BY COALESCE(update_time, create_time) DESC) = 1`,分区内取最新版本
|
|
|
|
|
|
|
+- **过滤**:`DATE_FORMAT(GREATEST(create_time, NULLIF(update_time, '')), 'yyyyMMdd') = '${dt}'`,只保留活动日落在 T-1 那天的记录
|
|
|
|
|
+- **写入分区**:动态分区 `PARTITION (dt)`,行级 dt = 活动日;过滤已保证活动日=T-1,实际只写到 dt=T-1 一个分区
|
|
|
|
|
+- **跨日漂移修正**:raw dt=T-2 因 48h 宽窗抓到的"活动日漂到 T-1"的部分,被 union 进 ods dt=T-1(详见 ADR-03)
|
|
|
|
|
+- **dedupe**:`ROW_NUMBER() OVER (PARTITION BY id ORDER BY GREATEST(create_time, NULLIF(update_time, '')) DESC) = 1`,过滤已锁定单一活动日,分区内每单取最新一版
|
|
|
- **跨 ods dt 不去重**:同 pk 多 dt 分区并存 = 上层拉链表(SCD Type 2)的底层
|
|
- **跨 ods dt 不去重**:同 pk 多 dt 分区并存 = 上层拉链表(SCD Type 2)的底层
|
|
|
- **重跑幂等**:动态分区 INSERT OVERWRITE 只覆盖 SELECT 出现的 dt,其他历史 dt 保留(实测见 §8 + tests/integration/spark/idempotence/)
|
|
- **重跑幂等**:动态分区 INSERT OVERWRITE 只覆盖 SELECT 出现的 dt,其他历史 dt 保留(实测见 §8 + tests/integration/spark/idempotence/)
|
|
|
|
|
|
|
@@ -100,7 +101,7 @@ dwd / dws 引入**真业务时间**作 dt 锚点(≠ ods 的 update_time 锚
|
|
|
|
|
|
|
|
DS 补数把 sched 设为补数目标日,所有时间变量按补数日重算。补数与定时的变量取值规则**完全一致**,无特殊处理。
|
|
DS 补数把 sched 设为补数目标日,所有时间变量按补数日重算。补数与定时的变量取值规则**完全一致**,无特殊处理。
|
|
|
|
|
|
|
|
-例:补 5/1 → sched=5/1 → cdt=5/1, dt=4/30, tdt=5/2, pdt=4/29 → raw 抓 `[4/30, 5/2)` 写 raw dt=4/30;ods 取 raw dt=4/30 + raw dt=4/29 → filter DATE(update_time)=4/30 → 写 ods dt=4/30。
|
|
|
|
|
|
|
+例:补 5/1 → sched=5/1 → cdt=5/1, dt=4/30, tdt=5/2, pdt=4/29 → raw 抓 `create ∈ [4/30, 5/2) OR update ∈ [4/30, 5/2)` 写 raw dt=4/30;ods 取 raw dt=4/30 + raw dt=4/29 → filter 活动日 `GREATEST(create,update)`=4/30 → 写 ods dt=4/30。
|
|
|
|
|
|
|
|
## 8. 串行重跑 / 日期递增幂等
|
|
## 8. 串行重跑 / 日期递增幂等
|
|
|
|
|
|