定时轮询会漏跑原因解析(热点监控的幂等与补偿设计)

定时任务的可靠性不等于「准时触发」,而是「执行可靠、结果正确」。漏跑与重跑集中在四个成因:进程重启错过触发点、单次执行超过调度间隔造成堆叠、单线程调度被长任务阻塞、多实例同时抢跑。对应的四件事必须做齐——用唯一键保证幂等、用执行记录支撑崩溃后补偿、用锁与实例数上限控制并发、用心跳监控发现静默失效。热点监控(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 状态却早已超过预估耗时的记录,说明执行期间进程异常退出。启动巡检把这类记录重置为失败,交给重试或补偿逻辑接手,避免任务永久卡在中间态。

五、重试与死信兜底

网络抖动、第三方接口超时、数据库瞬时不可用都会让单轮失败。重试要设边界:

  1. 只对可恢复异常重试——超时、连接失败、限流;参数错误与解析异常直接失败,重试无意义;
  2. 用指数退避加随机抖动,避免多任务同时重试形成二次冲击;
  3. 重试上限设 3 次,避免持续压垮下游;
  4. 仍失败的批次写入失败记录表,字段保留任务名、批次号、失败原因;
  5. 由一个低频补偿任务扫描失败记录,重跑或转人工,同时触发一条聚合通知。

死信兜底的价值在于「失败不丢」。失败信息落库后,可以离线分析失败集中在哪个数据源、哪个时段,这比逐条告警有用得多。

六、静默失效要靠监控发现

可靠性设计里容易被忽略的一环是:任务根本没触发时,系统不会报错,日志里也不会有异常。发现这种静默失效需要外部视角。

一是心跳判断,用「距上次成功执行的时长」做告警条件,超过两倍调度间隔就报警。二是引入外部看护——任务成功后主动上报一次心跳,外部服务在预期时间内没收到就通知。三是记录每轮真正的新增条数,长时间为零且无异常,通常意味着采集规则失效或登录态过期,而非平台没有热点。

上线顺序建议是:先做幂等,再配限流与并发上限,接着补执行记录与补偿,收尾接监控告警。顺序颠倒的代价通常是一次重复推送事故。

常见问题(FAQ)

Q1:锁过期了任务还没跑完怎么办?

按耗时上限设过期时间,并使用支持自动续期的锁,同时保留幂等校验兜底。

Q2:进程重启后要补跑所有错过的轮次吗?

热点类任务只补一次即可;统计类任务按缺失时间片分批补,且必须带幂等键。

Q3:任务没触发但系统不报错,如何发现?

用距上次成功时间做告警条件,配合成功后上报心跳的外部看护服务。

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 qiqicto@qq.com 举报,一经查实,本站将立刻删除。
赞 (0)
套餐观察官的头像套餐观察官普通用户

相关推荐

返回顶部