Spring Boot 3.0.6升级后RabbitMQ监听器无法正常工作求助
Spring Boot 3.0.6升级后RabbitMQ监听器无法正常消费问题
问题描述
将组件从Spring Boot 2.6.6升级至3.0.6后,RabbitMQ监听器无法正常作为消费者工作,但在Spring Boot 2.6.6环境下可正常运行。应用自身无相关异常日志输出,但RabbitMQ服务端日志出现以下错误。
RabbitMQ配置代码
@Configuration public class RabbitConfig implements RabbitListenerConfigurer { private static final Logger LOGGER = LogManager.getLogger(); @Value("${spring.rabbitmq.host:localhost}") private String host; @Value("${spring.rabbitmq.virtual-host:#{null}}") private String virtualHost; @Value("${spring.rabbitmq.port:5672}") private int port; @Value("${csn.traceability.component_name}") private String connectionName; @Value("${spring.rabbitmq.username}") private String username; @Value("${spring.rabbitmq.password:#{null}}") private String password; @Value("${spring.rabbitmq.passwordFile:#{null}}") private String passwordFile; @Value("${spring.rabbitmq.useSSL:false}") private boolean useSSL; @Value("${spring.rabbitmq.sslAlgorithm:TLSv1.1}") private String sslAlgorithm; @Bean MessageConverter messageConverter(ObjectMapper objectMapper) { return new Jackson2JsonMessageConverter(objectMapper, "*"); } @Bean public CachingConnectionFactory connectionFactory() { CachingConnectionFactory connectionFactory = new CachingConnectionFactory(host); connectionFactory.setUsername(username); connectionFactory.setPassword(password); connectionFactory.setPort(port); connectionFactory.setConnectionNameStrategy(cf -> connectionName); String amqpProtocol = useSSL ? "amqps" : "amqp"; if (LOGGER.isInfoEnabled()) { LOGGER.info("RabbitConfig values [{}] [{}] [{}] [{}] [{}]", amqpProtocol, username, host, port, virtualHost); } String connectionString = String.format("%s://%s:%s@%s:%s", amqpProtocol, username, password, host, port); if (StringUtils.isNotBlank(virtualHost)) { connectionString = connectionString + String.format("/%s", virtualHost); connectionFactory.setVirtualHost(virtualHost); } connectionFactory.setUri(connectionString); return connectionFactory; } @Bean(name = "rabbitListenerContainerFactory") public SimpleRabbitListenerContainerFactory listenerFactory(MessageConverter messageConverter, CachingConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory() { @Override protected void initializeContainer(SimpleMessageListenerContainer instance, RabbitListenerEndpoint endpoint) { super.initializeContainer(instance, endpoint); // set value appropriately for what the consumer does and how long they typically take to process (depends on prefetch count also) instance.setShutdownTimeout(30000); } }; factory.setConcurrentConsumers(1); factory.setPrefetchCount(10); factory.setBatchSize(1); factory.setConnectionFactory(connectionFactory); factory.setMessageConverter(messageConverter); return factory; } /** * Method to configure rabbit listener */ @Override public void configureRabbitListeners(RabbitListenerEndpointRegistrar registrar) { registrar.setMessageHandlerMethodFactory(validatingHandlerMethodFactory()); } @Bean DefaultMessageHandlerMethodFactory validatingHandlerMethodFactory() { DefaultMessageHandlerMethodFactory factory = new DefaultMessageHandlerMethodFactory(); factory.setValidator(amqpValidator()); return factory; } @Bean Validator amqpValidator() { return new OptionalValidatorFactoryBean(); } @Bean public RabbitAdmin rabbitAdmin(CachingConnectionFactory cachingConnectionFactory) { return new RabbitAdmin(cachingConnectionFactory); } /** * Static Bean for Property Configuration Place Holder * @return Static Property Source Place Holder */ @Bean public static PropertySourcesPlaceholderConfigurer propertyConfigInDev() { return new PropertySourcesPlaceholderConfigurer(); } @PostConstruct private void readPropertiesFromFiles() throws IOException { if (password == null && passwordFile != null) { byte[] encoded = Files.readAllBytes(Paths.get(passwordFile)); password = new String(encoded, StandardCharsets.UTF_8); } } }
监听器代码
@Component @EnableMongoRepositories("com.csn.filing.ors.update.repository") public class OrsUpdateListener { private static final Logger logger = LoggerFactory.getLogger(OrsUpdateListener.class); private static final String LISTENER_ID = "print-ors-xfr-ack"; private static final String PAYX_PREFIX = "x-payx-"; private final UpdateStatusTracker updateStatusTracker; private OrsPDFAvailabilityTrackerRepository orsPDFTrackerRepository; @Value("${csn.traceability.component_name}") private String componentName; /** * update listener constructor * @param orsPDFTrackerRepository MongoDB Repository Interface */ @Autowired public OrsUpdateListener(OrsPDFAvailabilityTrackerRepository orsPDFTrackerRepository, UpdateStatusTracker updateStatusTracker) { this.orsPDFTrackerRepository = orsPDFTrackerRepository; this.updateStatusTracker = updateStatusTracker; } /** * Incoming ors update messages * @param updateMessage Rabbit ORS Update Message */ @RabbitListener(id = LISTENER_ID, bindings = @QueueBinding( exchange = @Exchange(value = "${listener.exchange}", type = "topic"), value = @Queue(value = "#{'${spring.rabbitmq.virtual-host:}' == 'filing' ? '${listener.queue}-cd' : '${listener.queue}'}", durable = "true"), key = "${listener.routingKey}"), containerFactory= "rabbitListenerContainerFactory") public void receiveMessage(@Valid @Payload OrsUpdateMessage updateMessage, @Headers Map<String, Object> requestHeaders) { boolean success = false; long startTime = 0; // Extract Traceability Information from Rabbit Message Headers setThreadTraceability(requestHeaders); try { // Capture start time and inject into ThreadContext startTime = System.currentTimeMillis(); ThreadContext.put(MarkAttributeEnum.start_time.toString(), String.valueOf(startTime)); // Ensure Service Name is Populated if (ThreadContext.get(MarkAttributeEnum.service_name.toString()) == null) { ThreadContext.put(MarkAttributeEnum.service_name.toString(), "OrsUpdateListener"); } // Mark Request Accepted logger.info(String.format("%s=%s", MarkAttributeEnum.mark, MarkEnum.request_accepted)); // Mark Transaction Start logger.info(String.format("%s=%s", MarkAttributeEnum.mark, MarkEnum.transaction_start)); final OrsPDFAvailabilityTracker orsPDFTracker = MongoTraceability.capture(() -> orsPDFTrackerRepository.findById(updateMessage.getRequestId()).orElse(null)); if (orsPDFTracker != null ) { String orsErrorCode = ""; StatusType orsStatus = StatusType.COMPLETED; if (updateMessage.getErrorCode() == OrsErrorCode.SUCCESS) { /* * If a message is successfully sent, we want to update its status and remove it from the pdfTracker repository. */ updateStatus(orsPDFTracker, orsStatus, orsErrorCode); MongoTraceability.captureNoReturn(() -> orsPDFTrackerRepository.delete(orsPDFTracker)); success = true; } else { logger.error("Received ORS response with error code: {}", updateMessage.getErrorCode()); /* * If a message contains an error code we want to update its status (as failed) and save it back into the database */ orsErrorCode = updateMessage.getErrorCode().toString(); orsStatus = StatusType.FAILED; updateStatus(orsPDFTracker, orsStatus, orsErrorCode); orsPDFTracker.setStatus(OrsStatus.NOT_SENT); MongoTraceability.captureNoReturn(() -> orsPDFTrackerRepository.save(orsPDFTracker)); } } else { logger.error("ORS response invalid requestId: {}", updateMessage.toString()); } } catch (Exception e) { success = false; logger.error("Error Processing ORS Update", e); throw new AmqpRejectAndDontRequeueException(e.getMessage()); } finally { // Retrieve start time and calculate duration long endTime = System.currentTimeMillis(); long duration = endTime - startTime; // Mark and Measure Transaction End logger.info(String.format("%s=%s,%s=%s,%s=%d", MarkAttributeEnum.status.toString(), success ? "PASS": "FAIL", MarkAttributeEnum.mark.toString(), MarkEnum.transaction_end, MarkAttributeEnum.duration.toString(), duration)); // Clear ThreadContext ThreadContext.clearAll(); } } /** * Method responsible in updating ORS status on the taxfilingstatus records * @param orsPDFTracker * @param orsStatus status that need to updated on trackers * @param orsErrorCode error code for ors processing */ private void updateStatus(OrsPDFAvailabilityTracker orsPDFTracker, StatusType orsStatus, String orsErrorCode){ ReportType reportType = getReportType(orsPDFTracker.getState(), orsPDFTracker.getReportType()); StateType miscState = null; // If report doesn't exits for the state and returntype, then it could be a misc report by ReportType if(reportType == null){ reportType = getReportType(null, orsPDFTracker.getReportType()); miscState = orsPDFTracker.getState(); } if(StringUtils.isBlank(orsPDFTracker.getInternalEmployeeNumber())){ updateStatusTracker.updateOrsInfo(orsPDFTracker.getClientId(), orsPDFTracker.getBranchNumber(), reportType, orsPDFTracker.getExtractId(), orsPDFTracker.getYear(), orsPDFTracker.getQuarter(), orsPDFTracker.getLocalCode(), orsStatus, orsErrorCode, miscState); }else { updateStatusTracker.updateEmployeeOrsInfo(orsPDFTracker.getClientId(), orsPDFTracker.getBranchNumber(), orsPDFTracker.getInternalEmployeeNumber(), reportType, orsPDFTracker.getExtractId(), orsPDFTracker.getYear(), orsPDFTracker.getQuarter(), orsStatus, orsErrorCode); } } /** * Initiate Thread Context with Message Traceability Information * * @param rabbitHeaders * Rabbit Message Headers */ private void setThreadTraceability(Map<String, Object> rabbitHeaders) { // Clear Any Possible Traceability Remnants ThreadContext.clearAll(); if (rabbitHeaders != null) { // Store payx headers in ThreadContext rabbitHeaders.keySet().stream().filter(h -> h.toLowerCase().startsWith(PAYX_PREFIX)) .forEach(h -> ThreadContext.put(h.toLowerCase(), rabbitHeaders.get(h).toString())); // If transaction id was not received, create one along with a // transaction-start mark if (!ThreadContext.containsKey(MarkAttributeEnum.transaction_id.toString())) { ThreadContext.put(MarkAttributeEnum.transaction_id.toString(), UUID.randomUUID().toString()); ThreadContext.put(MarkAttributeEnum.transaction_unknown.toString(), "true"); ThreadContext.put(MarkAttributeEnum.business_process_name.toString(), "tax_print"); ThreadContext.put(MarkAttributeEnum.session_id.toString(), "0"); if (!ThreadContext.containsKey(MarkAttributeEnum.user.toString())) { ThreadContext.put(MarkAttributeEnum.user.toString(), "unk"); } } // Add Component Name ThreadContext.put(MarkAttributeEnum.component_name.toString(), componentName); } } }
RabbitMQ错误日志
- 多次报错:
operation basic.publish caused a channel exception not_found: no exchange 'tng.ors.service-fed' in vhost '/'(虚拟主机'/'中不存在交换器'tng.ors.service-fed') - 出现消费者交付确认超时警告及通道异常:
Consumer 51 on channel 1 has timed out waiting for delivery acknowledgement、delivery acknowledgement on channel 1 timed out
排查与修复建议
修正连接工厂配置覆盖问题
原connectionFactory()方法中,先设置了virtualHost等参数,随后调用setUri()会覆盖之前的配置。调整代码逻辑,要么统一使用URI配置,要么删除setUri()调用,保留单独设置参数的方式:// 移除setUri调用,保留单独参数设置 // connectionFactory.setUri(connectionString);或者确保URI拼接正确,且调用
setUri()后不再重复设置参数。验证虚拟主机与交换器配置
- 确认
${listener.exchange}配置值是否为tng.ors.service-fed,且该交换器存在于目标虚拟主机(如filing)而非默认的/。 - 通过RabbitMQ管理界面检查应用连接的虚拟主机,确保与配置文件中的
spring.rabbitmq.virtual-host一致。 - 若交换器未自动创建,手动在目标虚拟主机中创建
tng.ors.service-fed交换器,或检查RabbitMQ用户是否有交换器声明权限。
- 确认
解决交付确认超时
在监听器容器工厂中添加确认超时配置,匹配消息处理的实际耗时:factory.setAckTimeout(60000); // 设置为60秒,根据实际情况调整同时排查消息处理逻辑中的慢操作(如MongoDB查询/更新),优化数据库操作性能,减少单条消息处理时间。
检查配置属性加载
确认Spring Boot 3.0中配置属性的加载逻辑未发生变化,${listener.exchange}、${listener.queue}等属性是否正确注入。可以通过添加日志打印配置值,验证属性是否正确加载。
内容的提问来源于stack exchange,提问作者sekhar
相关产品推荐
相关产品推荐

