Why Does Reducing Cardinality of Keyspace Improves Throughput/Performance for Streaming Aggregation for Cloud Dataflow

Viewed 42

When writing an Apache Beam/Dataflow job, reducing cardinality of keyspace significantly improves the performance of grouping functions more than reducing the payload size. Does this have to do with shuffling cost, this means transferring more bytes for a single key has smaller cost than having to move around more keys? Any publications or documentation that explains the reasoning would be helpful. Thanks

0 Answers
Related