
本文介绍一种高效、可扩展的定时告警方案,用于在记录临近截止或已超期时精准触发用户通知,兼顾 kafka 事件驱动架构与静态数据场景,避免多实例轮询冲突。
本文介绍一种高效、可扩展的定时告警方案,用于在记录临近截止或已超期时精准触发用户通知,兼顾 kafka 事件驱动架构与静态数据场景,避免多实例轮询冲突。
在基于 Kafka 的事件驱动系统中,对动态行为(如创建、更新)触发通知较为自然;但面对静态表(如表 A 中大量含 created_time 和 deadline_time 的待处理记录),缺乏事件源时,传统“高频轮询 + 多进程并发扫描”极易引发重复告警、数据库压力激增及分布式竞态问题。
推荐方案:轻量级调度 + 内存优先队列 + 状态去重
核心思路是将时间维度从被动查询转为主动驱动,避免全表扫描,同时保证单点可靠、多实例安全:
-
使用分布式可靠调度器(如 Quartz Cluster、XXL-JOB 或 Kafka-based Scheduler)
配置固定间隔(如每 2 分钟)触发一次全局协调任务。该任务不直接查库发消息,而是执行以下原子操作:- 查询
WHERE deadline_time BETWEEN NOW() AND NOW() + INTERVAL '30 MINUTE'(即将到期)UNION ALLWHERE deadline_time (已超期且未处理) - 将结果集按
deadline_time升序排序,生成唯一批次 ID(如alert_batch_20241125_1430) - 向 Kafka 发送一条 控制消息(topic:
alert-schedule),携带批次 ID 与最小/最大 deadline 时间范围
- 查询
-
消费者端采用“Leader-Election + 有序队列”模式
所有告警服务实例订阅alert-schedule,通过 Kafka Consumer Group 自动实现负载均衡;利用 Kafka 的分区语义(如按批次 ID hash 到固定 partition)确保同一批次仅由一个实例消费。该实例执行:// 伪代码:内存优先队列驱动精准触发 PriorityQueue<AlertRecord> queue = new PriorityQueue<>((a, b) -> a.getDeadlineTime().compareTo(b.getDeadlineTime()) ); List<AlertRecord> candidates = queryFromDB(batchId); // 基于批次ID查缓存或DB queue.addAll(candidates); while (!queue.isEmpty()) { AlertRecord r = queue.poll(); if (r.getDeadlineTime().isBefore(Instant.now().plusSeconds(60))) { sendKafkaNotification(r); // 发送业务通知到 user-topic updateStatusToNotified(r.getId()); // 幂等更新状态 } else { // 重入队列尾部,等待下一轮调度(或放入 DelayQueue) break; } } -
关键保障机制
- ✅ 幂等性:每条记录更新
notified_at字段 + 唯一索引(record_id, alert_type),防止重复推送 - ✅ 容错性:消费者处理失败时,Kafka 自动重投(配合
enable.auto.commit=false+ 手动 commit offset) - ✅ 可伸缩性:调度器与消费者解耦;新增实例仅分担消费负载,不增加 DB 压力
- ⚠️ 注意:若
deadline_time可能动态变更,需在更新时同步发送deadline_update事件,刷新内存队列或触发重新调度
- ✅ 幂等性:每条记录更新
该方案相比 Cron 全表扫描更精准,比纯内存队列更可靠(持久化调度上下文),且天然适配 Kafka 生态。在千级 TPS 场景下实测,平均告警延迟

















