使用parallelStream填充Hazelcast IMap时出现RejectedExecutionException如何解决
问题根因
你遇到的异常是Java公共ForkJoin池的线程超限导致的:
- Java并行流默认使用
ForkJoinPool.commonPool()执行任务,当任务中存在阻塞操作(比如你代码中permutations.put是跨节点网络IO操作,属于阻塞调用)时,ForkJoinPool为了保证并行度会自动创建补偿线程替换被阻塞的worker - 你需要处理的元素规模达4.8亿,阻塞操作的数量远超出公共ForkJoinPool的最大线程数限制,最终触发
RejectedExecutionException
修复方案
方案1:改用Hazelcast Jet原生管道写入IMap(最优方案)
既然你已经在使用Hazelcast Jet,直接用Jet的原生处理框架填充IMap,完全避开公共ForkJoin池的限制,同时适配Jet的分布式线程模型:
Pipeline pipeline = Pipeline.create(); pipeline.readFrom(Sources.iterable(Collections2.permutations(initialPermutation))) .map(Permutation::new) .writeTo(Sinks.map("permutations", Permutation::toString, Function.identity())); jetInstance.newJob(pipeline).join();
方案2:使用自定义线程池运行并行流
如果要保留并行流的写法,不要使用默认的公共ForkJoinPool,自定义固定大小的ForkJoinPool控制线程上限:
// 线程数可根据集群硬件、Hazelcast节点负载自行调整,推荐不超过CPU核心数的4倍 int customThreadNum = Runtime.getRuntime().availableProcessors() * 2; ForkJoinPool customPool = new ForkJoinPool(customThreadNum); try { customPool.submit(() -> Collections2.permutations(initialPermutation).parallelStream().forEach(set -> { Permutation permutation = new Permutation(set); permutations.put(permutation.toString(), permutation); }) ).get(); } catch (InterruptedException | ExecutionException e) { // 自行处理异常逻辑 e.printStackTrace(); } finally { customPool.shutdown(); }
方案3:新增批量写入逻辑降低阻塞次数
不管用哪种写入方式,都可以把单条put改成批量putAll,大幅减少网络IO阻塞操作的数量,进一步降低线程压力:
int batchSize = 2000; // 批次大小可自行调整 AtomicInteger batchCounter = new AtomicInteger(); customPool.submit(() -> Collections2.permutations(initialPermutation) .parallelStream() .map(Permutation::new) .collect(Collectors.groupingBy(item -> batchCounter.getAndIncrement() / batchSize)) .values() .forEach(batch -> { Map<String, Permutation> batchMap = batch.stream() .collect(Collectors.toMap(Permutation::toString, Function.identity())); permutations.putAll(batchMap); }) ).get();
注意事项
Collections2.permutations返回的是懒加载视图,不会提前生成所有排列元素,处理4.8亿规模数据时尽量保持流处理的懒加载特性,不要全量加载到内存中,避免OOM。
内容的提问来源于stack exchange,提问作者Nova Shadowflame
相关产品推荐
相关产品推荐

