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 dbd21cf56af950045f75ca38b536ef06d4c8cb8f Author: Xu Huang <[email protected]> AuthorDate: Thu Nov 28 15:16:28 2024 +0800 [FLINK-36815][API] Introduce a NonPartitionedContext parameter in Open method of UserFunction for DSv2 --- .../api/function/OneInputStreamProcessFunction.java | 12 ++++++++++++ .../flink/datastream/api/function/ProcessFunction.java | 12 ------------ .../api/function/TwoInputBroadcastStreamProcessFunction.java | 12 ++++++++++++ .../function/TwoInputNonBroadcastStreamProcessFunction.java | 12 ++++++++++++ .../api/function/TwoOutputStreamProcessFunction.java | 12 ++++++++++++ .../flink/datastream/impl/operators/ProcessOperator.java | 2 +- .../impl/operators/TwoInputBroadcastProcessOperator.java | 2 +- .../impl/operators/TwoInputNonBroadcastProcessOperator.java | 2 +- .../datastream/impl/operators/TwoOutputProcessOperator.java | 2 +- .../flink/datastream/impl/functions/ProcessFunctionTest.java | 10 ++++++---- 10 files changed, 58 insertions(+), 20 deletions(-) diff --git a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/OneInputStreamProcessFunction.java b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/OneInputStreamProcessFunction.java index 342de0a6ed0..afe0b433f0e 100644 --- a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/OneInputStreamProcessFunction.java +++ b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/OneInputStreamProcessFunction.java @@ -26,6 +26,18 @@ import org.apache.flink.datastream.api.context.PartitionedContext; /** This contains all logical related to process records from single input. */ @Experimental public interface OneInputStreamProcessFunction<IN, OUT> extends ProcessFunction { + /** + * Initialization method for the function. It is called before the actual working methods (like + * processRecord) and thus suitable for one time setup work. + * + * <p>By default, this method does nothing. + * + * @throws Exception Implementations may forward exceptions, which are caught by the runtime. + * When the runtime catches an exception, it aborts the task and lets the fail-over logic + * decide whether to retry the task execution. + */ + default void open(NonPartitionedContext<OUT> ctx) throws Exception {} + /** * Process record and emit data through {@link Collector}. * diff --git a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/ProcessFunction.java b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/ProcessFunction.java index 9721aba35ec..95ab3f22543 100644 --- a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/ProcessFunction.java +++ b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/ProcessFunction.java @@ -28,18 +28,6 @@ import java.util.Set; /** Base class for all user defined process functions. */ @Experimental public interface ProcessFunction extends Function { - /** - * Initialization method for the function. It is called before the actual working methods (like - * processRecord) and thus suitable for one time setup work. - * - * <p>By default, this method does nothing. - * - * @throws Exception Implementations may forward exceptions, which are caught by the runtime. - * When the runtime catches an exception, it aborts the task and lets the fail-over logic - * decide whether to retry the task execution. - */ - default void open() throws Exception {} - /** * Explicitly declares states upfront. Each specific state must be declared in this method * before it can be used. diff --git a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputBroadcastStreamProcessFunction.java b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputBroadcastStreamProcessFunction.java index 32d7168c59c..4789c6bbc9b 100644 --- a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputBroadcastStreamProcessFunction.java +++ b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputBroadcastStreamProcessFunction.java @@ -29,6 +29,18 @@ import org.apache.flink.datastream.api.context.PartitionedContext; */ @Experimental public interface TwoInputBroadcastStreamProcessFunction<IN1, IN2, OUT> extends ProcessFunction { + /** + * Initialization method for the function. It is called before the actual working methods (like + * processRecord) and thus suitable for one time setup work. + * + * <p>By default, this method does nothing. + * + * @throws Exception Implementations may forward exceptions, which are caught by the runtime. + * When the runtime catches an exception, it aborts the task and lets the fail-over logic + * decide whether to retry the task execution. + */ + default void open(NonPartitionedContext<OUT> ctx) throws Exception {} + /** * Process record from non-broadcast input and emit data through {@link Collector}. * diff --git a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputNonBroadcastStreamProcessFunction.java b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputNonBroadcastStreamProcessFunction.java index 6fe67a28f2f..abb88089535 100644 --- a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputNonBroadcastStreamProcessFunction.java +++ b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputNonBroadcastStreamProcessFunction.java @@ -26,6 +26,18 @@ import org.apache.flink.datastream.api.context.PartitionedContext; /** This contains all logical related to process records from two non-broadcast input. */ @Experimental public interface TwoInputNonBroadcastStreamProcessFunction<IN1, IN2, OUT> extends ProcessFunction { + /** + * Initialization method for the function. It is called before the actual working methods (like + * processRecord) and thus suitable for one time setup work. + * + * <p>By default, this method does nothing. + * + * @throws Exception Implementations may forward exceptions, which are caught by the runtime. + * When the runtime catches an exception, it aborts the task and lets the fail-over logic + * decide whether to retry the task execution. + */ + default void open(NonPartitionedContext<OUT> ctx) throws Exception {} + /** * Process record from the first input and emit data through {@link Collector}. * diff --git a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoOutputStreamProcessFunction.java b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoOutputStreamProcessFunction.java index a9f76e70239..db895306ee0 100644 --- a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoOutputStreamProcessFunction.java +++ b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoOutputStreamProcessFunction.java @@ -26,6 +26,18 @@ import org.apache.flink.datastream.api.context.TwoOutputNonPartitionedContext; /** This contains all logical related to process and emit records to two output streams. */ @Experimental public interface TwoOutputStreamProcessFunction<IN, OUT1, OUT2> extends ProcessFunction { + /** + * Initialization method for the function. It is called before the actual working methods (like + * processRecord) and thus suitable for one time setup work. + * + * <p>By default, this method does nothing. + * + * @throws Exception Implementations may forward exceptions, which are caught by the runtime. + * When the runtime catches an exception, it aborts the task and lets the fail-over logic + * decide whether to retry the task execution. + */ + default void open(TwoOutputNonPartitionedContext<OUT1, OUT2> ctx) throws Exception {} + /** * Process and emit record to the first/second output through {@link Collector}s. * diff --git a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java index ce703012b25..28eeb84afc2 100644 --- a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java +++ b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java @@ -76,7 +76,7 @@ public class ProcessOperator<IN, OUT> getOperatorStateBackend()); outputCollector = getOutputCollector(); nonPartitionedContext = getNonPartitionedContext(); - userFunction.open(); + userFunction.open(nonPartitionedContext); } @Override diff --git a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java index cf707f8dedf..8c066145cd4 100644 --- a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java +++ b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java @@ -80,7 +80,7 @@ public class TwoInputBroadcastProcessOperator<IN1, IN2, OUT> operatorContext, getOperatorStateBackend()); this.nonPartitionedContext = getNonPartitionedContext(); - this.userFunction.open(); + this.userFunction.open(this.nonPartitionedContext); } @Override diff --git a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java index d952efed124..b4ca7cac2e1 100644 --- a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java +++ b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java @@ -83,7 +83,7 @@ public class TwoInputNonBroadcastProcessOperator<IN1, IN2, OUT> operatorContext, operatorStateBackend); this.nonPartitionedContext = getNonPartitionedContext(); - this.userFunction.open(); + this.userFunction.open(this.nonPartitionedContext); } @Override diff --git a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java index fb6209ef18f..a3c7d348ab0 100644 --- a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java +++ b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java @@ -92,7 +92,7 @@ public class TwoOutputProcessOperator<IN, OUT_MAIN, OUT_SIDE> operatorContext, operatorStateStore); this.nonPartitionedContext = getNonPartitionedContext(); - this.userFunction.open(); + this.userFunction.open(this.nonPartitionedContext); } @Override diff --git a/flink-datastream/src/test/java/org/apache/flink/datastream/impl/functions/ProcessFunctionTest.java b/flink-datastream/src/test/java/org/apache/flink/datastream/impl/functions/ProcessFunctionTest.java index 2d4330e204e..4bd1b4b0f5d 100644 --- a/flink-datastream/src/test/java/org/apache/flink/datastream/impl/functions/ProcessFunctionTest.java +++ b/flink-datastream/src/test/java/org/apache/flink/datastream/impl/functions/ProcessFunctionTest.java @@ -22,6 +22,7 @@ import org.apache.flink.api.common.typeinfo.Types; import org.apache.flink.datastream.api.common.Collector; import org.apache.flink.datastream.api.context.NonPartitionedContext; import org.apache.flink.datastream.api.context.PartitionedContext; +import org.apache.flink.datastream.api.context.TwoOutputNonPartitionedContext; import org.apache.flink.datastream.api.function.OneInputStreamProcessFunction; import org.apache.flink.datastream.api.function.ProcessFunction; import org.apache.flink.datastream.api.function.TwoInputBroadcastStreamProcessFunction; @@ -51,7 +52,7 @@ public class ProcessFunctionTest { OneInputStreamProcessFunction<Integer, Integer> processFunction = new OneInputStreamProcessFunction<>() { @Override - public void open() throws Exception { + public void open(NonPartitionedContext<Integer> ctx) throws Exception { openInvoked.set(true); } @@ -86,7 +87,7 @@ public class ProcessFunctionTest { new TwoInputBroadcastStreamProcessFunction<>() { @Override - public void open() throws Exception { + public void open(NonPartitionedContext<Integer> ctx) throws Exception { openInvoked.set(true); } @@ -126,7 +127,7 @@ public class ProcessFunctionTest { new TwoInputNonBroadcastStreamProcessFunction<>() { @Override - public void open() throws Exception { + public void open(NonPartitionedContext<Integer> ctx) throws Exception { openInvoked.set(true); } @@ -167,7 +168,8 @@ public class ProcessFunctionTest { new TwoOutputStreamProcessFunction<>() { @Override - public void open() throws Exception { + public void open(TwoOutputNonPartitionedContext<Integer, Integer> ctx) + throws Exception { openInvoked.set(true); }
