This is an automated email from the ASF dual-hosted git repository.

guoweijie pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git

commit abd288a6106dc6ab3aa38f11728766c51fe851f0
Author: Xu Huang <[email protected]>
AuthorDate: Thu Nov 28 14:35:58 2024 +0800

    [FLINK-36815][API] Make ExecutionEnvironment#fromSource return 
ProcessConfigurableAndNonKeyedPartitionStream
---
 .../org/apache/flink/datastream/api/ExecutionEnvironment.java    | 5 +++--
 .../apache/flink/datastream/impl/ExecutionEnvironmentImpl.java   | 9 ++++++---
 2 files changed, 9 insertions(+), 5 deletions(-)

diff --git 
a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/ExecutionEnvironment.java
 
b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/ExecutionEnvironment.java
index 5dd505e1678..c113552d5d3 100644
--- 
a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/ExecutionEnvironment.java
+++ 
b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/ExecutionEnvironment.java
@@ -21,7 +21,7 @@ package org.apache.flink.datastream.api;
 import org.apache.flink.annotation.Experimental;
 import org.apache.flink.api.common.RuntimeExecutionMode;
 import org.apache.flink.api.connector.dsv2.Source;
-import org.apache.flink.datastream.api.stream.NonKeyedPartitionStream;
+import 
org.apache.flink.datastream.api.stream.NonKeyedPartitionStream.ProcessConfigurableAndNonKeyedPartitionStream;
 
 /**
  * This is the context in which a program is executed.
@@ -51,5 +51,6 @@ public interface ExecutionEnvironment {
     /** Set the execution mode for this environment. */
     ExecutionEnvironment setExecutionMode(RuntimeExecutionMode runtimeMode);
 
-    <OUT> NonKeyedPartitionStream<OUT> fromSource(Source<OUT> source, String 
sourceName);
+    <OUT> ProcessConfigurableAndNonKeyedPartitionStream<OUT> fromSource(
+            Source<OUT> source, String sourceName);
 }
diff --git 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/ExecutionEnvironmentImpl.java
 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/ExecutionEnvironmentImpl.java
index b9a1fe045fb..c79804844ce 100644
--- 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/ExecutionEnvironmentImpl.java
+++ 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/ExecutionEnvironmentImpl.java
@@ -42,8 +42,9 @@ import org.apache.flink.core.execution.PipelineExecutor;
 import org.apache.flink.core.execution.PipelineExecutorFactory;
 import org.apache.flink.core.execution.PipelineExecutorServiceLoader;
 import org.apache.flink.datastream.api.ExecutionEnvironment;
-import org.apache.flink.datastream.api.stream.NonKeyedPartitionStream;
+import 
org.apache.flink.datastream.api.stream.NonKeyedPartitionStream.ProcessConfigurableAndNonKeyedPartitionStream;
 import org.apache.flink.datastream.impl.stream.NonKeyedPartitionStreamImpl;
+import org.apache.flink.datastream.impl.utils.StreamUtils;
 import org.apache.flink.streaming.api.environment.CheckpointConfig;
 import org.apache.flink.streaming.api.graph.StreamGraph;
 import org.apache.flink.streaming.api.graph.StreamGraphGenerator;
@@ -161,7 +162,8 @@ public class ExecutionEnvironmentImpl implements 
ExecutionEnvironment {
     }
 
     @Override
-    public <OUT> NonKeyedPartitionStream<OUT> fromSource(Source<OUT> source, 
String sourceName) {
+    public <OUT> ProcessConfigurableAndNonKeyedPartitionStream<OUT> fromSource(
+            Source<OUT> source, String sourceName) {
         if (source instanceof WrappedSource) {
             org.apache.flink.api.connector.source.Source<OUT, ?, ?> 
innerSource =
                     ((WrappedSource<OUT>) source).getWrappedSource();
@@ -176,7 +178,8 @@ public class ExecutionEnvironmentImpl implements 
ExecutionEnvironment {
                             resolvedTypeInfo,
                             getParallelism(),
                             false);
-            return new NonKeyedPartitionStreamImpl<>(this, 
sourceTransformation);
+            return StreamUtils.wrapWithConfigureHandle(
+                    new NonKeyedPartitionStreamImpl<>(this, 
sourceTransformation));
         } else if (source instanceof FromDataSource) {
             Collection<OUT> data = ((FromDataSource<OUT>) source).getData();
             TypeInformation<OUT> outType = extractTypeInfoFromCollection(data);

Reply via email to