【问题标题】:HornetQ: How to reuse XAConnection and XASessionHornetQ:如何重用 XAConnection 和 XASession
【发布时间】:2013-09-09 02:35:41
【问题描述】:

我在尝试在我的 JBoss 应用程序中的多个工作人员上重用 XAConnection 和 XASession 时遇到了一些问题。我已经设法将问题简化为一种方法。它应该能够ProduceConsumer 使用相同的连接和会话来发送消息。目前我的应用程序有很多队列和工作人员,每个工作人员当前正在启动和启动每个自己的连接和会话,而不是共享它。这不应该是可能的吗?

这是我的代码示例:

import org.apache.log4j.Logger;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import javax.ejb.Singleton;
import javax.ejb.Startup;
import javax.jms.*;
import javax.jms.Queue;
import javax.naming.InitialContext;

@Singleton
@Startup
public class QueueTest {

    private Logger logger = Logger.getLogger(QueueTest.class);

    @PostConstruct
    public void startup() {
        try {
            String queue = "queue/Queue1";
            String message = "test";

            //setting up connection
            InitialContext iniCtx = new InitialContext();
            XAConnectionFactory qcf = (XAConnectionFactory) iniCtx.lookup("java:/JmsXA");
            XAConnection connection = qcf.createXAConnection();
            connection.start();
            logger.debug("creating connection at " + new java.util.Date());

            //setting up session
            XASession session = connection.createXASession();
            logger.debug("creating session at " + new java.util.Date());

            //find the queue
            Object queueObj = iniCtx.lookup(queue);
            Queue jmsQueue = (javax.jms.Queue)queueObj;

            //adding message to queue
            javax.jms.MessageProducer producer = session.createProducer(jmsQueue);
            javax.jms.TextMessage textMessage = session.createTextMessage(message);
            producer.send(textMessage);
            producer.close();
            logger.debug("Message added to queue");

            //receiving message from queue
            javax.jms.MessageConsumer consumer = session.createConsumer(jmsQueue);
            javax.jms.TextMessage messageReceived = (javax.jms.TextMessage)consumer.receive(5000);

            if (messageReceived==null)
                throw new Exception("No message reveived");

            logger.debug("Got message:"+messageReceived.getText());
            consumer.close();
        }
        catch(Exception e) {
            logger.debug("Error: " + e.getMessage(), e);
        }
    }

    @PreDestroy
    public void shutdown() {

    }
}

结果如下:

11:47:17,905 DEBUG [QueueTest] (MSC service thread 1-8) creating connection at Thu Sep 05 11:47:17 CEST 2013
11:47:18,041 DEBUG [QueueTest] (MSC service thread 1-8) creating session at Thu Sep 05 11:47:18 CEST 2013
11:47:18,065 DEBUG [QueueTest] (MSC service thread 1-8) Message added to queue
11:47:23,081 DEBUG [QueueTest] (MSC service thread 1-8) Error: No message reveived

如您所见,消费者没有收到任何消息。为什么?

编辑 1:

package dk.energimidt.uapi.zigbee.services;

import org.apache.log4j.Logger;

import javax.ejb.Stateless;
import javax.ejb.TransactionAttribute;
import javax.ejb.TransactionAttributeType;
import javax.jms.Queue;
import javax.jms.XAConnection;
import javax.jms.XAConnectionFactory;
import javax.jms.XASession;
import javax.naming.InitialContext;

@TransactionAttribute(TransactionAttributeType.REQUIRED)
@Stateless
public class QueueTestWorkerBean implements QueueTestWorker {

    private Logger logger = Logger.getLogger(QueueTestWorkerBean.class);

    public void run() {
        try {
            String queue = "queue/Queue1";
            String message = "test";

            //setting up connection
            InitialContext iniCtx = new InitialContext();
            XAConnectionFactory qcf = (XAConnectionFactory) iniCtx.lookup("java:/JmsXA");
            XAConnection connection = qcf.createXAConnection();
            connection.start();
            logger.debug("creating connection at " + new java.util.Date());

            //setting up session
            XASession session = connection.createXASession();
            logger.debug("creating session at " + new java.util.Date());

            //find the queue
            Object queueObj = iniCtx.lookup(queue);
            Queue jmsQueue = (javax.jms.Queue)queueObj;

            //adding message to queue
            javax.jms.MessageProducer producer = session.createProducer(jmsQueue);
            javax.jms.TextMessage textMessage = session.createTextMessage(message);
            producer.send(textMessage);
            producer.close();
            session.commit();
            logger.debug("Message added to queue");

            //receiving message from queue
            javax.jms.MessageConsumer consumer = session.createConsumer(jmsQueue);
            javax.jms.TextMessage messageReceived = (javax.jms.TextMessage)consumer.receive(5000);

            if (messageReceived==null)
                throw new Exception("No message reveived");

            logger.debug("Got message:"+messageReceived.getText());
            consumer.close();

            connection.close();
        }
        catch(Exception e) {
            logger.debug("Error: " + e.getMessage(), e);
        }
    }
}

