【发布时间】:2021-03-19 07:42:27
【问题描述】:
DefaultMessageListenerContainer.shutdown 或 DefaultMessageListenerContainer.destroy 不会从队列中删除使用者。
这是一个类似的帖子:SpringJMS - How to Disconnect a MessageListenerContainer
(不知道怎么解决)
下面是我的代码:
public class MainProgram {
private static final AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(MessageConsumerFacade.class);
public static final DefaultMessageListenerContainer container = context.getBean(DefaultMessageListenerContainer.class);
public static void main(String[] args) throws InterruptedException {
boolean startListener = isStartListener(); // to start and stop listener at
will
if(startListener){
if (!container.isRunning()) {
container.start();
}
}else{
if (container.isRunning()) {
container.stop();
}
}
}
}
public class MessageConsumerFacade {
private ConnectionFactory connectionFactory() {
ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory();
connectionFactory.setBrokerURL(url);
connectionFactory.setUserName(userName);
connectionFactory.setPassword(password);
RedeliveryPolicy policy = connectionFactory.getRedeliveryPolicy();
policy.setInitialRedeliveryDelay(30000);
policy.setRedeliveryDelay(30000);
policy.setMaximumRedeliveries(2);
connectionFactory.setNonBlockingRedelivery(true);
return connectionFactory;
}
@Bean
public MessageListenerContainer listenerContainer() {
DefaultMessageListenerContainer container = new DefaultMessageListenerContainer();
container.setConnectionFactory(connectionFactory());
container.setDestinationName(queueName);
container.setMessageListener(new MessageJmsListener());
container.setCacheLevel(DefaultMessageListenerContainer.CACHE_NONE);
container.setErrorHandler(new MessageErrorHandler());
container.setSessionTransacted(true);
container.setAutoStartup(false);
container.shutdown();
return container;
}
}
public class MessageJmsListener implements MessageListener {
@Override
public void onMessage(Message message) {
if (message instanceof TextMessage) {
try {
//process the message and create record in Data Base
} catch (Exception e) {
throw new RuntimeException(e);
}
}
}
}
public class MessageErrorHandler implements ErrorHandler {
@Override
public void handleError(Throwable t) {
//log error
}
}```
【问题讨论】:
标签: activemq spring-jms shutdown consumer message-listener