ActiveMQ内存超限与Broken pipe问题排查求助
错误日志解读
你遇到的错误日志如下:
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

