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