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