如何配置MassTransit 8 Saga的无限延迟重传直至实例终结?
我们基于MassTransit 8和Saga状态机搭建业务流程,实现了多个Saga之间、Saga与外部服务的互通。目前已配置未处理异常的通用消息重试策略:重试5次,每次间隔30秒。
针对长期错误场景,我们需要实现延迟重传策略:要求每5分钟重试一次,无限次数,直到Saga终结。业务规则允许Saga的TTL为12小时,超时后通过Schedule功能终结并移除Saga。期望在Saga存活的12小时内持续重传,Saga终结移除后自动停止重传(同时丢弃DLQ中的对应消息),构建可靠的故障恢复流程。
补充说明1
实际需要重传的是Saga发起的请求,简化代码示例如下:
public class SurveySaga : MassTransitStateMachine<SurveyState> { public Event<RequestSurveyCommand> Requested { get; private set; } public Request<SurveyState, GetConfig, GetConfigResponse> GetConfigRequest { get; private set; } public SurveySaga() { Initially( When(Requested) .Request(GetConfigRequest, context => new GetSomeConfigRequest { Foo = "bar" }) ); } }
当GetConfigRequest失败时,本可以通过When(GetConfigRequest.Faulted)处理,但添加无限重传后该方法未触发。我们希望在Saga终结/移除时停止对该请求的重传。
补充说明2
我们了解到针对请求(非事件)无法直接配置缺失实例时的丢弃逻辑。当前请求配置如下:
Request(() => GetConfigRequest, x => x.GetConfigRequestId, cfg => { cfg.ServiceAddress = _endpointer.GetAddress(EndpointName.GetConfig); cfg.Timeout = TimeSpan.FromSeconds(_options.RequestTimeoutInSeconds); });
该请求通过GetConfigRequestId与Saga实例关联,但暂无配置实例缺失时的处理选项。请问如何实现请求与Saga的关联,并在实例缺失时自动丢弃请求?
解决方案
一、针对请求的延迟重传与Saga终结时停止重传
MassTransit中针对Saga请求的重传,不能直接依赖全局UseDelayedRedelivery配置,需要结合Saga状态机的事件逻辑手动实现:
绑定重试调度到Saga生命周期
定义一个调度任务,在请求故障时触发延迟重试,并在Saga终结时取消所有待执行的调度:public Schedule<SurveyState, RetryGetConfig> RetryGetConfigSchedule { get; private set; } public SurveySaga() { InstanceState(x => x.CurrentState); // 定义5分钟间隔的重试调度,绑定到Saga实例 Schedule(() => RetryGetConfigSchedule, x => x.RetryGetConfigTokenId, s => { s.Delay = TimeSpan.FromMinutes(5); s.ReceiveEndpoint(e => e.ConfigureSaga<SurveyState>(_context)); }); Initially( When(Requested) .Request(GetConfigRequest, context => new GetSomeConfigRequest { Foo = "bar" }) ); // 仅当Saga未终结时,触发重试调度 When(GetConfigRequest.Faulted) .If(context => !context.Saga.IsCompleted, b => b.Schedule(RetryGetConfigSchedule, context => new RetryGetConfig { SagaId = context.Saga.CorrelationId }) ); // 处理调度事件,重新发起请求 When(RetryGetConfigSchedule.Received) .Request(GetConfigRequest, context => new GetSomeConfigRequest { Foo = "bar" }); // Saga终结时取消所有未执行的重试 When(SagaCompleted) .Unschedule(RetryGetConfigSchedule); }禁用请求的内置重试
避免全局重试与手动调度冲突,关闭请求的内置重试逻辑:Request(() => GetConfigRequest, x => x.GetConfigRequestId, cfg => { cfg.ServiceAddress = _endpointer.GetAddress(EndpointName.GetConfig); cfg.Timeout = TimeSpan.FromSeconds(_options.RequestTimeoutInSeconds); cfg.Retry(r => r.None()); // 禁用内置重试 });
二、处理Saga实例缺失时的请求丢弃
针对请求关联的Saga实例不存在的场景,通过自定义中间件实现消息丢弃:
添加缺失实例过滤中间件
在Saga接收端点配置中加入过滤器,检测到实例不存在时直接丢弃消息:// 在Saga端点配置中注册过滤器 endpoint.ConfigureSaga<SurveyState>(context, cfg => { cfg.UseFilter(new MissingInstanceFilter<SurveyState>()); }); // 自定义缺失实例过滤器 public class MissingInstanceFilter<T> : IFilter<SagaConsumeContext<T>> where T : class, ISaga { public async Task Send(SagaConsumeContext<T> context, IPipe<SagaConsumeContext<T>> next) { if (context.Saga == null) { await context.DiscardAsync(); // 直接丢弃,不进入DLQ return; } await next.Send(context); } public void Probe(ProbeContext context) { } }在重试逻辑中提前校验实例存在性
在发起重试前,先查询Saga存储确认实例是否存活:When(GetConfigRequest.Faulted) .Async(context => context.Repository.GetInstance(context.CorrelationId) .ContinueWith(t => { if (t.Result == null) return Task.CompletedTask; return context.Schedule(RetryGetConfigSchedule, new RetryGetConfig { SagaId = context.CorrelationId }); }) );
三、关键注意事项
- Saga终结逻辑:确保12小时的TTL调度正确触发Saga终结事件,所有关联的重试调度会被自动取消
- 幂等性保障:重传的请求必须实现幂等,避免外部服务重复处理
- DLQ管理:通过缺失实例过滤器直接丢弃无效消息,减少DLQ的无效积压
内容的提问来源于stack exchange,提问作者alkesos

