【问题标题】:Error while Creating JMS Message Producer in jboss EAP 7在 jboss EAP 7 中创建 JMS 消息生产者时出错
【发布时间】:2019-12-21 04:01:18
【问题描述】:

我已经在 J​​BOSS_EAP_7.0 中配置了 JMS 主题,并编写了一个简单的 java 代码来创建一个消息生产者。我有以下无状态bean

@Stateless
public class ExchangeSenderFacadeWrapperBean {


    private static final OMSLogHandlerI logger = new Log4j2Handler("ClientSenderFacadeBean");
    @Resource(lookup = "java:/JmsXA")     // inject ConnectionFactory (more)
    protected ConnectionFactory  factory;


    @Resource(lookup = "java:/jms/topic/ORD_CLINT_PUSH")
    protected Topic target;

    private Connection  connection = null;
    private Session session = null;



    public void sendMessage(String message) {

        MessageProducer producer= null;
        try {
            if(connection==null){  //todo verify
                connection = factory.createConnection();
            }
            session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
            producer = session.createProducer(target);
            producer.setDisableMessageID(true);
            TextMessage outmsg = session.createTextMessage(message);
            producer.send(outmsg);
            logger.info("Message was sent to Topic");
            producer.setTimeToLive(900000);//15min  //todo
        } catch (Exception e) {
            logger.error(" Error when sending order to jboss:", e);
            throw new OMSCoreRuntimeException(e.getMessage(), e);
        } finally {
            try {
                if (producer != null)
                    producer.close();
            } catch (JMSException e) {
                logger.warn("\n jms producer close error:",e);
            }
            try {
                if (session != null)
                    session.close();
            } catch (JMSException e) {
                logger.warn("\n jms session close error:",e);
            }
        }
    }

这工作正常,直到我进行简单的更改以将 sendMessage(String message) 方法移动到 pojo 类,如下所示。

@Stateless(name = "ExchangeSenderFacadeBean")
@Local({ExchangeSenderFacadeLocalI.class})
public class ExchangeSenderFacadeWrapperBean implements ExchangeSenderFacadeLocalI {
    @Resource(lookup = "java:/JmsXA")     // inject ConnectionFactory (more)
    protected ConnectionFactory factory;

    @EJB(beanName = "BeanRegistryLoader")
    protected BeanRegistryLoader omsRegistryBean;

    protected BeanRegistryCore beanRegistryCore;

    @Resource(lookup = "java:/jms/queue/ToExchange")
    protected Queue target;

    private ExchangeSenderFacadeCoreI exchangeSenderFacadeCore;


    @Override
    public void sendToExchange(ExchangeMessage exchangeMessage) {
        exchangeSenderFacadeCore.sendToExchange(exchangeMessage);

    }

    @PostConstruct
    public void init() {
        beanRegistryCore = omsRegistryBean.registry();
        if (exchangeSenderFacadeCore == null) {
            exchangeSenderFacadeCore = ((BeanRegistryCore) omsRegistryBean.registry()).getExchangeSenderFacadeCoreI();
            exchangeSenderFacadeCore.setBeanRegistryCore(omsRegistryBean.registry());
            exchangeSenderFacadeCore.setFactory(factory);
            exchangeSenderFacadeCore.setTargetQueue(target);
        }
    }

}

ConnectionFactory 和目标 Queue 在 EJB 中设置的变量 PostConstruct 方法和 pojo 类如下所示,现在包含创建和发布方法到 EJB 队列的逻辑

public class ExchangeSenderFacadeCore implements ExchangeSenderFacadeCoreI {
    private static final OMSLogHandlerI logger = new Log4j2HndlAdaptor("ExchangeSenderFacadeCore");
    private BeanRegistryCore beanRegistryCore;
    private ConnectionFactory factory;
    private Connection connection = null;
    private Session session = null;
    private long ttl = 900000;
    protected Queue targetQueue;

