使用Spring StreamBridge批量发送PubSub消息时出现内存溢出问题
问题描述
我使用StreamBridge向PubSub发送约10万条消息,采用5线程的线程池执行任务,核心代码如下:
private void fetchAndPublish(List<MyObject> list) { int threads = 5 < list.size() ? 5 : list.size(); ExecutorService es = Executors.newFixedThreadPool(threads); CompletableFuture<?>[] futures = list.stream() .map(s -> new MyRunnable(s)) .map(task -> CompletableFuture.runAsync(task, es)) .toArray(CompletableFuture[]::new); CompletableFuture.allOf(futures).join(); es.shutdown(); } public class MyRunnable implements Runnable { private final MyObject myObject; public MyRunnable(MyObject myObject) { this.myObject = myObject; } @Override public void run() { try { Message<MyObject> message = MessageBuilder.withPayload(myObject) .build(); messagePublisherService.publish(message); } catch (Exception exception) { log.error("Failed ", exception); } } }
消息发布服务代码:
@Component @Slf4j public class MessagePublisherService { @Value("${publish.destination:myDes-out-0}") private String destination; @Autowired private StreamBridge streamBridge; public void publish(Message<MyObject> message) { log.info("publishing {}", message.getPayload().getId()); streamBridge.send(destination, message); } }
程序初期运行正常,但一段时间后出现**内存溢出(Out of Memory)**异常,推送到StreamBridge的数据内存未及时释放。堆内存配置为16GB,但很快被耗尽,且单条数据体积并不大。
通过Dynatrace监控可见:
- 初期PubSub发布请求耗时仅200+ms,随着时间推移,发布耗时逐渐增加,资源占用持续攀升。
- 一段时间后,发布请求耗时大幅增加,资源被持续占用。
问题分析与解决方案
核心原因
StreamBridge的send方法默认是异步非阻塞的,若下游PubSub的生产速度跟不上线程池的发送速度,会导致消息在内部队列积压,内存无法释放。同时当前实现存在以下问题:
- 每次调用
fetchAndPublish都创建新线程池,无复用,造成线程资源浪费。 - 一次性提交10万条任务,无流量控制,直接撑爆StreamBridge内部缓冲区。
- 未处理
streamBridge.send的返回值,无法感知发送状态与积压情况。
具体优化方案
复用全局线程池
避免重复创建销毁线程池,改为全局单例配置:// 在配置类中定义全局线程池 @Bean public ExecutorService fixedThreadPool() { return Executors.newFixedThreadPool(5); }注入后复用:
@Autowired private ExecutorService es; private void fetchAndPublish(List<MyObject> list) { CompletableFuture<?>[] futures = list.stream() .map(s -> new MyRunnable(s)) .map(task -> CompletableFuture.runAsync(task, es)) .toArray(CompletableFuture[]::new); CompletableFuture.allOf(futures).join(); // 全局线程池无需shutdown }添加流量控制
分批次处理任务,避免瞬间提交大量请求:private void fetchAndPublish(List<MyObject> list) { int batchSize = 100; // 每批处理100条 for (int i = 0; i < list.size(); i += batchSize) { int end = Math.min(i + batchSize, list.size()); List<MyObject> batch = list.subList(i, end); CompletableFuture<?>[] futures = batch.stream() .map(s -> new MyRunnable(s)) .map(task -> CompletableFuture.runAsync(task, es)) .toArray(CompletableFuture[]::new); CompletableFuture.allOf(futures).join(); } }或用信号量控制并发数:
private final Semaphore semaphore = new Semaphore(100); // 同时最多发送100条 public class MyRunnable implements Runnable { @Override public void run() { try { semaphore.acquire(); Message<MyObject> message = MessageBuilder.withPayload(myObject).build(); messagePublisherService.publish(message); } catch (Exception exception) { log.error("Failed ", exception); } finally { semaphore.release(); } } }处理StreamBridge发送返回值
等待发送完成再释放资源,避免消息积压:// 修改publish方法返回Future public CompletableFuture<Void> publish(Message<MyObject> message) { log.info("publishing {}", message.getPayload().getId()); return streamBridge.send(destination, message); }在Runnable中等待发送完成:
@Override public void run() { try { Message<MyObject> message = MessageBuilder.withPayload(myObject).build(); messagePublisherService.publish(message).join(); } catch (Exception exception) { log.error("Failed ", exception); } }调整StreamBridge生产者配置
开启批量发送并限制缓冲区大小:spring.cloud.stream.bindings.myDes-out-0.producer.batch-mode=true spring.cloud.stream.bindings.myDes-out-0.producer.batch-size=100 spring.cloud.stream.bindings.myDes-out-0.producer.max-buffer-size=1000内存快照分析
用jmap或jvisualvm生成堆快照,确认内存占用对象类型,排查是否存在未释放的连接、日志对象等问题。
内容的提问来源于stack exchange,提问作者Patan
相关产品推荐
相关产品推荐

