raw_usr_traces_apd_d_backfill.py 2.8 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273
  1. #!/usr/bin/env /usr/bin/python3
  2. # -*- coding:utf-8 -*-
  3. """
  4. 埋点历史一次性回算(backfill)包装脚本 —— 全量一波灌,比日常逐日循环快得多。
  5. 流程:
  6. 1. 按 -dt 区间,把每个存在的本地 gz put 到 HDFS 暂存 /tmp/raw_usr_traces_backfill/dt={dt}/
  7. (Hive 风格分区目录,Spark 读时自动识别 dt 为分区列;缺文件跳过)
  8. 2. 一个 Spark job 跑 raw backfill SQL(读全部暂存、脱敏、动态分区灌 raw)
  9. 3. 一个 Spark job 跑 ods 解析 SQL(range = 实际回算到的 [min,max],动态分区灌 ods)
  10. 4. 清暂存
  11. CLI:-dt 支持单日 / `20250801-` / 区间 / 离散(复用 get_date_range)。
  12. 日常增量用 raw_usr_traces_apd_d.py(单日);本脚本只供一次性历史回算。
  13. """
  14. import os
  15. import subprocess
  16. import sys
  17. PROJECT_ROOT = os.path.abspath(os.path.join(os.path.dirname(os.path.abspath(__file__)), '..', '..', '..'))
  18. sys.path.insert(0, PROJECT_ROOT)
  19. from dw_base.common.config_constants import K_DT
  20. from dw_base.utils.config_utils import parse_args
  21. from dw_base.utils.datetime_utils import get_date_range, get_yesterday
  22. LOCAL_DIR = '/data/upload/traces'
  23. STAGING = '/tmp/raw_usr_traces_backfill'
  24. RAW_SQL = 'jobs/raw/usr/raw_usr_traces_apd_d_backfill.sql'
  25. ODS_SQL = 'jobs/ods/usr/ods_usr_traces_apd_d.sql'
  26. UDF = 'dw_base/udf/business/spark_traces_udf.py'
  27. STARTER = 'bin/spark-sql-starter.py'
  28. def run(cmd, cwd=None):
  29. print('+ ' + ' '.join(cmd))
  30. return subprocess.call(cmd, cwd=cwd)
  31. def main():
  32. config, _ = parse_args(sys.argv[1:])
  33. date_range = get_date_range(config.get(K_DT, get_yesterday()))
  34. run(['hdfs', 'dfs', '-rm', '-r', '-f', '-skipTrash', STAGING])
  35. put = []
  36. for dt in date_range:
  37. gz = '%s/traces-%s-%s-%s.json.gz' % (LOCAL_DIR, dt[0:4], dt[4:6], dt[6:8])
  38. if not os.path.isfile(gz):
  39. print('跳过 %s:无本地 gz' % dt)
  40. continue
  41. run(['hdfs', 'dfs', '-mkdir', '-p', '%s/dt=%s' % (STAGING, dt)])
  42. if run(['hdfs', 'dfs', '-put', '-f', gz, '%s/dt=%s/' % (STAGING, dt)]) == 0:
  43. put.append(dt)
  44. if not put:
  45. print('无可回算文件,退出')
  46. sys.exit(0)
  47. start, stop = min(put), max(put)
  48. print('回算 %d 天,范围 [%s, %s]' % (len(put), start, stop))
  49. rc_raw = run(['python3', STARTER, '-f', RAW_SQL, '-u', UDF,
  50. '-sc', 'spark.executor.memoryOverhead=4g'], cwd=PROJECT_ROOT)
  51. rc_ods = run(['python3', STARTER, '-f', ODS_SQL,
  52. '-p', 'start_date=%s' % start, '-p', 'stop_date=%s' % stop,
  53. '-sc', 'spark.executor.memoryOverhead=4g'], cwd=PROJECT_ROOT)
  54. run(['hdfs', 'dfs', '-rm', '-r', '-f', '-skipTrash', STAGING])
  55. if rc_raw or rc_ods:
  56. print('!! raw rc=%d / ods rc=%d' % (rc_raw, rc_ods))
  57. sys.exit(rc_raw or rc_ods)
  58. if __name__ == '__main__':
  59. main()