You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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是每个消息处理请求的专属上下文容器,可直接用来存储自定义上下文数据,无需自行实现访问器。

实现步骤

  1. 修改上下文构建装饰器,将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;
        }
    }
}
  1. 其他装饰器直接从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内置上下文,可注册一个服务作用域的上下文访问器,同一消息处理请求内的所有装饰器会共享同一个实例。

实现步骤

  1. 定义上下文访问器接口与实现:
public interface IConsumerContextAccessor
{
    ConsumerContext? ConsumerContext { get; set; }
}

public class ConsumerContextAccessor : IConsumerContextAccessor
{
    public ConsumerContext? ConsumerContext { get; set; }
}
  1. 将访问器注册为作用域服务:
services.AddScoped<IConsumerContextAccessor, ConsumerContextAccessor>();
  1. 在上下文构建装饰器中初始化上下文:
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;
        }
    }
}
  1. 其他装饰器通过注入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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.11 16:37:01