Spring RabbitMQ:网络中断后客户端为何无法自动恢复?
版本信息
org.springframework.amqp:spring-rabbit:2.4.7 openjdk version "11.0.12" 2021-07-20
根据多篇技术文章,RabbitMQ客户端偶尔会与Broker断开连接且无法自动恢复。按照Spring AMQP文档说明,自动恢复功能是AMQP客户端内置的,当前客户端使用SimpleRabbitListenerContainerFactory,已遵循文档建议未设置automaticRecoveryEnabled和topologyRecoveryEnabled。
日志中出现如下异常:
[AMQP Connection 35.82.255.136:5671] ERROR com.rabbitmq.client.impl.ForgivingExceptionHandler - An unexpected connection driver error occurred com.rabbitmq.client.MissedHeartbeatException: Heartbeat missing with heartbeat = 60 seconds at com.rabbitmq.client.impl.AMQConnection.handleSocketTimeout(AMQConnection.java:847) ~[amqp-client-5.13.1.jar:5.13.1] at com.rabbitmq.client.impl.AMQConnection.readFrame(AMQConnection.java:747) [amqp-client-5.13.1.jar:5.13.1] at com.rabbitmq.client.impl.AMQConnection.access$300(AMQConnection.java:47) [amqp-client-5.13.1.jar:5.13.1] at com.rabbitmq.client.impl.AMQConnection$MainLoop.run(AMQConnection.java:666) [amqp-client-5.13.1.jar:5.13.1] at java.lang.Thread.run(Thread.java:829) [?:?]
异常发生后,客户端未执行自动恢复,且在重启服务前,日志中不再出现任何与RabbitMQ客户端相关的内容。
XML配置
<?xml version="1.0" encoding="UTF-8"?> <beans profile="staging,prod" xmlns="http://www.springframework.org/schema/beans" xmlns:context="http://www.springframework.org/schema/context" xmlns:rabbit="http://www.springframework.org/schema/rabbit" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://www.springframework.org/schema/beans https://www.springframework.org/schema/beans/spring-beans.xsd http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd http://www.springframework.org/schema/rabbit https://www.springframework.org/schema/rabbit/spring-rabbit.xsd"> <rabbit:annotation-driven/> <!-- Scan for components with @RabbitListener. --> <context:component-scan base-package="com.predicine.ingress.tasks"> <context:include-filter type="regex" expression="com.acme.tasks.ConsumeBarcodes"/> </context:component-scan> <rabbit:connection-factory id="connectionFactory" connection-factory="clientConnectionFactory"/> <bean id="clientConnectionFactory" class="org.springframework.amqp.rabbit.connection.RabbitConnectionFactoryBean"> <property name="username" value="${rabbitmq.user}" /> <property name="password" value="${rabbitmq.password}" /> <property name="host" value="${rabbitmq.host}" /> <property name="port" value="${rabbitmq.port}" /> <property name="useSSL" value="true" /> <property name="keyStore" value="classpath:/keycert.jks" /> <property name="keyStorePassphrase" value="changeit" /> <property name="keyStoreType" value="PKCS12" /> <property name="trustStore" value="classpath:/trustStore" /> <property name="trustStorePassphrase" value="changeit" /> <property name="trustStoreType" value="PKCS12" /> </bean> <rabbit:admin connection-factory="connectionFactory"/> <rabbit:queue name="${rabbitmq.queue}" durable="true" auto-delete="false"> <rabbit:queue-arguments> <!-- Let messages sit up to 48 hours in queue. Presumably within that timeframe we can bring our consumers back online. --> <entry key="x-message-ttl" value="172800000" value-type="java.lang.Integer"/> </rabbit:queue-arguments> </rabbit:queue> <rabbit:fanout-exchange name="${rabbitmq.exchange}" durable="true" auto-delete="false"> <rabbit:bindings> <rabbit:binding queue="${rabbitmq.queue}" /> </rabbit:bindings> </rabbit:fanout-exchange> <bean id="rmqMessageConverter" class="org.springframework.amqp.support.converter.Jackson2JsonMessageConverter"> <constructor-arg ref="myObjectMapper"/> </bean> <bean id="rabbitListenerContainerFactory" class="org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory"> <property name="connectionFactory" ref="connectionFactory"/> <property name="concurrentConsumers" value="1"/> <property name="maxConcurrentConsumers" value="1"/> <property name="receiveTimeout" value="316224000000"/> <property name="messageConverter" ref="rmqMessageConverter"/> </bean> </beans>
Java代码
package com.acme.tasks; import com.acme.io.NewBarcodes; import javax.validation.constraints.NotNull; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.springframework.validation.annotation.Validated; @Validated @Component public class ConsumeBarcodes { private static final Logger logger = LoggerFactory.getLogger(ConsumeBarcodes.class); @Value("${rabbitmq.queue}") private String queueName; public ConsumeBarcodes() {} @RabbitListener(queues = "${rabbitmq.queue}", ackMode = "AUTO") public void ingestNewBarcodes(@NotNull NewBarcodes newBarcodes) { logger.debug("RECEIVED message in ingestNewBarcodes from RabbitMQ queue: {}", queueName); // XXX Process newBarcodes here. } }
问题解答
1. 为何自动恢复功能未生效?
出现MissedHeartbeatException时,RabbitMQ客户端会直接关闭连接,默认的ForgivingExceptionHandler仅记录错误但不触发恢复逻辑。另外你设置的receiveTimeout值过大(约3660天),导致容器的消息接收线程长时间阻塞,无法及时响应连接状态变化,进而无法触发自动恢复流程。
同时,Spring AMQP的SimpleRabbitListenerContainer依赖连接工厂的自动恢复能力,但如果连接因心跳超时被强制关闭,容器可能无法感知连接状态变化,尤其是接收线程被长时间阻塞时。
2. 是否应在SimpleRabbitListenerContainerFactory中启用automaticRecoveryEnabled和topologyRecoveryEnabled?
不需要在容器工厂中直接设置这两个参数,它们是RabbitMQ原生客户端的配置项,需通过RabbitConnectionFactoryBean配置。注意:RabbitMQ原生客户端5.0+版本默认启用automaticRecoveryEnabled,但Spring AMQP会默认将topologyRecoveryEnabled设为false,因为Spring AMQP自身会负责队列、交换机、绑定等拓扑结构的恢复与管理。
你可以通过RabbitConnectionFactoryBean显式配置automaticRecoveryEnabled=true(默认即为true),同时保持topologyRecoveryEnabled=false(Spring AMQP推荐值),确保由Spring容器管理拓扑恢复。
3. 在当前配置下,如何实现消费者处理器handleShutdownSignal?
可以通过实现ChannelAwareMessageListener接口替代当前的@RabbitListener方法,该接口包含handleShutdownSignal方法,用于处理连接或通道关闭信号,示例代码如下:
package com.acme.tasks; import com.acme.io.NewBarcodes; import javax.validation.constraints.NotNull; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.springframework.validation.annotation.Validated; import com.rabbitmq.client.Channel; import com.rabbitmq.client.ShutdownSignalException; @Validated @Component public class ConsumeBarcodes implements ChannelAwareMessageListener { private static final Logger logger = LoggerFactory.getLogger(ConsumeBarcodes.class); @Value("${rabbitmq.queue}") private String queueName; @Autowired private Jackson2JsonMessageConverter rmqMessageConverter; public ConsumeBarcodes() {} @Override @RabbitListener(queues = "${rabbitmq.queue}", ackMode = "AUTO") public void onMessage(Message message, Channel channel) throws Exception { NewBarcodes newBarcodes = (NewBarcodes) rmqMessageConverter.fromMessage(message); logger.debug("RECEIVED message in ingestNewBarcodes from RabbitMQ queue: {}", queueName); // XXX Process newBarcodes here. } @Override public void handleShutdownSignal(String consumerTag, Channel channel, ShutdownSignalException sig) { logger.error("Received shutdown signal for consumer tag: {}", consumerTag, sig); // 此处可添加自定义恢复逻辑,比如记录告警、触发连接检查等 } }
也可以通过自定义RabbitListenerErrorHandler或配置容器的shutdownListener处理关闭信号,但实现ChannelAwareMessageListener是最直接的方式。
内容的提问来源于stack exchange,提问作者Kode Charlie

