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

使用Rebus ESB的IErrorHandler时出现循环序列化异常

Rebus ESB二级重试后死信消息多次序列化问题

在使用Rebus ESB的IErrorHandler实现二级重试和死信消息处理时,出现异常行为:通过Headers.DeferCount控制消息延迟次数,达到最大延迟次数后消息进入死信队列,但消息体被循环多次序列化,需要多次反序列化才能获取原始内容。本地环境正常,部署后出现该问题。

相关代码片段

自定义HandlePoisonMessage实现

public async Task HandlePoisonMessage(TransportMessage transportMessage, ITransactionContext transactionContext, ExceptionInfo exception)
{
    int MaxDeferCount = this.settings.Settings.MaxDeferCount;
    string messageID= transportMessage.Headers.GetValueOrDefault(Headers.MessageId, Guid.NewGuid().ToString());
    Dictionary<string, string> headers = transportMessage.Headers;
    int deferCount = 1;
    if (transportMessage.Headers.ContainsKey(Headers.DeferCount))
    {
        deferCount = int.Parse(transportMessage.Headers[Headers.DeferCount]);
    }

    if (deferCount <= MaxDeferCount)
    {
        deferCount++;
        var delay = TimeSpan.FromSeconds(2 * deferCount);

        if(headers.ContainsKey(Headers.DeferCount))
        {
            headers[Headers.DeferCount]=deferCount.ToString();
        }
        else
        {
            headers.Add(Headers.DeferCount,deferCount.ToString());
        }

        if (!headers.ContainsKey(CustomHeaders.ApplicationGroup))
        {
           headers.Add(CustomHeaders.ApplicationGroup, settings.Settings.ApplicationGroup);
        }
      
        await bus.Defer(delay, transportMessage, headers);

        log.Info($"Message with ID:{messageID} deferred {deferCount} time{(deferCount > 1 ? "s" : "")}");

    }
    else
    {
        if (!headers.ContainsKey(CustomHeaders.ApplicationGroup))
        {
            headers.Add(CustomHeaders.ApplicationGroup, settings.Settings.ApplicationGroup);
        }

        await errorHandler.HandlePoisonMessage(transportMessage,transactionContext,exception);
        
        log.Info($"Message sent to error queue Message ID:{messageID}");
    }
}

Rebus配置

builder.Services.AddRebus((configure, provider) =>
{
    return configure
        .Transport(t => t.UseRabbitMq($"amqp://{settings.Settings.UserName}:{settings.Settings.Password}@{settings.Settings.HostName}", $"{settings.Settings.EndpointQueueName}").Prefetch(settings.Settings.MaxPrefetchCount))
        .Options(o =>
        {
            o.RetryStrategy(appSettings.Settings.ErrorQueueName, maxDeliveryAttempts: settings.Settings.MaxDeliveryAttempts, secondLevelRetriesEnabled: true);
        })
        .OnCreated(async bus =>
        {
            // 省略其他代码
        });
}, isDefaultBus: true, key: "Default");

死信队列中的消息示例

