为何并行无序Java Stream的forEachOrdered仍按源顺序消费?
我近期基于阻塞队列实现了单消费者多生产者的算法,在实现简化版本时,尝试通过基础Java Stream来完成,但未达到预期效果。
复现代码
LongStream.range(0, 32) .unordered() .parallel() .map(value -> { // 'do some work', which can vary in duration try { Thread.sleep(new Random(value).nextLong(1000)); } catch (final InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException(e); } System.out.println("map " + value); return value; }).forEachOrdered(value -> { // actual implementation writes to output stream System.out.println(">>> " + value); System.out.flush(); });
运行输出
map 7 map 29 map 8 map 11 map 28 map 21 map 20 map 13 map 6 map 10 map 1 map 15 map 5 map 24 map 26 map 14 map 9 map 2 map 31 map 17 map 22 map 16 map 3 map 23 map 12 map 27 map 0 >>> 0 >>> 1 >>> 2 >>> 3 map 4 >>> 4 >>> 5 >>> 6 >>> 7 >>> 8 >>> 9 >>> 10 >>> 11 >>> 12 >>> 13 >>> 14 >>> 15 >>> 16 >>> 17 map 25 map 30 map 18 >>> 18 map 19 >>> 19 >>> 20 >>> 21 >>> 22 >>> 23 >>> 24 >>> 25 >>> 26 >>> 27 >>> 28 >>> 29 >>> 30 >>> 31
疑问与需求
我的初始想法是:每个map操作完成后即可被终端操作消费,且forEachOrdered会逐个处理元素(无需手动同步)。我已将流声明为unordered,认为在无序流中forEachOrdered不必保证按源顺序处理,但实际终端操作始终按流源顺序执行,即使让值0的处理耗时极长,其他值不耗时,结果依然如此。
更新说明:原实现是单消费者多生产者模型——生产者线程从任务队列取任务,将结果推至单消费者队列;消费者线程串行处理结果(需写入输出流,必须同步)。我希望用Stream替代:将任务转为流,map作为生产者,终端操作写入输出。选择forEachOrdered是因为它承诺逐个处理元素,且无序流无需保证顺序,觉得这样能实现无序并行处理,但实际不符。改用forEach并手动同步可正常工作。
我是否存在理解误区?该现象是否正常?若正常,成因是什么?Stream实现是否有优化空间?
解答
1. 核心理解误区
你对forEachOrdered的语义存在误解:无论流是否声明为unordered,forEachOrdered的核心职责就是严格按流的相遇顺序执行终端操作。unordered只是允许Stream API在中间操作阶段打乱顺序以提升并行效率,但不会改变forEachOrdered本身强制顺序消费的语义。
Javadoc明确规定:forEachOrdered会遵循流的相遇顺序(如果存在)执行动作。对于原本有序的流(比如LongStream.range生成的流),即使调用了unordered(),底层并未完全丢弃源的顺序信息,forEachOrdered依然会遵循这个顺序执行。
2. 现象正常的成因
LongStream.range生成的是天然有序流,调用unordered()仅表示允许流在处理时不维护顺序,但并没有抹除源的顺序属性。在并行处理时,map阶段可以乱序执行,但forEachOrdered为了保证顺序,会等待当前元素之前的所有元素都完成处理后,才会执行当前元素的终端操作——这就是为什么0的map完成后,终端操作从0开始连续输出,因为它要严格按源顺序消费所有元素。
简单来说:unordered允许中间操作乱序,但forEachOrdered强制终端操作按相遇顺序执行,而原流的相遇顺序并未被完全清除,因此终端操作始终按源顺序输出。
3. 满足需求的正确方案
你的需求是并行处理任务,结果完成一个就串行消费一个(不要求顺序),正确的实现方式是用forEach加同步机制,而非forEachOrdered:
forEach在并行流中会乱序执行终端操作,正好符合“完成即消费”的需求;- 由于输出流需要串行写入,需在forEach的动作逻辑中添加同步(比如
synchronized块,或利用PrintStream本身的线程安全特性)。
调整后的示例代码:
Object lock = new Object(); LongStream.range(0, 32) .unordered() .parallel() .map(value -> { try { Thread.sleep(new Random(value).nextLong(1000)); } catch (final InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException(e); } System.out.println("map " + value); return value; }).forEach(value -> { synchronized (lock) { System.out.println(">>> " + value); System.out.flush(); } });
4. Stream实现的优化空间
从JDK设计角度,forEachOrdered的语义就是强制顺序,因此不会为无序流做特殊优化。如果需要“完成即串行消费”的语义,目前只能通过forEach加同步实现,或者使用更底层的并发工具(比如你原本的阻塞队列方案)——因为Stream API的终端操作语义里,只提供了“无序并行消费(forEach)”和“有序串行消费(forEachOrdered)”两种选项,没有中间态的语义支持。
内容的提问来源于stack exchange,提问作者Koekje