    public ExchangeSenderFacadeCore() {
        if (System.getProperty(OMSConst.SYS_PROPERTY_JMS_TTL) != null && System.getProperty(OMSConst.SYS_PROPERTY_JMS_TTL).length() > 0) {
            ttl = Long.parseLong(System.getProperty(OMSConst.SYS_PROPERTY_JMS_TTL));
        }
        logger.info("LN:103", "==JMS Topic TTL:" + ttl);
    }

    @Override
    public void processSendToExchange(ExchangeMessage exchangeMessage) {
        sendToExchange(exchangeMessage);
    }

    public boolean isParallelRunEnabled() {
        Object isParallelRun = beanRegistryCore.getCacheAdaptorI().cacheGet(OMSConst.DEFAULT_TENANCY_CODE, OMSConst.APP_PARAM_IS_PARALLEL_RUN, CACHE_NAMES.SYS_PARAMS_CACHE_CORE);
        if (isParallelRun != null && String.valueOf(isParallelRun).equals(OMSConst.STRING_1)) {
            return true;
        }
        return false;
    }

    @Override
    public void sendToExchange(ExchangeMessage exchangeMessage) {
        MessageProducer producer = null;
        try {
            if (isParallelRunEnabled()) {
                logger.info("LN:66", "== Message send to exchange skipped,due to parallel run enabled");
                return;
            }
            if (connection == null) {
                connection = factory.createConnection();
            }
            session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
            producer = session.createProducer(targetQueue);
            producer.setDisableMessageID(true);
            Message message = beanRegistryCore.getJmsExchangeMsgTransformerI().transformToJMSMessage(session, exchangeMessage);
            producer.send(message);
            producer.setTimeToLive(ttl);//default 15min
            logger.elkLog("78", "-1", LogEventsEnum.SENT_TO_EXCHANGE, exchangeMessage.toString());
        } catch (Exception e) {
            logger.error("LN:80", " Error when sending order to exchange:", e);
            throw new OMSCoreRuntimeException(e.getMessage(), e);
        } finally {
            try {
                if (producer != null)
                    producer.close();
            } catch (JMSException e) {
                logger.error("LN:87", "JMS producer close error:", e);
            }
            try {
                if (session != null)
                    session.close();
            } catch (JMSException e) {
                logger.error("LN:93", "JMS session close error:", e);
            }
        }
    }

    @Override
    public void processSendToExchangeSync(ExchangeMessage exchangeMessage) {

    }

    @Override
    public BeanRegistryCore getBeanRegistryCore() {
        return beanRegistryCore;
    }

    @Override
    public void setBeanRegistryCore(BeanRegistryCore beanRegistryCore) {
        this.beanRegistryCore = beanRegistryCore;
    }

    @Override
    public ConnectionFactory getFactory() {
        return factory;
    }

    @Override
    public void setFactory(ConnectionFactory factory) {
        this.factory = factory;
    }

    @Override
    public Queue getTargetQueue() {
        return targetQueue;
    }

    @Override
    public void setTargetQueue(Queue targetQueue) {
        this.targetQueue = targetQueue;
    }
}

但是当我执行审核代码时,它给了我以下错误

javax.ejb.EJBTransactionRolledbackException:生产者已关闭

任何可能的修复?

【问题讨论】:

