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

ActiveMQ内存超限与Broken pipe问题排查求助

ActiveMQ内存耗尽与生产者限流问题排查

错误日志解读

你遇到的错误日志如下:

2023-06-09 12:57:57,650 | INFO  | Usage Manager Memory Limit (1336252826) reached on queue://VIBER, size 0. Producers will be throttled to the rate at which messages are removed from this destination to prevent flooding it. See http://activemq.apache.org/producer-flow-control.html for more info. | org.apache.activemq.broker.region.Queue | ActiveMQ Transport: tcp:///10.255.102.101:45200@61616
2023-06-09 12:58:27,735 | INFO  | Usage(default:memory:queue://VIBER:memory) percentUsage=48%, usage=646397829, limit=1336252826, percentUsageMinDelta=1%;Parent:Usage(default:memory) percentUsage=100%, usage=1336254815, limit=1336252826, percentUsageMinDelta=1%: Usage Manager Memory Limit reached. Producer (ID:49-bot-a.ukrposhta.loc-36283-1685673275857-17607:1:1:1) stopped to prevent flooding queue://VIBER. See http://activemq.apache.org/producer-flow-control.html for more info (blocking for: 8s) | org.apache.activemq.broker.region.Queue | ActiveMQ Transport: tcp:///10.255.102.101:45200@61616

日志核心含义:

  • VIBER队列的内存配额已耗尽,Broker开始限流生产者,强制生产速度与消息消费速度一致
  • 整个Broker的内存使用率达到100%,生产者被直接阻塞,避免队列被消息彻底淹没

后续出现的Broken pipe错误,是Broker内存耗尽后无法处理连接请求或数据发送,导致TCP连接中断。

代码逻辑问题分析

1. 生产者一次性推送超大量消息

produceNews方法会一次性生成50万条(Viber)或25万条(Telegram)消息,通过NewsProducer一次性提交到队列。使用SESSION_TRANSACTED会话模式时,所有消息会先缓存到Broker内存中,直到会话提交。这种批量生产方式会瞬间占满Broker内存配额,直接触发流控机制。

2. 消费者生命周期过短

consumeNews方法中,创建固定线程池提交消费者任务后,调用invokeAll立即关闭线程池。这意味着消费者线程仅会消费当前队列中的消息,执行完成后就终止,不会持续监听队列。若后续生产者继续发送消息,队列将无消费者处理,消息持续堆积最终耗尽内存。

3. 消费速度远低于生产速度

Telegram/Viber的API调用本身属于慢操作,而你仅配置8-10个消费者线程,面对几十万条消息,消费速度完全跟不上生产速度,消息会在Broker内存中越积越多。

4. 定时任务可能重复触发推送

每半小时运行的定时任务会检查数据库新闻并触发推送,若上一次消息未消费完成,新的推送任务会继续往队列塞消息,加剧堆积。

ActiveMQ配置问题

1. 队列未配置消息堆积限制

destinationPolicy仅给Topic配置了pendingMessageLimitStrategy,未给Queue设置任何消息数量限制。队列中的消息会无限制堆积在内存中,直到耗尽Broker内存。

2. 内存配额配置不合理

systemUsage中设置内存使用为JVM堆的70%,若JVM堆容量本身不大(例如仅2GB),70%的配额仅1.4GB,几十万条消息极易超过这个限制。

优化方案

1. 生产者优化

  • 分批发送消息:不要一次性提交几十万条消息,改为每次发送1000-5000条后提交会话,降低Broker内存压力。
  • 启用异步发送:通过producer.setAsyncSend(true)开启异步发送,提高生产效率,同时避免生产者被阻塞。
  • 压缩消息:对JSON格式的消息进行GZIP压缩,减少单条消息的内存占用。

2. 消费者优化

  • 持续监听队列:修改NewsConsumer逻辑,让消费者线程通过循环调用consumer.receive()或异步监听的方式持续监听队列,而非一次性消费完就退出;同时不要在消费完成后立即关闭线程池,保持消费者长期运行。
  • 合理调整线程池大小:根据第三方平台的API限流规则,适当增加消费者线程数(注意不要超过平台调用频率限制),提升消费速度。
  • 批量消费:采用批量消费模式,每次从队列拉取多条消息处理,减少Session交互次数,提高消费效率。

3. ActiveMQ配置优化

  • 给队列配置消息堆积限制:在destinationPolicy中添加Queue的策略,限制队列最大待处理消息数,超过后自动丢弃或降级:
<policyEntry=">" queue>
  <pendingMessageLimitStrategy>
    <constantPendingMessageLimitStrategy limit="10000"/>
  </pendingMessageLimitStrategy>
</policyEntry>
  • 调整内存配额:根据服务器实际内存情况,增大JVM堆内存,并合理设置memoryUsage的百分比或绝对值。例如JVM堆设置为8GB时,memoryUsage可设置为5GB左右。
  • 可选启用持久化:若允许消息持久化,可将队列配置为持久化模式,内存满后消息会写入磁盘,避免内存耗尽,但会带来一定性能开销。

4. 整体流程优化

  • 避免重复推送:在定时任务中添加状态检查,数据库标记新闻推送状态,仅未推送的新闻才触发推送。
  • 监控队列状态:添加ActiveMQ队列监控,实时查看消息堆积数量、内存使用率等指标,及时处理堆积问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 11:18:07