【问题标题】:WildFly JMS issue: "Consumer is closed"WildFly JMS 问题:“消费者已关闭”
【发布时间】:2021-06-19 22:54:03
【问题描述】:

我在 WildFly 19 中部署了一个应用程序和一个用于在队列中发送/接收消息的独立 Java 应用程序。此应用程序运行良好,但是当我尝试进行负载测试时,我从 Java 应用程序中收到如下所示的“消费者已关闭”异常。之后,我的 Java 应用程序无法从队列中发送和接收任何消息。

Caused by: org.apache.activemq.artemis.api.core.ActiveMQObjectClosedException: AMQ119017: Consumer is closed
        ... 11 common frames omitted
2021-03-18 08:46:16,949 [pool-2-thread-1] ERROR c.v.d.h.c.dip.jms.ResponseConsumer - Error while waiting for Response from Queue
javax.jms.IllegalStateException: AMQ119017: Consumer is closed
        at org.apache.activemq.artemis.core.client.impl.ClientConsumerImpl.checkClosed(ClientConsumerImpl.java:952) [artemis-core-client-2.6.3.jbossorg-00014.jar:2.6.3.jbossorg-00014]
        at org.apache.activemq.artemis.core.client.impl.ClientConsumerImpl.receive(ClientConsumerImpl.java:195) [artemis-core-client-2.6.3.jbossorg-00014.jar:2.6.3.jbossorg-00014]
        at org.apache.activemq.artemis.core.client.impl.ClientConsumerImpl.receive(ClientConsumerImpl.java:379) [artemis-core-client-2.6.3.jbossorg-00014.jar:2.6.3.jbossorg-00014]
        at org.apache.activemq.artemis.jms.client.ActiveMQMessageConsumer.getMessage(ActiveMQMessageConsumer.java:211) [artemis-jms-client-2.6.3.jbossorg-00014.jar:2.6.3.jbossorg-00014]
        at org.apache.activemq.artemis.jms.client.ActiveMQMessageConsumer.receive(ActiveMQMessageConsumer.java:132) [artemis-jms-client-2.6.3.jbossorg-00014.jar:2.6.3.jbossorg-00014]
        at com.verizon.delphi.hyperion.core.dip.jms.ResponseConsumer$1.run(ResponseConsumer.java:139) [classes/:na]
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) [na:1.8.0_112]
        at java.util.concurrent.FutureTask.run(FutureTask.java:266) [na:1.8.0_112]
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142) [na:1.8.0_112]
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617) [na:1.8.0_112]
        at java.lang.Thread.run(Thread.java:745) [na:1.8.0_112]
Caused by: org.apache.activemq.artemis.api.core.ActiveMQObjectClosedException: AMQ119017: Consumer is closed​

这是队列配置的standalone-full.xml