  • 您能否提供EJBTransactionRolledbackException 的完整堆栈跟踪?
  • 另外,没有必要缓存connection,因为您是从池化的JmsXA 连接工厂获取的。您每次都可以简单地创建/关闭连接。当然,如果连接工厂没有被池化,这将是一种反模式。
  • @JustinBertram 为了支持你的论点,我也得到了以下异常。 javax.ejb.EJBTransactionRolledbackException: Only allowed one session per connection 当我每次创建/关闭连接时,正如你提到的,我没有遇到异常并且工作正常,但问题是即使缓存连接消息生产者在场景 1 中也能正常工作。自从我移动 sendMessage() 后它就失败了POJO 类的方法,如场景 2。此场景背后的任何逻辑??
  • 如果它在您每次创建/关闭连接时都有效,那么我会这样做。这简化了您的代码并允许池执行其设计的任务。当您缓存来自池的连接时,尤其是在像 Java EE 这样的环境中,其中有关于允许您使用连接做什么的特定规则(例如,每个连接只有一个会话),那么奇怪的事情可能会发生。如果没有一个可以实际运行以查看发生了什么的示例,我真的无法提供更多。
  • @JustinBertram 感谢您在这种情况下提供的所有帮助,我做了一些研究并得到了以下答案,请您验证一下吗?这对我有很大的帮助。

标签: java jboss jms stateless-session-bean jms-queue


【解决方案1】:

在对问题进行了深入搜索后,我发现这篇 https://developer.jboss.org/wiki/ShouldICacheJMSConnectionsAndJMSSessions 文章发布在 JBOSS 开发者线程之一上。这清楚地解释了缓存连接和其他 JMS 相关资源作为 JMS 代码反模式的原因,该代码在 JEE 应用程序服务器中运行。

简而言之,JCA 层池化 JMS 连接和 JMS 会话。因此,当您调用 createConnection() 或 createSession() 时,在大多数情况下,它并没有真正调用实际的 JMS 实现来实际创建新的 JMS 连接或 JMS 会话,它只是从其自己的内部缓存中返回一个。

此外,JBOSS 服务器也管理无状态会话 bean 池。无状态会话 bean 仅在您完成其目的后才可在连接池上使用,而不是事先。同时连接(新创建的 JMS 或缓存的)用于在无状态会话 bean 中创建 JMS 会话(session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE)),也完成了它的用途,并且在 JCA 层连接池上也可用。因此,在无状态 EJB 类中调用缓存连接不会出现异常,即使 Oracle 不建议这样做。

public void sendToExchange(ExchangeMessage exchangeMessage) {
        MessageProducer producer = null;
        try {
            if (isParallelRunEnabled()) {
                logger.info("LN:66", "== Message send to exchange skipped,due to parallel run enabled");
                return;
            }
            if (connection == null) {
                connection = factory.createConnection();
            }
            session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
            producer = session.createProducer(targetQueue);
            producer.setDisableMessageID(true);
            Message message = beanRegistryCore.getJmsExchangeMsgTransformerI().transformToJMSMessage(session, exchangeMessage);
            producer.send(message);
            producer.setTimeToLive(ttl);//default 15min
            logger.elkLog("78", "-1", LogEventsEnum.SENT_TO_EXCHANGE, exchangeMessage.toString());
        } catch (Exception e) {
            logger.error("LN:80", " Error when sending order to exchange:", e);
            throw new OMSCoreRuntimeException(e.getMessage(), e);
        } finally {
            try {
                if (producer != null)
                    producer.close();
            } catch (JMSException e) {
                logger.error("LN:87", "JMS producer close error:", e);
            }
            try {
                if (session != null)
                    session.close();
            } catch (JMSException e) {
                logger.error("LN:93", "JMS session close error:", e);
            }
        }
    }

但是在这种情况下,由于同一个 POJO 类实例可以在多个场合使用,如下所示。它不保证连接在JCA层连接池中被释放和可用并给出异常。

@PostConstruct
    public void init() {
        beanRegistryCore = omsRegistryBean.registry();
        if (exchangeSenderFacadeCore == null) {
            exchangeSenderFacadeCore = ((BeanRegistryCore) omsRegistryBean.registry()).getExchangeSenderFacadeCoreI();
            exchangeSenderFacadeCore.setBeanRegistryCore(omsRegistryBean.registry());
            exchangeSenderFacadeCore.setFactory(factory);
            exchangeSenderFacadeCore.setTargetQueue(target);
        }
    }

【讨论】:

    猜你喜欢
    • 2019-12-14
    • 2020-03-09
    • 2019-09-06
    • 1970-01-01
    • 2023-03-17
    • 1970-01-01
    • 1970-01-01
    • 2017-02-22
    • 1970-01-01
    相关资源
    最近更新 更多