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

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}'
---&gt; 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 18:37:04