<subsystem xmlns="urn:jboss:domain:messaging-activemq:9.0">
    <server name="default">
        <statistics enabled="${wildfly.messaging-activemq.statistics-enabled:${wildfly.statistics-enabled:false}}"/>
        <security-setting name="#">
            <role name="guest" send="true" consume="true" create-non-durable-queue="true" delete-non-durable-queue="true"/>
        </security-setting>
        <address-setting name="#" dead-letter-address="jms.queue.DLQ" expiry-address="jms.queue.ExpiryQueue" max-size-bytes="10485760" page-size-bytes="2097152" message-counter-history-day-limit="10"/>
        <http-connector name="http-connector" socket-binding="http" endpoint="http-acceptor"/>
        <http-connector name="http-connector-throughput" socket-binding="http" endpoint="http-acceptor-throughput">
            <param name="batch-delay" value="50"/>
        </http-connector>
        <in-vm-connector name="in-vm" server-id="0">
            <param name="buffer-pooling" value="false"/>
        </in-vm-connector>
        <http-acceptor name="http-acceptor" http-listener="default"/>
        <http-acceptor name="http-acceptor-throughput" http-listener="default">
            <param name="batch-delay" value="50"/>
            <param name="direct-deliver" value="false"/>
        </http-acceptor>
        <in-vm-acceptor name="in-vm" server-id="0">
            <param name="buffer-pooling" value="false"/>
        </in-vm-acceptor>
        <jms-queue name="ExpiryQueue" entries="java:/jms/queue/ExpiryQueue"/>
        <jms-queue name="DLQ" entries="java:/jms/queue/DLQ"/>
        <jms-queue name="probeRequestQueue" entries="jms/queue/probeRequestQueue java:jboss/exported/jms/queue/probeRequestQueue queue/probeRequestQueue" durable="true"/>
        <jms-queue name="probeResponseQueue" entries="jms/queue/probeResponseQueue java:jboss/exported/jms/queue/probeResponseQueue queue/probeResponseQueue" durable="true"/>
        <jms-topic name="serverStateTopic" entries="java:jboss/exported/topic/serverStateTopic"/>
        <connection-factory name="InVmConnectionFactory" entries="java:/ConnectionFactory" connectors="in-vm"/>        
        <connection-factory name="RemoteConnectionFactory" entries="java:jboss/exported/jms/RemoteConnectionFactory java:/jms/RemoteConnectionFactory" connectors="http-connector" block-on-acknowledge="true" reconnect-attempts="-1" />
        <pooled-connection-factory name="activemq-ra" entries="java:/JmsXA java:jboss/DefaultJMSConnectionFactory" connectors="in-vm" reconnect-attempts="0" transaction="xa"/>
    </server>
</subsystem>​

目前我对这个问题感到震惊并寻求支持。谢谢!

编辑:

生产者类代码:

public class RequestValidator {
    
    private static final Logger L = LoggerFactory.getLogger(RequestValidator.class);
    private ServerStateListener listener;
    static final String REQUEST_QUEUE = "jms/queue/requestQueue";
    private long jmsProducerTimeToLive;
    
    public RequestValidator(ServerStateListener listener, Properties config) {
        this.listener = listener;
        this.jmsProducerTimeToLive = Long.parseLong(config.getProperty("jms.producer.mgs.timetolive", "5000"));
    }
    
    public void messageReceived(RequestDto request, Object callback) throws Exception {
        
        final Connection connection = Util.getConnection();
        
        if (Util.isAlive(connection)) {
            Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
            MessageProducer producer = session.createProducer(Util.getQueue(REQUEST_QUEUE));
            ObjectMessage requestMsg = session.createObjectMessage();
            requestMsg.setObject(request);
            producer.setTimeToLive(jmsProducerTimeToLive);
            producer.send(requestMsg);
            L.info("Request sent.....");
        } else {
            L.error("Unable to acquire Connection to MyApp");
            throw new OssException(OssExceptionType.ExceptFTOSSAppUnavailable, "Unable to acquire connection to MyApp");
        }
    }
}

消费者代码:

public class ResponseConsumer {
    
    private static final Logger L = LoggerFactory.getLogger(ResponseConsumer.class);
    private final ExecutorService es;
    private SimpleResponseTransmitter responseTransmitter;
    private Connection connection;
    
    static final String RESPONSE_QUEUE = "jms/queue/responseQueue";
    
    private RunnableAdapter runner;
    private boolean connected;
    
    public ResponseConsumer(SimpleResponseTransmitter responseTransmitter)       {
        this.responseTransmitter = responseTransmitter;
        this.es = new ThreadPoolExecutor(5, 60, 60L,
                TimeUnit.SECONDS, new LinkedBlockingQueue<Runnable>());
        reconnect();
    }
    
    public final void reconnect() {
        try {
            if (Util.isAlive(connection)) {
                return;
            }
            this.connected = false;
            this.connection = Util.getConnection();
            if (Util.isAlive(connection)) {
                connection.setExceptionListener(new ConsumerExceptionListener());
                L.info("Starting Response Consumer....");
                try {
                    Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
                    
                    final MessageConsumer responseConsumer = session.createConsumer(Util.getQueue(RESPONSE_QUEUE));
                    connection.start();
                    this.runner = new RunnableAdapter() {
                        
                        public void run() {
                            while (isAllowedToRun()) {
                                try {
                                    final ObjectMessage message;
                                    if ((message = (ObjectMessage) responseConsumer.receive(30000)) != null) {
                                        es.submit(new ResponseBuilder(message));
                                    }
                                } catch (JMSException ex) {
                                    L.error("Error while waiting for Response from Queue", ex);
                                }
                            }
                        }
                    };
                    this.connected = true;
                } catch (JMSException ex) {
                    L.error("Error in queue consumer", ex);
                }
                
            } else {
                L.error("Application NOT reachable for reading response queue");
            }
        } catch (JMSException ex) {
            L.error("Error creating JMS connection for queue consumer", ex);
        }
    }
}​