现在我在 Session.Commit() 上遇到异常:

10:46:03,697 DEBUG [QueueTestWorkerBean] (MSC service thread 1-14) creating connection at Tue Sep 17 10:46:03 CEST 2013
10:46:04,343 DEBUG [QueueTestWorkerBean] (MSC service thread 1-14) creating session at Tue Sep 17 10:46:04 CEST 2013
10:46:04,355 DEBUG [QueueTestWorkerBean] (MSC service thread 1-14) Error: XA connection: javax.jms.TransactionInProgressException: XA connection
    at org.hornetq.ra.HornetQRASession.commit(HornetQRASession.java:386)
    at QueueTestWorkerBean.run(QueueTestWorkerBean.java:45) [library-1.0.0.jar:]

【问题讨论】:

    标签: jakarta-ee jboss jms hornetq xa


    【解决方案1】:

    在您真正收到消息之前,您需要提交会话对象。在 producer.send 语句之后,需要添加 session.commit。

    另外我建议最后关闭生产者。

    另一件看起来不对的事情是你在生产者被销毁之后创建了消费者。

    【讨论】:

    • 感谢您的回复。在 producer.send 之后的行中添加提交,会导致 XA 连接:javax.jms.TransactionInProgressException
    • 我已经撤销了你的赏金,@Dennis。您可以根据需要重新提供并奖励它。这给了both另一个答案。
    【解决方案2】:

    我在那里看到了一些(实际上是 2 个)混淆:

    i - 您正在使用 XA 会话,但您没有声明任何事务边界……这通常在会话 Bean 和 MDB 上完成。我不确定你是否可以在这个无状态下做到这一点。

    如果您不使用任何声明性事务,则必须手动登记 XID。

    ii - jmsXA 是默认的资源适配器连接工厂。它上面已经有一个游泳池了。因此,每当您创建一个新会话时,您都会从池中取出。当您关闭它时,您会将其返回到池中。

    您可以使用常规的连接工厂。就在 InVMConnectionFactory 中(或者你在独立设备上定义的任何东西,假设你在 JBoss 上,在 PooledConnectionFactories 之外......然后只使用常规 JMS。

    即使是常规的连接工厂也可以与 XA 一起使用,但在这种情况下,您需要确保直接使用事务管理器的 api 来获取它。

    如果您使用常规连接工厂,则可以根据需要保持连接。

    请告诉我进展如何,我会帮助你。我知道你开始了赏金..但我会免费回答:)

    我在 EJB 教程中找不到任何关于将事务与单例一起使用的示例。

    我建议您通过 Statless 或 Stateful Session Bean 使用它,然后将 @TransactionAttribute 应用于 Bean。

    Java EE 6 教程有一些很好的信息:

    http://docs.oracle.com/javaee/6/tutorial/doc/bncij.html

    请注意,在您提交之前,一条消息将不可用。因此,如果您在交易中发送消息,您将无法在同一交易中接收它。

    在您的 edit1 示例中,您正在发送一条消息并在同一事务中使用它。这是行不通的,因为您首先需要提交生产方法,然后才能使用它。在这种情况下,您需要两个事务,因此 Edit1 已损坏。

    另外:确保最后关闭连接。由于您使用的是 JmsXA(或池连接工厂),因此应用程序服务器将自动完成轮询。

    【讨论】:

    • 感谢您的回复。我会试着看看它。你能给我一个第 1 部分中描述的解决方案的小例子吗?
    • 我没有任何用于 Single 类的示例。你应该在 EE6 上寻找事务注释
    • 非常抱歉,我刚刚将赏金放在了错误的答案上。 :( 有什么办法可以给你奖励吗?
    • 我用 EDIT1 编辑了主题。我已将代码转换为 bean 并添加了 @TransactionAttribute。现在我在 Session.Commit() 中得到一个 TransactionInProgressException
    猜你喜欢
    • 1970-01-01
    • 2012-10-29
    • 1970-01-01
    • 1970-01-01
    • 2018-10-21
    • 1970-01-01
    • 2012-07-22
    • 1970-01-01
    • 2013-01-25
    相关资源
    最近更新 更多