定时热点检查的核心是一套”调度触发 → 平台拉取 → 增量去重 → AI 分析 → 实时推送”的闭环。调度层用 APScheduler 的 cron 或 interval 触发器周期性驱动,去重层用内容指纹挡掉重复,分析层产出结构化标签,最后经 WebSocket 推给前端。下面给出完整流程与各层落地细节。
一、整体流程
一次检查任务的执行顺序固定如下:
- 调度器到点触发
check_job,载入当前监控词表与平台配置; - 对各平台(Twitter、B站等)并发拉取近窗口内的新内容;
- 用内容指纹做增量去重,过滤已处理过的重复项;
- 把增量内容送大模型做情感、话题、热点判定;
- 命中阈值的结果经 WebSocket 实时推送给订阅前端;
- 写入热点表与去重指纹库,更新本轮检查时间戳。
任意一步失败都不应阻断整轮,错误要落日志并进入重试,避免单次异常让监控空窗。
二、调度器选型与配置
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 天过期,既控内存,又允许周期性话题二次进入。
四、异常兜底与告警
定时任务怕”静默失败”——进程没崩,但某轮悄悄没跑。三道防线:
- 给
check_job包 try/except,异常写日志并抛回让调度器记录; - 接
EVENT_JOB_ERROR监听器,任务失败即发告警; - 监控
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 即可持久化恢复。