You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

排查与修复建议

  1. 修正连接工厂配置覆盖问题
    原connectionFactory()方法中,先设置了virtualHost等参数,随后调用setUri()会覆盖之前的配置。调整代码逻辑,要么统一使用URI配置,要么删除setUri()调用,保留单独设置参数的方式:

    // 移除setUri调用,保留单独参数设置
    // connectionFactory.setUri(connectionString);
    

    或者确保URI拼接正确,且调用setUri()后不再重复设置参数。

  2. 验证虚拟主机与交换器配置

    • 确认${listener.exchange}配置值是否为tng.ors.service-fed,且该交换器存在于目标虚拟主机(如filing)而非默认的/。
    • 通过RabbitMQ管理界面检查应用连接的虚拟主机,确保与配置文件中的spring.rabbitmq.virtual-host一致。
    • 若交换器未自动创建,手动在目标虚拟主机中创建tng.ors.service-fed交换器,或检查RabbitMQ用户是否有交换器声明权限。
  3. 解决交付确认超时
    在监听器容器工厂中添加确认超时配置,匹配消息处理的实际耗时:

    factory.setAckTimeout(60000); // 设置为60秒,根据实际情况调整
    

    同时排查消息处理逻辑中的慢操作(如MongoDB查询/更新),优化数据库操作性能,减少单条消息处理时间。

  4. 检查配置属性加载
    确认Spring Boot 3.0中配置属性的加载逻辑未发生变化,${listener.exchange}、${listener.queue}等属性是否正确注入。可以通过添加日志打印配置值,验证属性是否正确加载。

内容的提问来源于stack exchange,提问作者sekhar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 20:19:49