定时任务的可靠性不等于「准时触发」,而是「执行可靠、结果正确」。漏跑与重跑集中在四个成因:进程重启错过触发点、单次执行超过调度间隔造成堆叠、单线程调度被长任务阻塞、多实例同时抢跑。对应的四件事必须做齐——用唯一键保证幂等、用执行记录支撑崩溃后补偿、用锁与实例数上限控制并发、用心跳监控发现静默失效。热点监控(AI 热点监控工具)每分钟轮询多个平台,任何一环缺失都会表现为「面板上突然没有新热点」。下面给出各环节的实现。
一、漏跑与重跑的四个成因
| 现象 | 根因 | 对策 |
|---|---|---|
| 重启后中间几轮没数据 | 触发时刻进程不在,调度器无补偿 | 启动时对比上次成功时间,超期立即补跑 |
| 任务互相堆叠、数据错乱 | 执行时长超过调度间隔 | 限制同任务并发实例数,超时丢弃本轮 |
| 某个数据源卡住,其他源也停 | 调度器默认单线程被长任务阻塞 | 配置线程池,给每个任务设超时 |
| 同一条热点推送两次 | 多实例并发触发同一任务 | 分布式锁 + 唯一键幂等校验 |
值得先纠正一个常见误判:只加锁不做幂等是不够的。锁提前过期、人工重跑、消息重投都会绕过锁,幂等才是兜底那一层。
二、幂等:唯一键加状态机
幂等的含义是同一任务重复执行多次,结果与执行一次一致。
2.1 唯一索引拦住重复写入
热点监控里的幂等键取「平台 + 平台内条目 ID + 轮次日期」,靠数据库唯一索引拦重复写入:
CREATE TABLE IF NOT EXISTS hot_items (
id INTEGER PRIMARY KEY AUTOINCREMENT,
platform TEXT NOT NULL,
item_id TEXT NOT NULL,
round_date TEXT NOT NULL,
title TEXT,
analyzed INTEGER DEFAULT 0, -- 0 待分析 1 已分析
notified INTEGER DEFAULT 0, -- 0 未推送 1 已推送
created_at TEXT DEFAULT CURRENT_TIMESTAMP,
UNIQUE (platform, item_id, round_date)
);
写入侧用忽略冲突的语义,重复轮次不会产生脏数据:
def save_items(conn, platform: str, items: list, round_date: str) -> int:
rows = [(platform, it["id"], round_date, it["title"]) for it in items]
cur = conn.executemany(
"INSERT OR IGNORE INTO hot_items(platform, item_id, round_date, title) "
"VALUES (?, ?, ?, ?)",
rows,
)
conn.commit()
return cur.rowcount # 实际新增条数,等于本轮真正的新热点
2.2 状态位把流程切成可续跑的阶段
analyzed 与 notified 两个状态位把「采集 → 大模型分析 → 推送」拆成状态机。每步只处理符合前置状态的记录,任务重跑时已完成的记录自动跳过,天然幂等。大模型调用与消息推送都有成本,这一层直接决定账单和用户体验。
三、并发控制:实例数上限与阻塞策略
任务执行时间超过间隔时必须明确策略。三种取向对应不同业务:
| 策略 | 行为 | 适用任务 |
|---|---|---|
| 串行等待 | 上一轮未结束,本次排队 | 必须完整执行的统计汇总 |
| 丢弃本轮 | 上一轮未结束,跳过本次触发 | 状态巡检、热点轮询 |
| 允许并发 | 多实例同时运行 | 完全无状态、无数据竞争的任务 |
热点轮询选「丢弃本轮」——错过一轮等下一轮即可,堆叠反而会加重目标站压力。用 APScheduler 时对应三个参数:
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.executors.pool import ThreadPoolExecutor
scheduler = BackgroundScheduler(
executors={"default": ThreadPoolExecutor(8)}, # 避免长任务阻塞其他任务
job_defaults={
"coalesce": True, # 堆积的多次触发合并为一次
"max_instances": 1, # 同一任务同时只跑一个实例
"misfire_grace_time": 30, # 迟到 30 秒内仍执行,超过则跳过
},
)
scheduler.add_job(poll_platform, "interval", minutes=1, args=["bilibili"],
id="poll_bilibili", replace_existing=True)
coalesce 处理的是进程卡顿后攒下的多次触发:合并成一次,而不是补跑十几轮。多实例部署时再叠一层分布式锁,锁的过期时间要大于任务预估耗时上限,否则锁提前释放,另一个节点会并发执行同一批次。带自动续期的锁实现更稳妥。
四、崩溃恢复与错过补偿
调度器只管进程活着的时间。进程崩溃、容器重启、机器维护期间的触发点全部丢失,靠调度器自身找不回来。
4.1 启动时对比上次成功时间
补偿的做法是把「上次成功时间」持久化,启动时判断是否超期:
import time
def bootstrap_catch_up(conn, job_name: str, interval_sec: int, runner) -> None:
row = conn.execute(
"SELECT MAX(finished_at) FROM job_runs "
"WHERE job_name = ? AND status = 'success'",
(job_name,),
).fetchone()
last_success = row[0]
if last_success is None:
runner() # 首次部署,立即执行一次
return
gap = time.time() - float(last_success)
if gap > interval_sec: # 错过了至少一轮
runner() # 只补一次,不逐轮回放
补偿只补一次是刻意的:热点数据关心的是当前榜单,回放二十轮历史触发既无意义又会打满接口配额。对于必须逐轮补齐的任务(如按日统计),改为按缺失的时间片分批补跑,并让每个时间片带上幂等键。
4.2 执行记录表就是恢复依据
每轮执行都落一条记录,字段够用即可:
CREATE TABLE IF NOT EXISTS job_runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
job_name TEXT NOT NULL,
started_at REAL NOT NULL,
finished_at REAL,
duration_ms INTEGER,
status TEXT NOT NULL, -- running / success / failed
error TEXT
);
处于 running 状态却早已超过预估耗时的记录,说明执行期间进程异常退出。启动巡检把这类记录重置为失败,交给重试或补偿逻辑接手,避免任务永久卡在中间态。
五、重试与死信兜底
网络抖动、第三方接口超时、数据库瞬时不可用都会让单轮失败。重试要设边界:
- 只对可恢复异常重试——超时、连接失败、限流;参数错误与解析异常直接失败,重试无意义;
- 用指数退避加随机抖动,避免多任务同时重试形成二次冲击;
- 重试上限设 3 次,避免持续压垮下游;
- 仍失败的批次写入失败记录表,字段保留任务名、批次号、失败原因;
- 由一个低频补偿任务扫描失败记录,重跑或转人工,同时触发一条聚合通知。
死信兜底的价值在于「失败不丢」。失败信息落库后,可以离线分析失败集中在哪个数据源、哪个时段,这比逐条告警有用得多。
六、静默失效要靠监控发现
可靠性设计里容易被忽略的一环是:任务根本没触发时,系统不会报错,日志里也不会有异常。发现这种静默失效需要外部视角。
一是心跳判断,用「距上次成功执行的时长」做告警条件,超过两倍调度间隔就报警。二是引入外部看护——任务成功后主动上报一次心跳,外部服务在预期时间内没收到就通知。三是记录每轮真正的新增条数,长时间为零且无异常,通常意味着采集规则失效或登录态过期,而非平台没有热点。
上线顺序建议是:先做幂等,再配限流与并发上限,接着补执行记录与补偿,收尾接监控告警。顺序颠倒的代价通常是一次重复推送事故。
常见问题(FAQ)
Q1:锁过期了任务还没跑完怎么办?
按耗时上限设过期时间,并使用支持自动续期的锁,同时保留幂等校验兜底。
Q2:进程重启后要补跑所有错过的轮次吗?
热点类任务只补一次即可;统计类任务按缺失时间片分批补,且必须带幂等键。
Q3:任务没触发但系统不报错,如何发现?
用距上次成功时间做告警条件,配合成功后上报心跳的外部看护服务。