使用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处理逻辑,两者叠加可能导致消息被重复延迟处理,加剧了序列化嵌套的次数。
解决方案
避免直接传递TransportMessage
不要直接将TransportMessage传给Defer方法,应先反序列化原始消息体,再传递原始消息对象进行延迟发送:// 注入ISerializer实例 var serializer = provider.GetRequiredService<ISerializer>(); // 反序列化TransportMessage得到原始消息对象 var originalMessage = await serializer.Deserialize(transportMessage); // 使用原始消息对象进行延迟发送 await bus.Defer(delay, originalMessage, headers);调整重试逻辑
若Rebus内置的二级重试策略已满足需求,可考虑移除自定义的HandlePoisonMessage中手动控制延迟的逻辑,避免重复处理。若必须保留自定义逻辑,需确保与内置策略的重试次数不叠加。检查序列化配置
统一本地和部署环境的序列化器配置(如JSON序列化器的ReferenceLoopHandling、TypeNameHandling等设置),确保环境一致性,避免本地调试时掩盖问题。
内容的提问来源于stack exchange,提问作者Amour Rashid
相关产品推荐
相关产品推荐