其他:

class Util {
    private static final Logger L = LoggerFactory.getLogger(Util.class);
    private static final Object lock = new Object();
    private static TopicConnection topicConnection;
    private static Connection connection;
    protected static MyApplicationClient client;
    
    private static ConnectionFactory getConnectionFactory() {
    return client.lookup("java:/jms/RemoteConnectionFactory", ConnectionFactory.class);
     } 
     
    private static TopicConnectionFactory  getTopicConnectionFactory() {
    return client.lookup("java:/jms/RemoteConnectionFactory", TopicConnectionFactory.class);
     } 
     
     public static boolean isAlive(Connection connection) {
        try {
            return (connection != null && connection.getMetaData() != null);
        } catch (JMSException ex) {

            return false;
        }
    }
    public static Queue getQueue(String jndiName) {
        return client.lookup(jndiName, Queue.class);
    }

    public static Topic getTopic(String jndiName) {
        return client.lookup(jndiName, Topic.class);
    }
    public static Connection getConnection() throws JMSException {
        if (isAlive(connection)) {
            return connection;
        }
        synchronized (lock) {
            if (client.reconnect()) {
                final ConnectionFactory connectionFactory = getConnectionFactory();
                if (connectionFactory != null) {
                    connection = connectionFactory.createConnection(USERNAME,PASSWORD);
                    connection.start();
                }
            }
        }
        return connection;
    }
}

WildFly 管理控制台 - 请求队列详细信息: WildFly 管理控制台 - 响应队列详细信息:

【问题讨论】:

  • 你能详细说明你的“负载测试”吗?您如何进行负载测试?为什么您的应用程序不处理异常并重新初始化使用者或连接?日志中是否有任何消息表明存在任何问题?你能粘贴你的客户端代码吗?
  • @JustinBertram 负载测试不过是我将 n 个请求发送到队列中,它将从应用服务器端接收并将处理后的响应返回给提交者。我已按照您的要求粘贴了代码。
  • @JustinBertram 如何增加/减少消费者数量和设置每个队列的过期时间?

标签: java jms wildfly activemq-artemis


【解决方案1】:

我认为问题出在消费者的设计上。在reconnect 中,您将在try 块中创建javax.jms.Session 实例和javax.jms.MessageConsumer 实例。然后在RunnableAdapter 实例中使用MessageConsumer 实例。我假设这个RunnableAdapter 实例将在稍后运行。然而,问题在于,当RunnableAdapter 实例稍后实际运行时,用于创建消费者的Session 将超出范围,甚至可能已被垃圾回收。这意味着Consumer 实例将不起作用。您应该重构您的应用程序,以免发生这种情况。

【讨论】:

  • 嗨贾斯汀,我尝试关闭打开的对象,但仍然出现同样的问题。当我检查 Wildfly 管理控制台的队列状态时,发现以下情况。我向 request-q 发送了 95 个请求,其中 47 个请求仍处于处理状态,48 个已处理的响应发送回 response-q。
  • 你是什么意思,“试图关闭打开的对象”?它被关闭的事实的问题。
  • 47 个请求在那里被击中,请求和响应队列计数都没有得到更新
  • 据我所知,您的问题是由您的客户编写方式引起的(我在回答中对此进行了解释)。
  • 我的意思是这些 sessionObj.close() 和 MessageConsumerObj.close() 因为它没有在任何地方关闭
猜你喜欢
  • 2019-07-05
  • 2018-01-28
  • 2017-10-30
  • 1970-01-01
  • 2015-03-26
  • 1970-01-01
  • 2011-06-04
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多