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

如何修改MassTransit.Kafka消息JSON序列化:驼峰转帕斯卡

问题描述

我用MassTransit结合Kafka实现了Saga状态机,配置代码如下:

services.AddMassTransit(x =>
{
    // Configure the transport
    x.UsingInMemory();

    x.AddRider(rider =>
    {
        rider.AddSagaStateMachine<OrderStateMachine, OrderSagaState>()
            .EntityFrameworkRepository(r =>
            {
                r.ConcurrencyMode = ConcurrencyMode.Optimistic; // or use Pessimistic, which does not require RowVersion

                r.AddDbContext<DbContext, PostgresDbContext>((provider, builder) =>
                {
                    builder.UseNpgsql(configuration.GetConnectionString("PostgresSql"), m =>
                    {
                        m.MigrationsAssembly(Assembly.GetExecutingAssembly().GetName().Name);
                        m.MigrationsHistoryTable($"__{nameof(PostgresDbContext)}");
                    });
                });

                //This line is added to enable PostgreSQL features
                r.UsePostgres();
            });

        rider.UsingKafka(new ClientConfig() { 
            BootstrapServers = configuration["Kafka:BootstrapServers"],
            //SaslUsername = "", SaslPassword = "", SaslMechanism = SaslMechanism.Plain, 
            SecurityProtocol = SecurityProtocol.Plaintext,  
            ApiVersionRequest = true }, 
            (context, kafka ) => {
                // Configure saga endpoints for Kafka topics
                kafka.TopicEndpoint<CreateOrderMessage>(TopicNames.CreateOrder, TopicNames.CreateOrder + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });

                kafka.TopicEndpoint<PaymentCompletedEvent>(TopicNames.PaymentCompleted, TopicNames.PaymentCompleted + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });

                kafka.TopicEndpoint<PaymentFailedEvent>(TopicNames.PaymentFailed, TopicNames.PaymentFailed + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });

                kafka.TopicEndpoint<StockReservedEvent>(TopicNames.StockReserved, TopicNames.StockReserved + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });

                kafka.TopicEndpoint<StockReservationFailedEvent>(TopicNames.StockReservationFailed, TopicNames.StockReservationFailed + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });
            });

        rider.AddProducer<string, OrderCreatedEvent>(TopicNames.OrderCreated);
        rider.AddProducer<string, CompletePaymentMessage>(TopicNames.CompletePayment);
        rider.AddProducer<string, OrderCompletedEvent>(TopicNames.OrderCompleted);
        rider.AddProducer<string, OrderFailedEvent>(TopicNames.OrderFailed);
        rider.AddProducer<string, StockRollBackMessage>(TopicNames.StockRollback);
    });
});

当前发布到Kafka的消息采用驼峰式命名,示例如下:

{
    "orderItemList": [
        {
            "productId": 4,
            "count": 2
        }
    ],
    "correlationId": "c93b2d0f-0c87-427b-9744-d7026a677bce"
}

由于协作服务使用CAP库订阅这些消息时,反序列化后属性都是默认值,需要把消息改为帕斯卡式命名(首字母大写),请问如何修改MassTransit Kafka Rider的JSON序列化规则?

解决方案

要修改MassTransit Kafka Rider的JSON序列化规则,让消息采用帕斯卡式命名,需要配置JsonSerializerOptions,指定属性命名策略为保持C#类原属性名称(即帕斯卡命名),同时确保生产者和消费者使用相同的配置。

修改后的完整配置代码如下:

services.AddMassTransit(x =>
{
    x.UsingInMemory();

    x.AddRider(rider =>
    {
        rider.AddSagaStateMachine<OrderStateMachine, OrderSagaState>()
            .EntityFrameworkRepository(r =>
            {
                r.ConcurrencyMode = ConcurrencyMode.Optimistic;

                r.AddDbContext<DbContext, PostgresDbContext>((provider, builder) =>
                {
                    builder.UseNpgsql(configuration.GetConnectionString("PostgresSql"), m =>
                    {
                        m.MigrationsAssembly(Assembly.GetExecutingAssembly().GetName().Name);
                        m.MigrationsHistoryTable($"__{nameof(PostgresDbContext)}");
                    });
                });

                r.UsePostgres();
            });

        rider.UsingKafka(new ClientConfig() { 
            BootstrapServers = configuration["Kafka:BootstrapServers"],
            SecurityProtocol = SecurityProtocol.Plaintext,  
            ApiVersionRequest = true }, 
            (context, kafka ) => {
                // 配置JSON序列化选项,使用帕斯卡命名
                var jsonOptions = new JsonSerializerOptions
                {
                    PropertyNamingPolicy = null, // 不转换命名,保持C#类的帕斯卡命名
                    WriteIndented = false // 根据需求设置是否格式化输出
                };

                // 为生产者和消费者统一配置序列化/反序列化规则
                kafka.ConfigureJsonSerializer(options => jsonOptions);
                kafka.ConfigureJsonDeserializer(options => jsonOptions);

                // 原有TopicEndpoint配置保持不变
                kafka.TopicEndpoint<CreateOrderMessage>(TopicNames.CreateOrder, TopicNames.CreateOrder + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });

                kafka.TopicEndpoint<PaymentCompletedEvent>(TopicNames.PaymentCompleted, TopicNames.PaymentCompleted + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });

                kafka.TopicEndpoint<PaymentFailedEvent>(TopicNames.PaymentFailed, TopicNames.PaymentFailed + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });

                kafka.TopicEndpoint<StockReservedEvent>(TopicNames.StockReserved, TopicNames.StockReserved + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });

                kafka.TopicEndpoint<StockReservationFailedEvent>(TopicNames.StockReservationFailed, TopicNames.StockReservationFailed + "-group", e =>
                {
                    e.ConfigureSaga<OrderSagaState>(context);
                });
            });

        // 原有Producer配置保持不变
        rider.AddProducer<string, OrderCreatedEvent>(TopicNames.OrderCreated);
        rider.AddProducer<string, CompletePaymentMessage>(TopicNames.CompletePayment);
        rider.AddProducer<string, OrderCompletedEvent>(TopicNames.OrderCompleted);
        rider.AddProducer<string, OrderFailedEvent>(TopicNames.OrderFailed);
        rider.AddProducer<string, StockRollBackMessage>(TopicNames.StockRollback);
    });
});

关键说明

  • 设置PropertyNamingPolicy = null会让JSON序列化器直接使用C#类的属性名称(如OrderItemList、ProductId),不再自动转换为驼峰命名。
  • 必须同时配置ConfigureJsonSerializer(生产者使用)和ConfigureJsonDeserializer(消费者使用),保证两端序列化规则一致。
  • 如果需要更个性化的命名转换,可以自定义实现JsonNamingPolicy类,替换PropertyNamingPolicy的值。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 01:31:03