# 埋点数据同步实现(开发文档) > 设计/取舍见 `13-埋点同步-设计.md`;全量事件+字段+脱敏目录见 `16-埋点raw建模.md`。 ## 1. raw 层:薄表 单表 `raw.raw_usr_traces_apd_d`(不可变事件流 = 追加)。**只做脱敏、不拍平**——落脱敏后的 `_source` 整行 JSON,拍平归 ods。 ```sql CREATE EXTERNAL TABLE IF NOT EXISTS raw.raw_usr_traces_apd_d ( es_id STRING COMMENT 'ES 文档 _id,去重键', event_name STRING COMMENT '事件名(_source.event),便于 ods 路由', raw_json STRING COMMENT '脱敏后的 _source 整行 JSON' ) COMMENT '埋点 raw 层(脱敏后整行 JSON)' PARTITIONED BY (dt STRING COMMENT 'yyyymmdd,取自 gz 文件名') STORED AS ORC LOCATION '/user/hive/warehouse/raw.db/raw_usr_traces_apd_d'; ``` DDL:`manual/ddl/raw/usr/raw_usr_traces_apd_d_create.sql`。`es_id` / `event_name` 不敏感、拎出来便于 ods 路由与去重;其余全在 `raw_json`。 ## 2. 脱敏 ### 2.1 配置 `conf/tracking-mask.ini` 按事件分段,作用于 `properties` 顶层与 `params` 两层;`drop` 与 `mask` 同字段时 drop 优先。 ```ini [event:PayOrder] drop = receiverName, receiverTelephone, receiverArea, receiverAddress # mask = field:method method ∈ md5/month_trunc/mask_middle/keep_first_n/keep_last_n ``` 现状:5 个含敏事件(PayOrder / GroupPayOrder / SaveReceiverAddress / GroupOrderEditAddress / MemberLimitedPaymentClick)统一 drop 收货四要素。 ### 2.2 脱敏核心 `dw_base/tracking/mask.py`(纯 Python,可测) - `load_mask_conf(path)` → `{event: {'drop':[...], 'mask':{field:method}}}`;**文件缺失 fail-fast**(脱敏配置缺失会导致敏感数据原样入仓,不可静默) - `apply_mask(event_name, properties, conf)` → 脱敏后新 dict(不改入参);未声明事件原样返回(兜底) - `apply_method(method, value)` → 6 方法(语义对齐 `dw_base/datax/mask.py`,执行另写——埋点无源库) - 单测 `tests/unit/tracking/test_mask.py`:drop / 各 method / 兜底 / drop+mask 优先级 / fail-fast ### 2.3 Spark UDF `dw_base/udf/business/spark_traces_udf.py` `mask_source(line) -> STRING`:解析 ES hit、按事件脱敏 `properties`、回吐脱敏后 `_source` JSON。 - **返回 String 而非 struct**:框架 UDF 自动注册只认普通函数、默认 StringType(`@udf` struct 不被注册)。`es_id`/`event_name` 不敏感,SQL 侧 `get_json_object` 原生取。 - 配置经 SQL `ADD FILE conf/tracking-mask.ini` 分发,UDF 首次调用经 `SparkFiles` 懒加载。 - **前置依赖**:Python UDF 在 executor 上 import dw_base 需 `dw_base/__init__` executor-safe(`conf/env.sh` 缺失时跳过 driver 初始化,见 `33b468d`)——否则 executor 反序列化即崩。 ## 3. 入仓(历史 / 增量同一套) - SQL `jobs/raw/usr/raw_usr_traces_apd_d.sql`:`ADD FILE` 配置 → `CREATE TEMP VIEW USING text`(读 HDFS 临时 gz,自动解压)→ `mask_source` + 原生取 es_id/event_name → `INSERT OVERWRITE PARTITION(dt='${dt}')`。 - 包装脚本 `jobs/raw/usr/raw_usr_traces_apd_d.py`:逐日 `hdfs put` 本地 gz 到 `/tmp/raw_usr_traces/{dt}/` → 调 `bin/spark-sql-starter.py -f -u -dt {dt}` → 清临时目录。`-dt` 支持单日/区间/离散(复用 `get_date_range`);**缺文件跳过**(缺数据容忍)。 - 历史 = 补跑多 dt(同脚本传区间);增量 = 每日昨日。 ## 4. ods 层:解析拍平 `ods.ods_usr_traces_apd_d`:解析 `raw_json` JSON,公共属性 typed 拍平成列(35 列)+ `params_json` 半结构化兜底。DDL `manual/ddl/ods/usr/...`,解析 SQL `jobs/ods/usr/ods_usr_traces_apd_d.sql`。 - 埋点 ods 特例:非业务库类型恢复、是 JSON 解析;params 不 per-event 拍平(event explosion);web 端字段不拍平、回查走 raw。 - **dt = 文件日、N=1 不归位**:`WHERE dt='${dt}'` 读、`PARTITION(dt='${dt}')` 静态写;ES 按事件日分索引,文件日≈事件日 99.4%,~0.6% 迟到/未来小偏差按当天落(业务允许,见 `13` §8)。 - 字段清单见 `16-埋点raw建模.md` §1(公共属性)。 ## 5. 冒烟实证(2026-06-17,dt=20250906) - raw 入仓:**2,042,539 行**(与全量探查一致);PayOrder 9 条,收货四要素 **0 残留**(脱敏生效) - ods 解析:**2,042,539 行**(零丢失);event_time 东八区正常、user_id/lib/params_json 解析正确 - 入仓链端到端跑通:gz → text view → `mask_source`(executor)→ raw → ods ## 6. 调度 DolphinScheduler 每日 T+1 跑入仓脚本(昨日)。失败重跑(`INSERT OVERWRITE` 幂等)。dwd 待埋点重构后建(ADR-13)。