【问题标题】:MassTransit consume non MassTransit messageMassTransit 使用非 MassTransit 消息
【发布时间】:2017-11-18 11:11:42
【问题描述】:

我有一个控制台应用程序将消息发布到 RabbitMQ 交换。使用 MassTransit 构建的订阅者是否可以使用此消息?

这是发布者代码:

    public virtual void Send(LogEntryMessage message)
    {

        using (var connection = _factory.CreateConnection())
        using (var channel = connection.CreateModel())
        {
            var props = channel.CreateBasicProperties();
            props.CorrelationId = Guid.NewGuid().ToString();

            var body = Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(message));

            channel.BasicPublish(exchange: _endpointConfiguration.Exchange, routingKey: _endpointConfiguration.RoutingKey, basicProperties: null,
                body: body);
        }
    }

这是订阅者代码:

      IBusControl ConfigureBus()
      {
        return Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            var host = cfg.Host(new Uri("rabbitmq://localhost"), h =>
            {
                h.Username(username);
                h.Password(password);
            });

            cfg.ReceiveEndpoint(host, "LogEntryQueue", e =>
            {
                e.Handler<LogEntryMessage>(context =>
                Console.Out.WriteLineAsync($"Value was entered: {context.Message.MessageBody}"));
            });
        });
    }

这是消费者代码:

    public class LogEntryMessageProcessor : IConsumer<LogEntryMessage>
    {
        public Task Consume(ConsumeContext<LogEntryMessage> context)
        {
            Console.Out.WriteLineAsync($"Value was entered: 
                      {context.Message.Message.MessageBody}");
            return Task.FromResult(0);
        }
    }

【问题讨论】:

    标签: rabbitmq masstransit rabbitmq-exchange


    【解决方案1】:

    希望你能在Interoperability部分得到答案,特别是看example message。

    基本上,你需要按照一些简单的规则来构造一个 JSON 对象。

    示例消息如下所示:

    {
        "destinationAddress": "rabbitmq://localhost/input_queue",
        "headers": {},
        "message": {
            "value": "Some Value",
            "customerId": 27
        },
        "messageType": [
            "urn:message:MassTransit.Tests:ValueMessage"
        ]
    }
    

    您可以通过创建发布者和消费者轻松检查更复杂的消息的外观,运行程序以创建绑定,然后停止消费者并发布一些消息。它们将在订阅者队列中,因此您可以使用管理插件轻松阅读它们。

    【讨论】:

    • 感谢您的提示。我能够发布和订阅,但消息处理器无法序列化消息。 MassTransit 消息上下文已创建,它包含消息对象,但消息的所有属性均为空。
    • 如果你能在github上写一个示例repo,我大概可以找点时间看看。
    • 谢谢!非常感谢任何帮助!
    • 我创建了github.com/cksanjose/masstransit_pubsub 和github.com/cksanjose/masstransit_subscriber。第一个是普通的 RabbitMQ 发布者。订阅者是 MassTransit 和 RabbitMQ。
    • 似乎github.com/cksanjose/masstransit_pubsub/blob/master/… 应该是 LogEntryPayload 并且您的消费者应该使用该类型而不是外部类型,外部类型本质上是您的 MT 消息信封版本。
    【解决方案2】:

    为了让 MassTransit 处理由非 MassTransit 客户端发布的消息, 消息必须包含 MassTransit 所需的元数据,如Interoperability 页面中所述。 消息的消费者必须处理消息的有效负载。 在下面的代码中,payload 是 LogEntryPayload:

    public class LogEntryMessageProcessor : IConsumer<LogEntryPayload>
    {
        public Task Consume(ConsumeContext<LogEntryPayload> context)
        {
            //var payload = context.GetPayload<LogEntryPayload>();
            Console.Out.WriteLineAsync($"Value was entered: {context.Message.Id} - {context.Message.MessageBody}");
            return Task.FromResult(0);
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-03-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-12-13
      • 1970-01-01
      • 2023-02-26
      • 1970-01-01
      相关资源
      最近更新 更多