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

使用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的生产速度跟不上线程池的发送速度,会导致消息在内部队列积压,内存无法释放。同时当前实现存在以下问题:

  1. 每次调用fetchAndPublish都创建新线程池,无复用,造成线程资源浪费。
  2. 一次性提交10万条任务,无流量控制,直接撑爆StreamBridge内部缓冲区。
  3. 未处理streamBridge.send的返回值,无法感知发送状态与积压情况。

具体优化方案

  1. 复用全局线程池
    避免重复创建销毁线程池,改为全局单例配置:

    // 在配置类中定义全局线程池
    @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
    }
    
  2. 添加流量控制
    分批次处理任务,避免瞬间提交大量请求:

    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();
            }
        }
    }
    
  3. 处理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);
        }
    }
    
  4. 调整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
    
  5. 内存快照分析
    用jmap或jvisualvm生成堆快照,确认内存占用对象类型,排查是否存在未释放的连接、日志对象等问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 06:34:56