Skip to content

Commit 04195e9

Browse files
committed
Add bucketing to redistributeByKey and add option to redistribute by key.
1 parent 7392bb4 commit 04195e9

7 files changed

Lines changed: 138 additions & 28 deletions

File tree

sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java

Lines changed: 31 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -723,6 +723,9 @@ public abstract static class Read<K, V>
723723
@Pure
724724
public abstract @Nullable Boolean getOffsetDeduplication();
725725

726+
@Pure
727+
public abstract @Nullable Boolean getRedistributeByRecordKey();
728+
726729
@Pure
727730
public abstract @Nullable Duration getWatchTopicPartitionDuration();
728731

@@ -793,6 +796,8 @@ abstract Builder<K, V> setConsumerFactoryFn(
793796

794797
abstract Builder<K, V> setOffsetDeduplication(Boolean offsetDeduplication);
795798

799+
abstract Builder<K, V> setRedistributeByRecordKey(Boolean redistributeByRecordKey);
800+
796801
abstract Builder<K, V> setTimestampPolicyFactory(
797802
TimestampPolicyFactory<K, V> timestampPolicyFactory);
798803

@@ -908,11 +913,15 @@ static <K, V> void setupExternalBuilder(
908913
&& config.offsetDeduplication != null) {
909914
builder.setOffsetDeduplication(config.offsetDeduplication);
910915
}
916+
if (config.redistribute && config.redistributeByRecordKey != null) {
917+
builder.setRedistributeByRecordKey(config.redistributeByRecordKey);
918+
}
911919
} else {
912920
builder.setRedistributed(false);
913921
builder.setRedistributeNumKeys(0);
914922
builder.setAllowDuplicates(false);
915923
builder.setOffsetDeduplication(false);
924+
builder.setRedistributeByRecordKey(false);
916925
}
917926
}
918927

@@ -982,6 +991,7 @@ public static class Configuration {
982991
private Boolean redistribute;
983992
private Boolean allowDuplicates;
984993
private Boolean offsetDeduplication;
994+
private Boolean redistributeByRecordKey;
985995
private Long dynamicReadPollIntervalSeconds;
986996

987997
public void setConsumerConfig(Map<String, String> consumerConfig) {
@@ -1044,6 +1054,10 @@ public void setOffsetDeduplication(Boolean offsetDeduplication) {
10441054
this.offsetDeduplication = offsetDeduplication;
10451055
}
10461056

1057+
public void setRedistributeByRecordKey(Boolean redistributeByRecordKey) {
1058+
this.redistributeByRecordKey = redistributeByRecordKey;
1059+
}
1060+
10471061
public void setDynamicReadPollIntervalSeconds(Long dynamicReadPollIntervalSeconds) {
10481062
this.dynamicReadPollIntervalSeconds = dynamicReadPollIntervalSeconds;
10491063
}
@@ -1149,6 +1163,10 @@ public Read<K, V> withOffsetDeduplication(Boolean offsetDeduplication) {
11491163
return toBuilder().setOffsetDeduplication(offsetDeduplication).build();
11501164
}
11511165

1166+
public Read<K, V> withRedistributeByRecordKey(Boolean redistributeByRecordKey) {
1167+
return toBuilder().setRedistributeByRecordKey(redistributeByRecordKey).build();
1168+
}
1169+
11521170
/**
11531171
* Internally sets a {@link java.util.regex.Pattern} of topics to read from. All the partitions
11541172
* from each of the matching topics are read.
@@ -1667,6 +1685,11 @@ private void checkRedistributeConfiguration() {
16671685
LOG.warn(
16681686
"Offsets used for deduplication are available in WindowedValue's metadata. Combining, aggregating, mutating them may risk with data loss.");
16691687
}
1688+
if (getRedistributeByRecordKey() != null && getRedistributeByRecordKey()) {
1689+
checkState(
1690+
isRedistributed(),
1691+
"withRedistributeByRecordKey can only be used when withRedistribute is set.");
1692+
}
16701693
}
16711694

16721695
private void warnAboutUnsafeConfigurations(PBegin input) {
@@ -1847,11 +1870,15 @@ public PCollection<KafkaRecord<K, V>> expand(PBegin input) {
18471870
}
18481871

18491872
if (kafkaRead.getOffsetDeduplication() != null && kafkaRead.getOffsetDeduplication()) {
1850-
// TODO: Expose Kafka read options to control byOffsetShard vs byRecordKey.
1851-
return output.apply(
1852-
KafkaReadRedistribute.<K, V>byOffsetShard(kafkaRead.getRedistributeNumKeys()));
1873+
if (kafkaRead.getRedistributeByRecordKey() != null
1874+
&& kafkaRead.getRedistributeByRecordKey()) {
1875+
return output.apply(
1876+
KafkaReadRedistribute.<K, V>byRecordKey(kafkaRead.getRedistributeNumKeys()));
1877+
} else {
1878+
return output.apply(
1879+
KafkaReadRedistribute.<K, V>byOffsetShard(kafkaRead.getRedistributeNumKeys()));
1880+
}
18531881
}
1854-
18551882
RedistributeArbitrarily<KafkaRecord<K, V>> redistribute =
18561883
Redistribute.<KafkaRecord<K, V>>arbitrarily()
18571884
.withAllowDuplicates(kafkaRead.isAllowDuplicates());

sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistribute.java

Lines changed: 25 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,8 @@
1717
*/
1818
package org.apache.beam.sdk.io.kafka;
1919

20+
import static java.nio.charset.StandardCharsets.UTF_8;
21+
2022
import org.apache.beam.sdk.transforms.DoFn;
2123
import org.apache.beam.sdk.transforms.PTransform;
2224
import org.apache.beam.sdk.transforms.ParDo;
@@ -34,8 +36,8 @@ public static <K, V> KafkaReadRedistribute<K, V> byOffsetShard(@Nullable Integer
3436
return new KafkaReadRedistribute<>(numBuckets, false);
3537
}
3638

37-
public static <K, V> KafkaReadRedistribute<K, V> byRecordKey() {
38-
return new KafkaReadRedistribute<>(null, true);
39+
public static <K, V> KafkaReadRedistribute<K, V> byRecordKey(@Nullable Integer numBuckets) {
40+
return new KafkaReadRedistribute<>(numBuckets, true);
3941
}
4042

4143
// The number of buckets to shard into.
@@ -53,8 +55,8 @@ public PCollection<KafkaRecord<K, V>> expand(PCollection<KafkaRecord<K, V>> inpu
5355

5456
if (byRecordKey) {
5557
return input
56-
.apply("Pair with record key", ParDo.of(new AssignRecordKeyFn<K, V>()))
57-
.apply(Redistribute.<K, KafkaRecord<K, V>>byKey().withAllowDuplicates(false))
58+
.apply("Pair with record key", ParDo.of(new AssignRecordKeyFn<K, V>(numBuckets)))
59+
.apply(Redistribute.<Integer, KafkaRecord<K, V>>byKey().withAllowDuplicates(false))
5860
.apply(Values.create());
5961
}
6062

@@ -87,14 +89,29 @@ public void processElement(
8789
}
8890
}
8991

90-
static class AssignRecordKeyFn<K, V> extends DoFn<KafkaRecord<K, V>, KV<K, KafkaRecord<K, V>>> {
92+
static class AssignRecordKeyFn<K, V>
93+
extends DoFn<KafkaRecord<K, V>, KV<Integer, KafkaRecord<K, V>>> {
94+
95+
private @Nullable Integer numBuckets;
9196

92-
public AssignRecordKeyFn() {}
97+
public AssignRecordKeyFn(@Nullable Integer numBuckets) {
98+
this.numBuckets = numBuckets;
99+
}
93100

94101
@ProcessElement
95102
public void processElement(
96-
@Element KafkaRecord<K, V> element, OutputReceiver<KV<K, KafkaRecord<K, V>>> receiver) {
97-
receiver.output(KV.of(element.getKV().getKey(), element));
103+
@Element KafkaRecord<K, V> element,
104+
OutputReceiver<KV<Integer, KafkaRecord<K, V>>> receiver) {
105+
K key = element.getKV().getKey();
106+
String keyString = key == null ? "" : key.toString();
107+
int hash = Hashing.farmHashFingerprint64().hashBytes(keyString.getBytes(UTF_8)).asInt();
108+
109+
if (numBuckets != null && numBuckets > 0) {
110+
UnsignedInteger unsignedNumBuckets = UnsignedInteger.fromIntBits(numBuckets);
111+
hash = UnsignedInteger.fromIntBits(hash).mod(unsignedNumBuckets).intValue();
112+
}
113+
114+
receiver.output(KV.of(hash, element));
98115
}
99116
}
100117
}

sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibilityTest.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -117,7 +117,8 @@ private PipelineResult testReadTransformCreationWithImplementationBoundPropertie
117117
false, /*allowDuplicates*/
118118
0, /*numKeys*/
119119
null, /*offsetDeduplication*/
120-
null /*topics*/)));
120+
null /*topics*/,
121+
null /*redistributeByRecordKey*/)));
121122
return p.run();
122123
}
123124

sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaIOTest.java

Lines changed: 64 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -394,7 +394,8 @@ static KafkaIO.Read<Integer, Long> mkKafkaReadTransform(
394394
false, /*allowDuplicates*/
395395
0, /*numKeys*/
396396
null, /*offsetDeduplication*/
397-
null /*topics*/);
397+
null /*topics*/,
398+
null /*redistributeByRecordKey*/);
398399
}
399400

400401
static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithOffsetDedup(
@@ -407,7 +408,24 @@ static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithOffsetDedup(
407408
false, /*allowDuplicates*/
408409
100, /*numKeys*/
409410
true, /*offsetDeduplication*/
410-
null /*topics*/);
411+
null /*topics*/,
412+
null /*redistributeByRecordKey*/);
413+
}
414+
415+
static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithRedistributeByRecordKey(
416+
int numElements,
417+
@Nullable SerializableFunction<KV<Integer, Long>, Instant> timestampFn,
418+
boolean byRecordKey) {
419+
return mkKafkaReadTransform(
420+
numElements,
421+
numElements,
422+
timestampFn,
423+
true, /*redistribute*/
424+
false, /*allowDuplicates*/
425+
100, /*numKeys*/
426+
true, /*offsetDeduplication*/
427+
null /*topics*/,
428+
byRecordKey /*redistributeByRecordKey*/);
411429
}
412430

413431
static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithTopics(
@@ -422,7 +440,8 @@ static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithTopics(
422440
false, /*allowDuplicates*/
423441
0, /*numKeys*/
424442
null, /*offsetDeduplication*/
425-
topics /*topics*/);
443+
topics /*topics*/,
444+
null /*redistributeByRecordKey*/);
426445
}
427446

428447
/**
@@ -437,7 +456,8 @@ static KafkaIO.Read<Integer, Long> mkKafkaReadTransform(
437456
@Nullable Boolean withAllowDuplicates,
438457
@Nullable Integer numKeys,
439458
@Nullable Boolean offsetDeduplication,
440-
@Nullable List<String> topics) {
459+
@Nullable List<String> topics,
460+
@Nullable Boolean redistributeByRecordKey) {
441461

442462
KafkaIO.Read<Integer, Long> reader =
443463
KafkaIO.<Integer, Long>read()
@@ -722,7 +742,8 @@ public void warningsWithAllowDuplicatesEnabledAndCommitOffsets() {
722742
true, /*allowDuplicates*/
723743
0, /*numKeys*/
724744
null, /*offsetDeduplication*/
725-
null /*topics*/)
745+
null /*topics*/,
746+
null /*redistributeByRecordKey*/)
726747
.commitOffsetsInFinalize()
727748
.withConsumerConfigUpdates(
728749
ImmutableMap.of(ConsumerConfig.GROUP_ID_CONFIG, "group_id"))
@@ -750,7 +771,8 @@ public void noWarningsWithNoAllowDuplicatesAndCommitOffsets() {
750771
false, /*allowDuplicates*/
751772
0, /*numKeys*/
752773
null, /*offsetDeduplication*/
753-
null /*topics*/)
774+
null /*topics*/,
775+
null /*redistributeByRecordKey*/)
754776
.commitOffsetsInFinalize()
755777
.withConsumerConfigUpdates(
756778
ImmutableMap.of(ConsumerConfig.GROUP_ID_CONFIG, "group_id"))
@@ -779,7 +801,8 @@ public void testNumKeysIgnoredWithRedistributeNotEnabled() {
779801
false, /*allowDuplicates*/
780802
0, /*numKeys*/
781803
null, /*offsetDeduplication*/
782-
null /*topics*/)
804+
null /*topics*/,
805+
null /*redistributeByRecordKey*/)
783806
.withRedistributeNumKeys(100)
784807
.commitOffsetsInFinalize()
785808
.withConsumerConfigUpdates(
@@ -2152,7 +2175,8 @@ public void testUnboundedSourceStartReadTime() {
21522175
false, /*allowDuplicates*/
21532176
0, /*numKeys*/
21542177
null, /*offsetDeduplication*/
2155-
null /*topics*/)
2178+
null /*topics*/,
2179+
null /*redistributeByRecordKey*/)
21562180
.withStartReadTime(new Instant(startTime))
21572181
.withoutMetadata())
21582182
.apply(Values.create());
@@ -2175,6 +2199,36 @@ public void testOffsetDeduplication() {
21752199
p.run();
21762200
}
21772201

2202+
@Test
2203+
public void testRedistributeByRecordKeyOn() {
2204+
int numElements = 1000;
2205+
2206+
PCollection<Long> input =
2207+
p.apply(
2208+
mkKafkaReadTransformWithRedistributeByRecordKey(
2209+
numElements, new ValueAsTimestampFn(), true)
2210+
.withoutMetadata())
2211+
.apply(Values.create());
2212+
2213+
addCountingAsserts(input, numElements, numElements, 0, numElements - 1);
2214+
p.run();
2215+
}
2216+
2217+
@Test
2218+
public void testRedistributeByRecordKeyOff() {
2219+
int numElements = 1000;
2220+
2221+
PCollection<Long> input =
2222+
p.apply(
2223+
mkKafkaReadTransformWithRedistributeByRecordKey(
2224+
numElements, new ValueAsTimestampFn(), false)
2225+
.withoutMetadata())
2226+
.apply(Values.create());
2227+
2228+
addCountingAsserts(input, numElements, numElements, 0, numElements - 1);
2229+
p.run();
2230+
}
2231+
21782232
@Rule public ExpectedException noMessagesException = ExpectedException.none();
21792233

21802234
@Test
@@ -2198,7 +2252,8 @@ public void testUnboundedSourceStartReadTimeException() {
21982252
false, /*allowDuplicates*/
21992253
0, /*numKeys*/
22002254
null, /*offsetDeduplication*/
2201-
null /*topics*/)
2255+
null /*topics*/,
2256+
null /*redistributeByRecordKey*/)
22022257
.withStartReadTime(new Instant(startTime))
22032258
.withoutMetadata())
22042259
.apply(Values.create());

sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaReadRedistributeTest.java

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919

2020
import static org.apache.beam.sdk.io.kafka.KafkaTimestampType.LOG_APPEND_TIME;
2121
import static org.apache.beam.sdk.values.TypeDescriptors.integers;
22-
import static org.apache.beam.sdk.values.TypeDescriptors.strings;
2322
import static org.junit.Assert.assertEquals;
2423

2524
import java.io.Serializable;
@@ -102,7 +101,7 @@ public void testRedistributeByKey() {
102101
.withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of())));
103102

104103
PCollection<KafkaRecord<String, Integer>> output =
105-
input.apply(KafkaReadRedistribute.byRecordKey());
104+
input.apply(KafkaReadRedistribute.byRecordKey(10));
106105

107106
PAssert.that(output).containsInAnyOrder(INPUTS);
108107

@@ -148,13 +147,13 @@ public void testAssignRecordKeyFn() {
148147
Create.of(inputs)
149148
.withCoder(KafkaRecordCoder.of(StringUtf8Coder.of(), VarIntCoder.of())));
150149

151-
PCollection<String> output =
150+
PCollection<Integer> output =
152151
input
153-
.apply(ParDo.of(new AssignRecordKeyFn<String, Integer>()))
152+
.apply(ParDo.of(new AssignRecordKeyFn<String, Integer>(2)))
154153
.apply(GroupByKey.create())
155-
.apply(MapElements.into(strings()).via(KV::getKey));
154+
.apply(MapElements.into(integers()).via(KV::getKey));
156155

157-
PAssert.that(output).containsInAnyOrder(ImmutableList.of("k1", "k2", "k3", "k5"));
156+
PAssert.that(output).containsInAnyOrder(ImmutableList.of(0, 1));
158157

159158
pipeline.run();
160159
}

sdks/java/io/kafka/upgrade/src/main/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslation.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,7 @@ static class KafkaIOReadWithMetadataTranslator implements TransformPayloadTransl
102102
.addBooleanField("allows_duplicates")
103103
.addNullableInt32Field("redistribute_num_keys")
104104
.addNullableBooleanField("offset_deduplication")
105+
.addNullableBooleanField("redistribute_by_record_key")
105106
.addNullableLogicalTypeField("watch_topic_partition_duration", new NanosDuration())
106107
.addByteArrayField("timestamp_policy_factory")
107108
.addNullableMapField("offset_consumer_config", FieldType.STRING, FieldType.BYTES)
@@ -229,6 +230,9 @@ public Row toConfigRow(Read<?, ?> transform) {
229230
if (transform.getOffsetDeduplication() != null) {
230231
fieldValues.put("offset_deduplication", transform.getOffsetDeduplication());
231232
}
233+
if (transform.getRedistributeByRecordKey() != null) {
234+
fieldValues.put("redistribute_by_record_key", transform.getRedistributeByRecordKey());
235+
}
232236
return Row.withSchema(schema).withFieldValues(fieldValues).build();
233237
}
234238

@@ -363,6 +367,12 @@ public Row toConfigRow(Read<?, ?> transform) {
363367
transform = transform.withOffsetDeduplication(offsetDeduplication);
364368
}
365369
}
370+
if (TransformUpgrader.compareVersions(updateCompatibilityBeamVersion, "2.69.0") >= 0) {
371+
@Nullable Boolean byRecordKey = configRow.getValue("redistribute_by_record_key");
372+
if (byRecordKey != null) {
373+
transform = transform.withRedistributeByRecordKey(byRecordKey);
374+
}
375+
}
366376
Duration maxReadTime = configRow.getValue("max_read_time");
367377
if (maxReadTime != null) {
368378
transform =

sdks/java/io/kafka/upgrade/src/test/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslationTest.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,7 @@ public class KafkaIOTranslationTest {
6666
READ_TRANSFORM_SCHEMA_MAPPING.put("getStopReadTime", "stop_read_time");
6767
READ_TRANSFORM_SCHEMA_MAPPING.put("getRedistributeNumKeys", "redistribute_num_keys");
6868
READ_TRANSFORM_SCHEMA_MAPPING.put("getOffsetDeduplication", "offset_deduplication");
69+
READ_TRANSFORM_SCHEMA_MAPPING.put("getRedistributeByRecordKey", "redistribute_by_record_key");
6970
READ_TRANSFORM_SCHEMA_MAPPING.put(
7071
"isCommitOffsetsInFinalizeEnabled", "is_commit_offset_finalize_enabled");
7172
READ_TRANSFORM_SCHEMA_MAPPING.put("isDynamicRead", "is_dynamic_read");

0 commit comments

Comments
 (0)