This is an automated email from the ASF dual-hosted git repository.
guoweijie pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
from 3841f062255 [FLINK-34562][docs-zh] Port Debezium Avro Confluent
changes (FLINK-34509) to Chinese
new 5b2e923be0a [FLINK-34548][API] Initialize the datastream v2 related
modules
new 59525e460af [FLINK-34548][API] Create flink-core-api module and let
flink-core depend on it
new 13790e03207 [FLINK-34548][API] Move Function interface to
flink-core-api
new 13cfaa76b5e [FLINK-34548][API] Introduce ProcessFunction and
RuntimeContext related interfaces
new cedbcce6eff [FLINK-34548][API] Introduce variants of ProcessFunction
new 9fa74a8a706 [FLINK-34548][API] Introduce stream interface and move
KeySelector to flink-core-api
new e1147ca7e39 [FLINK-34548][API] Introduce ExecutionEnvironment
new 4f71c5b4660 [FLINK-34548][API] Implement process function's underlying
operators
new ceafa5a5705 [FLINK-34548][API] Implement datastream
new 056660e0b69 [FLINK-34548][API] Supports FLIP-27 Source
new 28762497bdf [FLINK-34548][API] Supports sink-v2 Sink
The 11 revisions listed above as "new" are entirely new to this
repository and will be described in separate emails. The revisions
listed as "add" were already present in the repository and have only
been added to this reference.
Summary of changes:
flink-clients/pom.xml | 6 +
.../java/org/apache/flink/client/ClientUtils.java | 7 +
flink-core-api/pom.xml | 43 +++
.../flink/api/common/RuntimeExecutionMode.java | 0
.../flink/api/common/functions/Function.java | 0
.../org/apache/flink/api/connector/dsv2/Sink.java | 14 +-
.../apache/flink/api/connector/dsv2/Source.java | 14 +-
.../flink/api/java/functions/KeySelector.java | 0
flink-core/pom.xml | 6 +
.../dsv2/DataStreamV2SinkUtils.java} | 25 +-
.../connector/dsv2/DataStreamV2SourceUtils.java | 52 +++
.../dsv2/FromDataSource.java} | 27 +-
.../dsv2/WrappedSink.java} | 25 +-
.../dsv2/WrappedSource.java} | 25 +-
flink-datastream-api/pom.xml | 44 +++
.../flink/datastream/api/ExecutionEnvironment.java | 55 +++
.../flink/datastream/api/common/Collector.java | 30 +-
.../api/context/NonPartitionedContext.java | 21 +-
.../datastream/api/context/RuntimeContext.java | 19 +-
.../context/TwoOutputNonPartitionedContext.java | 22 +-
.../api/function/ApplyPartitionFunction.java | 27 +-
.../function/OneInputStreamProcessFunction.java | 45 +++
.../datastream/api/function/ProcessFunction.java | 50 +++
.../TwoInputBroadcastStreamProcessFunction.java | 68 ++++
.../TwoInputNonBroadcastStreamProcessFunction.java | 64 ++++
.../function/TwoOutputApplyPartitionFunction.java | 39 ++
.../function/TwoOutputStreamProcessFunction.java | 48 +++
.../datastream/api/stream/BroadcastStream.java | 76 ++++
.../flink/datastream/api/stream/DataStream.java | 15 +-
.../flink/datastream/api/stream/GlobalStream.java | 96 +++++
.../api/stream/KeyedPartitionStream.java | 212 +++++++++++
.../api/stream/NonKeyedPartitionStream.java | 118 +++++++
flink-datastream/pom.xml | 73 ++++
.../impl/ExecutionContextEnvironment.java | 58 +++
.../impl/ExecutionEnvironmentFactory.java | 23 +-
.../datastream/impl/ExecutionEnvironmentImpl.java | 359 +++++++++++++++++++
.../impl/common/KeyCheckedOutputCollector.java | 72 ++++
.../datastream/impl/common/OutputCollector.java | 33 +-
.../datastream/impl/common/TimestampCollector.java | 33 +-
.../impl/context/DefaultNonPartitionedContext.java | 21 +-
.../impl/context/DefaultRuntimeContext.java | 16 +-
.../DefaultTwoOutputNonPartitionedContext.java | 22 +-
.../impl/operators/KeyedProcessOperator.java | 53 +++
.../KeyedTwoInputBroadcastProcessOperator.java | 54 +++
.../KeyedTwoInputNonBroadcastProcessOperator.java | 56 +++
.../operators/KeyedTwoOutputProcessOperator.java | 78 ++++
.../datastream/impl/operators/ProcessOperator.java | 71 ++++
.../TwoInputBroadcastProcessOperator.java | 86 +++++
.../TwoInputNonBroadcastProcessOperator.java | 86 +++++
.../impl/operators/TwoOutputProcessOperator.java | 113 ++++++
.../datastream/impl/stream/AbstractDataStream.java | 83 +++++
.../impl/stream/BroadcastStreamImpl.java | 127 +++++++
.../datastream/impl/stream/GlobalStreamImpl.java | 186 ++++++++++
.../impl/stream/KeyedPartitionStreamImpl.java | 393 +++++++++++++++++++++
.../impl/stream/NonKeyedPartitionStreamImpl.java | 199 +++++++++++
.../flink/datastream/impl/utils/StreamUtils.java | 288 +++++++++++++++
.../DataStreamV2SinkTransformation.java | 96 +++++
.../DataStreamV2SinkTransformationTranslator.java | 367 +++++++++++++++++++
.../impl/ExecutionEnvironmentImplTest.java | 124 +++++++
.../impl/TestingExecutionEnvironmentFactory.java | 39 ++
.../datastream/impl/TestingTransformation.java | 33 +-
.../impl/common/KeyCheckedOutputCollectorTest.java | 72 ++++
.../impl/common/OutputCollectorTest.java | 53 +++
.../impl/common/TestingTimestampCollector.java | 72 ++++
.../impl/operators/KeyedProcessOperatorTest.java | 123 +++++++
.../KeyedTwoInputBroadcastProcessOperatorTest.java | 155 ++++++++
...yedTwoInputNonBroadcastProcessOperatorTest.java | 156 ++++++++
.../KeyedTwoOutputProcessOperatorTest.java | 158 +++++++++
.../impl/operators/ProcessOperatorTest.java | 83 +++++
.../TwoInputBroadcastProcessOperatorTest.java | 112 ++++++
.../TwoInputNonBroadcastProcessOperatorTest.java | 115 ++++++
.../operators/TwoOutputProcessOperatorTest.java | 108 ++++++
.../impl/stream/BroadcastStreamImplTest.java | 86 +++++
.../impl/stream/GlobalStreamImplTest.java | 91 +++++
.../impl/stream/KeyedPartitionStreamImplTest.java | 168 +++++++++
.../stream/NonKeyedPartitionStreamImplTest.java | 134 +++++++
.../datastream/impl/stream/StreamTestUtils.java | 119 +++++++
.../datastream/impl/utils/StreamUtilsTest.java | 197 +++++++++++
.../util/TwoInputStreamOperatorTestHarness.java | 13 +
pom.xml | 5 +
80 files changed, 6377 insertions(+), 178 deletions(-)
create mode 100644 flink-core-api/pom.xml
rename {flink-core =>
flink-core-api}/src/main/java/org/apache/flink/api/common/RuntimeExecutionMode.java
(100%)
copy {flink-core =>
flink-core-api}/src/main/java/org/apache/flink/api/common/functions/Function.java
(100%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-core-api/src/main/java/org/apache/flink/api/connector/dsv2/Sink.java (68%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-core-api/src/main/java/org/apache/flink/api/connector/dsv2/Source.java
(68%)
rename {flink-core =>
flink-core-api}/src/main/java/org/apache/flink/api/java/functions/KeySelector.java
(100%)
copy
flink-core/src/main/java/org/apache/flink/api/{common/functions/Function.java
=> connector/dsv2/DataStreamV2SinkUtils.java} (59%)
create mode 100644
flink-core/src/main/java/org/apache/flink/api/connector/dsv2/DataStreamV2SourceUtils.java
copy
flink-core/src/main/java/org/apache/flink/api/{common/functions/Function.java
=> connector/dsv2/FromDataSource.java} (64%)
copy
flink-core/src/main/java/org/apache/flink/api/{common/functions/Function.java
=> connector/dsv2/WrappedSink.java} (60%)
copy
flink-core/src/main/java/org/apache/flink/api/{common/functions/Function.java
=> connector/dsv2/WrappedSource.java} (58%)
create mode 100644 flink-datastream-api/pom.xml
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/ExecutionEnvironment.java
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/common/Collector.java
(57%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/context/NonPartitionedContext.java
(56%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/context/RuntimeContext.java
(57%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/context/TwoOutputNonPartitionedContext.java
(53%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/ApplyPartitionFunction.java
(53%)
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/OneInputStreamProcessFunction.java
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/ProcessFunction.java
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputBroadcastStreamProcessFunction.java
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoInputNonBroadcastStreamProcessFunction.java
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoOutputApplyPartitionFunction.java
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/function/TwoOutputStreamProcessFunction.java
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/stream/BroadcastStream.java
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/stream/DataStream.java
(67%)
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/stream/GlobalStream.java
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/stream/KeyedPartitionStream.java
create mode 100644
flink-datastream-api/src/main/java/org/apache/flink/datastream/api/stream/NonKeyedPartitionStream.java
create mode 100644 flink-datastream/pom.xml
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/ExecutionContextEnvironment.java
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream/src/main/java/org/apache/flink/datastream/impl/ExecutionEnvironmentFactory.java
(62%)
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/ExecutionEnvironmentImpl.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/common/KeyCheckedOutputCollector.java
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream/src/main/java/org/apache/flink/datastream/impl/common/OutputCollector.java
(50%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream/src/main/java/org/apache/flink/datastream/impl/common/TimestampCollector.java
(50%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultNonPartitionedContext.java
(62%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultRuntimeContext.java
(67%)
copy
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultTwoOutputNonPartitionedContext.java
(57%)
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/KeyedProcessOperator.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/KeyedTwoInputBroadcastProcessOperator.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/KeyedTwoInputNonBroadcastProcessOperator.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/KeyedTwoOutputProcessOperator.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/stream/AbstractDataStream.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/stream/BroadcastStreamImpl.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/stream/GlobalStreamImpl.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/stream/KeyedPartitionStreamImpl.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/stream/NonKeyedPartitionStreamImpl.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/datastream/impl/utils/StreamUtils.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/streaming/api/transformations/DataStreamV2SinkTransformation.java
create mode 100644
flink-datastream/src/main/java/org/apache/flink/streaming/runtime/translators/DataStreamV2SinkTransformationTranslator.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/ExecutionEnvironmentImplTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/TestingExecutionEnvironmentFactory.java
rename
flink-core/src/main/java/org/apache/flink/api/common/functions/Function.java =>
flink-datastream/src/test/java/org/apache/flink/datastream/impl/TestingTransformation.java
(53%)
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/common/KeyCheckedOutputCollectorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/common/OutputCollectorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/common/TestingTimestampCollector.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/operators/KeyedProcessOperatorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/operators/KeyedTwoInputBroadcastProcessOperatorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/operators/KeyedTwoInputNonBroadcastProcessOperatorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/operators/KeyedTwoOutputProcessOperatorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/operators/ProcessOperatorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperatorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperatorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperatorTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/stream/BroadcastStreamImplTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/stream/GlobalStreamImplTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/stream/KeyedPartitionStreamImplTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/stream/NonKeyedPartitionStreamImplTest.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/stream/StreamTestUtils.java
create mode 100644
flink-datastream/src/test/java/org/apache/flink/datastream/impl/utils/StreamUtilsTest.java