← 返回博客列表

【魔码量化工程实战进阶 #06】定时调度与任务编排:让数据自己每天0650到位

2026年08月31日 18:02 · 魔码数服 · 魔码量化工程实战进阶

【魔码量化工程实战进阶 #06】定时调度与任务编排:让数据自己每天 06:50 到位

入门系列第 18 篇讲了"定时拉数"(单只 time.sleep 轮询)。本篇把它升级成生产级调度:每天盘后自动跑、任务之间有依赖、失败会告警、而且绝对不踩"盘中禁跑"红线。这是把前 5 篇(管道 / 存储 / 并发 / 重试 / 增量)串成一个自动运转系统的最后一环。

本文你将得到什么

  1. 为什么 while True: sleep(60) 不是调度(进程一挂全停)
  2. crontab(生产部署)与 APScheduler(代码内调度)两种选法
  3. 一个任务编排骨架:拉数 → 落库 → 校验 → 告警,依赖串行
  4. 四个调度红线坑(时区 / 任务重叠 / 盘中禁跑 / 无日志)

一、痛点:人工每天点"运行"

前 5 篇把单次拉数做成了可靠流水线,但还得你每天手动跑。人的问题:忘跑、周末漏、出差断更、盘中手滑点了一下(红线!)。调度要解决"到点自动跑、不用人在"。

新手爱写 while True: time.sleep(86400),问题:进程被 OOM 杀掉就永远停了;没有"任务依赖",校验和拉数无法编排;没有告警,失败了你一周后才发现。


二、工程方案:两种调度器

2.1 crontab(生产部署首选)

把拉数脚本注册成系统定时任务,进程挂了 cron 下次照跑:

# 每天 18:40(盘后)跑全市场日 K 增量更新;周日下午回补复权最近 5 日
40 18 * * 1-5  /usr/bin/python /opt/moma/pull_daily.py >> /var/log/moma/pull.log 2>&1
30 18 * * 0     /usr/bin/python /opt/moma/pull_backfill.py >> /var/log/moma/backfill.log 2>&1

2.2 APScheduler(代码内调度,适合单机小服务)

from apscheduler.schedulers.blocking import BlockingScheduler

def job_daily():
    # 复用前 5 篇的 Puller + 增量 + 重试
    Puller().run(all_codes)        # #01 流水线
    backfill_recent(5)             # #05 增量回补
    if not validate():             # 校验
        alert("日 K 更新校验失败")  # 告警

sched = BlockingScheduler()
sched.add_job(job_daily, "cron", day_of_week="mon-fri", hour=18, minute=40, timezone="Asia/Shanghai")
sched.start()

三、任务编排:把"四步"串成一条链

单跑拉数不够,拉完必须校验、校验失败必须告警,否则脏数据静默进库更可怕。

def pipeline():
    steps = [("拉数",   pull_daily),
             ("落库",   upsert_to_sqlite),
             ("校验",   validate_freshness),   # 检查今日数据行数/最新日期
             ("告警",   alert_on_fail)]
    for name, fn in steps:
        try:
            ok = fn()
            log(f"[OK] {name}")
            if not ok:
                break                       # 校验不过,停止并告警
        except Exception as e:
            log(f"[FAIL] {name}: {e}")
            alert(f"{name} 失败: {e}")
            break

validate_freshness 示例:检查 max(t) == 今天,且今日行数 ≈ 昨日行数(突增/突降都报警)。


四、本机实测(编排骨架可跑;调度需运行时环境)

上面 pipeline() 是纯 Python,可直接 python -c "pipeline()" 验证四步串联逻辑(把 pull_daily 换成你 #01–#05 的实现)。调度本身(cron/APScheduler)在服务器上 7×24 运行,本机只需确认"逻辑对、配置对"。

一个最小健康检查脚本(放进 cron,每小时报一次库的新鲜度):

def health_check():
    latest = conn.execute("SELECT MAX(t) FROM kline").fetchone()[0]
    age = (datetime.date.today() - parse(latest)).days
    if age > 1:
        alert(f"数据已 {age} 天未更新,疑似调度中断")

五、原理深挖:四个调度红线坑

  1. 时区错:cron 默认 UTC。你写 40 18 想盘后跑,结果 UTC 18 点是北京时间凌晨 2 点。务必显式 timezone="Asia/Shanghai" 或用 TZ=Asia/Shanghai
  2. 任务重叠:上一次跑了 50 分钟还没完,下一次 18:40 又起,两只抢同一张表、重复拉。解决:BlockingSchedulercoalesce + max_instances=1,或 cron 前置加锁文件。
  3. 盘中禁跑(红线):调度时间必须完全避开 9:30–11:30、13:00–15:00。盘后选 18:40 安全;若做分钟级,更要在非交易时段。运维红线,违反即事故。
  4. 无日志无告警:调度跑在后台,没日志你不知道它死没死。每条 log + 失败 alert(邮件/企微机器人)是标配,否则"自动"等于"裸奔"。

六、小结

调度不是 sleep(86400),而是"系统级定时 + 任务编排 + 健康检查 + 告警"。把前 5 篇的流水线接进 crontab/APScheduler,你的数据每天 18:40 自己到位、自己校验、自己报警——你只需要偶尔看一眼告警群。到这一篇,模块一(数据工程基石)完整闭环:稳拉、会存、能并发、不怕抖、只增量、自动跑

当调度从"日 K"升级到"分钟级全市场",单任务时长、存储放大、盘中红线都会变苛刻——这正是魔码量化 Pro 包(完整历史分钟 K 线 + 全市场实时)与更高配额要支撑的场景。详见文末。


免责声明:本文所有示例数据仅用于接口演示,不构成任何投资建议;市场有风险,投资需谨慎。

系列持续更新中。 想要亲手跑通上面的代码?前往 魔码证书申请页 免费领取你的专属证书,复制即用、按次计费、稳定可用。

想亲自试一下?免费获取证书
客服微信
客服微信二维码