Java parallel stream - order of invoking the parallel() method

Viewed 845
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?

3 Answers

This explains why parallelism happened at original input source (File lines) level rather at the result of map (i.e Record pojo).

The entire stream is either parallel or sequential. We don't select a subset of operations to run sequentially or in parallel.

When the terminal operation is initiated, the stream pipeline is executed sequentially or in parallel depending on the orientation of the stream on which it is invoked. [...] When the terminal operation is initiated, the stream pipeline is executed sequentially or in parallel depending on the mode of the stream on which it is invoked. same source

As you mention, parallel streams use split iterators. Clearly, this is to partition the data before operations start running.


In my case, the input is a file IO stream. Which split iterator will be used?

Looking at the source, I see it uses java.nio.file.FileChannelLinesSpliterator


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.

Right. You can even call parallel() and sequential() multiple times. The one invoked last will win. When we call parallel(), we set that for the stream that's returned; and as stated above, all operations run either sequentially or in parallel.


In this case, Java shouldn't allow users to place parallel operation anywhere in the pipeline except at the original source...

This becomes a matter of opinions. I think Zabuza gives a good reason to support the JDK designers' choice.


The one way to achieve is to write my own customized split iterator. Is there any other way?

This depends on your operations

  • If findFirst() is your real terminal operation, then you don't even need to worry about parallel execution, because there won't be many calls to doSomething() anyway (findFirst() is short-circuiting). .parallel() in fact may cause more than one element to be processed, while findFirst() on a sequential stream would prevent that.
  • If your terminal operation doesn't create much data, then maybe you can create your Record objects using a sequential stream, then process the result in parallel:

    List<Record> smallData = Files.lines(inputFile.toPath(), 
                                         StandardCharsets.UTF_8)
      .map(record -> new Record(recordNumber.incrementAndGet(), record)) 
      .collect(Collectors.toList())
      .parallelStream()     
      .filter(record -> doSomeOperation())
      .collect(Collectors.toList());
    
  • If your pipeline would load a lot of data in memory (which may be the reason you're using Files.lines()), then perhaps you'll need a custom split iterator. Before I go there, though, I'd look into other options (such saving lines with an id column to start with - that's just my opinion).
    I'd also attempt to process records in smaller batches, like this:

    AtomicInteger recordNumber = new AtomicInteger();
    final int batchSize = 10;
    
    try(BufferedReader reader = Files.newBufferedReader(inputFile.toPath(), 
            StandardCharsets.UTF_8);) {
        Supplier<List<Record>> batchSupplier = () -> {
            List<Record> batch = new ArrayList<>();
            for (int i = 0; i < batchSize; i++) {
                String nextLine;
                try {
                    nextLine = reader.readLine();
                } catch (IOException e) {
                    //hanlde exception
                    throw new RuntimeException(e);
                }
    
                if(null == nextLine) 
                    return batch;
                batch.add(new Record(recordNumber.getAndIncrement(), nextLine));
            }
            System.out.println("next batch");
    
            return batch;
        };
    
        Stream.generate(batchSupplier)
            .takeWhile(list -> list.size() >= batchSize)
            .map(list -> list.parallelStream()
                             .filter(record -> doSomeOperation())
                             .collect(Collectors.toList()))
            .flatMap(List::stream)
            .forEach(System.out::println);
    }
    

    This executes doSomeOperation() in parallel without loading all the data into memory. But note that batchSize will need to be given a thought.

The original Stream design included the idea to support subsequent pipeline stages with different parallel execution settings, but this idea has been abandoned. The API may stem from this time, but on the other hand, an API design that forces the caller to make a single unambiguous decision for parallel or sequential execution would be much more complicated.

The actual Spliterator in use by Files.lines(…) is implementation-dependent. In Java 8 (Oracle or OpenJDK), you always get the same as with BufferedReader.lines(). In more recent JDKs, if the Path belongs to the default filesystem and the charset is one of the supported for this feature, you get a Stream with a dedicated Spliterator implementation, the java.nio.file.FileChannelLinesSpliterator. If the preconditions are not met, you get the same as with BufferedReader.lines(), which still is based on an Iterator implemented within BufferedReader and wrapped via Spliterators.spliteratorUnknownSize.

Your specific task is best handled with a custom Spliterator which can perform the line numbering right at the source, before the parallel processing, to allow subsequent parallel processing without restrictions.

public static Stream<Record> records(Path p) throws IOException {
    LineNoSpliterator sp = new LineNoSpliterator(p);
    return StreamSupport.stream(sp, false).onClose(sp);
}

private static class LineNoSpliterator implements Spliterator<Record>, Runnable {
    int chunkSize = 100;
    SeekableByteChannel channel;
    LineNumberReader reader;

    LineNoSpliterator(Path path) throws IOException {
        channel = Files.newByteChannel(path, StandardOpenOption.READ);
        reader=new LineNumberReader(Channels.newReader(channel,StandardCharsets.UTF_8));
    }

    @Override
    public void run() {
        try(Closeable c1 = reader; Closeable c2 = channel) {}
        catch(IOException ex) { throw new UncheckedIOException(ex); }
        finally { reader = null; channel = null; }
    }

    @Override
    public boolean tryAdvance(Consumer<? super Record> action) {
        try {
            String line = reader.readLine();
            if(line == null) return false;
            action.accept(new Record(reader.getLineNumber(), line));
            return true;
        } catch (IOException ex) {
            throw new UncheckedIOException(ex);
        }
    }

    @Override
    public Spliterator<Record> trySplit() {
        Record[] chunks = new Record[chunkSize];
        int read;
        for(read = 0; read < chunks.length; read++) {
            int pos = read;
            if(!tryAdvance(r -> chunks[pos] = r)) break;
        }
        return Spliterators.spliterator(chunks, 0, read, characteristics());
    }

    @Override
    public long estimateSize() {
        try {
            return (channel.size() - channel.position()) / 60;
        } catch (IOException ex) {
            return 0;
        }
    }

    @Override
    public int characteristics() {
        return ORDERED | NONNULL | DISTINCT;
    }
}

And the following is a simple demonstration of when application of parallel is applied. The output from peek clearly shows the difference between the two examples. Note: The map call is just tossed in to add another method prior to parallel.

IntStream.rangeClosed (1,20).peek(a->System.out.print(a+" "))
        .map(a->a + 200).sum();
System.out.println();
IntStream.rangeClosed(1,20).peek(a->System.out.print(a+" "))
        .map(a->a + 200).parallel().sum();
Related