【问题标题】:Multiple queues receiving same message from virtual topic creates a deadletter entry for one queue only从虚拟主题接收相同消息的多个队列仅为一个队列创建一个死信条目
【发布时间】:2022-01-19 22:42:24
【问题描述】:

我正在使用虚拟目的地在 ActiveMQ 5.15.13 中实现发布订阅模型。

我有一个虚拟主题VirtualTopic,并且绑定了两个队列。每个队列都有自己的重新传递策略。假设Queue 1 将重试消息 2 次,以防在处理消息时出现异常,Queue 2 将重试消息 3 次。发布重试消息将被发送到死信队列。我也在使用单独的死信队列策略,这样每个队列都有自己的死信队列。

我观察到,当向VirtualTopic 发送消息时,具有相同消息ID 的消息会被传递到两个队列。我面临一个问题,如果两个队列的消费者都无法成功处理消息。发往Queue 1 的消息在重试 2 次后被移至死信队列。但是Queue 2 没有死信队列,虽然队列 2 中的消息被重试了 3 次。

这是预期的行为吗?

代码:

public class ActiveMQRedelivery {

    private final ActiveMQConnectionFactory factory;

    public ActiveMQRedelivery(String brokerUrl) {
        factory = new ActiveMQConnectionFactory(brokerUrl);
        factory.setUserName("admin");
        factory.setPassword("password");
        factory.setAlwaysSyncSend(false);
    }

    public void publish(String topicAddress, String message) {
        final String topicName = "VirtualTopic." + topicAddress;
        try {
            final Connection producerConnection = factory.createConnection();
            producerConnection.start();
            final Session producerSession = producerConnection.createSession(false, AUTO_ACKNOWLEDGE);
            final MessageProducer producer = producerSession.createProducer(null);
            final TextMessage textMessage = producerSession.createTextMessage(message);
            final Topic topic = producerSession.createTopic(topicName);
            producer.send(topic, textMessage, PERSISTENT, DEFAULT_PRIORITY, DEFAULT_TIME_TO_LIVE);
        } catch (JMSException e) {
            throw new RuntimeException("Message could not be published", e);
        }

    }

    public void initializeConsumer(String queueName, String topicAddress, int numOfRetry) throws JMSException {
        factory.getRedeliveryPolicyMap().put(new ActiveMQQueue("*." + queueName + ".>"),
                getRedeliveryPolicy(numOfRetry));
        Connection connection = factory.createConnection();
        connection.start();
        final Session consumerSession = connection.createSession(false, CLIENT_ACKNOWLEDGE);
        final Queue queue = consumerSession.createQueue("Consumer." + queueName +
                ".VirtualTopic." + topicAddress);
        final MessageConsumer consumer = consumerSession.createConsumer(queue);
        consumer.setMessageListener(message -> {
            try {
                    System.out.println("in listener --- " + ((ActiveMQDestination)message.getJMSDestination()).getPhysicalName());
                    consumerSession.recover();
            } catch (JMSException e) {
                e.printStackTrace();
            }
        });
    }

    private RedeliveryPolicy getRedeliveryPolicy(int numOfRetry) {
        final RedeliveryPolicy redeliveryPolicy = new RedeliveryPolicy();
        redeliveryPolicy.setInitialRedeliveryDelay(0);
        redeliveryPolicy.setMaximumRedeliveries(numOfRetry);
        redeliveryPolicy.setMaximumRedeliveryDelay(-1);
        redeliveryPolicy.setRedeliveryDelay(0);
        return redeliveryPolicy;
    }

}

测试:

public class ActiveMQRedeliveryTest {

    private static final String brokerUrl = "tcp://0.0.0.0:61616";
    private ActiveMQRedelivery activeMQRedelivery;

    @Before
    public void setUp() throws Exception {
        activeMQRedelivery = new ActiveMQRedelivery(brokerUrl);
    }

    @Test
    public void testMessageRedeliveries() throws Exception {
        String topicAddress = "testTopic";
        activeMQRedelivery.initializeConsumer("queue1", topicAddress, 2);
        activeMQRedelivery.initializeConsumer("queue2", topicAddress, 3);
        activeMQRedelivery.publish(topicAddress, "TestMessage");
        Thread.sleep(3000);
    }

    @After
    public void tearDown() throws Exception {
    }
}

【问题讨论】:

  • 您找到解决此问题的方法了吗?

标签: java activemq dead-letter virtual-topic


【解决方案1】:

我最近遇到了这个问题。要解决此问题,需要将 2 个属性添加到 individualDeadLetterStrategy,如下所示

<deadLetterStrategy>
        <individualDeadLetterStrategy destinationPerDurableSubscriber="true" enableAudit="false" queuePrefix="DLQ." useQueueForQueueMessages="true"/>
</deadLetterStrategy>

属性说明:

destinationPerDurableSubscriber - 为每个持久订阅者启用单独的目标。

enableAudit - 死信策略具有默认启用的消息审核。这可以防止重复的消息被添加到配置的 DLQ。启用该属性后,当destinationPerDurableSubscriber 属性设置为true 时,未为多个订阅者传递到主题的同一消息将仅放置在其中一个订阅者DLQ 上,即假设两个消费者未能确认同一消息的同一消息。主题,该消息只会被放置在一个消费者而不是另一个消费者的 DLQ 上。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-03-17
    • 2012-05-24
    • 2013-02-07
    • 1970-01-01
    • 2014-04-13
    • 1970-01-01
    • 2017-02-13
    • 2019-07-23
    相关资源
    最近更新 更多