Rebus无法反序列化dotnet CAP发布的RabbitMQ消息
CAP发布RabbitMQ事件,Rebus消费反序列化失败问题
我使用CAP库发布事件,通过Rebus从RabbitMQ读取时出现反序列化异常,但CAP配合RabbitMQ.Client读取完全正常。
发布端CAP配置
builder.Services.AddCap(options => { options.FailedRetryInterval = capConfig.FailedRetryIntervalSeconds; options.FailedRetryCount = capConfig.FailedRetryCount; options.UseEntityFramework<MyDbContext>(); options.UseRabbitMQ(options => { options.HostName = rabbitMQConfig.Host; }); }); builder.Services.AddScoped<IMyPublisher, MyPublisher>();
事件发布逻辑
MyEvent.NAME为包含命名空间的类全名:
public async Task Publish(MyEvent myEvent) { using (var transaction = _dbContext.Database.BeginTransaction(_capPublisher, autoCommit: false)) { try { // 事务逻辑 await _capPublisher.PublishAsync(MyEvent.NAME, myEvent); await transaction.CommitAsync(); } catch (Exception ex) { _logger.LogError(ex, ex.Message); await transaction.RollbackAsync(); } } }
通过注入IMyPublisher调用Publish方法发布事件。
CAP正常消费示例
[HttpPost("CAPUNUSED")] [CapSubscribe(MyEvent.NAME)] public async Task TestCAPSubscribe(MyEvent myEvent) { await Console.Out.WriteLineAsync($"{myEvent.Id}"); }
Rebus消费端配置(版本8.2.2)
// 其他非Rebus配置 services.AddRebus(configure => configure .Transport(c => { c.UseRabbitMq(rabbitMqConfig.ConnectionString, rabbitMqConfig.QueueName) .InputQueueOptions(o => { o.SetAutoDelete(false); o.AddArgument("x-message-ttl", rabbitMqConfig.TimeToLive); }); }) //.Serialization(s => s.Register(c => new CustomJsonSerializer())) //.Serialization(s => s.UseNewtonsoftJson()) ); services.AutoRegisterHandlersFromAssemblyOf<MyEventHandler>(); var bus = services.BuildServiceProvider().GetRequiredService<IBus>(); await bus.Subscribe<MyEvent>(); // 其他非Rebus配置
Rebus事件处理器
public class MyEventHandler : IHandleMessages<MyEvent> { public Task Handle(MyEvent message) { Console.WriteLine("IT WORKS"); return Task.CompletedTask; } }
三种测试结果
1. 默认使用System.Text.Json序列化器
消息进入错误队列,抛出异常:
5 unhandled exceptions: 2/5/2024 11:43:03 AM +01:00: System.FormatException: Could not deserialize JSON text: '{"Id":34282,"CollabInfoId":132,"Application":"application","CompanyId":"d6b4f3bc-0f15-46a1-af86-b3dd369158c4","Name":"Name","Extension":"Extension","SizeQuantity":123,"Path":"path","Type":"Type","CreatedBy":"aecdbdc0-8cbc-48e7-a97c-6d91c1b01a6c","DateCreated":"0001-01-01T00:00:00","DateModified":"0001-01-01T00:00:00","FailedIterationsCounter":0,"MaxIterationsNumber":0}' ---> System.ArgumentNullException: Value cannot be null. (Parameter 'returnType') at System.Text.Json.ThrowHelper.ThrowArgumentNullException(String parameterName) at System.Text.Json.JsonSerializer.Deserialize(String json, Type returnType, JsonSerializerOptions options) at Rebus.Serialization.Json.SystemTextJsonSerializer.Deserialize(String bodyString, Type type) --- End of inner exception stack trace --- at Rebus.Serialization.Json.SystemTextJsonSerializer.Deserialize(String bodyString, Type type) at Rebus.Serialization.Json.SystemTextJsonSerializer.GetMessage(TransportMessage transportMessage, Encoding bodyEncoding) at Rebus.Serialization.Json.SystemTextJsonSerializer.Deserialize(TransportMessage transportMessage) at Rebus.Compression.UnzippingSerializerDecorator.Deserialize(TransportMessage transportMessage) at Rebus.Pipeline.Receive.DeserializeIncomingMessageStep.Process(IncomingStepContext context, Func`1 next) at Rebus.DataBus.ClaimCheck.HydrateIncomingMessageStep.Process(IncomingStepContext context, Func`1 next) at Rebus.Pipeline.Receive.HandleDeferredMessagesStep.Process(IncomingStepContext context, Func`1 next) at Rebus.Retry.Simple.DefaultRetryStep.Process(IncomingStepContext context, Func`1 next) // 重复异常信息省略
2. 使用Newtonsoft.Json序列化器
抛出异常:
1 unhandled exceptions: 2/5/2024 11:46:54 AM +01:00: Rebus.Exceptions.MessageCouldNotBeDispatchedToAnyHandlersException: Message with ID knuth-13090383562881270373 and type Newtonsoft.Json.Linq.JObject, Newtonsoft.Json could not be dispatched to any handlers (and will not be retried under the default fail-fast settings) at Rebus.Pipeline.Receive.DispatchIncomingMessageStep.Process(IncomingStepContext context, Func`1 next) at Rebus.Sagas.LoadSagaDataStep.Process(IncomingStepContext context, Func`1 next) at Rebus.Pipeline.Receive.ActivateHandlersStep.Process(IncomingStepContext context, Func`1 next) at Rebus.Pipeline.Receive.HandleRoutingSlipsStep.Process(IncomingStepContext context, Func`1 next) at Rebus.Pipeline.Receive.DeserializeIncomingMessageStep.Process(IncomingStepContext context, Func`1 next) at Rebus.DataBus.ClaimCheck.HydrateIncomingMessageStep.Process(IncomingStepContext context, Func`1 next) at Rebus.Pipeline.Receive.HandleDeferredMessagesStep.Process(IncomingStepContext context, Func`1 next) at Rebus.Retry.Simple.DefaultRetryStep.Process(IncomingStepContext context, Func`1 next)
3. 自定义序列化器(可正常工作,但不想硬编码)
public class CustomJsonSerializer : ISerializer { private readonly JsonSerializerOptions _jsonSerializerOptions; public CustomJsonSerializer() { _jsonSerializerOptions = new JsonSerializerOptions { PropertyNameCaseInsensitive = true, PropertyNamingPolicy = JsonNamingPolicy.CamelCase, }; } public Task<Message> Deserialize(TransportMessage transportMessage) { if (transportMessage == null) throw new ArgumentNullException(nameof(transportMessage)); var jsonString = Encoding.UTF8.GetString(transportMessage.Body); // 使用或不使用序列化选项均可正常工作 var body = JsonSerializer.Deserialize(jsonString, typeof(MyEvent), _jsonSerializerOptions); //var body = JsonSerializer.Deserialize(jsonString, typeof(MyEvent)); var message = new Message(transportMessage.Headers, body); return Task.FromResult(message); } public Task<TransportMessage> Serialize(Message message) { throw new NotImplementedException(); } }
需求
已尝试System.Text.Json、Newtonsoft.Json和自定义序列化器三种方案,仅自定义序列化器可正常工作,但希望避免硬编码序列化逻辑,让其他系统无需显式定义序列化器即可使用事件,求可行解决方案。
内容的提问来源于stack exchange,提问作者Dušan Dimitrov
相关产品推荐
相关产品推荐

