This is an automated email from the ASF dual-hosted git repository. xtsong pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit c90f30f40bda0874bfdb16b0ea20f8a556947abd Author: sxnan <[email protected]> AuthorDate: Fri May 31 23:50:37 2024 +0800 [FLINK-35475][runtime] Introduce isInternalSorterSupport to OperatorAttributes This closes #24874 --- .../api/operators/OperatorAttributes.java | 12 +- .../api/operators/OperatorAttributesBuilder.java | 8 +- .../sortpartition/KeyedSortPartitionOperator.java | 5 +- .../AbstractMultipleInputTransformation.java | 4 + .../transformations/OneInputTransformation.java | 4 + .../transformations/TwoInputTransformation.java | 4 + .../MultiInputTransformationTranslator.java | 3 +- .../OneInputTransformationTranslator.java | 14 +- .../TwoInputTransformationTranslator.java | 5 +- ...obGraphGeneratorWithOperatorAttributesTest.java | 174 +++++++++++++++++++++ 10 files changed, 214 insertions(+), 19 deletions(-) diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributes.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributes.java index 6c699e15f9f..db883c03631 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributes.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributes.java @@ -30,9 +30,11 @@ public class OperatorAttributes implements Serializable { private static final long serialVersionUID = 1L; private final boolean outputOnlyAfterEndOfStream; + private final boolean internalSorterSupported; - OperatorAttributes(boolean outputOnlyAfterEndOfStream) { + OperatorAttributes(boolean outputOnlyAfterEndOfStream, boolean internalSorterSupported) { this.outputOnlyAfterEndOfStream = outputOnlyAfterEndOfStream; + this.internalSorterSupported = internalSorterSupported; } /** @@ -49,4 +51,12 @@ public class OperatorAttributes implements Serializable { public boolean isOutputOnlyAfterEndOfStream() { return outputOnlyAfterEndOfStream; } + + /** + * Returns true iff the operator uses an internal sorter to sort inputs by key. When it is true, + * the input records will not to be sorted externally before being fed into this operator. + */ + public boolean isInternalSorterSupported() { + return internalSorterSupported; + } } diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributesBuilder.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributesBuilder.java index 3239444477f..a504aaa49cd 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributesBuilder.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributesBuilder.java @@ -23,6 +23,7 @@ import org.apache.flink.annotation.Experimental; @Experimental public class OperatorAttributesBuilder { private boolean outputOnlyAfterEndOfStream = false; + private boolean internalSorterSupported = false; /** * Set to true if and only if the operator only emits records after all its inputs have ended. @@ -34,7 +35,12 @@ public class OperatorAttributesBuilder { return this; } + public OperatorAttributesBuilder setInternalSorterSupported(boolean internalSorterSupported) { + this.internalSorterSupported = internalSorterSupported; + return this; + } + public OperatorAttributes build() { - return new OperatorAttributes(outputOnlyAfterEndOfStream); + return new OperatorAttributes(outputOnlyAfterEndOfStream, internalSorterSupported); } } diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/sortpartition/KeyedSortPartitionOperator.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/sortpartition/KeyedSortPartitionOperator.java index 8c45a54ed5a..dc4458d7176 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/sortpartition/KeyedSortPartitionOperator.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/sortpartition/KeyedSortPartitionOperator.java @@ -250,7 +250,10 @@ public class KeyedSortPartitionOperator<INPUT, KEY> extends AbstractStreamOperat @Override public OperatorAttributes getOperatorAttributes() { - return new OperatorAttributesBuilder().setOutputOnlyAfterEndOfStream(true).build(); + return new OperatorAttributesBuilder() + .setOutputOnlyAfterEndOfStream(true) + .setInternalSorterSupported(true) + .build(); } /** diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/AbstractMultipleInputTransformation.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/AbstractMultipleInputTransformation.java index 154ba63b0f3..7b97bc8ecbd 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/AbstractMultipleInputTransformation.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/AbstractMultipleInputTransformation.java @@ -94,4 +94,8 @@ public abstract class AbstractMultipleInputTransformation<OUT> extends PhysicalT public boolean isOutputOnlyAfterEndOfStream() { return operatorFactory.getOperatorAttributes().isOutputOnlyAfterEndOfStream(); } + + public boolean isInternalSorterSupported() { + return operatorFactory.getOperatorAttributes().isInternalSorterSupported(); + } } diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/OneInputTransformation.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/OneInputTransformation.java index ba8162d8e5b..6aeaf2d96ee 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/OneInputTransformation.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/OneInputTransformation.java @@ -188,4 +188,8 @@ public class OneInputTransformation<IN, OUT> extends PhysicalTransformation<OUT> public boolean isOutputOnlyAfterEndOfStream() { return operatorFactory.getOperatorAttributes().isOutputOnlyAfterEndOfStream(); } + + public boolean isInternalSorterSupported() { + return operatorFactory.getOperatorAttributes().isInternalSorterSupported(); + } } diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/TwoInputTransformation.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/TwoInputTransformation.java index f7b60c15d2f..007f908c809 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/TwoInputTransformation.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/transformations/TwoInputTransformation.java @@ -235,4 +235,8 @@ public class TwoInputTransformation<IN1, IN2, OUT> extends PhysicalTransformatio public boolean isOutputOnlyAfterEndOfStream() { return operatorFactory.getOperatorAttributes().isOutputOnlyAfterEndOfStream(); } + + public boolean isInternalSorterSupported() { + return operatorFactory.getOperatorAttributes().isInternalSorterSupported(); + } } diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/MultiInputTransformationTranslator.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/MultiInputTransformationTranslator.java index b2a39d1acf4..be2980f3ab6 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/MultiInputTransformationTranslator.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/MultiInputTransformationTranslator.java @@ -132,7 +132,8 @@ public class MultiInputTransformationTranslator<OUT> private void maybeApplyBatchExecutionSettings( final AbstractMultipleInputTransformation<OUT> transformation, final Context context) { - if (transformation instanceof KeyedMultipleInputTransformation) { + if (transformation instanceof KeyedMultipleInputTransformation + && !transformation.isInternalSorterSupported()) { KeyedMultipleInputTransformation<OUT> keyedTransformation = (KeyedMultipleInputTransformation<OUT>) transformation; List<Transformation<?>> inputs = transformation.getInputs(); diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/OneInputTransformationTranslator.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/OneInputTransformationTranslator.java index 804d1f7c09a..6d7ae8103f1 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/OneInputTransformationTranslator.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/OneInputTransformationTranslator.java @@ -22,9 +22,6 @@ import org.apache.flink.annotation.Internal; import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.streaming.api.graph.StreamConfig; import org.apache.flink.streaming.api.graph.TransformationTranslator; -import org.apache.flink.streaming.api.operators.SimpleOperatorFactory; -import org.apache.flink.streaming.api.operators.StreamOperatorFactory; -import org.apache.flink.streaming.api.operators.sortpartition.KeyedSortPartitionOperator; import org.apache.flink.streaming.api.transformations.OneInputTransformation; import java.util.Collection; @@ -80,16 +77,7 @@ public final class OneInputTransformationTranslator<IN, OUT> private void maybeApplyBatchExecutionSettings( final OneInputTransformation<IN, OUT> transformation, final Context context) { KeySelector<IN, ?> keySelector = transformation.getStateKeySelector(); - if (keySelector != null) { - // KeyedSortPartitionOperator doesn't need sorted input because it sorts data - // internally. - StreamOperatorFactory<OUT> operatorFactory = transformation.getOperatorFactory(); - if (operatorFactory instanceof SimpleOperatorFactory - && operatorFactory.getStreamOperatorClass( - Thread.currentThread().getContextClassLoader()) - == KeyedSortPartitionOperator.class) { - return; - } + if (keySelector != null && !transformation.isInternalSorterSupported()) { BatchExecutionUtils.applyBatchExecutionSettings( transformation.getId(), context, StreamConfig.InputRequirement.SORTED); } diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/TwoInputTransformationTranslator.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/TwoInputTransformationTranslator.java index 13dca02c4e1..8313985e001 100644 --- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/TwoInputTransformationTranslator.java +++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/translators/TwoInputTransformationTranslator.java @@ -89,8 +89,9 @@ public class TwoInputTransformationTranslator<IN1, IN2, OUT> ? StreamConfig.InputRequirement.SORTED : StreamConfig.InputRequirement.PASS_THROUGH; - if (input1Requirement == StreamConfig.InputRequirement.SORTED - || input2Requirement == StreamConfig.InputRequirement.SORTED) { + if ((input1Requirement == StreamConfig.InputRequirement.SORTED + || input2Requirement == StreamConfig.InputRequirement.SORTED) + && !transformation.isInternalSorterSupported()) { BatchExecutionUtils.applyBatchExecutionSettings( transformation.getId(), context, input1Requirement, input2Requirement); } diff --git a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGeneratorWithOperatorAttributesTest.java b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGeneratorWithOperatorAttributesTest.java index 77536de7bc0..94c9e2ec173 100644 --- a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGeneratorWithOperatorAttributesTest.java +++ b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGeneratorWithOperatorAttributesTest.java @@ -19,6 +19,7 @@ package org.apache.flink.streaming.api.graph; import org.apache.flink.api.common.functions.MapFunction; +import org.apache.flink.api.common.typeinfo.BasicTypeInfo; import org.apache.flink.api.common.typeinfo.Types; import org.apache.flink.configuration.Configuration; import org.apache.flink.runtime.io.network.partition.ResultPartitionType; @@ -28,10 +29,15 @@ import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.co.CoProcessFunction; import org.apache.flink.streaming.api.functions.sink.v2.DiscardingSink; +import org.apache.flink.streaming.api.operators.ChainingStrategy; import org.apache.flink.streaming.api.operators.OperatorAttributes; import org.apache.flink.streaming.api.operators.OperatorAttributesBuilder; import org.apache.flink.streaming.api.operators.StreamMap; +import org.apache.flink.streaming.api.operators.StreamOperator; +import org.apache.flink.streaming.api.operators.StreamOperatorFactory; +import org.apache.flink.streaming.api.operators.StreamOperatorParameters; import org.apache.flink.streaming.api.operators.co.CoProcessOperator; +import org.apache.flink.streaming.api.transformations.KeyedMultipleInputTransformation; import org.apache.flink.util.Collector; import org.junit.jupiter.api.Test; @@ -144,6 +150,139 @@ public class StreamingJobGraphGeneratorWithOperatorAttributesTest { assertThat(node.getManagedMemoryOperatorScopeUseCaseWeights()).hasSize(weightSize); } + @Test + void testOneInputOperatorWithInternalSorterSupported() { + final StreamExecutionEnvironment env = + StreamExecutionEnvironment.getExecutionEnvironment(new Configuration()); + + final DataStream<Integer> source1 = env.fromData(1, 2, 3).name("source1"); + source1.keyBy(x -> x) + .transform( + "internalSorter", + Types.INT, + new StreamOperatorWithConfigurableOperatorAttributes<>( + (MapFunction<Integer, Integer>) value -> value, + new OperatorAttributesBuilder() + .setOutputOnlyAfterEndOfStream(true) + .setInternalSorterSupported(true) + .build())) + .keyBy(x -> x) + .transform( + "noInternalSorter", + Types.INT, + new StreamOperatorWithConfigurableOperatorAttributes<>( + (MapFunction<Integer, Integer>) value -> value, + new OperatorAttributesBuilder() + .setOutputOnlyAfterEndOfStream(true) + .build())) + .sinkTo(new DiscardingSink<>()) + .name("sink"); + + final StreamGraph streamGraph = env.getStreamGraph(false); + Map<String, StreamNode> nodeMap = new HashMap<>(); + for (StreamNode node : streamGraph.getStreamNodes()) { + nodeMap.put(node.getOperatorName(), node); + } + + assertThat(nodeMap.get("internalSorter").getInputRequirements()).isEmpty(); + assertThat(nodeMap.get("noInternalSorter").getInputRequirements().get(0)) + .isEqualTo(StreamConfig.InputRequirement.SORTED); + } + + @Test + void testTwoInputOperatorWithInternalSorterSupported() { + final StreamExecutionEnvironment env = + StreamExecutionEnvironment.getExecutionEnvironment(new Configuration()); + + final DataStream<Integer> source1 = env.fromData(1, 2, 3).name("source1"); + final DataStream<Integer> source2 = env.fromData(1, 2, 3).name("source2"); + source1.keyBy(x -> x) + .connect(source2.keyBy(x -> x)) + .transform( + "internalSorter", + Types.INT, + new TwoInputStreamOperatorWithConfigurableOperatorAttributes<>( + new OperatorAttributesBuilder() + .setOutputOnlyAfterEndOfStream(true) + .setInternalSorterSupported(true) + .build())) + .keyBy(x -> x) + .connect(source2.keyBy(x -> x)) + .transform( + "noInternalSorter", + Types.INT, + new TwoInputStreamOperatorWithConfigurableOperatorAttributes<>( + new OperatorAttributesBuilder() + .setOutputOnlyAfterEndOfStream(true) + .build())) + .sinkTo(new DiscardingSink<>()) + .name("sink"); + + final StreamGraph streamGraph = env.getStreamGraph(false); + Map<String, StreamNode> nodeMap = new HashMap<>(); + for (StreamNode node : streamGraph.getStreamNodes()) { + nodeMap.put(node.getOperatorName(), node); + } + + assertThat(nodeMap.get("internalSorter").getInputRequirements()).isEmpty(); + assertThat(nodeMap.get("noInternalSorter").getInputRequirements().get(0)) + .isEqualTo(StreamConfig.InputRequirement.SORTED); + assertThat(nodeMap.get("noInternalSorter").getInputRequirements().get(1)) + .isEqualTo(StreamConfig.InputRequirement.SORTED); + } + + @Test + void testMultipleInputOperatorWithInternalSorterSupported() { + final StreamExecutionEnvironment env = + StreamExecutionEnvironment.getExecutionEnvironment(new Configuration()); + + final DataStream<Integer> source1 = env.fromData(1, 2, 3).name("source1"); + final DataStream<Integer> source2 = env.fromData(1, 2, 3).name("source2"); + + KeyedMultipleInputTransformation<Integer> transform = + new KeyedMultipleInputTransformation<>( + "internalSorter", + new OperatorAttributesConfigurableOperatorFactory<>( + new OperatorAttributesBuilder() + .setOutputOnlyAfterEndOfStream(true) + .setInternalSorterSupported(true) + .build()), + BasicTypeInfo.INT_TYPE_INFO, + 3, + BasicTypeInfo.INT_TYPE_INFO); + + transform.addInput(source1.keyBy(x -> x).getTransformation(), x -> x); + transform.addInput(source2.getTransformation(), null); + + KeyedMultipleInputTransformation<Integer> transform2 = + new KeyedMultipleInputTransformation<>( + "noInternalSorter", + new OperatorAttributesConfigurableOperatorFactory<>( + new OperatorAttributesBuilder() + .setOutputOnlyAfterEndOfStream(true) + .build()), + BasicTypeInfo.INT_TYPE_INFO, + 3, + BasicTypeInfo.INT_TYPE_INFO); + + transform2.addInput(transform, null); + transform2.addInput(source2.keyBy(x -> x).getTransformation(), x -> x); + + new DataStream<>(env, transform2).sinkTo(new DiscardingSink<>()); + + final StreamGraph streamGraph = env.getStreamGraph(false); + Map<String, StreamNode> nodeMap = new HashMap<>(); + for (StreamNode node : streamGraph.getStreamNodes()) { + nodeMap.put(node.getOperatorName(), node); + } + + assertThat(nodeMap.get("internalSorter").getInputRequirements()).isEmpty(); + assertThat(nodeMap.get("noInternalSorter").getInputRequirements().get(0)) + .isEqualTo(StreamConfig.InputRequirement.PASS_THROUGH); + assertThat(nodeMap.get("noInternalSorter").getInputRequirements().get(1)) + .isEqualTo(StreamConfig.InputRequirement.SORTED); + } + private static class StreamOperatorWithConfigurableOperatorAttributes<IN, OUT> extends StreamMap<IN, OUT> { private final OperatorAttributes attributes; @@ -176,6 +315,41 @@ public class StreamingJobGraphGeneratorWithOperatorAttributesTest { } } + private static class OperatorAttributesConfigurableOperatorFactory<OUT> + implements StreamOperatorFactory<OUT> { + + private final OperatorAttributes operatorAttributes; + + public OperatorAttributesConfigurableOperatorFactory( + OperatorAttributes operatorAttributes) { + this.operatorAttributes = operatorAttributes; + } + + @Override + public <T extends StreamOperator<OUT>> T createStreamOperator( + StreamOperatorParameters<OUT> parameters) { + throw new UnsupportedOperationException(); + } + + @Override + public void setChainingStrategy(ChainingStrategy strategy) {} + + @Override + public ChainingStrategy getChainingStrategy() { + return ChainingStrategy.DEFAULT_CHAINING_STRATEGY; + } + + @Override + public Class<? extends StreamOperator> getStreamOperatorClass(ClassLoader classLoader) { + return StreamMap.class; + } + + @Override + public OperatorAttributes getOperatorAttributes() { + return operatorAttributes; + } + } + private static class NoOpCoProcessFunction<IN1, IN2, OUT> extends CoProcessFunction<IN1, IN2, OUT> { @Override
