Rebus管道步骤中用作用域服务传递Correlation Id遇阻
解决Rebus中使用RebusTransactionScopeSuppressor时Correlation Id转发失效的问题
针对你遇到的问题——使用RebusTransactionScopeSuppressor立即发送消息时,原出站步骤无法加载入站消息导致Correlation Id转发失效,这里提供两种可行的解决方案:
方案一:利用Rebus的MessageContext传递Correlation Id
Rebus的MessageContext.Current基于AsyncLocal实现,同一个异步处理流内均可访问,适合在入站和出站步骤间传递Correlation Id,无需额外作用域服务:
- 添加入站步骤捕获Correlation Id
public class CaptureCorrelationIdIncomingStep : IIncomingStep { public async Task Process(IncomingStepContext context, Func<Task> next) { var inboundMessage = context.Load<Message>(); if (inboundMessage.Headers.TryGetValue(Rebus.Messages.Headers.CorrelationId, out var correlationId)) { MessageContext.Current.Items["CurrentCorrelationId"] = correlationId; } await next(); } }
- 添加出站步骤附着Correlation Id
public class AttachCorrelationIdOutgoingStep : IOutgoingStep { public async Task Process(OutgoingStepContext context, Func<Task> next) { if (MessageContext.Current.Items.TryGetValue("CurrentCorrelationId", out var correlationIdObj) && correlationIdObj is string correlationId && !string.IsNullOrWhiteSpace(correlationId)) { var outboundMessage = context.Load<Message>(); outboundMessage.Headers[Rebus.Messages.Headers.CorrelationId] = correlationId; } await next(); } }
- 注册步骤到Rebus配置
Configure.With(serviceProviderAdapter) .Transport(t => t.Use...()) // 替换为你的Transport配置 .Options(options => { options.AddIncomingStep<CaptureCorrelationIdIncomingStep>(); options.AddOutgoingStep<AttachCorrelationIdOutgoingStep>(); }) .Start();
方案二:使用作用域服务存储Correlation Id
如果你更倾向于用自定义作用域服务,可通过Rebus的Advanced.ServiceProvider在出站步骤中直接获取作用域服务,无需手动加载IServiceScope:
- 定义作用域存储服务
public interface ICorrelationIdStore { string? CorrelationId { get; set; } } public class CorrelationIdStore : ICorrelationIdStore { public string? CorrelationId { get; set; } }
- 注册作用域服务
services.AddScoped<ICorrelationIdStore>();
- 入站步骤捕获Correlation Id
public class CaptureCorrelationIdStep : IIncomingStep { private readonly ICorrelationIdStore _correlationIdStore; public CaptureCorrelationIdStep(ICorrelationIdStore correlationIdStore) { _correlationIdStore = correlationIdStore; } public async Task Process(IncomingStepContext context, Func<Task> next) { var inboundMessage = context.Load<Message>(); if (inboundMessage.Headers.TryGetValue(Rebus.Messages.Headers.CorrelationId, out var correlationId)) { _correlationIdStore.CorrelationId = correlationId; } await next(); } }
- 出站步骤获取并附着Correlation Id
public class AttachCorrelationIdStep : IOutgoingStep { public async Task Process(OutgoingStepContext context, Func<Task> next) { var correlationIdStore = context.Advanced.ServiceProvider.GetService<ICorrelationIdStore>(); if (correlationIdStore != null && !string.IsNullOrWhiteSpace(correlationIdStore.CorrelationId)) { var outboundMessage = context.Load<Message>(); outboundMessage.Headers[Rebus.Messages.Headers.CorrelationId] = correlationIdStore.CorrelationId; } await next(); } }
- 注册步骤到Rebus配置
和方案一的注册方式一致,将两个步骤添加到Rebus的入站和出站管道中。
关键说明
- 方案一无需额外依赖Rebus的服务注入机制,利用AsyncLocal特性确保异步流内的上下文传递,完美适配
RebusTransactionScopeSuppressor场景。 - 方案二通过
context.Advanced.ServiceProvider直接获取作用域服务,避免手动加载IServiceScope的问题,因为Rebus在入站处理的作用域内发送消息时,出站步骤的ServiceProvider会自动绑定到当前作用域。
内容的提问来源于stack exchange,提问作者Nikola Vacic
相关产品推荐
相关产品推荐

