定时热点检查跑方法详解(监控工具的轮询调度与增量去重)

定时热点检查的核心是一套”调度触发 → 平台拉取 → 增量去重 → AI 分析 → 实时推送”的闭环。调度层用 APScheduler 的 cron 或 interval 触发器周期性驱动,去重层用内容指纹挡掉重复,分析层产出结构化标签,最后经 WebSocket 推给前端。下面给出完整流程与各层落地细节。

一、整体流程

一次检查任务的执行顺序固定如下:

  1. 调度器到点触发 check_job,载入当前监控词表与平台配置;
  2. 对各平台(Twitter、B站等)并发拉取近窗口内的新内容;
  3. 用内容指纹做增量去重,过滤已处理过的重复项;
  4. 把增量内容送大模型做情感、话题、热点判定;
  5. 命中阈值的结果经 WebSocket 实时推送给订阅前端;
  6. 写入热点表与去重指纹库,更新本轮检查时间戳。

任意一步失败都不应阻断整轮,错误要落日志并进入重试,避免单次异常让监控空窗。

二、调度器选型与配置

APScheduler 由触发器、执行器、作业存储、调度器四部分构成。监控常驻服务选 BackgroundScheduler,不阻塞主进程;纯脚本才用 BlockingScheduler。

调度器 是否阻塞主线程 适用场景
BlockingScheduler 是 独立脚本、无其它主逻辑
BackgroundScheduler 否 Web 应用、常驻服务
AsyncIOScheduler 否 asyncio 异步环境

触发器有三种:date(一次性)、interval(固定间隔)、cron(类 crontab 周期)。热点检查多用后两种:

from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.cron import CronTrigger

sched = BackgroundScheduler(timezone="Asia/Shanghai")

# 每 10 分钟轮询一次(interval)
sched.add_job(
    check_job, "interval", minutes=10,
    id="hot_check_interval",
    max_instances=1,            # 禁止上一轮未完又起新一轮
    coalesce=True,              # 错过多次只补一次
    misfire_grace_time=300,     # 允许延迟上限 300 秒执行
)

# 或每天整点跑(cron)
sched.add_job(
    check_job, CronTrigger(hour="*", minute=0),
    id="hot_check_cron", coalesce=True,
)

sched.start()

max_instances=1 防止慢查询堆积,coalesce=True 在调度器重启后只补一次而非补齐所有漏跑,misfire_grace_time 容忍短暂延迟。这三项是生产定时任务的标准保险。

三、增量去重的实现

去重的目的是同一内容只分析、只推送一次。两种常用指纹:

  • 平台 ID 去重:Twitter 的 id、B站的 bvid 天然不重复,直接存 Redis Set,O(1) 判断;
  • 内容哈希去重:对正文做归一化(去空白、转小写)后取 SHA-256,挡住改写搬运。
import hashlib
import redis

r = redis.Redis(host="localhost", port=6379, db=0)

def is_new(post_id: str, text: str) -> bool:
    # 双指纹:ID 与内容哈希任一命中即视为重复
    fid = f"post:id:{post_id}"
    fhash = f"post:hash:{hashlib.sha256(text.strip().encode()).hexdigest()}"
    if r.exists(fid) or r.exists(fhash):
        return False
    r.set(fid, 1, ex=60 * 60 * 24 * 7)   # 7 天过期
    r.set(fhash, 1, ex=60 * 60 * 24 * 7)
    return True

ID 去重快但挡不住搬运,内容哈希补上这一层。指纹设 7 天过期,既控内存,又允许周期性话题二次进入。

四、异常兜底与告警

定时任务怕”静默失败”——进程没崩,但某轮悄悄没跑。三道防线:

  1. 给 check_job 包 try/except,异常写日志并抛回让调度器记录;
  2. 接 EVENT_JOB_ERROR 监听器,任务失败即发告警;
  3. 监控 next_run_time,长时间未更新则触发人工介入。
from apscheduler.events import EVENT_JOB_ERROR

def on_error(event):
    # 这里可接邮件/Webhook/IM 告警
    print(f"任务 {event.job_id} 执行失败: {event.exception}")

sched.add_listener(on_error, EVENT_JOB_ERROR)

轮询间隔要与平台限流匹配。Twitter 的 recent search 在 App-Only 下约 450 次/15 分钟,把检查频率与每轮请求数乘积控制在此预算内,再配合指数退避,才能长期稳定运行。

常见问题(FAQ)

Q1:interval 和 cron 怎么选?

固定间隔用 interval;需”每天 8 点”这类日历规则用 cron。

Q2:任务执行时间超过间隔怎么办?

设 max_instances=1 并 coalesce=True,避免堆叠与重复触发。

Q3:重启后任务会丢吗?

默认内存存储会丢,用 SQLAlchemy/Redis 作 JobStore 即可持久化恢复。

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

相关推荐

返回顶部