Prechádzať zdrojové kódy

feat(raw): 埋点历史回算 backfill — 全量一波灌(gz 入 dt= 子目录分区发现)+ ods 改日期范围动态分区

tianyu.chu 1 mesiac pred
rodič
commit
5cd73d58ec

+ 9 - 6
jobs/ods/usr/ods_usr_traces_apd_d.sql

@@ -1,12 +1,14 @@
 -- 作者:tianyu.chu
 -- 日期:2026-06-10
 -- 工单:(无)
--- 目的:埋点 raw → ods,解析脱敏后 _source JSON 拍平公共属性 + 保留 params_json;按文件日 dt 静态分区写入(N=1 不归位)
+-- 目的:埋点 raw → ods,解析脱敏后 _source JSON 拍平公共属性 + 保留 params_json;按文件日 dt 动态分区写入
 -- 状态:[待执行]
--- 备注:dt 静态写 ${dt} = 文件/上传日(N=1 不归位);实测文件内 ~99.4% 事件日=文件日,~0.6% 迟到/未来小偏差按当天落,业务允许(窗口决策见 workspace/20260610/埋点迟到漂移分布-窗口决策.md);
---       事件不可变,无双源 union / 无去重;es_id 单文件内唯一;时区随集群(东八区)
+-- 备注:日期范围 [${start_date}, ${stop_date}] 含两端、动态分区按 raw.dt 写——日常增量传 start=stop=单日,
+--       历史回算传全区间,一份 SQL 两用(避免解析逻辑复制漂移,对齐 ADR-11);
+--       dt = 文件日(N=1 不归位):实测 ~99.4% 事件日=文件日,~0.6% 迟到/未来小偏差按当天落,业务允许(见 kb/13 §8);
+--       动态覆盖本环境只覆盖 SELECT 出现的分区(kb/21 §8);事件不可变,无双源 union / 无去重;时区随集群(东八区)
 
-INSERT OVERWRITE TABLE ods.ods_usr_traces_apd_d PARTITION (dt = '${dt}')
+INSERT OVERWRITE TABLE ods.ods_usr_traces_apd_d PARTITION (dt)
 SELECT
     es_id                                                                              AS es_id,
     event_name                                                                         AS event_name,
@@ -42,6 +44,7 @@ SELECT
     CAST(get_json_object(raw_json, '$.properties.isFirstTime')          AS BOOLEAN)     AS is_first_time,
     CAST(get_json_object(raw_json, '$.properties.resumeFromBackground') AS BOOLEAN)     AS resume_from_background,
     CAST(get_json_object(raw_json, '$.properties.eventDuration')        AS BIGINT)      AS event_duration,
-    get_json_object(raw_json, '$.properties.params')                                   AS params_json
+    get_json_object(raw_json, '$.properties.params')                                   AS params_json,
+    dt                                                                                 AS dt
 FROM raw.raw_usr_traces_apd_d
-WHERE dt = '${dt}';
+WHERE dt BETWEEN '${start_date}' AND '${stop_date}';

+ 71 - 0
jobs/raw/usr/raw_usr_traces_apd_d_backfill.py

