【问题标题】:Cannot concurrently consume from ActiveMQ embedded broker无法同时从 ActiveMQ 嵌入式代理消费
【发布时间】:2015-04-01 23:35:46
【问题描述】:

我正在寻找一种将 10 条消息发布到 ActiveMQ 嵌入式代理的方法,并在同一个 VM 上使用 JMS API 同时使用它们。

下面的代码存在某种竞争,因为有时它会并行消耗 2、4、8 条消息,并在 latch.await 调用时挂起直到超时。

public final class ActiveMQJMSParallelTest {
    private static final Logger logger = LoggerFactory.getLogger(ActiveMQJMSParallelTest.class);
    private static final int numberOfMessages = 10;

    public static void main(final String[] args) throws Exception {
        final Properties props = new Properties();
        props.setProperty(Context.INITIAL_CONTEXT_FACTORY, "org.apache.activemq.jndi.ActiveMQInitialContextFactory");
        props.setProperty(Context.PROVIDER_URL, "vm://localhost?broker.persistent=false");
        props.setProperty("queue.parallelQueue", "parallelQueue");
        final Context jndiContext = new InitialContext(props);
        final ConnectionFactory connectionFactory = (ConnectionFactory) jndiContext.lookup("ConnectionFactory");
        final Destination destination = (Destination) jndiContext.lookup("parallelQueue");
        final Connection connection = connectionFactory.createConnection();
        Session session = null;
        try {
            session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE);
            final MessageProducer producer = session.createProducer(destination);
            for (int i = 0; i < numberOfMessages; i++) {
                final TextMessage message = session.createTextMessage();
                message.setText("This is message " + (i + 1));
                producer.send(message);
                logger.info("Produced message: {}", message);
            }
            session.commit();
        } finally {
            if (session != null)
                session.close();
        }

        final CountDownLatch latch = new CountDownLatch(numberOfMessages);
        final ExecutorService pool = Executors.newFixedThreadPool(numberOfMessages);
        for (int i = 0; i < numberOfMessages; i++) {
            pool.submit(new Runnable() {
                @Override public void run() {
                    try {
                        final Connection connection = connectionFactory.createConnection();
                        connection.start();
                        final Session session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE);
                        final Queue destination = session.createQueue("parallelQueue");
                        final MessageConsumer consumer = session.createConsumer(destination);
                        final Message received = consumer.receive();
                        logger.info("Consuming message: {}", received);
                        latch.countDown();
                        latch.await(1, TimeUnit.MINUTES);
                        logger.info("Consumed message: {}", received);
                        session.close();
                        connection.close();
                    } catch(Exception e) {
                        e.printStackTrace();
                    }
                }
            });
        }
        latch.await(10, TimeUnit.MINUTES);
        jndiContext.close();
    }
}

有人可以为这个任务变出工作代码吗?

【问题讨论】:

    标签: java jms activemq


    【解决方案1】:

    如果您想确保每个消费者有机会一次获取一条消息,那么您应该使用零预取值,这样代理就不会尝试调度第一个消费者的预取限制等当他们到达时。

    看看documentation 页面上的预取是如何工作的。

    【讨论】:

    • 谢谢你在这里回答,蒂莫西。
    • 为了让这个例子工作,下面的行改变是必要的:"vm://localhost?broker.persistent=false&amp;jms.prefetchPolicy.all=0"
    猜你喜欢
    • 2011-05-28
    • 1970-01-01
    • 2012-05-13
    • 2013-01-06
    • 2012-10-19
    • 2015-03-28
    • 2018-02-16
    • 1970-01-01
    • 2012-10-06
    相关资源
    最近更新 更多