如何在Rebus中为每个消息处理器实现共享上下文?
Rebus装饰器间共享上下文的解决方案
问题背景
Rebus为整个消息处理管道使用单个服务作用域,我希望通过装饰器给消息处理器添加横切关注点(比如消费者级幂等性),但目前无法让同一处理器执行链中的多个装饰器共享上下文状态。每个装饰器能独立解析依赖,但同一消息处理流程里,装饰器之间没有明显的共享状态方式。
尝试过但无效的代码示例
services.Decorate(typeof(IHandleMessages<>), typeof(IdempotentHandlerDecorator<>)); services.Decorate(typeof(IHandleMessages<>), typeof(MessageLogDecorator<>)); services.Decorate(typeof(IHandleMessages<>), typeof(ConsumerContextBuilderDecorator<>)); internal interface IConsumerContextAccesor { ConsumerContext ConsumerContext { get; } } public class ConsumerContext { public string MessageName { get; set; } = string.Empty; public string ConsumerName { get; set; } = string.Empty; public long EllapsedMilliseconds { get; set; } public MessageConsumptionResult ConsumptionResult { get; set; } = MessageConsumptionResult.Unknown; public Exception? Exception { get; set; } public DateTimeOffset ProcessedAt { get; set; } } public class ConsumerContextBuilderDecorator<TMessage>( IHandleMessages<TMessage> innerHandler) : IHandleMessages<TMessage>, IConsumerContextAccesor { private ConsumerContext _consumerContext = null!; public ConsumerContext ConsumerContext => _consumerContext; public async Task Handle(TMessage message) { if (message == null) return; _consumerContext = new(); _consumerContext.MessageName = message.GetType().Name; _consumerContext.ConsumerName = innerHandler.GetType().Name; var sw = Stopwatch.StartNew(); try { await innerHandler.Handle(message); _consumerContext.ConsumptionResult = MessageConsumptionResult.Success; } catch (Exception ex) { _consumerContext.ConsumptionResult = MessageConsumptionResult.Failed; _consumerContext.Exception = ex; throw; } finally { sw.Stop(); _consumerContext.ProcessedAt = DateTime.UtcNow; _consumerContext.EllapsedMilliseconds = sw.ElapsedMilliseconds; } } } public class IdempotentHandlerDecorator<TMessage>( IHandleMessages<TMessage> innerHandler, IMessageContext messageContext, IMessageTracker messageTracker) : IHandleMessages<TMessage>, IConsumerContextAccesor { public ConsumerContext ConsumerContext => (innerHandler as IConsumerContextAccesor)?.ConsumerContext!; public async Task Handle(TMessage message) { if (message == null) return; var messageId = messageContext.Headers[Headers.MessageId]; if (await messageTracker.IsProcessed(messageId, this.ConsumerContext.ConsumerName)) { this.ConsumerContext.ConsumptionResult = Core.Entities.MessageConsumptionResult.Skipped; return; } await innerHandler.Handle(message); this.ConsumerContext.ConsumptionResult = Core.Entities.MessageConsumptionResult.Success; await messageTracker.MarkAsProcessed( messageId: messageId, consumerName: this.ConsumerContext.ConsumerName, messageName: this.ConsumerContext.MessageName); } } public class MessageLogDecorator<TMessage>( IHandleMessages<TMessage> innerHandler, IMessageContext messageContext, IDbContextFactory<MessageConsumptionDbContext> dbContextFactory, ILogger<MessageLogDecorator<TMessage>> logger) : IHandleMessages<TMessage>, IConsumerContextAccesor { public ConsumerContext ConsumerContext => (innerHandler as IConsumerContextAccesor)?.ConsumerContext!; public async Task Handle(TMessage message) { if (message == null) return; try { await innerHandler.Handle(message); } finally { if (this.ConsumerContext is not null) { await RecordMessageConsumptionAsync(message); } else { logger.LogError($"Configuration error. {nameof(MessageLogDecorator<TMessage>)} used in a pipeline with no pipeline context"); } } } private async Task RecordMessageConsumptionAsync(TMessage message) { if (message == null) return; var messageId = messageContext.Headers[Headers.MessageId]; using var scope = new TransactionScope( TransactionScopeOption.Suppress, TransactionScopeAsyncFlowOption.Enabled); using var dbContext = dbContextFactory.CreateDbContext(); dbContext.MessageConsumptions.Add( new MessageConsumption { MessageId = messageId, ConsumerName = this.ConsumerContext!.ConsumerName!, Headers = JsonSerializer.Serialize(messageContext.Headers), Payload = JsonSerializer.Serialize(message), Result = this.ConsumerContext.ConsumptionResult, Exception = this.ConsumerContext.Exception?.ToString(), EllapsedMs = this.ConsumerContext.EllapsedMilliseconds, OccurredAt = this.ConsumerContext.ProcessedAt }); await dbContext.SaveChangesAsync(); scope.Complete(); } }
可行解决方案
方案1:复用Rebus内置的IMessageContext
Rebus的IMessageContext是每个消息处理请求的专属上下文容器,可直接用来存储自定义上下文数据,无需自行实现访问器。
实现步骤
- 修改上下文构建装饰器,将
ConsumerContext存入IMessageContext.Items:
public class ConsumerContextBuilderDecorator<TMessage>( IHandleMessages<TMessage> innerHandler, IMessageContext messageContext) : IHandleMessages<TMessage> { public async Task Handle(TMessage message) { if (message == null) return; var consumerContext = new ConsumerContext { MessageName = message.GetType().Name, ConsumerName = innerHandler.GetType().Name, ProcessedAt = DateTime.UtcNow }; // 将上下文存入Rebus消息上下文的Items集合 messageContext.Items["ConsumerContext"] = consumerContext; var sw = Stopwatch.StartNew(); try { await innerHandler.Handle(message); consumerContext.ConsumptionResult = MessageConsumptionResult.Success; } catch (Exception ex) { consumerContext.ConsumptionResult = MessageConsumptionResult.Failed; consumerContext.Exception = ex; throw; } finally { sw.Stop(); consumerContext.EllapsedMilliseconds = sw.ElapsedMilliseconds; } } }
- 其他装饰器直接从
IMessageContext.Items获取上下文:
public class IdempotentHandlerDecorator<TMessage>( IHandleMessages<TMessage> innerHandler, IMessageContext messageContext, IMessageTracker messageTracker) : IHandleMessages<TMessage> { public async Task Handle(TMessage message) { if (message == null) return; // 从Rebus上下文获取自定义ConsumerContext if (!messageContext.Items.TryGetValue("ConsumerContext", out var contextObj) || contextObj is not ConsumerContext consumerContext) { throw new InvalidOperationException("ConsumerContext未在消息上下文中初始化"); } var messageId = messageContext.Headers[Headers.MessageId]; if (await messageTracker.IsProcessed(messageId, consumerContext.ConsumerName)) { consumerContext.ConsumptionResult = MessageConsumptionResult.Skipped; return; } await innerHandler.Handle(message); consumerContext.ConsumptionResult = MessageConsumptionResult.Success; await messageTracker.MarkAsProcessed( messageId: messageId, consumerName: consumerContext.ConsumerName, messageName: consumerContext.MessageName); } }
方案2:使用请求作用域的上下文访问器
若不想依赖Rebus内置上下文,可注册一个服务作用域的上下文访问器,同一消息处理请求内的所有装饰器会共享同一个实例。
实现步骤
- 定义上下文访问器接口与实现:
public interface IConsumerContextAccessor { ConsumerContext? ConsumerContext { get; set; } } public class ConsumerContextAccessor : IConsumerContextAccessor { public ConsumerContext? ConsumerContext { get; set; } }
- 将访问器注册为作用域服务:
services.AddScoped<IConsumerContextAccessor, ConsumerContextAccessor>();
- 在上下文构建装饰器中初始化上下文:
public class ConsumerContextBuilderDecorator<TMessage>( IHandleMessages<TMessage> innerHandler, IConsumerContextAccessor contextAccessor) : IHandleMessages<TMessage> { public async Task Handle(TMessage message) { if (message == null) return; var consumerContext = new ConsumerContext { MessageName = message.GetType().Name, ConsumerName = innerHandler.GetType().Name, ProcessedAt = DateTime.UtcNow }; contextAccessor.ConsumerContext = consumerContext; var sw = Stopwatch.StartNew(); try { await innerHandler.Handle(message); consumerContext.ConsumptionResult = MessageConsumptionResult.Success; } catch (Exception ex) { consumerContext.ConsumptionResult = MessageConsumptionResult.Failed; consumerContext.Exception = ex; throw; } finally { sw.Stop(); consumerContext.EllapsedMilliseconds = sw.ElapsedMilliseconds; } } }
- 其他装饰器通过注入
IConsumerContextAccessor获取上下文:
public class MessageLogDecorator<TMessage>( IHandleMessages<TMessage> innerHandler, IMessageContext messageContext, IDbContextFactory<MessageConsumptionDbContext> dbContextFactory, ILogger<MessageLogDecorator<TMessage>> logger, IConsumerContextAccessor contextAccessor) : IHandleMessages<TMessage> { public async Task Handle(TMessage message) { if (message == null) return; var consumerContext = contextAccessor.ConsumerContext ?? throw new InvalidOperationException("ConsumerContext未初始化"); try { await innerHandler.Handle(message); } finally { await RecordMessageConsumptionAsync(message, consumerContext); } } private async Task RecordMessageConsumptionAsync(TMessage message, ConsumerContext consumerContext) { if (message == null) return; var messageId = messageContext.Headers[Headers.MessageId]; using var scope = new TransactionScope( TransactionScopeOption.Suppress, TransactionScopeAsyncFlowOption.Enabled); using var dbContext = dbContextFactory.CreateDbContext(); dbContext.MessageConsumptions.Add( new MessageConsumption { MessageId = messageId, ConsumerName = consumerContext.ConsumerName!, Headers = JsonSerializer.Serialize(messageContext.Headers), Payload = JsonSerializer.Serialize(message), Result = consumerContext.ConsumptionResult, Exception = consumerContext.Exception?.ToString(), EllapsedMs = consumerContext.EllapsedMilliseconds, OccurredAt = consumerContext.ProcessedAt }); await dbContext.SaveChangesAsync(); scope.Complete(); } }
关键注意事项
- 装饰器注册顺序:
ConsumerContextBuilderDecorator必须是最外层装饰器,确保上下文在其他装饰器执行前完成初始化 - 作用域安全:Rebus为每个消息处理请求创建单独的服务作用域,因此作用域服务会自动在同一请求内共享,无需担心并发问题
- 避免单例存储:不要用单例依赖存储上下文,会导致不同请求的上下文互相覆盖
内容的提问来源于stack exchange,提问作者martincardosotoledo
相关产品推荐
相关产品推荐

