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

如何让基于Spring Integration的ActiveMQ Artemis生产者感知服务器阻塞?

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 08:25:05