【问题标题】:Only consuming messages with certain headers using RabbitMQ and SpringAMQP仅使用 RabbitMQ 和 SpringAMQP 使用具有某些标头的消息
【发布时间】:2014-10-18 19:17:23
【问题描述】:

我正在尝试将消息发布到队列,然后让某些消费者仅在它包含某个标头时才使用它,如果它包含另一个标头,则另一个消费者使用它。

到目前为止,我所做的是设置一个标头交换,仅当消息包含该标头时才将消息路由到某个队列。

这是我用来设置交换、队列和监听器的配置:

<!-- Register Queue Listener Beans -->
<bean id="ActionMessageListener" class="com.mycee.Action" />    

<!-- Register RabbitMQ Connections -->
<rabbit:connection-factory 
    id="connectionFactory"
    port="${rabbit.port}"
    virtual-host="${rabbit.virtual}" 
    host="${rabbit.host}" 
    username="${rabbit.username}"
    password="${rabbit.password}" 
    connection-factory="nativeConnectionFactory" />

<!-- Register RabbitMQ Listeners -->
<rabbit:listener-container          
    connection-factory="connectionFactory"
    channel-transacted="true"
    requeue-rejected="true"
    concurrency="${rabbit.consumers}">
    <rabbit:listener queues="${queue.myqueue}" ref="ActionMessageListener" method="handle"/>        
</rabbit:listener-container>    

<!-- Setup RabbitMQ headers exchange -->
<rabbit:headers-exchange id="${exchange.myexchange}" name="${exchange.myexchange}">
    <rabbit:bindings>
        <rabbit:binding queue="${queue.myqueue}" key="action" value="action3" />            
    </rabbit:bindings>
</rabbit:headers-exchange>

<rabbit:admin connection-factory="connectionFactory"/>  
<rabbit:queue name="${queue.myqueue}" />

所以我使用 action 的键和 action3 的值将 myqueue 绑定到 myexchange。

现在当我在交易所发布时:

即使操作设置为 action1 而不是 action3,ChannelAwareMessageListener 仍在使用它

public class Action implements ChannelAwareMessageListener {

    @Override
    public void onMessage(Message message, Channel channel) throws Exception {

        System.out.println(message.toString());

    }

}

要么我没有正确使用 headers-exchange,要么我没有正确配置它 - 有什么建议吗?

【问题讨论】:

    标签: java spring rabbitmq spring-amqp


    【解决方案1】:

    这样不行;每个消费者都需要一个单独的队列。见the tutorial

    当多个消费者从同一个队列消费时,他们会竞争所有消息;您不能在消费者端选择消息; “选择”是由交换机通过将消息路由到特定队列来完成的。

    【讨论】:

    • 好吧,只是为了确认我今天学到的东西,不管我使用的是标题、主题、直接还是扇出交换,所有这些交换所做的都是将消息路由到队列。如果我有 100 个消费者,每个消费者都使用不同的操作,我应该有一个名为 action00 到 action99 的 100 个队列,而不是像上面问题中那样将它们全部路由到“myqueue”。此外,没有办法根据某些参数(你总是得到第一个参数)从队列中消费消息,或者直接绑定到交换器并直接从交换器消费,交换器的目的是路由到队列。
    • 为了巩固这个想法,如果我通过扇出交换从发布者进行 RPC 调用,假设消息被扇出到 10 个队列,被 10 个消费者消费,然后他们都响应,是否会消耗第一个响应,而其余响应会一直处于等待状态,直到超时?
    • 是的(对第一条评论中的所有人)。第二条评论 - 这取决于您如何配置您的回复,但基本上,是的,所有消费者都可以回复。困难在于知道要消费多少回复(您必须知道您希望回复多少消费者)。您还必须以某种方式关联回复。实现这两者的一种简单方法是为回复声明一个临时自动删除队列并收听它,直到发生超时。在 Spring-AMQP 中没有对这种场景的开箱即用支持,但使用 RabbitTemplate.execute() 推出自己的场景并不难
    猜你喜欢
    • 2019-08-24
    • 2014-10-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-06-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多