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

Reply via email to