| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273 |
- #!/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,
- '-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,
- '-sc', 'spark.executor.memoryOverhead=4g'], 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()
|