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

使用MassTransit从Azure EventHub消费自定义格式非信封消息遇阻

解决MassTransit消费Azure EventHub自定义格式消息的问题

问题原因

使用UseRawJsonSerializer后,MassTransit无法将消息中自定义的type字段与你的EventHubMessage类型关联,导致ConsumeContextMessageTypeFilter找不到匹配的消费管道,消息无法到达消费者。

解决方案

添加自定义消息类型解析逻辑,让MassTransit能根据消息中的type字段识别对应的CLR消费类型,同时确保序列化配置正确。

修改后的完整配置代码

var builder = Host.CreateApplicationBuilder();
builder.Services
    .AddMassTransit(x =>
    {
        x.UsingInMemory();
        x.AddRider(rider =>
        {
            rider.AddConsumer<EventHubMessageConsumer>();

            rider.UsingEventHub((context, k) =>
            {
                k.Host(
                    "Endpoint=sb://localhost;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=SAS_KEY_VALUE;UseDevelopmentEmulator=true;");
                k.Storage("UseDevelopmentStorage=true");

                k.ReceiveEndpoint("eh1", c =>
                {
                    // 使用原始JSON序列化器
                    c.UseRawJsonSerializer(RawSerializerOptions.AnyMessageType, true);
                    
                    // 添加自定义类型解析逻辑:根据消息的type字段映射到EventHubMessage
                    c.UseMessageTypeConverter(async context =>
                    {
                        // 将消息体读取为JObject(以Newtonsoft.Json为例)
                        var jsonBody = await context.Body.ReadAsStringAsync();
                        var message = Newtonsoft.Json.JsonConvert.DeserializeObject<Newtonsoft.Json.Linq.JObject>(jsonBody);
                        var typeValue = message["type"]?.ToString();
                        
                        // 当type为ExampleMessage时,返回对应的CLR类型
                        return typeValue == "ExampleMessage" ? typeof(EventHubMessage) : null;
                    });
                    
                    // 配置消费者
                    c.ConfigureConsumer<EventHubMessageConsumer>(context);
                });
            });
        });
    });

var host = builder.Build();
await host.RunAsync();

public class EventHubMessageConsumer :
    IConsumer<EventHubMessage>
{
    public Task Consume(ConsumeContext<EventHubMessage> context)
    {
        // 处理消息逻辑
        Console.WriteLine($"Received message type: {context.Message.Type}");
        return Task.CompletedTask;
    }
}

public record EventHubMessage
{
    public string Type { get; init; }
    public object Data { get; init; } // 对应消息中的data字段,可根据实际业务调整类型
}

关键说明

  • UseMessageTypeConverter:自定义类型解析逻辑,从消息体中读取type字段,映射到你的EventHubMessage类型,让MassTransit明确该用哪个消费者处理消息。
  • 确保EventHubMessage的属性与消息JSON结构完全匹配,比如添加Data属性对应消息中的data字段。
  • 如果使用System.Text.Json替代Newtonsoft.Json,只需调整反序列化的代码逻辑即可。

另一种简化方案(固定消息类型)

如果你的EventHub中所有消息都符合EventHubMessage结构,可以直接指定序列化器的目标类型,跳过类型解析步骤:

k.ReceiveEndpoint("eh1", c =>
{
    // 直接指定反序列化为EventHubMessage类型
    c.UseRawJsonSerializer(RawSerializerOptions.SpecificMessageType, true);
    c.ConfigureConsumer<EventHubMessageConsumer>(context);
    // 绑定消息类型与事件中心
    c.Message<EventHubMessage>(m => m.SetEntityName("eh1"));
});

这种方案适用于所有消息都是同一类型的场景,无需解析type字段。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 23:14:53