本文目录导读:

批量任务分发避免重复的核心思路是 “确保每个任务只被一个消费者处理一次”,这通常被称为 “幂等性”(Idempotency)或 “去重”(Deduplication)。
根据不同的架构和场景,有几种主流且有效的解决方案:
使用消息队列的“精准投递”特性
这是最直接的方法,现代消息队列(如 RabbitMQ、Kafka、Pulsar)本身就提供了避免消息被重复消费的机制。
- 关键机制: Consumer 的 Offset / Ack 确认机制 + 手动提交。
- 消费者在处理完任务后才向消息队列发送“确认(Ack)”。
- 如果消费者在处理过程中崩溃,消息队列会认为该任务未处理,并将其重新分发给另一个消费者(或重启后的同一消费者)。
- 避免重复的关键: 消费者必须在业务逻辑完全执行成功后再提交 Ack,如果业务逻辑里包含多个步骤(如写数据库、调用API),确保所有步骤完成后,最后才提交 Offset/Ack。
任务目标系统的幂等性(最根本方法)
这是解决重复消费的最佳实践,无论分发机制如何,重复最终都能被系统自身消化。
- 原理: 任务的处理逻辑(如写数据库、调用API)本身支持重复执行多次,结果与执行一次相同。
- 实现方式:
- 数据库唯一键: 任务处理结果写入数据库时,使用任务的唯一ID(如
job_id)作为主键或唯一索引,如果重复插入,会直接报错或忽略。 - 业务逻辑幂等: 将用户余额设置为100元”是幂等的(设置10次也是100元);“给用户余额增加10元”则不是幂等的(增加10次变100元),后者需要改为“如果用户余额未到100元,则将其设置为100元”或使用版本号乐观锁。
- 数据库唯一键: 任务处理结果写入数据库时,使用任务的唯一ID(如
引入“任务状态机”(分布式锁配合)
在任务调度器(如 Celery、Quartz、XXL-JOB)层面,维护任务的状态。
- 实现流程:
- 任务创建时: 在中央存储(如 Redis、数据库)中记录任务ID和状态(如
pending)。 - 调度器获取任务: 使用分布式锁(如 Redis SETNX)或数据库乐观锁(
UPDATE tasks SET status='processing' WHERE id=? AND status='pending')来原子地抢任务。 - 原子性状态转移: 只有成功将状态从
pending改为processing的调度器,才获得该任务。 - 处理完成后: 将状态更新为
done或failed。 - 定期清理/重试: 长时间处于
processing状态的任务(可能因为节点崩溃)会被重新标记为pending并被再次分发。
- 任务创建时: 在中央存储(如 Redis、数据库)中记录任务ID和状态(如
利用去重表/缓存(短时间窗口)
如果任务ID天然唯一(如UUID),可以使用一个简单的去重缓存。
- 实现方式:
- 消费前查询: 在任务处理器开始工作前,先查询一个高可用且支持过期时间的存储(如 Redis),检查该任务ID是否出现过。
- SetNX 原子操作: 使用 Redis 的
SET key NX EX 600(不存在则设置,设置成功返回OK,设置失败表示已存在),如果设置成功,说明这是第一次见到这个任务,允许处理;如果返回false,说明任务已处理过,直接跳过。 - 注意: 这个方法的有效期取决于Redis Key的过期时间,如果任务处理时间很长,需要动态延长过期时间,否则容易误判。
业务层自定义过滤(兜底策略)
在极端情况(如网络重试导致重复分发)下,任务处理方需要有能力自行过滤。
- 实现方式: 在处理任务入口处,维护一个本地或全局的已处理任务ID集合。
- 小型系统: 可以使用内存中的 Set,配合定时清理(例如每10分钟清理一次,只保留最近处理的ID)。
- 大型系统: 使用 Bloom Filter(布隆过滤器)进行快速判断,但存在一定的误判可能(判断为“可能已处理”其实是“未处理”),需要结合数据库确认。
总结与最佳实践
| 方法 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| 消息队列 Ack 控制 | 原生支持,实现简单 | 需要正确处理好Ack时机,仍有极小概率因网络中断导致重复 | 大多数消息驱动架构 |
| 目标系统幂等 | 最彻底,能防御所有来源的重复 | 业务改造相对复杂,需设计业务逻辑 | 所有关键业务(如支付、下单) |
| 任务状态机 + 锁 | 保证强一致性,无重复 | 依赖分布式锁,可能导致系统瓶颈 | 调度中间件(如 Celery)、定时任务 |
| 去重表/缓存 | 性能高,适合短时间窗口 | 有数据丢失风险(缓存丢了),无法保证长期不重复 | 短时效任务(如验证码发送) |
推荐的最佳组合实践:
- 源头控制: 使用消息队列(如 Kafka/RabbitMQ)的手动 Ack,确保业务处理完再确认。
- 底层防御: 目标系统(数据库、API)实现幂等性(如唯一键、版本号)。
- 中间加固: 如果是高可用系统,任务调度器配合分布式锁或状态机。
- 兜底方案: 处理器自身使用 Redis + SetNX 或 布隆过滤器 进行快速去重。
这样,即使某个环节出现网络波动、节点重启导致的重复分发,系统也能安全地避免重复处理。