如何避免MongoDB多节点场景下同一邮件数据被重复读取发送
MongoDB多节点邮件任务防重复消费实现方案
核心依赖MongoDB单文档读写的原子性特性,不需要额外引入分布式锁组件,就能实现多节点下任务的互斥抢占,避免重复发送。
1. 邮件集合字段改造
先给原有邮件集合补充任务调度相关字段:
status:整型状态标记,0=待发送,1=发送中,2=发送成功,3=发送失败lockedBy:字符串类型,记录抢占当前任务的节点唯一标识(建议用机器名+进程ID+服务启动时间戳生成,保证全局唯一)lockedAt:UTC时间类型,记录任务被抢占的时间,用于异常节点的任务超时回收retryCount:整型,记录任务被抢占重试的次数,避免失败任务无限循环发送lastError:字符串类型,记录发送失败的错误信息,方便排查问题
注意:需要给status、lockedAt字段创建复合索引,避免任务抢占查询时全表扫描,提升高并发下的性能。
2. 核心原子抢占逻辑
绝对不要使用「先查询待发送任务、再更新任务状态」的两步操作,这两步之间存在并发窗口,会导致多个节点拿到同一条任务。
直接使用MongoDB提供的FindOneAndUpdate原子操作,把「匹配待执行任务」和「标记任务被当前节点锁定」两个动作合并成单个数据库操作,同一时间只会有一个节点能成功更新同一条文档,从根源上避免并发冲突。
.NET驱动下的实现代码如下:
// 服务启动时生成当前节点全局唯一标识,整个服务生命周期内固定不变 private readonly string _currentNodeId = $"{Environment.MachineName}_{Process.GetCurrentProcess().Id}_{DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()}"; private readonly IMongoCollection<EmailTask> _emailTaskCollection; // 任务锁超时时间,根据单封邮件最大发送耗时调整,建议留3-5倍冗余,比如单封最多发1分钟就设5分钟 private readonly TimeSpan _taskLockTimeout = TimeSpan.FromMinutes(5); // 最大重试次数 private const int MaxRetryCount = 3; /// <summary> /// 尝试抢占1条可执行的邮件任务 /// </summary> /// <returns>抢占到的任务,无可用任务时返回null</returns> public async Task<EmailTask> TryAcquireEmailTaskAsync() { var utcNow = DateTime.UtcNow; // 过滤可抢占的任务:要么是待发送状态,要么是发送中但已经超过锁超时时间(原节点可能宕机/异常) var filter = Builders<EmailTask>.Filter.Or( Builders<EmailTask>.Filter.Eq(task => task.Status, 0), Builders<EmailTask>.Filter.And( Builders<EmailTask>.Filter.Eq(task => task.Status, 1), Builders<EmailTask>.Filter.Lt(task => task.LockedAt, utcNow.Subtract(_taskLockTimeout)) ) ); // 原子更新:把匹配到的任务标记为当前节点锁定 var update = Builders<EmailTask>.Update .Set(task => task.Status, 1) .Set(task => task.LockedBy, _currentNodeId) .Set(task => task.LockedAt, utcNow) .Inc(task => task.RetryCount, 1); var options = new FindOneAndUpdateOptions<EmailTask> { ReturnDocument = ReturnDocument.After, // 排序规则:优先抢占重试次数少、等待时间久的任务,避免老任务积压 Sort = Builders<EmailTask>.Sort.Ascending(task => task.RetryCount).Ascending(task => task.LockedAt) }; return await _emailTaskCollection.FindOneAndUpdateAsync(filter, update, options); }
3. 任务执行完成后的状态更新
节点抢占到任务、执行完邮件发送逻辑后,必须更新任务状态,更新时要带上lockedBy作为过滤条件,保证只能更新自己抢到的任务,不会误改其他节点正在处理的任务:
/// <summary> /// 标记任务发送结果 /// </summary> public async Task MarkTaskResultAsync(string taskId, bool sendSuccess, string errorMsg = null) { // 只允许更新当前节点自己锁定的任务 var filter = Builders<EmailTask>.Filter.Eq(task => task.Id, taskId) & Builders<EmailTask>.Filter.Eq(task => task.LockedBy, _currentNodeId); UpdateDefinition<EmailTask> update; if (sendSuccess) { // 发送成功,标记为成功状态,清理锁相关字段 update = Builders<EmailTask>.Update .Set(task => task.Status, 2) .Set(task => task.LastError, null) .Unset(task => task.LockedBy) .Unset(task => task.LockedAt); } else { // 发送失败,判断是否超过最大重试次数 var currentTask = await _emailTaskCollection.Find(Builders<EmailTask>.Filter.Eq(task => task.Id, taskId)).FirstOrDefaultAsync(); if (currentTask != null && currentTask.RetryCount >= MaxRetryCount) { // 超过重试次数,标记为永久失败,后续人工排查 update = Builders<EmailTask>.Update .Set(task => task.Status, 3) .Set(task => task.LastError, errorMsg) .Unset(task => task.LockedBy) .Unset(task => task.LockedAt); } else { // 未到重试上限,重置为待发送状态,等待下次抢占 update = Builders<EmailTask>.Update .Set(task => task.Status, 0) .Set(task => task.LastError, errorMsg) .Unset(task => task.LockedBy) .Unset(task => task.LockedAt); } } await _emailTaskCollection.UpdateOneAsync(filter, update); }
4. 生产环境优化建议
- 批量抢占优化:如果单节点需要批量拉取任务处理,不要循环调用单条抢占方法,可以基于
FindOneAndUpdate配合数组过滤器一次抢占多条任务,减少数据库交互次数,核心原子逻辑和单条抢占一致。 - 幂等兜底:哪怕做了任务抢占,建议发送邮件前先基于业务唯一键(比如邮件唯一ID+收件人地址)做本地幂等校验,极端异常场景下也能100%避免重复发送。
- 监控告警:对状态为发送失败(status=3)、重试次数超过阈值的任务配置告警,及时处理发送异常的邮件。
- 超时时间校准:上线前压测单封邮件的最大发送耗时,锁超时时间必须大于这个值,避免正常发送中的任务被其他节点误抢占。
内容的提问来源于stack exchange,提问作者ClemB
相关产品推荐
相关产品推荐