eyJIZWFkZXJzIjp7ImN1c3RvbS1hcHBsaWNhdGlvbi1ncm91cCI6IkNsYWltcyIsInJiczItY29udGVudC10eXBlIjoiYXBwbGljYXRpb24vanNvbjtjaGFyc2V0PXV0Zi04IiwicmJzMi1jb3JyLWlkIjoiYzFmYmFkYTctMmIxZi00NjEzLWE5OWMtMzU2MmY5ZjlmOTdjIiwicmJzMi1jb3JyLXNlcSI6IjEzIiwicmJzMi1kZWZlci1jb3VudCI6IjQiLCJyYnMyLWRlZmVyLXJlY2lwaWVudCI6ImNsYWltc19wcm9jZXNzaW5nX2F1dG9tYXRpb24iLCJyYnMyLWludGVudCI6InAycCIsInJiczItbXNnLWlkIjoiY2UxZDI0NWItZDFiYi00Yzc1LTk1NDgtMDdlNDY2MzlhZmFhIiwicmJzMi1tc2ctdHlwZSI6Ik5ISUYuU2hhcmVkLk1lc3NhZ2VzLlJldmlld1ByaWNlcywgTkhJRi5TaGFyZWQuTWVzc2FnZXMiLCJyYnMyLXJldHVybi1hZGRyZXNzIjoiY2xhaW1zX3Byb2Nlc3NpbmdfYXV0b21hdGlvbiIsInJiczItc2VuZGVyLWFkZHJlc3MiOiJjbGFpbXNfcHJvY2Vzc2luZ19hdXRvbWF0aW9uIiwicmJzMi1zZW50dGltZSI6IjIwMjYtMDQtMTBUMjI6MzI6MDguMTAwMDQwOVx1MDAyQjAzOjAwIiwieC1kb3RuZXQtcHViLXNlcS1ubyI6IjMyMjIiLCJyYnMyLWRlZmVycmVkLXVudGlsIjoiMjAyNi0wNC0xMVQwMjo1Njo1Ny4xNzQxMDk4XHUwMDJCMDM6MDAifSwiQm9keSI6ImV5SklaV0ZrWlhKeklqcDdJbkppY3pJdFkyOXVkR1Z1ZEMxMGVYQmxJam9pWVhCd2JHbGpZWFJwYjI0dmFuTnZianRqYUdGeWMyVjBQWFYwWmkwNElpd2ljbUp6TWkxamIzSnlMV2xrSWpvaVl6Rm1ZbUZrWVRjdE1tSXhaaTAwTmpFekxXRTVPV010TXpVMk1tWTVaamxtT1Rkaklpd2ljbUp6TWkxamIzSnlMWE5sY1NJNklqRXpJaXdpY21Kek1pMXBiblJsYm5RaU9pSndNbkFpTENKeVluTXlMVzF6WnkxcFpDSTZJbU5sTVdReU5EVmlMV1F4WW1JdE5HTTNOUzA1TlRRNExUQTNaVFEyTmpNNVlXWmhZU0lzSW5KaWN6SXRiWE5uTFhSNWNHVWlPaUpPU0VsR0xsTm9ZWEpsWkM1TlpYTnpZV2RsY3k1U1pYWnBaWGRRY21salpYTXNJRTVJU1VZdVUyaGhjbVZrTGsxbGMzTmhaMlZ6SWl3aWNtSnpNaTF5WlhSMWNtNHRZV1JrY21WemN5STZJbU5zWVdsdGMxOXdjbTlqWlhOemFXNW5YMkYxZEc5dFlYUnBiMjRpTENKeVluTXlMWE5sYm1SbGNpMWhaR1J5WlhOeklqb2lZMnhoYVcxelgzQnliMk5sYzNOcGJtZGZZWFYwYjIxaGRHbHZiaUlzSW5KaWN6SXRjMlZ1ZEhScGJXVWlPaUl5TURJMkxUQTBMVEV3VkRJeU9qTXlPakE0TGpFd01EQTBNRGxjZFRBd01rSXdNem93TUNJc0luZ3RaRzkwYm1WMExYQjFZaTF6WlhFdGJtOGlPaUk1TkRBMU1qVWlMQ0p5WW5NeUxXUmxabVZ5TFdOdmRXNTBJam9pTWlJc0luSmljekl0WkdWbVpYSnlaV1F0ZFc1MGFXd2lPaUl5TURJMkxUQTBMVEV3VkRJek9qSTBPakkwTGpBeU1EUXpOVGxjZFRBd01rSXdNem93TUNJc0luSmljekl0WkdWbVpYSXRjbVZqYVhCcFpXNTBJam9pWTJ4aGFXMXpYM0J5YjJObGMzTnBibWRmWVhWMGIyMWhkR2x2YmlJc0ltTjFjM1J2YlMxaGNIQnNhV05oZEdsdmJpMW5jbTkxY0NJNklrTnNZV2x0Y3lKOUxDSkNiMlI1SWpvaVpYbEtSMkl5ZUhCaU1HeEZTV3B2YVZsVVFUQlplbFpzVFRKWmRFMUhVVFZaZVRBd1dsUkpkMHhVYXpSUFJHZDBXWHBPYWxwRVl6SmFSR3N6VDBkR2FrbHVNRDBpZlE9PSJ9

问题成因分析

  • 嵌套序列化触发:调用bus.Defer(delay, transportMessage, headers)时,直接传递了TransportMessage实例。Rebus的Defer方法会将传入的对象整体序列化,而TransportMessage本身已经包含了序列化后的消息体,导致二次(甚至多次)序列化,最终死信消息的Body是多层嵌套的序列化结果。
  • 环境差异原因:本地环境可能因调试模式下序列化器的配置(如允许循环引用、宽松的反序列化规则)或消息传递的本地特性,使得嵌套序列化的问题未暴露;部署环境的序列化器配置更严格,或RabbitMQ的消息存储/传递方式放大了该问题。
  • 重试策略冲突:Rebus配置中已经启用了内置的RetryStrategy(二级重试),同时自定义了HandlePoisonMessage处理逻辑,两者叠加可能导致消息被重复延迟处理,加剧了序列化嵌套的次数。

解决方案

  1. 避免直接传递TransportMessage
    不要直接将TransportMessage传给Defer方法,应先反序列化原始消息体,再传递原始消息对象进行延迟发送:

    // 注入ISerializer实例
    var serializer = provider.GetRequiredService<ISerializer>();
    // 反序列化TransportMessage得到原始消息对象
    var originalMessage = await serializer.Deserialize(transportMessage);
    // 使用原始消息对象进行延迟发送
    await bus.Defer(delay, originalMessage, headers);
    
  2. 调整重试逻辑
    若Rebus内置的二级重试策略已满足需求,可考虑移除自定义的HandlePoisonMessage中手动控制延迟的逻辑,避免重复处理。若必须保留自定义逻辑,需确保与内置策略的重试次数不叠加。

  3. 检查序列化配置
    统一本地和部署环境的序列化器配置(如JSON序列化器的ReferenceLoopHandling、TypeNameHandling等设置),确保环境一致性,避免本地调试时掩盖问题。

内容的提问来源于stack exchange,提问作者Amour Rashid

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 15:04:54