Skip to content

Commit

Permalink
propagate targetParallelism from SortedBucketIO (#2997)
Browse files Browse the repository at this point in the history
😅
  • Loading branch information
clairemcginty authored May 21, 2020
1 parent f853138 commit 2c1b71b
Showing 1 changed file with 1 addition and 1 deletion.
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ public <V> CoGbkTransform<K, V> to(SortedBucketIO.Write<K, V> write) {
public PCollection<KV<K, CoGbkResult>> expand(PBegin input) {
List<BucketedInput<?, ?>> bucketedInputs =
reads.stream().map(Read::toBucketedInput).collect(Collectors.toList());
return input.apply(new SortedBucketSource<>(keyClass, bucketedInputs));
return input.apply(new SortedBucketSource<>(keyClass, bucketedInputs, targetParallelism));
}
}

Expand Down

0 comments on commit 2c1b71b

Please sign in to comment.