AtomicInteger recordNumber = new AtomicInteger();
Files.lines(inputFile.toPath(), StandardCharsets.UTF_8)
.map(record -> new Record(recordNumber.incrementAndGet(), record))
.parallel()
.filter(record -> doSomeOperation())
.findFirst()
When I wrote this I assumed that the threads will be spawned only the map call since parallel is placed after map. But some lines in the file were getting different record numbers for every execution.
I read the official Java stream documentation and a few web sites to understand how streams work under the hood.
A few questions:
Java parallel stream works based on SplitIterator, which is implemented by every collection like ArrayList,LinkedList etc. When we construct a parallel stream out of those collections, the corresponding split iterator will be used to split and iterate the collection. This explains why parallelism happened at original input source (File lines) level rather at the result of map (i.e Record pojo). Is my understanding correct?
In my case, the input is a file IO stream. Which split iterator will be used?
It doesn't matter where we place
parallel()in the pipeline. The original input source will always be split and remaining intermediate operations will be applied.In this case, Java shouldn't allow users to place parallel operation anywhere in the pipeline except at the original source. Because, it is giving wrong understanding for those who doesn't know how java stream works internally. I know
parallel()operation would have been defined for Stream object type and so, it is working this way. But, it is better to provide some alternate solution.In the above code snippet, I am trying to add a line number to every record in the input file and so it should be ordered. However, I want to apply
doSomeOperation()in parallel as it is heavy weight logic. The one way to achieve is to write my own customized split iterator. Is there any other way?