From 727a25d391e73e2c31df0357d6622c18bcc1f89b Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Thu, 11 Sep 2025 13:31:52 -0700 Subject: [PATCH 1/4] Set default redistribute keys for KafkaIO read. --- .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java index e048a996a8c7..b90fe336c1e6 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java @@ -655,6 +655,8 @@ public static WriteRecords writeRecords() { ///////////////////////// Read Support \\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\ + private static final int DEFAULT_REDISTRIBUTE_NUM_KEYS = 32768; + /** * A {@link PTransform} to read from Kafka topics. See {@link KafkaIO} for more information on * usage and configuration. @@ -1099,7 +1101,11 @@ public Read withTopicPartitions(List topicPartitions) { * @return an updated {@link Read} transform. */ public Read withRedistribute() { - return toBuilder().setRedistributed(true).build(); + Builder builder = toBuilder().setRedistributed(true); + if (getRedistributeNumKeys() == 0) { + builder = builder.setRedistributeNumKeys(DEFAULT_REDISTRIBUTE_NUM_KEYS); + } + return builder.build(); } /** @@ -2667,7 +2673,11 @@ public ReadSourceDescriptors withProcessingTime() { /** Enable Redistribute. */ public ReadSourceDescriptors withRedistribute() { - return toBuilder().setRedistribute(true).build(); + Builder builder = toBuilder().setRedistribute(true); + if (getRedistributeNumKeys() == 0) { + builder = builder.setRedistributeNumKeys(DEFAULT_REDISTRIBUTE_NUM_KEYS); + } + return builder.build(); } public ReadSourceDescriptors withAllowDuplicates() { From 478f82e905a95299812485ee94bdc84e940ecc6b Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Thu, 11 Sep 2025 14:21:43 -0700 Subject: [PATCH 2/4] Add test for Kafka redistribute default keys. --- .../apache/beam/sdk/io/kafka/KafkaIOTest.java | 44 +++++++++++++++++++ 1 file changed, 44 insertions(+) diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java index 3d441f8dc521..bbf1cce02b16 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java @@ -792,6 +792,50 @@ public void testNumKeysIgnoredWithRedistributeNotEnabled() { p.run(); } + @Test + public void testDefaultRedistributeNumKeys() { + int numElements = 1000; + // Redistribute is not used. + KafkaIO.Read read = + mkKafkaReadTransform( + numElements, + numElements, + new ValueAsTimestampFn(), + false, /*redistribute*/ + false, /*allowDuplicates*/ + null, /*numKeys*/ + null, /*offsetDeduplication*/ + null /*topics*/); + assertEquals(0, read.getRedistributeNumKeys()); + + // Redistribute is used and defaulted the number of keys. + read = + mkKafkaReadTransform( + numElements, + numElements, + new ValueAsTimestampFn(), + true, /*redistribute*/ + false, /*allowDuplicates*/ + null, /*numKeys*/ + null, /*offsetDeduplication*/ + null /*topics*/); + // Default is by DEFAULT_REDISTRIBUTE_NUM_KEYS in KafkaIO. + assertEquals(32768, read.getRedistributeNumKeys()); + + // Redistribute is used and specified the number of keys. + read = + mkKafkaReadTransform( + numElements, + numElements, + new ValueAsTimestampFn(), + true, /*redistribute*/ + false, /*allowDuplicates*/ + 10, /*numKeys*/ + null, /*offsetDeduplication*/ + null /*topics*/); + assertEquals(10, read.getRedistributeNumKeys()); + } + @Test public void testDisableRedistributeKafkaOffsetLegacy() { thrown.expect(Exception.class); From 2cd4f976fa42d7eb71809868b63dacfbbbc78e6b Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Thu, 11 Sep 2025 14:33:22 -0700 Subject: [PATCH 3/4] Clarify test comments. --- .../org/apache/beam/sdk/io/kafka/KafkaIOTest.java | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java index bbf1cce02b16..83c2e1b38826 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java @@ -30,6 +30,7 @@ import static org.hamcrest.Matchers.matchesPattern; import static org.hamcrest.Matchers.not; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; @@ -795,7 +796,7 @@ public void testNumKeysIgnoredWithRedistributeNotEnabled() { @Test public void testDefaultRedistributeNumKeys() { int numElements = 1000; - // Redistribute is not used. + // Redistribute is not used and does not modify the read transform further. KafkaIO.Read read = mkKafkaReadTransform( numElements, @@ -806,9 +807,10 @@ public void testDefaultRedistributeNumKeys() { null, /*numKeys*/ null, /*offsetDeduplication*/ null /*topics*/); + assertFalse(read.isRedistributed()); assertEquals(0, read.getRedistributeNumKeys()); - // Redistribute is used and defaulted the number of keys. + // Redistribute is used and defaulted the number of keys due to no user setting. read = mkKafkaReadTransform( numElements, @@ -819,10 +821,11 @@ public void testDefaultRedistributeNumKeys() { null, /*numKeys*/ null, /*offsetDeduplication*/ null /*topics*/); - // Default is by DEFAULT_REDISTRIBUTE_NUM_KEYS in KafkaIO. + assertTrue(read.isRedistributed()); + // Default is defined by DEFAULT_REDISTRIBUTE_NUM_KEYS in KafkaIO. assertEquals(32768, read.getRedistributeNumKeys()); - // Redistribute is used and specified the number of keys. + // Redistribute is set with user-specified the number of keys. read = mkKafkaReadTransform( numElements, @@ -833,6 +836,7 @@ public void testDefaultRedistributeNumKeys() { 10, /*numKeys*/ null, /*offsetDeduplication*/ null /*topics*/); + assertTrue(read.isRedistributed()); assertEquals(10, read.getRedistributeNumKeys()); } From 21e91856165e08c4292c7edfb93c736fdae9dc5c Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Tue, 16 Sep 2025 14:30:51 -0500 Subject: [PATCH 4/4] Update java doc comments to match the default keys update. --- .../org/apache/beam/sdk/io/kafka/KafkaIO.java | 26 ++++++++++++++++--- 1 file changed, 23 insertions(+), 3 deletions(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java index b90fe336c1e6..045a74a8507e 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java @@ -655,6 +655,12 @@ public static WriteRecords writeRecords() { ///////////////////////// Read Support \\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\\ + /** + * Default number of keys to redistribute Kafka inputs into. + * + *

This value is used when {@link Read#withRedistribute()} is used without {@link + * Read#withRedistributeNumKeys(int redistributeNumKeys)}. + */ private static final int DEFAULT_REDISTRIBUTE_NUM_KEYS = 32768; /** @@ -1127,10 +1133,11 @@ public Read withAllowDuplicates(Boolean allowDuplicates) { * Redistributes Kafka messages into a distinct number of keys for processing in subsequent * steps. * - *

Specifying an explicit number of keys is generally recommended over redistributing into an - * unbounded key space. + *

If unset, defaults to {@link KafkaIO#DEFAULT_REDISTRIBUTE_NUM_KEYS}. * - *

Must be used with {@link KafkaIO#withRedistribute()}. + *

Use zero to disable bucketing into a distinct number of keys. + * + *

Must be used with {@link Read#withRedistribute()}. * * @param redistributeNumKeys specifies the total number of keys for redistributing inputs. * @return an updated {@link Read} transform. @@ -2684,6 +2691,19 @@ public ReadSourceDescriptors withAllowDuplicates() { return toBuilder().setAllowDuplicates(true).build(); } + /** + * Redistributes Kafka messages into a distinct number of keys for processing in subsequent + * steps. + * + *

If unset, defaults to {@link KafkaIO#DEFAULT_REDISTRIBUTE_NUM_KEYS}. + * + *

Use zero to disable bucketing into a distinct number of keys. + * + *

Must be used with {@link ReadSourceDescriptors#withRedistribute()}. + * + * @param redistributeNumKeys specifies the total number of keys for redistributing inputs. + * @return an updated {@link Read} transform. + */ public ReadSourceDescriptors withRedistributeNumKeys(int redistributeNumKeys) { return toBuilder().setRedistributeNumKeys(redistributeNumKeys).build(); }