From e570870ee930ca92e486f5b82c9664dbe1fec407 Mon Sep 17 00:00:00 2001 From: Heejong Lee Date: Thu, 13 May 2021 11:21:38 -0700 Subject: [PATCH] [BEAM-10670] Use non-SDF based translation for Read by default --- .../beam/runners/core/construction/SplittableParDo.java | 7 +++---- .../main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 5 +---- 2 files changed, 4 insertions(+), 8 deletions(-) diff --git a/runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SplittableParDo.java b/runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SplittableParDo.java index 639c1551525c..45b5e11e66bd 100644 --- a/runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SplittableParDo.java +++ b/runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SplittableParDo.java @@ -646,14 +646,13 @@ public void tearDown() { /** * Converts {@link Read} based Splittable DoFn expansions to primitive reads implemented by {@link - * PrimitiveBoundedRead} and {@link PrimitiveUnboundedRead} if either the experiment {@code - * use_deprecated_read} or {@code beam_fn_api_use_deprecated_read} are specified. + * PrimitiveBoundedRead} and {@link PrimitiveUnboundedRead} if the experiment {@code use_sdf_read} + * is not specified. * *

TODO(BEAM-10670): Remove the primitive Read and make the splittable DoFn the only option. */ public static void convertReadBasedSplittableDoFnsToPrimitiveReadsIfNecessary(Pipeline pipeline) { - if (ExperimentalOptions.hasExperiment(pipeline.getOptions(), "beam_fn_api_use_deprecated_read") - || ExperimentalOptions.hasExperiment(pipeline.getOptions(), "use_deprecated_read")) { + if (!ExperimentalOptions.hasExperiment(pipeline.getOptions(), "use_sdf_read")) { convertReadBasedSplittableDoFnsToPrimitiveReads(pipeline); } } 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 499a1f8818e8..2665ccf11602 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 @@ -1223,10 +1223,7 @@ public PCollection> expand(PBegin input) { Coder valueCoder = getValueCoder(coderRegistry); // For read from unbounded in a bounded manner, we actually are not going through Read or SDF. - if (ExperimentalOptions.hasExperiment( - input.getPipeline().getOptions(), "beam_fn_api_use_deprecated_read") - || ExperimentalOptions.hasExperiment( - input.getPipeline().getOptions(), "use_deprecated_read") + if (!ExperimentalOptions.hasExperiment(input.getPipeline().getOptions(), "use_sdf_read") || getMaxNumRecords() < Long.MAX_VALUE || getMaxReadTime() != null) { return input.apply(new ReadFromKafkaViaUnbounded<>(this, keyCoder, valueCoder));