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

使用CompletableFuture与线程池处理百万级发送任务的优化问题

问题影响评估

现有实现的影响不止内存占用层面,具体问题如下:

  • 单看内存开销:单个CompletableFuture实例在64位JVM中约占3264字节,100万个实例总内存占用约30MB60MB,数值上不算特别夸张,但存在多个隐性严重问题
  • 线程安全Bug:你存储Future用到的ArrayList不是线程安全类,多线程调用send方法向同一个ID对应的List添加元素时,会出现并发修改异常、元素丢失等问题
  • 调用逻辑不匹配:你提供的客户端调用代码无参调用notifySendComplete(),但方法定义要求传入ID参数,实际运行会直接报错返回,无法等待任务完成
  • 内存泄漏风险:如果某批次任务的notifySendComplete被漏调用,对应List里的所有Future会一直留在jobMap中无法被回收,随着调用次数增多最终会触发OOM

更优实现方案

核心思路是不存储所有Future实例,用原子计数器跟踪任务完成进度,完全满足「不能关闭线程池、不能使用CountDownLatch」的约束:

实现代码

public class MySender {
    private final MyPublisher myPublisher;
    private final ExecutorService threadPool;
    // 按批次ID存储剩余未完成任务数
    private final Map<String, AtomicLong> jobCounter = Maps.newConcurrentMap();
    // 按批次ID存储完成通知信号
    private final Map<String, CompletableFuture<Void>> finishSignalMap = Maps.newConcurrentMap();

    public MySender (final MyPublisher myPublisher,
                     ExecutorService threadPool) {
        this.myPublisher= myPublisher;
        this.threadPool = threadPool;
    }

    // 批量提交任务前先初始化批次
    public void initBatch(String batchId, int totalTaskNum) {
        jobCounter.put(batchId, new AtomicLong(totalTaskNum));
        finishSignalMap.put(batchId, new CompletableFuture<>());
    }

    public void send(final MyData event) {
        CompletableFuture.runAsync(() -> {
            try {
                doSubmit(event);
            } finally {
                String batchId = event.getID();
                long leftTaskNum = jobCounter.get(batchId).decrementAndGet();
                // 所有任务执行完成后触发完成信号、清理资源
                if (leftTaskNum == 0) {
                    finishSignalMap.get(batchId).complete(null);
                    jobCounter.remove(batchId);
                    finishSignalMap.remove(batchId);
                }
            }
        }, threadPool);
    }

    public void notifySendComplete(final String batchId) {
        CompletableFuture<Void> finishSignal = finishSignalMap.get(batchId);
        if (finishSignal == null) {
            return;
        }
        // 阻塞等待批次全部发送完成
        finishSignal.join();
    }

    private void doSubmit(final MyData event) {
         try {
             myPublisher.send(event);
         } catch(Exception e) {
             // 日志打印错误
         }
    }
}

客户端调用方式

String batchId = "自定义批次ID";
// 提交任务前先初始化批次,传入当前批次总任务数
mySender.initBatch(batchId, myInputList.size());
myInputList.forEach(input -> {
    MyData event = createData(input);
    event.setID(batchId);
    mySender.send(event);
});
// 阻塞等待当前批次所有消息发送完成
mySender.notifySendComplete(batchId);

方案优势

  • 内存占用极低:不管任务量是10万还是1000万,每个批次仅存储2个对象(计数器+完成信号),完全不会出现大量Future无法回收的问题
  • 线程安全:所有计数操作都是原子操作,不存在并发修改问题
  • 无内存泄漏风险:任务完成后自动清理Map中的对应记录
  • 符合所有约束:不需要关闭线程池,也没有用到CountDownLatch

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 21:00:07