@@ -0,0 +1,71 @@
+#!/usr/bin/env /usr/bin/python3
+# -*- coding:utf-8 -*-
+"""
+埋点历史一次性回算(backfill)包装脚本 —— 全量一波灌,比日常逐日循环快得多。
+
+流程:
+  1. 按 -dt 区间,把每个存在的本地 gz put 到 HDFS 暂存 /tmp/raw_usr_traces_backfill/dt={dt}/
+     (Hive 风格分区目录,Spark 读时自动识别 dt 为分区列;缺文件跳过)
+  2. 一个 Spark job 跑 raw backfill SQL(读全部暂存、脱敏、动态分区灌 raw)
+  3. 一个 Spark job 跑 ods 解析 SQL(range = 实际回算到的 [min,max],动态分区灌 ods)
+  4. 清暂存
+
+CLI:-dt 支持单日 / `20250801-` / 区间 / 离散(复用 get_date_range)。
+日常增量用 raw_usr_traces_apd_d.py(单日);本脚本只供一次性历史回算。
+"""
+import os
+import subprocess
+import sys
+
+PROJECT_ROOT = os.path.abspath(os.path.join(os.path.dirname(os.path.abspath(__file__)), '..', '..', '..'))
+sys.path.insert(0, PROJECT_ROOT)
+
+from dw_base.common.config_constants import K_DT
+from dw_base.utils.config_utils import parse_args
+from dw_base.utils.datetime_utils import get_date_range, get_yesterday
+
+LOCAL_DIR = '/data/upload/traces'
+STAGING = '/tmp/raw_usr_traces_backfill'
+RAW_SQL = 'jobs/raw/usr/raw_usr_traces_apd_d_backfill.sql'
+ODS_SQL = 'jobs/ods/usr/ods_usr_traces_apd_d.sql'
+UDF = 'dw_base/udf/business/spark_traces_udf.py'
+STARTER = 'bin/spark-sql-starter.py'
+
+
+def run(cmd, cwd=None):
+    print('+ ' + ' '.join(cmd))
+    return subprocess.call(cmd, cwd=cwd)
+
+
+def main():
+    config, _ = parse_args(sys.argv[1:])
+    date_range = get_date_range(config.get(K_DT, get_yesterday()))
+
+    run(['hdfs', 'dfs', '-rm', '-r', '-f', '-skipTrash', STAGING])
+    put = []
+    for dt in date_range:
+        gz = '%s/traces-%s-%s-%s.json.gz' % (LOCAL_DIR, dt[0:4], dt[4:6], dt[6:8])
+        if not os.path.isfile(gz):
+            print('跳过 %s:无本地 gz' % dt)
+            continue
+        run(['hdfs', 'dfs', '-mkdir', '-p', '%s/dt=%s' % (STAGING, dt)])
+        if run(['hdfs', 'dfs', '-put', '-f', gz, '%s/dt=%s/' % (STAGING, dt)]) == 0:
+            put.append(dt)
+    if not put:
+        print('无可回算文件,退出')
+        sys.exit(0)
+    start, stop = min(put), max(put)
+    print('回算 %d 天,范围 [%s, %s]' % (len(put), start, stop))
+
+    rc_raw = run(['python3', STARTER, '-f', RAW_SQL, '-u', UDF], cwd=PROJECT_ROOT)
+    rc_ods = run(['python3', STARTER, '-f', ODS_SQL,
+                  '-p', 'start_date=%s' % start, '-p', 'stop_date=%s' % stop], cwd=PROJECT_ROOT)
+    run(['hdfs', 'dfs', '-rm', '-r', '-f', '-skipTrash', STAGING])
+
+    if rc_raw or rc_ods:
+        print('!! raw rc=%d / ods rc=%d' % (rc_raw, rc_ods))
+    sys.exit(rc_raw or rc_ods)
+
+
+if __name__ == '__main__':
+    main()

+ 23 - 0
jobs/raw/usr/raw_usr_traces_apd_d_backfill.sql

@@ -0,0 +1,23 @@
+-- 作者:tianyu.chu
+-- 日期:2026-06-17
+-- 工单:(无)
+-- 目的:埋点历史一次性回算 —— 全部 gz 一个 Spark job 解析脱敏、动态分区灌 raw(dt 走目录分区发现)
+-- 状态:[待执行]
+-- 备注:配批量 put —— 每个 gz 放到 /tmp/raw_usr_traces_backfill/dt={yyyymmdd}/,Spark 读时自动把 dt 当分区列
+--       (比 input_file_name() 稳:后者在 text view + Python UDF 查询里返回空,会把行全落到 null 分区);
+--       日常增量仍用 raw_usr_traces_apd_d.py/.sql(单日静态分区)。本文件只供一次性回算。
+--       动态覆盖本环境只覆盖 SELECT 出现的分区(kb/21 §8)——全史回算写全部分区,一波灌满。
+
+ADD FILE conf/tracking-mask.ini;
+
+CREATE OR REPLACE TEMPORARY VIEW traces_gz_backfill
+USING text
+OPTIONS (path '/tmp/raw_usr_traces_backfill/');
+
+INSERT OVERWRITE TABLE raw.raw_usr_traces_apd_d PARTITION (dt)
+SELECT
+    get_json_object(value, '$._id')           AS es_id,
+    get_json_object(value, '$._source.event') AS event_name,
+    mask_source(value)                        AS raw_json,
+    CAST(dt AS STRING)                        AS dt
+FROM traces_gz_backfill;