Cron 脚本迁到 DAG:先封装任务,再做双跑比对

📅 2026/8/12 18:26:17
Cron 脚本迁到 DAG:先封装任务,再做双跑比对
Cron 脚本迁到 DAG先封装任务再做双跑比对在运维与数据工程体系中长期依赖操作系统原生的crontab运行独立 Python 脚本容易积累技术债务。此类脚本通常缺乏显式依赖控制与全局状态感知当上游同步任务异常时下游分析脚本仍可能盲目执行导致产生无效或错误数据。直接下线 Cron 把代码重写一遍切到 Airflow 或 Prefect风险极其高昂。我们要采取分阶段平滑迁移的策略通过任务封装解耦、新旧系统双跑Dual-Run校验以及基于数据指纹的自动断路器完成无缝过渡。flowchart TD subgraph LegacyCron[旧架构: 隐式依赖的 Cron 单体乱象] C1[Cron 02:00 sync_db.py] --|无状态无校验| DB[(MySQL)] C2[Cron 03:00 calc_report.py] --|假定上游已完成| DB end subgraph MigrationStrategy[渐进式平滑迁移路径] Step1[1. Task 装饰器封装 幂等状态落库] Step2[2. 构造 Dual-Run 双跑比较器] Step3[3. 数据指纹 Checksum 自动比对] Step4[4. 优雅切流与旧 Cron 完全下线] Step1 -- Step2 -- Step3 -- Step4 end subgraph NewDAG[新架构: 显式 DAG 拓扑 熔断防线] D1[Prefect/Airflow DAG Task A] --|执行成功 产出 Checksum| D2[DAG Task B: 自动触发] D1 --|报错或数据异常| Breaker[断路器熔断: 阻止下游执行并告警] end1. 拆解痛点 Cron 架构的核心缺陷Cron 定时任务最核心的问题在于缺乏显式的依赖关系与状态感知。没有状态机Cron 只管在特定的时间触发python run.py至于上个周期的任务是成功了还是还在卡死僵死Cron 完全不关心缺乏数据质量校验Data Quality Check脚本执行成功Exit Code 0只代表 Python 语法没报错不代表洗出来的数据符合质量逻辑如清洗后数据条数为 0重试机制原始一旦脚本崩掉要么靠人肉半夜爬起来手动跑要么在 Shell 脚本里硬写死循环 retry极其容易造成数据库死锁。重构的目标不是推翻原有业务逻辑而是用结构化的 DAG 调度流Directed Acyclic Graph将逻辑包裹起来并赋予其重试、可观测性与依赖控制能力。2. 第一阶段Python 逻辑提取与 Task 状态增强迁移的第一步是在不修改旧 Cron 触发的前提下用包装函数Wrapper Pattern改造原有 Python 脚本。把原来的单体过程式代码重构为拥有输入参数、输出返回值以及状态记录的函数。我们使用轻量级的装饰器将原有的业务代码无缝升级为具备上下文感知能力的 Task。import sys import time import hashlib import functools import structlog from typing import Callable, Any, Dict logger structlog.get_logger() # 模拟全局调度状态持久化表 TASK_EXECUTION_STATE: Dict[str, Dict[str, Any]] {} def pipeline_task(task_name: str, max_retries: int 3, retry_delay: float 1.0): 用于平滑改造老旧 Cron 脚本的防护装饰器 def decorator(func: Callable): functools.wraps(func) def wrapper(*args, **kwargs): retry_count 0 start_time time.time() while retry_count max_retries: try: logger.info(task_execution_start, tasktask_name, attemptretry_count 1) # 执行原老旧业务逻辑 result func(*args, **kwargs) # 计算产出数据的指纹校验和 (Data Checksum) data_str str(result) checksum hashlib.md5(data_str.encode(utf-8)).hexdigest() cost_ms (time.time() - start_time) * 1000 # 记录成功状态 TASK_EXECUTION_STATE[task_name] { status: SUCCESS, checksum: checksum, cost_ms: cost_ms, timestamp: time.time() } logger.info(task_execution_success, tasktask_name, checksumchecksum, cost_mscost_ms) return result except Exception as e: retry_count 1 logger.warning(task_execution_failed, tasktask_name, attemptretry_count, errorstr(e)) if retry_count max_retries: TASK_EXECUTION_STATE[task_name] { status: FAILED, error: str(e), timestamp: time.time() } # 抛出异常阻断可能存在的错误下游 raise RuntimeError(fTask [{task_name}] 在 {max_retries} 次重试后依然失败: {str(e)}) time.sleep(retry_delay) return wrapper return decorator # --- 改造前的老脚本逻辑 --- pipeline_task(task_namesync_user_orders, max_retries2) def legacy_sync_orders_task(date_str: str): # 模拟原来的 ETL 同步逻辑 if date_str 2026-08-12: return [{order_id: 1001, amount: 99.5}, {order_id: 1002, amount: 200.0}] else: raise ValueError(数据库网络抖动超时)通过这一层轻量包装原有的老脚本瞬间拥有了自动重试、异常告警与产出物 MD5 指纹计算能力。3. 第二阶段Dual-Run 双跑比对与渐进式切流在第二阶段我们在新的 DAG 调度引擎如 Airflow 或 Prefect中创建镜像管道让新调度框架与旧 Cron 并行跑同一份数据。但是新管道的洗表结果不直接覆盖生产数据库而是写入临时 Shadow 表。两边跑完后由比较器DualRunComparator自动拉取两边结果的checksum指纹与数据条数Row Count。class DualRunComparator: 新旧管道双跑数据对比组件 staticmethod def compare_outputs(cron_output: Any, dag_output: Any) - bool: cron_str str(cron_output) dag_str str(dag_output) cron_md5 hashlib.md5(cron_str.encode(utf-8)).hexdigest() dag_md5 hashlib.md5(dag_str.encode(utf-8)).hexdigest() is_matched (cron_md5 dag_md5) if not is_matched: logger.error(dual_run_mismatch_detected, cron_md5cron_md5, dag_md5dag_md5) else: logger.info(dual_run_matched_perfectly, checksumcron_md5) return is_matched # --- 运行双跑校验测试 --- cron_res legacy_sync_orders_task(2026-08-12) dag_res legacy_sync_orders_task(2026-08-12) # 新 DAG Task 调用的逻辑 matched DualRunComparator.compare_outputs(cron_res, dag_res) assert matched, 新旧管道输出数据指纹必须完全一致连续跑完 7 天的数据双跑比对如果每日的指纹匹配率达到 100%我们就可以自信地注释掉服务器crontab里的旧配置将新 DAG 调度升级为主路。将隐式依赖重构为可被监控与熔断的 DAG 拓扑图通过平滑迁移策略既能保障存量业务稳定又能提升系统的整体可观测性。5. 双跑比对要先定义允许差异迁移期间的新旧任务不必要求字节级完全相同但必须先写明哪些字段允许因时间、排序或外部数据变化而不同哪些字段一旦不同就要阻断切换。比对结果要按任务类型聚合不能只看总成功率某个低频但高价值任务持续偏差仍值得单独处置。发现问题时保留输入版本、依赖版本和中间产物摘要才能复现差异并确定是代码迁移、数据质量还是调度时序造成的。确认切换后也不要立刻删掉旧任务。保留一个受控的只读回放窗口并定期抽样复核关键产物等依赖、数据和调度节奏都稳定后再回收旧链路。这样出现问题时还有可比较的基线。