如何让基于Spring Integration的ActiveMQ Artemis生产者感知服务器阻塞?
这个问题我之前在生产环境碰到过类似情况——Artemis在磁盘不足这类资源耗尽场景下的默认行为确实容易让生产者“蒙在鼓里”。结合Spring Integration的使用场景,咱们可以从主动配置服务器行为、客户端超时与异常处理、客户端事件监听、提前指标预警这几个维度来解决:
1. 配置Artemis服务器,让阻塞场景主动抛出异常
Artemis默认在磁盘满时会静默阻塞生产者(仅日志警告),不会主动拒绝请求。我们可以通过修改broker.xml的参数,让服务器在资源不足时直接向生产者抛出异常:
<!-- 设置磁盘使用率阈值,比如90%时触发保护机制 --> <max-disk-usage>90</max-disk-usage> <!-- 开启磁盘满时拒绝发送请求,直接向生产者抛出异常 --> <send-fail-on-disk-full>true</send-fail-on-disk-full> <!-- 可选:设置阻塞生产者的超时时间,超时后同样抛出异常 --> <blocking-timeout-millis>30000</blocking-timeout-millis>
解释:开启send-fail-on-disk-full后,当磁盘使用率超过阈值,服务器会直接返回错误,而不是无限阻塞生产者,这样客户端就能立刻感知到异常。
2. 配置Spring Integration生产者的超时与异常处理
Spring Integration的JMS发送组件默认没有设置超时,会一直等待服务器响应。我们需要给生产者添加超时限制,并配置异常处理链路:
方式一:通过JmsTemplate配置超时
@Bean public JmsTemplate jmsTemplate(ConnectionFactory connectionFactory) { JmsTemplate template = new JmsTemplate(connectionFactory); template.setSendTimeout(5000); // 设置5秒发送超时,超时后抛出JmsException return template; }
方式二:在XML配置的出站适配器中设置超时
<int-jms:outbound-channel-adapter id="artemisOutbound" channel="messageInputChannel" connection-factory="artemisConnectionFactory" destination="targetQueue" send-timeout="5000"/> <!-- 5秒超时 -->
添加异常处理逻辑
用Spring Integration的ExpressionEvaluatingRequestHandlerAdvice捕获发送异常,实现告警或重试:
@Bean public ExpressionEvaluatingRequestHandlerAdvice artemisErrorAdvice() { ExpressionEvaluatingRequestHandlerAdvice advice = new ExpressionEvaluatingRequestHandlerAdvice(); advice.setFailureChannel("artemisErrorChannel"); // 异常消息发送到错误通道 advice.setThrowExceptionOnFailure(true); // 可选:是否继续向上抛出异常 return advice; }
然后在出站适配器中引用这个Advice,就能在发送失败时触发自定义处理逻辑(比如告警、写入死信队列)。
3. 利用Artemis客户端的事件监听
Artemis客户端提供了会话监听和异步发送回调,可以更细粒度地感知服务器状态变化:
注册会话监听器
监听连接失败、会话异常等事件,比如磁盘满导致服务器断开连接时触发告警:
@Bean public ConnectionFactory artemisConnectionFactory() { ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory("tcp://artemis-server:61616"); factory.setClientSessionListener(new ClientSessionListener() { @Override public void connectionFailed(ActiveMQException exception, boolean failedOver) { log.error("Artemis连接异常,原因:{}", exception.getMessage(), exception); // 这里可以触发告警,比如调用运维通知接口 } @Override public void connectionCreated(ClientSession session) { log.info("Artemis连接恢复成功"); } }); return factory; }
异步发送的回调处理
如果用异步发送方式,可以通过SendCallback捕获异常:
ClientProducer producer = session.createProducer(targetQueue); producer.send(message, new SendCallback() { @Override public void onException(ActiveMQException exception) { log.error("异步发送消息到Artemis失败:{}", exception.getMessage(), exception); // 处理失败逻辑,比如重试或记录到异常表 } @Override public void onSuccess() { log.debug("消息发送成功"); } });
4. 主动监控Artemis指标,提前预警
除了被动感知异常,最好能在资源耗尽前提前预警。可以通过Artemis的JMX接口监控关键指标:
- 磁盘使用率(
diskUsage) - 队列待处理消息数
- 生产者阻塞数(
blockedProducers)
结合Spring Boot Actuator,可以把这些指标暴露出来,用Prometheus+Grafana设置阈值告警(比如磁盘使用率超过85%时触发告警),在磁盘满之前就介入处理,避免阻塞场景发生。
总结一下:通过服务器配置主动抛异常+客户端超时与异常处理+事件监听+提前指标预警的组合方式,既能让生产者及时感知服务器阻塞,也能从根源上避免这类场景的发生。另外要注意,开启send-fail-on-disk-full后,需要确保应用有重试机制或死信队列处理,避免消息丢失。
内容的提问来源于stack exchange,提问作者madarinho

