apache / apache/beam

Support zero-shuffle grouping operations

Open
#19,044 0 comments 0 reactions 0 assignees View on GitHub
dsl euphoria improvement P3 performance
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

On some occasions input dataset might be already correctly shuffled (i.e. as a result of previous operation(s)), which means that subsequent grouping operation could leverage this and remove the unneeded shuffle. Example (pseudocode):
```

 Dataset input = ...

 Dataset> counts1 = CountByKey.of(input)

  
.keyBy(e -> e)

   .windowBy( /* some small window */ )

   .output();

 Dataset> counts2 = SumByKey.of(counts1)

   .keyBy(Pair::getFirst)

   .windowBy( /* larger window
*/ )

   .output();

```

Now, the second `ReduceByKey` already might have correct shuffle (depends on runner), but isn't able to leverage this, because it isn't aware that the key grouping key has not changed from the previous operation.

Proposed change:
```

 Dataset input = ...

 Dataset> counts1 = CountByKey.of(input)

  
.keyBy(e -> e)

   .windowBy( /* some small window */ )

   .output();

 Dataset> counts2 = SumByKey.of(counts1)

   .keyByLocally(Pair::getFirst)

   .windowBy( /* larger
window */ )

   .output();

```

Introduce `keyByLocally` to keyed operations, which will tell the runner that the grouping is preserved from one keyed operator to the other.

This will probably require some support on Beam SDK side, because this information has to be passed to the runner (so that i.e. FlinkRunner can make use of something like `DataStreamUtils#reinterpretAsKeyedStream`.

Imported from Jira [BEAM-5330](https://issues.apache.org/jira/browse/BEAM-5330). Original Jira may contain additional context.
Reported by: janl.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.