如何禁止NServiceBus Saga完成后以相同关联ID重复执行?
当然可以实现,但NServiceBus默认确实不会保留已完成Saga的Correlation ID记录(完成后会删除Saga数据行),所以需要我们自己添加额外的跟踪逻辑来实现这个需求。下面是两种常见的实现思路:
方法一:单独跟踪已完成的Correlation ID
这种思路是创建一个独立的存储(比如数据库表、缓存等),专门记录已经完成的Saga的Correlation ID和对应的Saga类型,在Saga启动前先检查这个存储:
创建跟踪存储
可以新建一个类似CompletedSagaCorrelations的数据库表,字段至少包含:CorrelationId:要跟踪的关联IDSagaType:Saga的类型名称(区分不同Saga)CompletedAt:完成时间(用于后续清理过期记录)
在Saga启动前检查
在处理启动Saga的消息时,先查询这个跟踪存储,如果已存在相同的Correlation ID和Saga类型记录,就跳过Saga的启动逻辑:public async Task Handle(StartMySagaMessage message, IMessageHandlerContext context) { // 注入自定义的跟踪器服务 var tracker = context.Extensions.Get<ICompletedSagaTracker>(); var hasCompleted = await tracker.HasCompletedAsync(nameof(MySaga), message.CorrelationId); if (hasCompleted) { _logger.LogInformation("Saga with Correlation ID {CorrelationId} has already completed - skipping execution.", message.CorrelationId); return; } // 正常启动Saga的逻辑 Data.CorrelationId = message.CorrelationId; // ...其他业务处理 }Saga完成时记录跟踪信息
在处理Saga完成的消息时,将当前的Correlation ID和Saga类型写入跟踪存储:public async Task Handle(CompleteMySagaMessage message, IMessageHandlerContext context) { var tracker = context.Extensions.Get<ICompletedSagaTracker>(); await tracker.TrackCompletionAsync(nameof(MySaga), Data.CorrelationId); // 完成Saga,默认会删除Saga数据行 MarkAsComplete(); }定期清理跟踪数据
为了避免跟踪存储无限膨胀,可以添加定时任务,清理掉超过一定时间(比如30天)的已完成记录。
方法二:保留已完成的Saga实例并标记状态
这种思路是修改Saga的持久化行为,不删除已完成的Saga数据,而是标记为“已完成”状态,后续相同Correlation ID的消息进来时,会找到这个已完成的实例,再判断是否跳过处理:
修改Saga数据类
添加一个IsCompleted字段来标记Saga状态:public class MySagaData : ContainSagaData { public string CorrelationId { get; set; } public bool IsCompleted { get; set; } }配置Saga关联规则
确保相同Correlation ID的消息能找到已完成的Saga实例:protected override void ConfigureHowToFindSaga(SagaPropertyMapper<MySagaData> mapper) { mapper.MapSaga(saga => saga.CorrelationId) .ToMessage(message => message.CorrelationId); }处理启动消息时检查状态
在处理启动消息时,如果找到的Saga实例已经标记为完成,就跳过逻辑:public async Task Handle(StartMySagaMessage message, IMessageHandlerContext context) { if (Data.IsCompleted) { _logger.LogInformation("Saga with Correlation ID {CorrelationId} is already completed - skipping.", message.CorrelationId); return; } // 如果是新实例,初始化数据 if (Data.Id == Guid.Empty) { Data.CorrelationId = message.CorrelationId; } // 正常的业务处理逻辑 }完成Saga时标记状态
注意不要调用默认的MarkAsComplete()(会删除数据),而是设置IsCompleted状态:public async Task Handle(CompleteMySagaMessage message, IMessageHandlerContext context) { Data.IsCompleted = true; // 不要调用MarkAsComplete(),保留Saga数据 }
这种方法的缺点是Saga表会保留所有已完成的实例,需要定期清理这些数据来避免数据库性能问题。
注意事项
- 无论哪种方法,都要确保跟踪逻辑的可靠性(比如使用事务保证记录的原子性)。
- 根据业务需求选择合适的清理策略,避免存储无限增长。
- 可以结合日志记录,方便排查重复触发的问题。
内容的提问来源于stack exchange,提问作者l1pton17

