14-埋点同步-开发.md 4.6 KB

埋点数据同步实现(开发文档)

设计/取舍见 13-埋点同步-设计.md;全量事件+字段+脱敏目录见 16-埋点raw建模.md

1. raw 层:薄表

单表 raw.raw_usr_traces_apd_d(不可变事件流 = 追加)。只做脱敏、不拍平——落脱敏后的 _source 整行 JSON,拍平归 ods。

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.sqles_id / event_name 不敏感、拎出来便于 ods 路由与去重;其余全在 raw_json

2. 脱敏

2.1 配置 conf/tracking-mask.ini

按事件分段,作用于 properties 顶层与 params 两层;dropmask 同字段时 drop 优先。

[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.sqlADD 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 <sql> -u <udf> -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)。