Pārlūkot izejas kodu

fix(ods): 埋点 ods 加 DISTRIBUTE BY dt 防动态分区写 OOM + backfill ods 调用补 memoryOverhead

ods 全量回算动态写 276 分区时每 task 持有过多 ORC writer,
off-heap 超默认 512MB overhead 被 YARN 杀,task 累计 4 次失败 abort。
DISTRIBUTE BY dt 按分区聚合大幅降低单 task 并发 writer 数;
日常单分区无副作用。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
tianyu.chu 1 mēnesi atpakaļ
vecāks
revīzija
8b3d7ee63b

+ 2 - 1
jobs/ods/usr/ods_usr_traces_apd_d.sql

@@ -47,4 +47,5 @@ SELECT
     get_json_object(raw_json, '$.properties.params')                                   AS params_json,
     dt                                                                                 AS dt
 FROM raw.raw_usr_traces_apd_d
-WHERE dt BETWEEN '${start_date}' AND '${stop_date}';
+WHERE dt BETWEEN '${start_date}' AND '${stop_date}'
+DISTRIBUTE BY dt;

+ 2 - 1
jobs/raw/usr/raw_usr_traces_apd_d_backfill.py

@@ -60,7 +60,8 @@ def main():
     rc_raw = run(['python3', STARTER, '-f', RAW_SQL, '-u', UDF,
                   '-sc', 'spark.executor.memoryOverhead=4g'], cwd=PROJECT_ROOT)
     rc_ods = run(['python3', STARTER, '-f', ODS_SQL,
-                  '-p', 'start_date=%s' % start, '-p', 'stop_date=%s' % stop], cwd=PROJECT_ROOT)
+                  '-p', 'start_date=%s' % start, '-p', 'stop_date=%s' % stop,
+                  '-sc', 'spark.executor.memoryOverhead=4g'], cwd=PROJECT_ROOT)
     run(['hdfs', 'dfs', '-rm', '-r', '-f', '-skipTrash', STAGING])
 
     if rc_raw or rc_ods: