Skip to content

Commit d07d2e8

Browse files
committed
Add byRecordKey property to Kafka read compatibility.
1 parent a7632c0 commit d07d2e8

3 files changed

Lines changed: 17 additions & 12 deletions

File tree

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -139,6 +139,12 @@ Object getDefaultValue() {
139139
},
140140
OFFSET_DEDUPLICATION(LEGACY),
141141
LOG_TOPIC_VERIFICATION,
142+
REDISTRIBUTE_BY_RECORD_KEY {
143+
@Override
144+
Object getDefaultValue() {
145+
return false;
146+
}
147+
},
142148
;
143149

144150
private final @NonNull ImmutableSet<KafkaIOReadImplementation> supportedImplementations;

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

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

20-
import static org.apache.beam.sdk.io.kafka.KafkaIOTest.mkKafkaReadTransform;
2120
import static org.apache.beam.sdk.io.kafka.KafkaIOTest.mkKafkaReadTransformWithOffsetDedup;
2221
import static org.hamcrest.MatcherAssert.assertThat;
2322
import static org.hamcrest.Matchers.containsInAnyOrder;
@@ -109,7 +108,7 @@ private PipelineResult testReadTransformCreationWithImplementationBoundPropertie
109108
Function<KafkaIO.Read<Integer, Long>, KafkaIO.Read<Integer, Long>> kafkaReadDecorator) {
110109
p.apply(
111110
kafkaReadDecorator.apply(
112-
mkKafkaReadTransform(
111+
KafkaIOTest.mkKafkaReadTransformBase(
113112
1000,
114113
null,
115114
new ValueAsTimestampFn(),

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

Lines changed: 10 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -386,7 +386,7 @@ public Consumer<byte[], byte[]> apply(Map<String, Object> config) {
386386

387387
static KafkaIO.Read<Integer, Long> mkKafkaReadTransform(
388388
int numElements, @Nullable SerializableFunction<KV<Integer, Long>, Instant> timestampFn) {
389-
return mkKafkaReadTransform(
389+
return mkKafkaReadTransformBase(
390390
numElements,
391391
numElements,
392392
timestampFn,
@@ -400,7 +400,7 @@ static KafkaIO.Read<Integer, Long> mkKafkaReadTransform(
400400

401401
static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithOffsetDedup(
402402
int numElements, @Nullable SerializableFunction<KV<Integer, Long>, Instant> timestampFn) {
403-
return mkKafkaReadTransform(
403+
return mkKafkaReadTransformBase(
404404
numElements,
405405
numElements,
406406
timestampFn,
@@ -416,7 +416,7 @@ static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithRedistributeByRecordK
416416
int numElements,
417417
@Nullable SerializableFunction<KV<Integer, Long>, Instant> timestampFn,
418418
boolean byRecordKey) {
419-
return mkKafkaReadTransform(
419+
return mkKafkaReadTransformBase(
420420
numElements,
421421
numElements,
422422
timestampFn,
@@ -432,7 +432,7 @@ static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithTopics(
432432
int numElements,
433433
@Nullable SerializableFunction<KV<Integer, Long>, Instant> timestampFn,
434434
List<String> topics) {
435-
return mkKafkaReadTransform(
435+
return mkKafkaReadTransformBase(
436436
numElements,
437437
numElements,
438438
timestampFn,
@@ -448,7 +448,7 @@ static KafkaIO.Read<Integer, Long> mkKafkaReadTransformWithTopics(
448448
* Creates a consumer with two topics, with 10 partitions each. numElements are (round-robin)
449449
* assigned all the 20 partitions.
450450
*/
451-
static KafkaIO.Read<Integer, Long> mkKafkaReadTransform(
451+
static KafkaIO.Read<Integer, Long> mkKafkaReadTransformBase(
452452
int numElements,
453453
@Nullable Integer maxNumRecords,
454454
@Nullable SerializableFunction<KV<Integer, Long>, Instant> timestampFn,
@@ -737,7 +737,7 @@ public void warningsWithAllowDuplicatesEnabledAndCommitOffsets() {
737737

738738
PCollection<Long> input =
739739
p.apply(
740-
mkKafkaReadTransform(
740+
mkKafkaReadTransformBase(
741741
numElements,
742742
numElements,
743743
new ValueAsTimestampFn(),
@@ -766,7 +766,7 @@ public void noWarningsWithNoAllowDuplicatesAndCommitOffsets() {
766766

767767
PCollection<Long> input =
768768
p.apply(
769-
mkKafkaReadTransform(
769+
mkKafkaReadTransformBase(
770770
numElements,
771771
numElements,
772772
new ValueAsTimestampFn(),
@@ -796,7 +796,7 @@ public void testNumKeysIgnoredWithRedistributeNotEnabled() {
796796

797797
PCollection<Long> input =
798798
p.apply(
799-
mkKafkaReadTransform(
799+
mkKafkaReadTransformBase(
800800
numElements,
801801
numElements,
802802
new ValueAsTimestampFn(),
@@ -2170,7 +2170,7 @@ public void testUnboundedSourceStartReadTime() {
21702170

21712171
PCollection<Long> input =
21722172
p.apply(
2173-
mkKafkaReadTransform(
2173+
mkKafkaReadTransformBase(
21742174
numElements,
21752175
maxNumRecords,
21762176
new ValueAsTimestampFn(),
@@ -2247,7 +2247,7 @@ public void testUnboundedSourceStartReadTimeException() {
22472247
int startTime = numElements / 20;
22482248

22492249
p.apply(
2250-
mkKafkaReadTransform(
2250+
mkKafkaReadTransformBase(
22512251
numElements,
22522252
numElements,
22532253
new ValueAsTimestampFn(),

0 commit comments

Comments
 (0)