【发布时间】:2016-11-24 11:40:34
【问题描述】:
一旦我使用来自rabbitmq 的消息,我已经实现了多线程以在jira 中执行一些操作。我正在使用spring amqp(1.6.1版)
一旦线程捕获到异常,我会将状态设置为错误并输出我将来引用的对象。在将此输出对象发送到队列时。我正面临上述执行
代码:
连接工厂:
@Configuration
@PropertySources({ @PropertySource("classpath:application.properties") })
public class RabbitMQConfiguration {
@Autowired
private Environment environment;
@Bean
public ConnectionFactory connectionFactory() {
// TODO make it possible to customize in subclasses.
CachingConnectionFactory connectionFactory = new CachingConnectionFactory(environment.getProperty("bip.rabbitmq.url"));
connectionFactory.setUsername(environment.getProperty("bip.rabbitmq.username"));
connectionFactory.setPassword(environment.getProperty("bip.rabbitmq.password"));
return connectionFactory;
}
@Bean
public MessageConverter jsonMessageConverter() {
return new Jackson2JsonMessageConverter();
}
/**
* @return the admin bean that can declare queues etc.
*/
@Bean
public AmqpAdmin amqpAdmin() {
RabbitAdmin rabbitAdmin = new RabbitAdmin(connectionFactory());
return rabbitAdmin;
}
@Bean
public RabbitTemplate rabbitTemplate() {
RabbitTemplate template = new RabbitTemplate(connectionFactory());
template.setMessageConverter(jsonMessageConverter());
return template;
}
@Bean(name = "jiraQueueListenerContainerFactory")
public SimpleRabbitListenerContainerFactory jiraQueueListenerContainerFactory() {
SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
factory.setConnectionFactory(connectionFactory());
factory.setMessageConverter(new Jackson2JsonMessageConverter());
factory.setReceiveTimeout(10L);
return factory;
}
}
消息处理程序:
@Component
@PropertySources({ @PropertySource("classpath:application.properties") })
public class JiraMessageHandler {
@Autowired
private RabbitTemplate rabbitTemplate;
private static ExecutorService executor = Executors.newFixedThreadPool(Constants.THREAD_SIZE);
private static Logger logger = LogManager.getLogger(JiraMessageHandler.class);
@RabbitListener(containerFactory = "jiraQueueListenerContainerFactory", queues = Constants.QUEUE_NAME)
public void handleMessage(HashMap<String, Object> jiraMessage) {
logger.info(jiraMessage.toString() + System.currentTimeMillis());
logger.info("Jira Message Handler");
BaseServerAdapter jiraProcessingAdapter = new JiraProcessingAdapter();
Future future = executor.submit(jiraProcessingAdapter);
JiraAdapterOutput jiraAdapterOutput = new JiraAdapterOutput();
Future future = executor.submit(jiraProcessingAdapter);
jiraAdapterOutput = (JiraAdapterOutput) future.get();
try {
jiraAdapterOutput = (JiraAdapterOutput) future.get();
if (jiraAdapterOutput.getOutputMap().get("activityStatus") == "SUCCESS") {
logger.info("Successfully Executed Jira ::: " + new Date() + "::: "
+ jiraAdapterOutput.getOutputMap().get("jiraId"));
rabbitTemplate.convertAndSend(Constants.ADAPTER_OUTPUT_QUEUE, senderMap);
}else if (jiraAdapterOutput.getOutputMap().get("activityStatus").equalsIgnoreCase("FAIL")) {
logger.info("Successfully Executed Jira ::: " + new Date() + "::: "
+ jiraAdapterOutput.getOutputMap().get("jiraId"));
sendMessageForProcessingToBIP(senderMap);
}
private boolean sendMessageForProcessingToBIP(HashMap<String, ExchangeDTO> senderMap) {
try {
rabbitTemplate.convertAndSend(Constants.WFM_ERROR_QUEUE, senderMap);
return true;
} catch (Exception e) {
**logger.info("Message sending failed, try again:::::::" + e.getMessage());**
}
return false;
}
它显示“应用程序上下文已关闭,ConnectionFactory 无法再创建连接。”
我做错了什么。 我在这里也提到了:https://jira.spring.io/browse/AMQP-546
【问题讨论】: