From 403e2b40f2c638b1ac238c88b9e3c1604a4801b1 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Wed, 10 Sep 2025 09:53:13 -0700 Subject: [PATCH 01/19] Add deterministic redistribute sharding for KafkaIO read. --- .../beam/sdk/transforms/Redistribute.java | 42 ++++++++++++++----- .../org/apache/beam/sdk/io/kafka/KafkaIO.java | 21 +++++----- 2 files changed, 42 insertions(+), 21 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java index a01b5f570a57..6f9f4a2180c2 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java @@ -48,7 +48,7 @@ public class Redistribute { /** @return a {@link RedistributeArbitrarily} transform with default configuration. */ public static RedistributeArbitrarily arbitrarily() { - return new RedistributeArbitrarily<>(null, false); + return new RedistributeArbitrarily<>(null, false, false); } /** @return a {@link RedistributeByKey} transform with default configuration. */ @@ -131,23 +131,33 @@ public void processElement( public static class RedistributeArbitrarily extends PTransform, PCollection> { // The number of buckets to shard into. - // A runner is free to ignore this (a runner may ignore the transorm + // A runner is free to ignore this (a runner may ignore the transform // entirely!) This is a performance optimization to prevent having // unit sized bundles on the output. If unset, uses a random integer key. private @Nullable Integer numBuckets = null; private boolean allowDuplicates = false; + private boolean deterministicSharding = false; - private RedistributeArbitrarily(@Nullable Integer numBuckets, boolean allowDuplicates) { + private RedistributeArbitrarily( + @Nullable Integer numBuckets, boolean allowDuplicates, boolean deterministicSharding) { this.numBuckets = numBuckets; this.allowDuplicates = allowDuplicates; + this.deterministicSharding = deterministicSharding; } public RedistributeArbitrarily withNumBuckets(@Nullable Integer numBuckets) { - return new RedistributeArbitrarily<>(numBuckets, this.allowDuplicates); + return new RedistributeArbitrarily<>( + numBuckets, this.allowDuplicates, this.deterministicSharding); } public RedistributeArbitrarily withAllowDuplicates(boolean allowDuplicates) { - return new RedistributeArbitrarily<>(this.numBuckets, allowDuplicates); + return new RedistributeArbitrarily<>( + this.numBuckets, allowDuplicates, this.deterministicSharding); + } + + public RedistributeArbitrarily withDeterministicSharding(boolean deterministicSharding) { + return new RedistributeArbitrarily<>( + this.numBuckets, this.allowDuplicates, deterministicSharding); } public boolean getAllowDuplicates() { @@ -157,7 +167,9 @@ public boolean getAllowDuplicates() { @Override public PCollection expand(PCollection input) { return input - .apply("Pair with random key", ParDo.of(new AssignShardFn<>(numBuckets))) + .apply( + "Pair with key", + ParDo.of(new AssignShardFn<>(numBuckets, this.deterministicSharding))) .apply(Redistribute.byKey().withAllowDuplicates(this.allowDuplicates)) .apply(Values.create()); } @@ -191,21 +203,31 @@ public void processElement( } static class AssignShardFn extends DoFn> { - private int shard; + private int randomShard; private @Nullable Integer numBuckets; + private boolean deterministicSharding; - public AssignShardFn(@Nullable Integer numBuckets) { + public AssignShardFn(@Nullable Integer numBuckets, boolean deterministicSharding) { this.numBuckets = numBuckets; + this.deterministicSharding = deterministicSharding; + this.randomShard = 0; } @Setup public void setup() { - shard = ThreadLocalRandom.current().nextInt(); + if (deterministicSharding) { + randomShard = ThreadLocalRandom.current().nextInt(); + } } @ProcessElement public void processElement(@Element T element, OutputReceiver> r) { - ++shard; + int shard = 0; + if (deterministicSharding && element != null) { + shard = element.hashCode(); + } else { + shard = ++randomShard; + } // Smear the shard into something more random-looking, to avoid issues // with runners that don't properly hash the key being shuffled, but rely // on it being random-looking. E.g. Spark takes the Java hashCode() of keys, 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 045a74a8507e..f43a8610f317 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 @@ -79,6 +79,7 @@ import org.apache.beam.sdk.transforms.PTransform; import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.transforms.Redistribute; +import org.apache.beam.sdk.transforms.Redistribute.RedistributeArbitrarily; import org.apache.beam.sdk.transforms.Reshuffle; import org.apache.beam.sdk.transforms.SerializableFunction; import org.apache.beam.sdk.transforms.SimpleFunction; @@ -1858,18 +1859,16 @@ public PCollection> expand(PBegin input) { "Offsets committed due to usage of commitOffsetsInFinalize() and may not capture all work processed due to use of withRedistribute() with duplicates enabled"); } - if (kafkaRead.getRedistributeNumKeys() == 0) { - return output.apply( - "Insert Redistribute", - Redistribute.>arbitrarily() - .withAllowDuplicates(kafkaRead.isAllowDuplicates())); - } else { - return output.apply( - "Insert Redistribute with Shards", - Redistribute.>arbitrarily() - .withAllowDuplicates(kafkaRead.isAllowDuplicates()) - .withNumBuckets((int) kafkaRead.getRedistributeNumKeys())); + RedistributeArbitrarily> redistribute = + Redistribute.>arbitrarily() + .withAllowDuplicates(kafkaRead.isAllowDuplicates()); + if (kafkaRead.getRedistributeNumKeys() != 0) { + redistribute = redistribute.withNumBuckets((int) kafkaRead.getRedistributeNumKeys()); + } + if (kafkaRead.getOffsetDeduplication() != null && kafkaRead.getOffsetDeduplication()) { + redistribute = redistribute.withDeterministicSharding(true); } + return output.apply("Redistribute", redistribute); } return output; } From a84ab497a771ae546eef0dd0c7ddf5d3fb55d9f5 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Thu, 11 Sep 2025 09:10:07 -0700 Subject: [PATCH 02/19] Address PR feedback. --- .../beam/sdk/transforms/Redistribute.java | 84 ++++++++++++------- 1 file changed, 52 insertions(+), 32 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java index 6f9f4a2180c2..6ee75174ce9d 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java @@ -36,6 +36,7 @@ import org.apache.beam.sdk.values.ValueInSingleWindow; import org.apache.beam.sdk.values.WindowingStrategy; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedInteger; import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Duration; @@ -166,10 +167,16 @@ public boolean getAllowDuplicates() { @Override public PCollection expand(PCollection input) { - return input - .apply( - "Pair with key", - ParDo.of(new AssignShardFn<>(numBuckets, this.deterministicSharding))) + PCollection> sharded; + if (deterministicSharding) { + sharded = + input.apply( + "Pair with deterministic key", + ParDo.of(new AssignDeterministicShardFn(numBuckets))); + } else { + sharded = input.apply("Pair with random key", ParDo.of(new AssignShardFn(numBuckets))); + } + return sharded .apply(Redistribute.byKey().withAllowDuplicates(this.allowDuplicates)) .apply(Values.create()); } @@ -202,46 +209,59 @@ public void processElement( } } + static class Sharding { + static int Smear(int shard) { + // Smear the shard into something more random-looking, to avoid issues + // with runners that don't properly hash the key being shuffled, but rely + // on it being random-looking. E.g. Spark takes the Java hashCode() of keys, + // which for Integer is a no-op, and it is an issue: + // http://hydronitrogen.com/poor-hash-partitioning-of-timestamps-integers-and-longs-in- + // spark.html + // This hashing strategy is copied from + // org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Hashing.smear(). + return 0x1b873593 * Integer.rotateLeft(shard * 0xcc9e2d51, 15); + } + + static int Bucket(int hash, @Nullable Integer numBuckets) { + if (numBuckets == null) return hash; + UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); + return UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue(); + } + } + static class AssignShardFn extends DoFn> { - private int randomShard; + private int shard; private @Nullable Integer numBuckets; - private boolean deterministicSharding; - public AssignShardFn(@Nullable Integer numBuckets, boolean deterministicSharding) { + public AssignShardFn(@Nullable Integer numBuckets) { this.numBuckets = numBuckets; - this.deterministicSharding = deterministicSharding; - this.randomShard = 0; } @Setup public void setup() { - if (deterministicSharding) { - randomShard = ThreadLocalRandom.current().nextInt(); - } + shard = ThreadLocalRandom.current().nextInt(); } @ProcessElement public void processElement(@Element T element, OutputReceiver> r) { - int shard = 0; - if (deterministicSharding && element != null) { - shard = element.hashCode(); - } else { - shard = ++randomShard; - } - // Smear the shard into something more random-looking, to avoid issues - // with runners that don't properly hash the key being shuffled, but rely - // on it being random-looking. E.g. Spark takes the Java hashCode() of keys, - // which for Integer is a no-op and it is an issue: - // http://hydronitrogen.com/poor-hash-partitioning-of-timestamps-integers-and-longs-in- - // spark.html - // This hashing strategy is copied from - // org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Hashing.smear(). - int hashOfShard = 0x1b873593 * Integer.rotateLeft(shard * 0xcc9e2d51, 15); - if (numBuckets != null) { - UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); - hashOfShard = UnsignedInteger.fromIntBits(hashOfShard).mod(unsignedNumBuckets).intValue(); - } - r.output(KV.of(hashOfShard, element)); + int hash = Sharding.Smear(++shard); + hash = Sharding.Bucket(hash, numBuckets); + r.output(KV.of(hash, element)); + } + } + + static class AssignDeterministicShardFn extends DoFn> { + private @Nullable Integer numBuckets; + + public AssignDeterministicShardFn(@Nullable Integer numBuckets) { + this.numBuckets = numBuckets; + } + + @ProcessElement + public void processElement(ProcessContext context) { + int hash = Hashing.farmHashFingerprint64().hashLong(context.currentRecordOffset()).asInt(); + hash = Sharding.Bucket(hash, numBuckets); + context.output(KV.of(hash, context.element())); } } From 0627d9483831a0829036cea016badf1e44d5436f Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Thu, 11 Sep 2025 09:15:42 -0700 Subject: [PATCH 03/19] Provide more detailed transform name for the redistribute. --- .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 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 f43a8610f317..dee432a84460 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 @@ -1862,13 +1862,18 @@ public PCollection> expand(PBegin input) { RedistributeArbitrarily> redistribute = Redistribute.>arbitrarily() .withAllowDuplicates(kafkaRead.isAllowDuplicates()); - if (kafkaRead.getRedistributeNumKeys() != 0) { - redistribute = redistribute.withNumBuckets((int) kafkaRead.getRedistributeNumKeys()); - } + StringBuilder redistributeName = new StringBuilder("Redistribute"); if (kafkaRead.getOffsetDeduplication() != null && kafkaRead.getOffsetDeduplication()) { redistribute = redistribute.withDeterministicSharding(true); + redistributeName.append(" deterministically"); + } else { + redistributeName.append(" randomly"); + } + if (kafkaRead.getRedistributeNumKeys() != 0) { + redistribute = redistribute.withNumBuckets((int) kafkaRead.getRedistributeNumKeys()); + redistributeName.append(" with bucketing"); } - return output.apply("Redistribute", redistribute); + return output.apply(redistributeName.toString(), redistribute); } return output; } From bc961654c056af85bd4fc41f80cb03fb66c14e51 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Thu, 11 Sep 2025 10:06:58 -0700 Subject: [PATCH 04/19] Address spotless precommit findings. --- .../org/apache/beam/sdk/transforms/Redistribute.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java index 6ee75174ce9d..16ba7138c37e 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java @@ -210,7 +210,7 @@ public void processElement( } static class Sharding { - static int Smear(int shard) { + static int smear(int shard) { // Smear the shard into something more random-looking, to avoid issues // with runners that don't properly hash the key being shuffled, but rely // on it being random-looking. E.g. Spark takes the Java hashCode() of keys, @@ -222,8 +222,8 @@ static int Smear(int shard) { return 0x1b873593 * Integer.rotateLeft(shard * 0xcc9e2d51, 15); } - static int Bucket(int hash, @Nullable Integer numBuckets) { - if (numBuckets == null) return hash; + static int bucket(int hash, @Nullable Integer numBuckets) { + if (numBuckets == null) { return hash; } UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); return UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue(); } @@ -244,8 +244,8 @@ public void setup() { @ProcessElement public void processElement(@Element T element, OutputReceiver> r) { - int hash = Sharding.Smear(++shard); - hash = Sharding.Bucket(hash, numBuckets); + int hash = Sharding.smear(++shard); + hash = Sharding.bucket(hash, numBuckets); r.output(KV.of(hash, element)); } } @@ -260,7 +260,7 @@ public AssignDeterministicShardFn(@Nullable Integer numBuckets) { @ProcessElement public void processElement(ProcessContext context) { int hash = Hashing.farmHashFingerprint64().hashLong(context.currentRecordOffset()).asInt(); - hash = Sharding.Bucket(hash, numBuckets); + hash = Sharding.bucket(hash, numBuckets); context.output(KV.of(hash, context.element())); } } From 65621c8770f1abb21d71b96d2aaa413c5f1c9f88 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Thu, 11 Sep 2025 11:57:37 -0700 Subject: [PATCH 05/19] Address spotless precommit findings. --- .../java/org/apache/beam/sdk/transforms/Redistribute.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java index 16ba7138c37e..fe1f672f5a93 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java @@ -223,7 +223,9 @@ static int smear(int shard) { } static int bucket(int hash, @Nullable Integer numBuckets) { - if (numBuckets == null) { return hash; } + if (numBuckets == null) { + return hash; + } UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); return UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue(); } From a611d9b4c91445cbfbed7ba77fadaa6b23e33db3 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Tue, 16 Sep 2025 13:30:16 -0500 Subject: [PATCH 06/19] Keep redistribute transform name the same. --- .../main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 7 +------ 1 file changed, 1 insertion(+), 6 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 dee432a84460..dcc881fc1851 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 @@ -1862,18 +1862,13 @@ public PCollection> expand(PBegin input) { RedistributeArbitrarily> redistribute = Redistribute.>arbitrarily() .withAllowDuplicates(kafkaRead.isAllowDuplicates()); - StringBuilder redistributeName = new StringBuilder("Redistribute"); if (kafkaRead.getOffsetDeduplication() != null && kafkaRead.getOffsetDeduplication()) { redistribute = redistribute.withDeterministicSharding(true); - redistributeName.append(" deterministically"); - } else { - redistributeName.append(" randomly"); } if (kafkaRead.getRedistributeNumKeys() != 0) { redistribute = redistribute.withNumBuckets((int) kafkaRead.getRedistributeNumKeys()); - redistributeName.append(" with bucketing"); } - return output.apply(redistributeName.toString(), redistribute); + return output.apply("Insert Redistribute", redistribute); } return output; } From 1ba7206a6baa659fd148b0a29f3a20ce77d87840 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Tue, 16 Sep 2025 16:51:52 -0500 Subject: [PATCH 07/19] Add deterministic sharding unit test. --- .../beam/sdk/transforms/Redistribute.java | 4 ++++ .../beam/sdk/transforms/RedistributeTest.java | 23 +++++++++++++++++++ 2 files changed, 27 insertions(+) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java index fe1f672f5a93..e6dd9d6b819c 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java @@ -165,6 +165,10 @@ public boolean getAllowDuplicates() { return allowDuplicates; } + public boolean getDeterministicSharding() { + return deterministicSharding; + } + @Override public PCollection expand(PCollection input) { PCollection> sharded; diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/RedistributeTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/RedistributeTest.java index ea46ffec4496..3f69968d5f33 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/RedistributeTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/RedistributeTest.java @@ -17,6 +17,7 @@ */ package org.apache.beam.sdk.transforms; +import static junit.framework.TestCase.assertTrue; import static org.apache.beam.sdk.TestUtils.KvMatcher.isKv; import static org.apache.beam.sdk.util.construction.PTransformTranslation.REDISTRIBUTE_ARBITRARILY_URN; import static org.apache.beam.sdk.util.construction.PTransformTranslation.REDISTRIBUTE_BY_KEY_URN; @@ -39,6 +40,7 @@ import org.apache.beam.sdk.testing.UsesTestStream; import org.apache.beam.sdk.testing.ValidatesRunner; import org.apache.beam.sdk.transforms.Redistribute.AssignShardFn; +import org.apache.beam.sdk.transforms.Redistribute.RedistributeArbitrarily; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; import org.apache.beam.sdk.transforms.windowing.FixedWindows; import org.apache.beam.sdk.transforms.windowing.GlobalWindow; @@ -123,6 +125,27 @@ public void testJustRedistribute() { pipeline.run(); } + @Test + @Category(ValidatesRunner.class) + public void testDeterministicRedistribute() { + List inputValues = Lists.newArrayList(); + for (int i = 0; i < 10; i++) { + inputValues.add(i); + } + + PCollection input = pipeline.apply(Create.of(inputValues).withCoder(VarIntCoder.of())); + RedistributeArbitrarily redistribute = + Redistribute.arbitrarily().withDeterministicSharding(true); + assertTrue(redistribute.getDeterministicSharding()); + + PCollection output = input.apply(redistribute); + PAssert.that(output).containsInAnyOrder(inputValues); + + assertEquals(input.getWindowingStrategy(), output.getWindowingStrategy()); + + pipeline.run(); + } + /** * Tests that timestamps are preserved after applying a {@link Redistribute} with the default * {@link WindowingStrategy}. From b4abac100b611be9b47a0877b46c805af76593b4 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Fri, 19 Sep 2025 10:22:49 -0700 Subject: [PATCH 08/19] Refactor to specific deterministic Kafka redistribute method. --- .../beam/sdk/transforms/Redistribute.java | 90 +++++-------------- .../beam/sdk/transforms/RedistributeTest.java | 23 ----- .../org/apache/beam/sdk/io/kafka/KafkaIO.java | 13 ++- .../sdk/io/kafka/KafkaReadRedistribute.java | 79 ++++++++++++++++ .../io/kafka/KafkaReadRedistributeTest.java | 84 +++++++++++++++++ 5 files changed, 193 insertions(+), 96 deletions(-) create mode 100644 sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java create mode 100644 sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java index e6dd9d6b819c..3a8bef28839a 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java @@ -36,7 +36,6 @@ import org.apache.beam.sdk.values.ValueInSingleWindow; import org.apache.beam.sdk.values.WindowingStrategy; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedInteger; import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Duration; @@ -49,7 +48,7 @@ public class Redistribute { /** @return a {@link RedistributeArbitrarily} transform with default configuration. */ public static RedistributeArbitrarily arbitrarily() { - return new RedistributeArbitrarily<>(null, false, false); + return new RedistributeArbitrarily<>(null, false); } /** @return a {@link RedistributeByKey} transform with default configuration. */ @@ -137,50 +136,28 @@ public static class RedistributeArbitrarily // unit sized bundles on the output. If unset, uses a random integer key. private @Nullable Integer numBuckets = null; private boolean allowDuplicates = false; - private boolean deterministicSharding = false; - private RedistributeArbitrarily( - @Nullable Integer numBuckets, boolean allowDuplicates, boolean deterministicSharding) { + private RedistributeArbitrarily(@Nullable Integer numBuckets, boolean allowDuplicates) { this.numBuckets = numBuckets; this.allowDuplicates = allowDuplicates; - this.deterministicSharding = deterministicSharding; } public RedistributeArbitrarily withNumBuckets(@Nullable Integer numBuckets) { - return new RedistributeArbitrarily<>( - numBuckets, this.allowDuplicates, this.deterministicSharding); + return new RedistributeArbitrarily<>(numBuckets, this.allowDuplicates); } public RedistributeArbitrarily withAllowDuplicates(boolean allowDuplicates) { - return new RedistributeArbitrarily<>( - this.numBuckets, allowDuplicates, this.deterministicSharding); - } - - public RedistributeArbitrarily withDeterministicSharding(boolean deterministicSharding) { - return new RedistributeArbitrarily<>( - this.numBuckets, this.allowDuplicates, deterministicSharding); + return new RedistributeArbitrarily<>(this.numBuckets, allowDuplicates); } public boolean getAllowDuplicates() { return allowDuplicates; } - public boolean getDeterministicSharding() { - return deterministicSharding; - } - @Override public PCollection expand(PCollection input) { - PCollection> sharded; - if (deterministicSharding) { - sharded = - input.apply( - "Pair with deterministic key", - ParDo.of(new AssignDeterministicShardFn(numBuckets))); - } else { - sharded = input.apply("Pair with random key", ParDo.of(new AssignShardFn(numBuckets))); - } - return sharded + return input + .apply("Pair with random key", ParDo.of(new AssignShardFn<>(numBuckets))) .apply(Redistribute.byKey().withAllowDuplicates(this.allowDuplicates)) .apply(Values.create()); } @@ -213,28 +190,6 @@ public void processElement( } } - static class Sharding { - static int smear(int shard) { - // Smear the shard into something more random-looking, to avoid issues - // with runners that don't properly hash the key being shuffled, but rely - // on it being random-looking. E.g. Spark takes the Java hashCode() of keys, - // which for Integer is a no-op, and it is an issue: - // http://hydronitrogen.com/poor-hash-partitioning-of-timestamps-integers-and-longs-in- - // spark.html - // This hashing strategy is copied from - // org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Hashing.smear(). - return 0x1b873593 * Integer.rotateLeft(shard * 0xcc9e2d51, 15); - } - - static int bucket(int hash, @Nullable Integer numBuckets) { - if (numBuckets == null) { - return hash; - } - UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); - return UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue(); - } - } - static class AssignShardFn extends DoFn> { private int shard; private @Nullable Integer numBuckets; @@ -250,24 +205,21 @@ public void setup() { @ProcessElement public void processElement(@Element T element, OutputReceiver> r) { - int hash = Sharding.smear(++shard); - hash = Sharding.bucket(hash, numBuckets); - r.output(KV.of(hash, element)); - } - } - - static class AssignDeterministicShardFn extends DoFn> { - private @Nullable Integer numBuckets; - - public AssignDeterministicShardFn(@Nullable Integer numBuckets) { - this.numBuckets = numBuckets; - } - - @ProcessElement - public void processElement(ProcessContext context) { - int hash = Hashing.farmHashFingerprint64().hashLong(context.currentRecordOffset()).asInt(); - hash = Sharding.bucket(hash, numBuckets); - context.output(KV.of(hash, context.element())); + ++shard; + // Smear the shard into something more random-looking, to avoid issues + // with runners that don't properly hash the key being shuffled, but rely + // on it being random-looking. E.g. Spark takes the Java hashCode() of keys, + // which for Integer is a no-op and it is an issue: + // http://hydronitrogen.com/poor-hash-partitioning-of-timestamps-integers-and-longs-in- + // spark.html + // This hashing strategy is copied from + // org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Hashing.smear(). + int hashOfShard = 0x1b873593 * Integer.rotateLeft(shard * 0xcc9e2d51, 15); + if (numBuckets != null) { + UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); + hashOfShard = UnsignedInteger.fromIntBits(hashOfShard).mod(unsignedNumBuckets).intValue(); + } + r.output(KV.of(hashOfShard, element)); } } diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/RedistributeTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/RedistributeTest.java index 3f69968d5f33..ea46ffec4496 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/RedistributeTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/RedistributeTest.java @@ -17,7 +17,6 @@ */ package org.apache.beam.sdk.transforms; -import static junit.framework.TestCase.assertTrue; import static org.apache.beam.sdk.TestUtils.KvMatcher.isKv; import static org.apache.beam.sdk.util.construction.PTransformTranslation.REDISTRIBUTE_ARBITRARILY_URN; import static org.apache.beam.sdk.util.construction.PTransformTranslation.REDISTRIBUTE_BY_KEY_URN; @@ -40,7 +39,6 @@ import org.apache.beam.sdk.testing.UsesTestStream; import org.apache.beam.sdk.testing.ValidatesRunner; import org.apache.beam.sdk.transforms.Redistribute.AssignShardFn; -import org.apache.beam.sdk.transforms.Redistribute.RedistributeArbitrarily; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; import org.apache.beam.sdk.transforms.windowing.FixedWindows; import org.apache.beam.sdk.transforms.windowing.GlobalWindow; @@ -125,27 +123,6 @@ public void testJustRedistribute() { pipeline.run(); } - @Test - @Category(ValidatesRunner.class) - public void testDeterministicRedistribute() { - List inputValues = Lists.newArrayList(); - for (int i = 0; i < 10; i++) { - inputValues.add(i); - } - - PCollection input = pipeline.apply(Create.of(inputValues).withCoder(VarIntCoder.of())); - RedistributeArbitrarily redistribute = - Redistribute.arbitrarily().withDeterministicSharding(true); - assertTrue(redistribute.getDeterministicSharding()); - - PCollection output = input.apply(redistribute); - PAssert.that(output).containsInAnyOrder(inputValues); - - assertEquals(input.getWindowingStrategy(), output.getWindowingStrategy()); - - pipeline.run(); - } - /** * Tests that timestamps are preserved after applying a {@link Redistribute} with the default * {@link WindowingStrategy}. 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 dcc881fc1851..e0c87ba2d582 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 @@ -1859,16 +1859,21 @@ public PCollection> expand(PBegin input) { "Offsets committed due to usage of commitOffsetsInFinalize() and may not capture all work processed due to use of withRedistribute() with duplicates enabled"); } + if (kafkaRead.getOffsetDeduplication() != null && kafkaRead.getOffsetDeduplication()) { + return output.apply( + KafkaReadRedistribute.redistribute() + .withNumBuckets(kafkaRead.getRedistributeNumKeys())); + } + RedistributeArbitrarily> redistribute = Redistribute.>arbitrarily() .withAllowDuplicates(kafkaRead.isAllowDuplicates()); - if (kafkaRead.getOffsetDeduplication() != null && kafkaRead.getOffsetDeduplication()) { - redistribute = redistribute.withDeterministicSharding(true); - } + String redistributeName = "Insert Redistribute"; if (kafkaRead.getRedistributeNumKeys() != 0) { redistribute = redistribute.withNumBuckets((int) kafkaRead.getRedistributeNumKeys()); + redistributeName = "Insert Redistribute with Shards"; } - return output.apply("Insert Redistribute", redistribute); + return output.apply(redistributeName, redistribute); } return output; } diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java new file mode 100644 index 000000000000..c6b0d19d150f --- /dev/null +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java @@ -0,0 +1,79 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.sdk.io.kafka; + +import org.apache.beam.sdk.transforms.DoFn; +import org.apache.beam.sdk.transforms.PTransform; +import org.apache.beam.sdk.transforms.ParDo; +import org.apache.beam.sdk.transforms.Redistribute; +import org.apache.beam.sdk.transforms.Values; +import org.apache.beam.sdk.values.KV; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedInteger; +import org.checkerframework.checker.nullness.qual.Nullable; + +public class KafkaReadRedistribute + extends PTransform>, PCollection>> { + public static KafkaReadRedistribute redistribute() { + return new KafkaReadRedistribute<>(null); + } + + // The number of buckets to shard into. + private @Nullable Integer numBuckets = null; + + private KafkaReadRedistribute(@Nullable Integer numBuckets) { + this.numBuckets = numBuckets; + } + + public KafkaReadRedistribute withNumBuckets(@Nullable Integer numBuckets) { + return new KafkaReadRedistribute<>(numBuckets); + } + + @Override + public PCollection> expand(PCollection> input) { + PCollection>> sharded = + input.apply("Pair with deterministic key", ParDo.of(new AssignShardFn(numBuckets))); + + return sharded + .apply(Redistribute.>byKey().withAllowDuplicates(false)) + .apply(Values.create()); + } + + static class AssignShardFn extends DoFn, KV>> { + private @Nullable Integer numBuckets; + + public AssignShardFn(@Nullable Integer numBuckets) { + this.numBuckets = numBuckets; + } + + @ProcessElement + public void processElement( + @Element KafkaRecord element, + OutputReceiver>> receiver) { + int hash = Hashing.farmHashFingerprint64().hashLong(element.getOffset()).asInt(); + + if (numBuckets != null && numBuckets > 0) { + UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); + hash = UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue(); + } + + receiver.output(KV.of(hash, element)); + } + } +} diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java new file mode 100644 index 000000000000..789ed670a060 --- /dev/null +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java @@ -0,0 +1,84 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.sdk.io.kafka; + +import static org.apache.beam.sdk.io.kafka.KafkaTimestampType.LOG_APPEND_TIME; +import static org.junit.Assert.assertEquals; + +import java.io.Serializable; +import org.apache.beam.sdk.coders.StringUtf8Coder; +import org.apache.beam.sdk.coders.VarIntCoder; +import org.apache.beam.sdk.testing.PAssert; +import org.apache.beam.sdk.testing.TestPipeline; +import org.apache.beam.sdk.testing.ValidatesRunner; +import org.apache.beam.sdk.transforms.Create; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.junit.Rule; +import org.junit.Test; +import org.junit.experimental.categories.Category; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +/** Tests for {@link KafkaReadRedistribute}. */ +@RunWith(JUnit4.class) +public class KafkaReadRedistributeTest implements Serializable { + + private static final ImmutableList> INPUTS = + ImmutableList.of( + MakeKafkaRecord("k1", 3, 1), + MakeKafkaRecord("k5", Integer.MAX_VALUE, 2), + MakeKafkaRecord("k5", Integer.MIN_VALUE, 3), + MakeKafkaRecord("k2", 66, 4), + MakeKafkaRecord("k1", 4, 5), + MakeKafkaRecord("k2", -33, 6), + MakeKafkaRecord("k3", 0, 7)); + + static KafkaRecord MakeKafkaRecord(String key, Integer value, Integer offset) { + return new KafkaRecord( + /*topic*/ "kafka", + /*partition*/ 1, + /*offset*/ offset, + /*timestamp*/ 123, + /*timestampType*/ LOG_APPEND_TIME, + /*headers*/ null, + key, + value); + } + + @Rule public final transient TestPipeline pipeline = TestPipeline.create(); + + @Test + @Category(ValidatesRunner.class) + public void testJustRedistribute() { + + PCollection> input = + pipeline.apply( + Create.of(INPUTS) + .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); + + PCollection> output = + input.apply(KafkaReadRedistribute.redistribute()); + + PAssert.that(output).containsInAnyOrder(INPUTS); + + assertEquals(input.getWindowingStrategy(), output.getWindowingStrategy()); + + pipeline.run(); + } +} From b36e2ec8b0c804b7950d7c05838206a038d1b7b4 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Fri, 19 Sep 2025 11:33:55 -0700 Subject: [PATCH 09/19] Add redistribute by key variant. --- .../org/apache/beam/sdk/io/kafka/KafkaIO.java | 4 +- .../sdk/io/kafka/KafkaReadRedistribute.java | 45 ++++++++++++++----- .../io/kafka/KafkaReadRedistributeTest.java | 23 +++++++++- 3 files changed, 56 insertions(+), 16 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 e0c87ba2d582..413529ee2b29 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 @@ -1860,9 +1860,9 @@ public PCollection> expand(PBegin input) { } if (kafkaRead.getOffsetDeduplication() != null && kafkaRead.getOffsetDeduplication()) { + // TODO: Expose Kafka read options to control byOffsetShard vs byRecordKey. return output.apply( - KafkaReadRedistribute.redistribute() - .withNumBuckets(kafkaRead.getRedistributeNumKeys())); + KafkaReadRedistribute.byOffsetShard(kafkaRead.getRedistributeNumKeys())); } RedistributeArbitrarily> redistribute = diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java index c6b0d19d150f..3e1048718f6f 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java @@ -30,35 +30,45 @@ public class KafkaReadRedistribute extends PTransform>, PCollection>> { - public static KafkaReadRedistribute redistribute() { - return new KafkaReadRedistribute<>(null); + public static KafkaReadRedistribute byOffsetShard(@Nullable Integer numBuckets) { + return new KafkaReadRedistribute<>(numBuckets, false); + } + + public static KafkaReadRedistribute byRecordKey() { + return new KafkaReadRedistribute<>(null, true); } // The number of buckets to shard into. private @Nullable Integer numBuckets = null; + // When redistributing, group records by the Kafka record's key instead of by offset hash. + private boolean byRecordKey = false; - private KafkaReadRedistribute(@Nullable Integer numBuckets) { + private KafkaReadRedistribute(@Nullable Integer numBuckets, boolean byRecordKey) { this.numBuckets = numBuckets; - } - - public KafkaReadRedistribute withNumBuckets(@Nullable Integer numBuckets) { - return new KafkaReadRedistribute<>(numBuckets); + this.byRecordKey = byRecordKey; } @Override public PCollection> expand(PCollection> input) { - PCollection>> sharded = - input.apply("Pair with deterministic key", ParDo.of(new AssignShardFn(numBuckets))); - return sharded + if (byRecordKey) { + return input + .apply("Pair with record key", ParDo.of(new AssignKeyFn())) + .apply(Redistribute.>byKey().withAllowDuplicates(false)) + .apply(Values.create()); + } + + return input + .apply("Pair with offset shard", ParDo.of(new AssignOffsetShardFn(numBuckets))) .apply(Redistribute.>byKey().withAllowDuplicates(false)) .apply(Values.create()); } - static class AssignShardFn extends DoFn, KV>> { + static class AssignOffsetShardFn + extends DoFn, KV>> { private @Nullable Integer numBuckets; - public AssignShardFn(@Nullable Integer numBuckets) { + public AssignOffsetShardFn(@Nullable Integer numBuckets) { this.numBuckets = numBuckets; } @@ -76,4 +86,15 @@ public void processElement( receiver.output(KV.of(hash, element)); } } + + static class AssignKeyFn extends DoFn, KV>> { + + public AssignKeyFn() {} + + @ProcessElement + public void processElement( + @Element KafkaRecord element, OutputReceiver>> receiver) { + receiver.output(KV.of(element.getKV().getKey(), element)); + } + } } diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java index 789ed670a060..2a429c45cb46 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java @@ -65,7 +65,7 @@ static KafkaRecord MakeKafkaRecord(String key, Integer value, I @Test @Category(ValidatesRunner.class) - public void testJustRedistribute() { + public void testRedistributeByOffsetShard() { PCollection> input = pipeline.apply( @@ -73,7 +73,26 @@ public void testJustRedistribute() { .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); PCollection> output = - input.apply(KafkaReadRedistribute.redistribute()); + input.apply(KafkaReadRedistribute.byOffsetShard(/*numBuckets*/ 10)); + + PAssert.that(output).containsInAnyOrder(INPUTS); + + assertEquals(input.getWindowingStrategy(), output.getWindowingStrategy()); + + pipeline.run(); + } + + @Test + @Category(ValidatesRunner.class) + public void testRedistributeByKey() { + + PCollection> input = + pipeline.apply( + Create.of(INPUTS) + .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); + + PCollection> output = + input.apply(KafkaReadRedistribute.byRecordKey()); PAssert.that(output).containsInAnyOrder(INPUTS); From 449603062cf9a5dbad1b277dc61dc8564cb753f5 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Fri, 19 Sep 2025 11:44:24 -0700 Subject: [PATCH 10/19] Add test of sharding fns. --- .../sdk/io/kafka/KafkaReadRedistribute.java | 6 +- .../io/kafka/KafkaReadRedistributeTest.java | 58 +++++++++++++++++++ 2 files changed, 61 insertions(+), 3 deletions(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java index 3e1048718f6f..5d58b5ee7e82 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java @@ -53,7 +53,7 @@ public PCollection> expand(PCollection> inpu if (byRecordKey) { return input - .apply("Pair with record key", ParDo.of(new AssignKeyFn())) + .apply("Pair with record key", ParDo.of(new AssignRecordKeyFn())) .apply(Redistribute.>byKey().withAllowDuplicates(false)) .apply(Values.create()); } @@ -87,9 +87,9 @@ public void processElement( } } - static class AssignKeyFn extends DoFn, KV>> { + static class AssignRecordKeyFn extends DoFn, KV>> { - public AssignKeyFn() {} + public AssignRecordKeyFn() {} @ProcessElement public void processElement( diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java index 2a429c45cb46..68b8130806c4 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java @@ -18,17 +18,27 @@ package org.apache.beam.sdk.io.kafka; import static org.apache.beam.sdk.io.kafka.KafkaTimestampType.LOG_APPEND_TIME; +import static org.apache.beam.sdk.values.TypeDescriptors.integers; +import static org.apache.beam.sdk.values.TypeDescriptors.strings; import static org.junit.Assert.assertEquals; import java.io.Serializable; +import java.util.List; import org.apache.beam.sdk.coders.StringUtf8Coder; import org.apache.beam.sdk.coders.VarIntCoder; +import org.apache.beam.sdk.io.kafka.KafkaReadRedistribute.AssignOffsetShardFn; +import org.apache.beam.sdk.io.kafka.KafkaReadRedistribute.AssignRecordKeyFn; import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; import org.apache.beam.sdk.testing.ValidatesRunner; import org.apache.beam.sdk.transforms.Create; +import org.apache.beam.sdk.transforms.GroupByKey; +import org.apache.beam.sdk.transforms.MapElements; +import org.apache.beam.sdk.transforms.ParDo; +import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.PCollection; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists; import org.junit.Rule; import org.junit.Test; import org.junit.experimental.categories.Category; @@ -100,4 +110,52 @@ public void testRedistributeByKey() { pipeline.run(); } + + @Test + @Category({ValidatesRunner.class}) + public void testAssignOutputShardFn() { + List> inputs = Lists.newArrayList(); + for (int i = 0; i < 10; i++) { + inputs.addAll(INPUTS); + } + + PCollection> input = + pipeline.apply( + Create.of(inputs) + .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); + + PCollection output = + input + .apply(ParDo.of(new AssignOffsetShardFn(2))) + .apply(GroupByKey.create()) + .apply(MapElements.into(integers()).via(KV::getKey)); + + PAssert.that(output).containsInAnyOrder(ImmutableList.of(0, 1)); + + pipeline.run(); + } + + @Test + @Category({ValidatesRunner.class}) + public void testAssignRecordKeyFn() { + List> inputs = Lists.newArrayList(); + for (int i = 0; i < 10; i++) { + inputs.addAll(INPUTS); + } + + PCollection> input = + pipeline.apply( + Create.of(inputs) + .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); + + PCollection output = + input + .apply(ParDo.of(new AssignRecordKeyFn())) + .apply(GroupByKey.create()) + .apply(MapElements.into(strings()).via(KV::getKey)); + + PAssert.that(output).containsInAnyOrder(ImmutableList.of("k1", "k2", "k3", "k5")); + + pipeline.run(); + } } From 99d8b787da95fe65726adc7652d77a9353570358 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Mon, 22 Sep 2025 15:40:04 -0700 Subject: [PATCH 11/19] Add bucketing to redistributeByKey and add option to redistribute by key. --- .../org/apache/beam/sdk/io/kafka/KafkaIO.java | 35 ++++++++- .../sdk/io/kafka/KafkaReadRedistribute.java | 33 +++++++-- ...IOReadImplementationCompatibilityTest.java | 3 +- .../apache/beam/sdk/io/kafka/KafkaIOTest.java | 73 ++++++++++++++++--- .../io/kafka/KafkaReadRedistributeTest.java | 11 ++- .../io/kafka/upgrade/KafkaIOTranslation.java | 10 +++ .../kafka/upgrade/KafkaIOTranslationTest.java | 1 + 7 files changed, 138 insertions(+), 28 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 413529ee2b29..dcd0ac3daaf0 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 @@ -731,6 +731,9 @@ public abstract static class Read @Pure public abstract @Nullable Boolean getOffsetDeduplication(); + @Pure + public abstract @Nullable Boolean getRedistributeByRecordKey(); + @Pure public abstract @Nullable Duration getWatchTopicPartitionDuration(); @@ -801,6 +804,8 @@ abstract Builder setConsumerFactoryFn( abstract Builder setOffsetDeduplication(Boolean offsetDeduplication); + abstract Builder setRedistributeByRecordKey(Boolean redistributeByRecordKey); + abstract Builder setTimestampPolicyFactory( TimestampPolicyFactory timestampPolicyFactory); @@ -916,11 +921,15 @@ static void setupExternalBuilder( && config.offsetDeduplication != null) { builder.setOffsetDeduplication(config.offsetDeduplication); } + if (config.redistribute && config.redistributeByRecordKey != null) { + builder.setRedistributeByRecordKey(config.redistributeByRecordKey); + } } else { builder.setRedistributed(false); builder.setRedistributeNumKeys(0); builder.setAllowDuplicates(false); builder.setOffsetDeduplication(false); + builder.setRedistributeByRecordKey(false); } } @@ -990,6 +999,7 @@ public static class Configuration { private Boolean redistribute; private Boolean allowDuplicates; private Boolean offsetDeduplication; + private Boolean redistributeByRecordKey; private Long dynamicReadPollIntervalSeconds; public void setConsumerConfig(Map consumerConfig) { @@ -1052,6 +1062,10 @@ public void setOffsetDeduplication(Boolean offsetDeduplication) { this.offsetDeduplication = offsetDeduplication; } + public void setRedistributeByRecordKey(Boolean redistributeByRecordKey) { + this.redistributeByRecordKey = redistributeByRecordKey; + } + public void setDynamicReadPollIntervalSeconds(Long dynamicReadPollIntervalSeconds) { this.dynamicReadPollIntervalSeconds = dynamicReadPollIntervalSeconds; } @@ -1162,6 +1176,10 @@ public Read withOffsetDeduplication(Boolean offsetDeduplication) { return toBuilder().setOffsetDeduplication(offsetDeduplication).build(); } + public Read withRedistributeByRecordKey(Boolean redistributeByRecordKey) { + return toBuilder().setRedistributeByRecordKey(redistributeByRecordKey).build(); + } + /** * Internally sets a {@link java.util.regex.Pattern} of topics to read from. All the partitions * from each of the matching topics are read. @@ -1680,6 +1698,11 @@ private void checkRedistributeConfiguration() { LOG.warn( "Offsets used for deduplication are available in WindowedValue's metadata. Combining, aggregating, mutating them may risk with data loss."); } + if (getRedistributeByRecordKey() != null && getRedistributeByRecordKey()) { + checkState( + isRedistributed(), + "withRedistributeByRecordKey can only be used when withRedistribute is set."); + } } private void warnAboutUnsafeConfigurations(PBegin input) { @@ -1860,11 +1883,15 @@ public PCollection> expand(PBegin input) { } if (kafkaRead.getOffsetDeduplication() != null && kafkaRead.getOffsetDeduplication()) { - // TODO: Expose Kafka read options to control byOffsetShard vs byRecordKey. - return output.apply( - KafkaReadRedistribute.byOffsetShard(kafkaRead.getRedistributeNumKeys())); + if (kafkaRead.getRedistributeByRecordKey() != null + && kafkaRead.getRedistributeByRecordKey()) { + return output.apply( + KafkaReadRedistribute.byRecordKey(kafkaRead.getRedistributeNumKeys())); + } else { + return output.apply( + KafkaReadRedistribute.byOffsetShard(kafkaRead.getRedistributeNumKeys())); + } } - RedistributeArbitrarily> redistribute = Redistribute.>arbitrarily() .withAllowDuplicates(kafkaRead.isAllowDuplicates()); diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java index 5d58b5ee7e82..3943d399cf24 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java @@ -17,6 +17,8 @@ */ package org.apache.beam.sdk.io.kafka; +import static java.nio.charset.StandardCharsets.UTF_8; + import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.PTransform; import org.apache.beam.sdk.transforms.ParDo; @@ -34,8 +36,8 @@ public static KafkaReadRedistribute byOffsetShard(@Nullable Integer return new KafkaReadRedistribute<>(numBuckets, false); } - public static KafkaReadRedistribute byRecordKey() { - return new KafkaReadRedistribute<>(null, true); + public static KafkaReadRedistribute byRecordKey(@Nullable Integer numBuckets) { + return new KafkaReadRedistribute<>(numBuckets, true); } // The number of buckets to shard into. @@ -53,8 +55,8 @@ public PCollection> expand(PCollection> inpu if (byRecordKey) { return input - .apply("Pair with record key", ParDo.of(new AssignRecordKeyFn())) - .apply(Redistribute.>byKey().withAllowDuplicates(false)) + .apply("Pair with record key", ParDo.of(new AssignRecordKeyFn(numBuckets))) + .apply(Redistribute.>byKey().withAllowDuplicates(false)) .apply(Values.create()); } @@ -87,14 +89,29 @@ public void processElement( } } - static class AssignRecordKeyFn extends DoFn, KV>> { + static class AssignRecordKeyFn + extends DoFn, KV>> { + + private @Nullable Integer numBuckets; - public AssignRecordKeyFn() {} + public AssignRecordKeyFn(@Nullable Integer numBuckets) { + this.numBuckets = numBuckets; + } @ProcessElement public void processElement( - @Element KafkaRecord element, OutputReceiver>> receiver) { - receiver.output(KV.of(element.getKV().getKey(), element)); + @Element KafkaRecord element, + OutputReceiver>> receiver) { + K key = element.getKV().getKey(); + String keyString = key == null ? "" : key.toString(); + int hash = Hashing.farmHashFingerprint64().hashBytes(keyString.getBytes(UTF_8)).asInt(); + + if (numBuckets != null && numBuckets > 0) { + UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); + hash = UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue(); + } + + receiver.output(KV.of(hash, element)); } } } diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java index 26682946afca..fb44be8db4b8 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java @@ -117,7 +117,8 @@ private PipelineResult testReadTransformCreationWithImplementationBoundPropertie false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/))); + null /*topics*/, + null /*redistributeByRecordKey*/))); return p.run(); } 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 83c2e1b38826..e07f067c50a5 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 @@ -395,7 +395,8 @@ static KafkaIO.Read mkKafkaReadTransform( false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/); + null /*topics*/, + null /*redistributeByRecordKey*/); } static KafkaIO.Read mkKafkaReadTransformWithOffsetDedup( @@ -408,7 +409,24 @@ static KafkaIO.Read mkKafkaReadTransformWithOffsetDedup( false, /*allowDuplicates*/ 100, /*numKeys*/ true, /*offsetDeduplication*/ - null /*topics*/); + null /*topics*/, + null /*redistributeByRecordKey*/); + } + + static KafkaIO.Read mkKafkaReadTransformWithRedistributeByRecordKey( + int numElements, + @Nullable SerializableFunction, Instant> timestampFn, + boolean byRecordKey) { + return mkKafkaReadTransform( + numElements, + numElements, + timestampFn, + true, /*redistribute*/ + false, /*allowDuplicates*/ + 100, /*numKeys*/ + true, /*offsetDeduplication*/ + null /*topics*/, + byRecordKey /*redistributeByRecordKey*/); } static KafkaIO.Read mkKafkaReadTransformWithTopics( @@ -423,7 +441,8 @@ static KafkaIO.Read mkKafkaReadTransformWithTopics( false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - topics /*topics*/); + topics /*topics*/, + null /*redistributeByRecordKey*/); } /** @@ -438,7 +457,8 @@ static KafkaIO.Read mkKafkaReadTransform( @Nullable Boolean withAllowDuplicates, @Nullable Integer numKeys, @Nullable Boolean offsetDeduplication, - @Nullable List topics) { + @Nullable List topics, + @Nullable Boolean redistributeByRecordKey) { KafkaIO.Read reader = KafkaIO.read() @@ -723,7 +743,8 @@ public void warningsWithAllowDuplicatesEnabledAndCommitOffsets() { true, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/) + null /*topics*/, + null /*redistributeByRecordKey*/) .commitOffsetsInFinalize() .withConsumerConfigUpdates( ImmutableMap.of(ConsumerConfig.GROUP_ID_CONFIG, "group_id")) @@ -751,7 +772,8 @@ public void noWarningsWithNoAllowDuplicatesAndCommitOffsets() { false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/) + null /*topics*/, + null /*redistributeByRecordKey*/) .commitOffsetsInFinalize() .withConsumerConfigUpdates( ImmutableMap.of(ConsumerConfig.GROUP_ID_CONFIG, "group_id")) @@ -780,7 +802,8 @@ public void testNumKeysIgnoredWithRedistributeNotEnabled() { false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/) + null /*topics*/, + null /*redistributeByRecordKey*/) .withRedistributeNumKeys(100) .commitOffsetsInFinalize() .withConsumerConfigUpdates( @@ -2200,7 +2223,8 @@ public void testUnboundedSourceStartReadTime() { false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/) + null /*topics*/, + null /*redistributeByRecordKey*/) .withStartReadTime(new Instant(startTime)) .withoutMetadata()) .apply(Values.create()); @@ -2223,6 +2247,36 @@ public void testOffsetDeduplication() { p.run(); } + @Test + public void testRedistributeByRecordKeyOn() { + int numElements = 1000; + + PCollection input = + p.apply( + mkKafkaReadTransformWithRedistributeByRecordKey( + numElements, new ValueAsTimestampFn(), true) + .withoutMetadata()) + .apply(Values.create()); + + addCountingAsserts(input, numElements, numElements, 0, numElements - 1); + p.run(); + } + + @Test + public void testRedistributeByRecordKeyOff() { + int numElements = 1000; + + PCollection input = + p.apply( + mkKafkaReadTransformWithRedistributeByRecordKey( + numElements, new ValueAsTimestampFn(), false) + .withoutMetadata()) + .apply(Values.create()); + + addCountingAsserts(input, numElements, numElements, 0, numElements - 1); + p.run(); + } + @Rule public ExpectedException noMessagesException = ExpectedException.none(); @Test @@ -2246,7 +2300,8 @@ public void testUnboundedSourceStartReadTimeException() { false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/) + null /*topics*/, + null /*redistributeByRecordKey*/) .withStartReadTime(new Instant(startTime)) .withoutMetadata()) .apply(Values.create()); diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java index 68b8130806c4..fccad4d06a44 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java @@ -19,7 +19,6 @@ import static org.apache.beam.sdk.io.kafka.KafkaTimestampType.LOG_APPEND_TIME; import static org.apache.beam.sdk.values.TypeDescriptors.integers; -import static org.apache.beam.sdk.values.TypeDescriptors.strings; import static org.junit.Assert.assertEquals; import java.io.Serializable; @@ -102,7 +101,7 @@ public void testRedistributeByKey() { .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); PCollection> output = - input.apply(KafkaReadRedistribute.byRecordKey()); + input.apply(KafkaReadRedistribute.byRecordKey(10)); PAssert.that(output).containsInAnyOrder(INPUTS); @@ -148,13 +147,13 @@ public void testAssignRecordKeyFn() { Create.of(inputs) .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); - PCollection output = + PCollection output = input - .apply(ParDo.of(new AssignRecordKeyFn())) + .apply(ParDo.of(new AssignRecordKeyFn(2))) .apply(GroupByKey.create()) - .apply(MapElements.into(strings()).via(KV::getKey)); + .apply(MapElements.into(integers()).via(KV::getKey)); - PAssert.that(output).containsInAnyOrder(ImmutableList.of("k1", "k2", "k3", "k5")); + PAssert.that(output).containsInAnyOrder(ImmutableList.of(0, 1)); pipeline.run(); } diff --git a/sdks/java/io/kafka/upgrade/src/main/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslation.java b/sdks/java/io/kafka/upgrade/src/main/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslation.java index 2ebdbf29e230..51d9b028bab0 100644 --- a/sdks/java/io/kafka/upgrade/src/main/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslation.java +++ b/sdks/java/io/kafka/upgrade/src/main/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslation.java @@ -102,6 +102,7 @@ static class KafkaIOReadWithMetadataTranslator implements TransformPayloadTransl .addBooleanField("allows_duplicates") .addNullableInt32Field("redistribute_num_keys") .addNullableBooleanField("offset_deduplication") + .addNullableBooleanField("redistribute_by_record_key") .addNullableLogicalTypeField("watch_topic_partition_duration", new NanosDuration()) .addByteArrayField("timestamp_policy_factory") .addNullableMapField("offset_consumer_config", FieldType.STRING, FieldType.BYTES) @@ -229,6 +230,9 @@ public Row toConfigRow(Read transform) { if (transform.getOffsetDeduplication() != null) { fieldValues.put("offset_deduplication", transform.getOffsetDeduplication()); } + if (transform.getRedistributeByRecordKey() != null) { + fieldValues.put("redistribute_by_record_key", transform.getRedistributeByRecordKey()); + } return Row.withSchema(schema).withFieldValues(fieldValues).build(); } @@ -363,6 +367,12 @@ public Row toConfigRow(Read transform) { transform = transform.withOffsetDeduplication(offsetDeduplication); } } + if (TransformUpgrader.compareVersions(updateCompatibilityBeamVersion, "2.69.0") >= 0) { + @Nullable Boolean byRecordKey = configRow.getValue("redistribute_by_record_key"); + if (byRecordKey != null) { + transform = transform.withRedistributeByRecordKey(byRecordKey); + } + } Duration maxReadTime = configRow.getValue("max_read_time"); if (maxReadTime != null) { transform = diff --git a/sdks/java/io/kafka/upgrade/src/test/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslationTest.java b/sdks/java/io/kafka/upgrade/src/test/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslationTest.java index b5848b316baf..845e89b3b659 100644 --- a/sdks/java/io/kafka/upgrade/src/test/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslationTest.java +++ b/sdks/java/io/kafka/upgrade/src/test/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslationTest.java @@ -66,6 +66,7 @@ public class KafkaIOTranslationTest { READ_TRANSFORM_SCHEMA_MAPPING.put("getStopReadTime", "stop_read_time"); READ_TRANSFORM_SCHEMA_MAPPING.put("getRedistributeNumKeys", "redistribute_num_keys"); READ_TRANSFORM_SCHEMA_MAPPING.put("getOffsetDeduplication", "offset_deduplication"); + READ_TRANSFORM_SCHEMA_MAPPING.put("getRedistributeByRecordKey", "redistribute_by_record_key"); READ_TRANSFORM_SCHEMA_MAPPING.put( "isCommitOffsetsInFinalizeEnabled", "is_commit_offset_finalize_enabled"); READ_TRANSFORM_SCHEMA_MAPPING.put("isDynamicRead", "is_dynamic_read"); From 89a124462716d2fa5ba5bd8dc69a504a6f7c3520 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Mon, 22 Sep 2025 17:07:28 -0700 Subject: [PATCH 12/19] Actually enable withRedistributeByRecordKey in KafkaIOTest. --- .../test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java | 3 +++ 1 file changed, 3 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 e07f067c50a5..e20ffe8b962f 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 @@ -495,6 +495,9 @@ static KafkaIO.Read mkKafkaReadTransform( if (offsetDeduplication != null && offsetDeduplication) { reader.withOffsetDeduplication(offsetDeduplication); } + if (redistributeByRecordKey != null && redistributeByRecordKey) { + reader.withRedistributeByRecordKey(redistributeByRecordKey); + } } return reader; } From 760f297717523884b37d8c387eacf18485c46953 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Mon, 22 Sep 2025 20:51:42 -0700 Subject: [PATCH 13/19] Add byRecordKey property to Kafka read compatibility. --- ...afkaIOReadImplementationCompatibility.java | 6 ++++++ ...IOReadImplementationCompatibilityTest.java | 3 +-- .../apache/beam/sdk/io/kafka/KafkaIOTest.java | 20 +++++++++---------- 3 files changed, 17 insertions(+), 12 deletions(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibility.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibility.java index 81a1de9b872b..8c5efb066d6e 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibility.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibility.java @@ -139,6 +139,12 @@ Object getDefaultValue() { }, OFFSET_DEDUPLICATION(LEGACY), LOG_TOPIC_VERIFICATION, + REDISTRIBUTE_BY_RECORD_KEY { + @Override + Object getDefaultValue() { + return false; + } + }, ; private final @NonNull ImmutableSet supportedImplementations; diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java index fb44be8db4b8..521f77bf4f3d 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java @@ -17,7 +17,6 @@ */ package org.apache.beam.sdk.io.kafka; -import static org.apache.beam.sdk.io.kafka.KafkaIOTest.mkKafkaReadTransform; import static org.apache.beam.sdk.io.kafka.KafkaIOTest.mkKafkaReadTransformWithOffsetDedup; import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.Matchers.containsInAnyOrder; @@ -109,7 +108,7 @@ private PipelineResult testReadTransformCreationWithImplementationBoundPropertie Function, KafkaIO.Read> kafkaReadDecorator) { p.apply( kafkaReadDecorator.apply( - mkKafkaReadTransform( + KafkaIOTest.mkKafkaReadTransformBase( 1000, null, new ValueAsTimestampFn(), 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 e20ffe8b962f..4a8d9e02298b 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 @@ -387,7 +387,7 @@ public Consumer apply(Map config) { static KafkaIO.Read mkKafkaReadTransform( int numElements, @Nullable SerializableFunction, Instant> timestampFn) { - return mkKafkaReadTransform( + return mkKafkaReadTransformBase( numElements, numElements, timestampFn, @@ -401,7 +401,7 @@ static KafkaIO.Read mkKafkaReadTransform( static KafkaIO.Read mkKafkaReadTransformWithOffsetDedup( int numElements, @Nullable SerializableFunction, Instant> timestampFn) { - return mkKafkaReadTransform( + return mkKafkaReadTransformBase( numElements, numElements, timestampFn, @@ -417,7 +417,7 @@ static KafkaIO.Read mkKafkaReadTransformWithRedistributeByRecordK int numElements, @Nullable SerializableFunction, Instant> timestampFn, boolean byRecordKey) { - return mkKafkaReadTransform( + return mkKafkaReadTransformBase( numElements, numElements, timestampFn, @@ -433,7 +433,7 @@ static KafkaIO.Read mkKafkaReadTransformWithTopics( int numElements, @Nullable SerializableFunction, Instant> timestampFn, List topics) { - return mkKafkaReadTransform( + return mkKafkaReadTransformBase( numElements, numElements, timestampFn, @@ -449,7 +449,7 @@ static KafkaIO.Read mkKafkaReadTransformWithTopics( * Creates a consumer with two topics, with 10 partitions each. numElements are (round-robin) * assigned all the 20 partitions. */ - static KafkaIO.Read mkKafkaReadTransform( + static KafkaIO.Read mkKafkaReadTransformBase( int numElements, @Nullable Integer maxNumRecords, @Nullable SerializableFunction, Instant> timestampFn, @@ -738,7 +738,7 @@ public void warningsWithAllowDuplicatesEnabledAndCommitOffsets() { PCollection input = p.apply( - mkKafkaReadTransform( + mkKafkaReadTransformBase( numElements, numElements, new ValueAsTimestampFn(), @@ -767,7 +767,7 @@ public void noWarningsWithNoAllowDuplicatesAndCommitOffsets() { PCollection input = p.apply( - mkKafkaReadTransform( + mkKafkaReadTransformBase( numElements, numElements, new ValueAsTimestampFn(), @@ -797,7 +797,7 @@ public void testNumKeysIgnoredWithRedistributeNotEnabled() { PCollection input = p.apply( - mkKafkaReadTransform( + mkKafkaReadTransformBase( numElements, numElements, new ValueAsTimestampFn(), @@ -2218,7 +2218,7 @@ public void testUnboundedSourceStartReadTime() { PCollection input = p.apply( - mkKafkaReadTransform( + mkKafkaReadTransformBase( numElements, maxNumRecords, new ValueAsTimestampFn(), @@ -2295,7 +2295,7 @@ public void testUnboundedSourceStartReadTimeException() { int startTime = numElements / 20; p.apply( - mkKafkaReadTransform( + mkKafkaReadTransformBase( numElements, numElements, new ValueAsTimestampFn(), From 638feabca80525b3d26fc0ead99b0a084ad3d231 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Mon, 22 Sep 2025 21:06:39 -0700 Subject: [PATCH 14/19] Fix comma formatting for mkKafkaReadTransform. --- ...aIOReadImplementationCompatibilityTest.java | 2 +- .../apache/beam/sdk/io/kafka/KafkaIOTest.java | 18 +++++++++--------- 2 files changed, 10 insertions(+), 10 deletions(-) diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java index 521f77bf4f3d..9d7881e3c3f2 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java @@ -116,7 +116,7 @@ private PipelineResult testReadTransformCreationWithImplementationBoundPropertie false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ null /*redistributeByRecordKey*/))); return p.run(); } 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 4a8d9e02298b..cc17626c3a90 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 @@ -395,7 +395,7 @@ static KafkaIO.Read mkKafkaReadTransform( false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ null /*redistributeByRecordKey*/); } @@ -409,7 +409,7 @@ static KafkaIO.Read mkKafkaReadTransformWithOffsetDedup( false, /*allowDuplicates*/ 100, /*numKeys*/ true, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ null /*redistributeByRecordKey*/); } @@ -425,7 +425,7 @@ static KafkaIO.Read mkKafkaReadTransformWithRedistributeByRecordK false, /*allowDuplicates*/ 100, /*numKeys*/ true, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ byRecordKey /*redistributeByRecordKey*/); } @@ -441,7 +441,7 @@ static KafkaIO.Read mkKafkaReadTransformWithTopics( false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - topics /*topics*/, + topics, /*topics*/ null /*redistributeByRecordKey*/); } @@ -746,7 +746,7 @@ public void warningsWithAllowDuplicatesEnabledAndCommitOffsets() { true, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ null /*redistributeByRecordKey*/) .commitOffsetsInFinalize() .withConsumerConfigUpdates( @@ -775,7 +775,7 @@ public void noWarningsWithNoAllowDuplicatesAndCommitOffsets() { false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ null /*redistributeByRecordKey*/) .commitOffsetsInFinalize() .withConsumerConfigUpdates( @@ -805,7 +805,7 @@ public void testNumKeysIgnoredWithRedistributeNotEnabled() { false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ null /*redistributeByRecordKey*/) .withRedistributeNumKeys(100) .commitOffsetsInFinalize() @@ -2226,7 +2226,7 @@ public void testUnboundedSourceStartReadTime() { false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ null /*redistributeByRecordKey*/) .withStartReadTime(new Instant(startTime)) .withoutMetadata()) @@ -2303,7 +2303,7 @@ public void testUnboundedSourceStartReadTimeException() { false, /*allowDuplicates*/ 0, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/, + null, /*topics*/ null /*redistributeByRecordKey*/) .withStartReadTime(new Instant(startTime)) .withoutMetadata()) From 14478900cb2cd8d9997fd0f20eeb1a5d4d0f8eae Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Mon, 22 Sep 2025 21:18:17 -0700 Subject: [PATCH 15/19] Fix cases where reader was not overwritten when building Kafka reader. --- .../test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 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 cc17626c3a90..5de47b4205fd 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 @@ -493,10 +493,10 @@ static KafkaIO.Read mkKafkaReadTransformBase( reader = reader.withRedistributeNumKeys(numKeys); } if (offsetDeduplication != null && offsetDeduplication) { - reader.withOffsetDeduplication(offsetDeduplication); + reader = reader.withOffsetDeduplication(offsetDeduplication); } if (redistributeByRecordKey != null && redistributeByRecordKey) { - reader.withRedistributeByRecordKey(redistributeByRecordKey); + reader = reader.withRedistributeByRecordKey(redistributeByRecordKey); } } return reader; From 6489053f32ab12e8c6cd049af639ab53178dd861 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Tue, 23 Sep 2025 08:33:51 -0700 Subject: [PATCH 16/19] Rebase and revert method rename for debugging. --- ...IOReadImplementationCompatibilityTest.java | 2 +- .../apache/beam/sdk/io/kafka/KafkaIOTest.java | 29 ++++++++++--------- 2 files changed, 17 insertions(+), 14 deletions(-) diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java index 9d7881e3c3f2..dd74f07cafab 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java @@ -108,7 +108,7 @@ private PipelineResult testReadTransformCreationWithImplementationBoundPropertie Function, KafkaIO.Read> kafkaReadDecorator) { p.apply( kafkaReadDecorator.apply( - KafkaIOTest.mkKafkaReadTransformBase( + KafkaIOTest.mkKafkaReadTransform( 1000, null, new ValueAsTimestampFn(), 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 5de47b4205fd..7637b14e1d8d 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 @@ -387,7 +387,7 @@ public Consumer apply(Map config) { static KafkaIO.Read mkKafkaReadTransform( int numElements, @Nullable SerializableFunction, Instant> timestampFn) { - return mkKafkaReadTransformBase( + return mkKafkaReadTransform( numElements, numElements, timestampFn, @@ -401,7 +401,7 @@ static KafkaIO.Read mkKafkaReadTransform( static KafkaIO.Read mkKafkaReadTransformWithOffsetDedup( int numElements, @Nullable SerializableFunction, Instant> timestampFn) { - return mkKafkaReadTransformBase( + return mkKafkaReadTransform( numElements, numElements, timestampFn, @@ -417,7 +417,7 @@ static KafkaIO.Read mkKafkaReadTransformWithRedistributeByRecordK int numElements, @Nullable SerializableFunction, Instant> timestampFn, boolean byRecordKey) { - return mkKafkaReadTransformBase( + return mkKafkaReadTransform( numElements, numElements, timestampFn, @@ -433,7 +433,7 @@ static KafkaIO.Read mkKafkaReadTransformWithTopics( int numElements, @Nullable SerializableFunction, Instant> timestampFn, List topics) { - return mkKafkaReadTransformBase( + return mkKafkaReadTransform( numElements, numElements, timestampFn, @@ -449,7 +449,7 @@ static KafkaIO.Read mkKafkaReadTransformWithTopics( * Creates a consumer with two topics, with 10 partitions each. numElements are (round-robin) * assigned all the 20 partitions. */ - static KafkaIO.Read mkKafkaReadTransformBase( + static KafkaIO.Read mkKafkaReadTransform( int numElements, @Nullable Integer maxNumRecords, @Nullable SerializableFunction, Instant> timestampFn, @@ -738,7 +738,7 @@ public void warningsWithAllowDuplicatesEnabledAndCommitOffsets() { PCollection input = p.apply( - mkKafkaReadTransformBase( + mkKafkaReadTransform( numElements, numElements, new ValueAsTimestampFn(), @@ -767,7 +767,7 @@ public void noWarningsWithNoAllowDuplicatesAndCommitOffsets() { PCollection input = p.apply( - mkKafkaReadTransformBase( + mkKafkaReadTransform( numElements, numElements, new ValueAsTimestampFn(), @@ -797,7 +797,7 @@ public void testNumKeysIgnoredWithRedistributeNotEnabled() { PCollection input = p.apply( - mkKafkaReadTransformBase( + mkKafkaReadTransform( numElements, numElements, new ValueAsTimestampFn(), @@ -832,7 +832,8 @@ public void testDefaultRedistributeNumKeys() { false, /*allowDuplicates*/ null, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/); + null, /*topics*/ + null /*redistributeByRecordKey*/); assertFalse(read.isRedistributed()); assertEquals(0, read.getRedistributeNumKeys()); @@ -846,7 +847,8 @@ public void testDefaultRedistributeNumKeys() { false, /*allowDuplicates*/ null, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/); + null, /*topics*/ + null /*redistributeByRecordKey*/); assertTrue(read.isRedistributed()); // Default is defined by DEFAULT_REDISTRIBUTE_NUM_KEYS in KafkaIO. assertEquals(32768, read.getRedistributeNumKeys()); @@ -861,7 +863,8 @@ public void testDefaultRedistributeNumKeys() { false, /*allowDuplicates*/ 10, /*numKeys*/ null, /*offsetDeduplication*/ - null /*topics*/); + null, /*topics*/ + null /*redistributeByRecordKey*/); assertTrue(read.isRedistributed()); assertEquals(10, read.getRedistributeNumKeys()); } @@ -2218,7 +2221,7 @@ public void testUnboundedSourceStartReadTime() { PCollection input = p.apply( - mkKafkaReadTransformBase( + mkKafkaReadTransform( numElements, maxNumRecords, new ValueAsTimestampFn(), @@ -2295,7 +2298,7 @@ public void testUnboundedSourceStartReadTimeException() { int startTime = numElements / 20; p.apply( - mkKafkaReadTransformBase( + mkKafkaReadTransform( numElements, numElements, new ValueAsTimestampFn(), From cf907b7acb8e8844f750e0b71a21105c20f14f94 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Tue, 23 Sep 2025 08:57:22 -0700 Subject: [PATCH 17/19] Address spotless finding for makeKafkaRecord. --- .../io/kafka/KafkaReadRedistributeTest.java | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java index fccad4d06a44..d1839b46c78e 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java @@ -50,15 +50,15 @@ public class KafkaReadRedistributeTest implements Serializable { private static final ImmutableList> INPUTS = ImmutableList.of( - MakeKafkaRecord("k1", 3, 1), - MakeKafkaRecord("k5", Integer.MAX_VALUE, 2), - MakeKafkaRecord("k5", Integer.MIN_VALUE, 3), - MakeKafkaRecord("k2", 66, 4), - MakeKafkaRecord("k1", 4, 5), - MakeKafkaRecord("k2", -33, 6), - MakeKafkaRecord("k3", 0, 7)); - - static KafkaRecord MakeKafkaRecord(String key, Integer value, Integer offset) { + makeKafkaRecord("k1", 3, 1), + makeKafkaRecord("k5", Integer.MAX_VALUE, 2), + makeKafkaRecord("k5", Integer.MIN_VALUE, 3), + makeKafkaRecord("k2", 66, 4), + makeKafkaRecord("k1", 4, 5), + makeKafkaRecord("k2", -33, 6), + makeKafkaRecord("k3", 0, 7)); + + static KafkaRecord makeKafkaRecord(String key, Integer value, Integer offset) { return new KafkaRecord( /*topic*/ "kafka", /*partition*/ 1, From 06469c5b2b8c056b5f34f7fe5be719bea83840f0 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Wed, 24 Sep 2025 09:16:35 -0700 Subject: [PATCH 18/19] Add tests for deterministic sharding. --- .../io/kafka/KafkaReadRedistributeTest.java | 75 ++++++++++++++++++- 1 file changed, 73 insertions(+), 2 deletions(-) diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java index d1839b46c78e..a14c6e3232e5 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java @@ -30,6 +30,7 @@ import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; import org.apache.beam.sdk.testing.ValidatesRunner; +import org.apache.beam.sdk.transforms.Count; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.transforms.GroupByKey; import org.apache.beam.sdk.transforms.MapElements; @@ -58,6 +59,26 @@ public class KafkaReadRedistributeTest implements Serializable { makeKafkaRecord("k2", -33, 6), makeKafkaRecord("k3", 0, 7)); + private static final ImmutableList> SAME_OFFSET_INPUTS = + ImmutableList.of( + makeKafkaRecord("k1", 3, 1), + makeKafkaRecord("k5", Integer.MAX_VALUE, 1), + makeKafkaRecord("k5", Integer.MIN_VALUE, 1), + makeKafkaRecord("k2", 66, 1), + makeKafkaRecord("k1", 4, 1), + makeKafkaRecord("k2", -33, 1), + makeKafkaRecord("k3", 0, 1)); + + private static final ImmutableList> SAME_KEY_INPUTS = + ImmutableList.of( + makeKafkaRecord("k1", 3, 1), + makeKafkaRecord("k1", Integer.MAX_VALUE, 2), + makeKafkaRecord("k1", Integer.MIN_VALUE, 3), + makeKafkaRecord("k1", 66, 4), + makeKafkaRecord("k1", 4, 5), + makeKafkaRecord("k1", -33, 6), + makeKafkaRecord("k1", 0, 7)); + static KafkaRecord makeKafkaRecord(String key, Integer value, Integer offset) { return new KafkaRecord( /*topic*/ "kafka", @@ -112,7 +133,7 @@ public void testRedistributeByKey() { @Test @Category({ValidatesRunner.class}) - public void testAssignOutputShardFn() { + public void testAssignOutputShardFnBucketing() { List> inputs = Lists.newArrayList(); for (int i = 0; i < 10; i++) { inputs.addAll(INPUTS); @@ -136,7 +157,7 @@ public void testAssignOutputShardFn() { @Test @Category({ValidatesRunner.class}) - public void testAssignRecordKeyFn() { + public void testAssignRecordKeyFnBucketing() { List> inputs = Lists.newArrayList(); for (int i = 0; i < 10; i++) { inputs.addAll(INPUTS); @@ -157,4 +178,54 @@ public void testAssignRecordKeyFn() { pipeline.run(); } + + @Test + @Category({ValidatesRunner.class}) + public void testAssignOutputShardFnDeterministic() { + List> inputs = Lists.newArrayList(); + for (int i = 0; i < 10; i++) { + inputs.addAll(SAME_OFFSET_INPUTS); + } + + PCollection> input = + pipeline.apply( + Create.of(inputs) + .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); + + PCollection output = + input + .apply(ParDo.of(new AssignOffsetShardFn(1024))) + .apply(GroupByKey.create()) + .apply(MapElements.into(integers()).via(KV::getKey)); + + PCollection count = output.apply("CountElements", Count.globally()); + PAssert.that(count).containsInAnyOrder(1L); + + pipeline.run(); + } + + @Test + @Category({ValidatesRunner.class}) + public void testAssignRecordKeyFnDeterministic() { + List> inputs = Lists.newArrayList(); + for (int i = 0; i < 10; i++) { + inputs.addAll(SAME_KEY_INPUTS); + } + + PCollection> input = + pipeline.apply( + Create.of(inputs) + .withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of()))); + + PCollection output = + input + .apply(ParDo.of(new AssignRecordKeyFn(1024))) + .apply(GroupByKey.create()) + .apply(MapElements.into(integers()).via(KV::getKey)); + + PCollection count = output.apply("CountElements", Count.globally()); + PAssert.that(count).containsInAnyOrder(1L); + + pipeline.run(); + } } From ebef84849a38337cfc5f6d2cd1bbadea372648b7 Mon Sep 17 00:00:00 2001 From: Tom Stepp Date: Thu, 25 Sep 2025 07:59:14 -0700 Subject: [PATCH 19/19] numBuckets as UnsignedInteger to reduce conversion overhead, and clarify sharding Fn display name. --- .../sdk/io/kafka/KafkaReadRedistribute.java | 31 ++++++++++++------- 1 file changed, 19 insertions(+), 12 deletions(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java index 3943d399cf24..61c0b671f292 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java @@ -28,6 +28,7 @@ import org.apache.beam.sdk.values.PCollection; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedInteger; +import org.checkerframework.checker.nullness.qual.NonNull; import org.checkerframework.checker.nullness.qual.Nullable; public class KafkaReadRedistribute @@ -55,23 +56,27 @@ public PCollection> expand(PCollection> inpu if (byRecordKey) { return input - .apply("Pair with record key", ParDo.of(new AssignRecordKeyFn(numBuckets))) + .apply("Pair with shard from key", ParDo.of(new AssignRecordKeyFn(numBuckets))) .apply(Redistribute.>byKey().withAllowDuplicates(false)) .apply(Values.create()); } return input - .apply("Pair with offset shard", ParDo.of(new AssignOffsetShardFn(numBuckets))) + .apply("Pair with shard from offset", ParDo.of(new AssignOffsetShardFn(numBuckets))) .apply(Redistribute.>byKey().withAllowDuplicates(false)) .apply(Values.create()); } static class AssignOffsetShardFn extends DoFn, KV>> { - private @Nullable Integer numBuckets; + private @NonNull UnsignedInteger numBuckets; public AssignOffsetShardFn(@Nullable Integer numBuckets) { - this.numBuckets = numBuckets; + if (numBuckets != null && numBuckets > 0) { + this.numBuckets = UnsignedInteger.fromIntBits(numBuckets); + } else { + this.numBuckets = UnsignedInteger.valueOf(0); + } } @ProcessElement @@ -80,9 +85,8 @@ public void processElement( OutputReceiver>> receiver) { int hash = Hashing.farmHashFingerprint64().hashLong(element.getOffset()).asInt(); - if (numBuckets != null && numBuckets > 0) { - UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); - hash = UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue(); + if (numBuckets != null) { + hash = UnsignedInteger.fromIntBits(hash).mod(numBuckets).intValue(); } receiver.output(KV.of(hash, element)); @@ -92,10 +96,14 @@ public void processElement( static class AssignRecordKeyFn extends DoFn, KV>> { - private @Nullable Integer numBuckets; + private @NonNull UnsignedInteger numBuckets; public AssignRecordKeyFn(@Nullable Integer numBuckets) { - this.numBuckets = numBuckets; + if (numBuckets != null && numBuckets > 0) { + this.numBuckets = UnsignedInteger.fromIntBits(numBuckets); + } else { + this.numBuckets = UnsignedInteger.valueOf(0); + } } @ProcessElement @@ -106,9 +114,8 @@ public void processElement( String keyString = key == null ? "" : key.toString(); int hash = Hashing.farmHashFingerprint64().hashBytes(keyString.getBytes(UTF_8)).asInt(); - if (numBuckets != null && numBuckets > 0) { - UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets); - hash = UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue(); + if (numBuckets != null) { + hash = UnsignedInteger.fromIntBits(hash).mod(numBuckets).intValue(); } receiver.output(KV.of(hash, element));