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));