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

如何使用MassTransit以纯文本形式消费RabbitMQ消息?

MassTransit 纯文本消息消费实现方案

MassTransit 本身没有内置的纯文本 ISerializerFactory,但可以通过自定义序列化器或直接消费原始消息的方式满足你的需求,以下是具体实现方案:

方案1:自定义纯文本序列化器(实现ISerializerFactory)

你可以自行实现 ISerializerFactory 及配套的序列化/反序列化逻辑,直接处理纯文本消息:

自定义序列化器代码

public class PlainTextSerializerFactory : ISerializerFactory
{
    public ISerializer CreateSerializer() => new PlainTextSerializer();
    public IDeserializer CreateDeserializer() => new PlainTextDeserializer();
    public string ContentType => "text/plain";
}

public class PlainTextSerializer : ISerializer
{
    public Task Serialize<T>(Stream stream, T message, CancellationToken cancellationToken)
    {
        if (message is string text)
            return stream.WriteAsync(Encoding.UTF8.GetBytes(text), cancellationToken);
        
        throw new InvalidOperationException("仅支持字符串类型消息");
    }
}

public class PlainTextDeserializer : IDeserializer
{
    public async Task<ConsumeContext> Deserialize(Stream stream, Headers headers, CancellationToken cancellationToken)
    {
        var text = await new StreamReader(stream).ReadToEndAsync(cancellationToken);
        return new ConsumeContextProxy(headers, text);
    }
}

配置MassTransit使用自定义序列化器

services.AddMassTransit(x =>
{
    x.AddConsumer<PlainTextConsumer>();

    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host("rabbitmq://localhost");
        
        // 注册自定义序列化器工厂
        cfg.AddSerializer(new PlainTextSerializerFactory());
        cfg.AddDeserializer(new PlainTextSerializerFactory());

        cfg.ReceiveEndpoint("plain-text-queue", e =>
        {
            e.ConfigureConsumer<PlainTextConsumer>(context);
            // 指定端点使用纯文本序列化/反序列化器
            e.UseSerializer<PlainTextSerializer>();
            e.UseDeserializer<PlainTextDeserializer>();
        });
    });
});

纯文本消费者实现

public class PlainTextConsumer : IConsumer<string>
{
    public Task Consume(ConsumeContext<string> context)
    {
        var plainTextContent = context.Message;
        // 此处处理纯文本消息逻辑
        return Task.CompletedTask;
    }
}

方案2:直接消费原始消息(跳过序列化)

如果不想自定义序列化器,也可以直接消费 RawMessage 绕过MassTransit的序列化流程:

原始消息消费者代码

public class RawTextConsumer : IConsumer<RawMessage>
{
    public async Task Consume(ConsumeContext<RawMessage> context)
    {
        var rawPlainText = await context.Message.Body.ReadAsStringAsync();
        // 处理原始纯文本内容
        return Task.CompletedTask;
    }
}

配置端点启用原始消息消费

cfg.ReceiveEndpoint("raw-text-queue", e =>
{
    e.ConfigureConsumer<RawTextConsumer>(context);
    e.ConsumeMessageOnly = true;
    e.UseRawJsonSerializer(); // 跳过默认序列化逻辑,直接读取原始内容
});

关键注意事项

  • 确保自定义序列化器的 ContentType 与生产者发送消息时的 content-type 头完全匹配,否则MassTransit无法正确选择你的反序列化器。
  • 未来对接其他队列系统时,只需在对应端点复用上述自定义序列化器配置或原始消息消费逻辑,即可保持跨队列的一致性处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 06:40:09