This is an automated email from the ASF dual-hosted git repository.
zakelly pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
from f5c44b683cd [FLINK-36783][table-planner] Use validated SqlNode in
`asQuery` while converting `SqlCreateTableAS` to `CreateTableASOperation`
add 734e736a429 [hotfix] Remove unnecessary 'V2' in names of state-related
classes
add faefd9ede07 [FLINK-35156][State] Support creating operator state from
the state v2 descriptors
add 7c03e9383ec [hotfix][Runtime/State] Consider
`isAsyncStateProcessingEnabled` during `initializeState`
add 1360064d747 [FLINK-35156][Runtime] Rework: Make operators of
DataStream V2 integrate with async state processing framework
add 22cc255ad99 [FLINK-35156] Use State V2 API in
`StatefulDataStreamV2ITCase`
No new revisions were added by this update.
Summary of changes:
.../flink/datastream/api/context/StateManager.java | 10 +--
.../impl/context/DefaultPartitionedContext.java | 8 +--
.../impl/context/DefaultStateManager.java | 78 ++++++++++------------
.../impl/operators/KeyedProcessOperator.java | 28 +++-----
.../KeyedTwoInputBroadcastProcessOperator.java | 28 +++-----
.../KeyedTwoInputNonBroadcastProcessOperator.java | 33 +++------
.../operators/KeyedTwoOutputProcessOperator.java | 29 +++-----
.../datastream/impl/operators/ProcessOperator.java | 29 +++++++-
.../TwoInputBroadcastProcessOperator.java | 29 +++++++-
.../TwoInputNonBroadcastProcessOperator.java | 29 +++++++-
.../impl/operators/TwoOutputProcessOperator.java | 31 +++++++--
.../context/DefaultNonPartitionedContextTest.java | 13 +++-
.../impl/context/DefaultStateManagerTest.java | 10 ++-
.../DefaultTwoOutputNonPartitionedContextTest.java | 13 +++-
.../operators/MockFreqCountProcessFunction.java | 2 +-
.../MockGlobalListAppenderProcessFunction.java | 2 +-
.../operators/MockListAppenderProcessFunction.java | 2 +-
.../operators/MockMultiplierProcessFunction.java | 2 +-
.../MockRecudingMultiplierProcessFunction.java | 2 +-
.../operators/MockSumAggregateProcessFunction.java | 2 +-
.../AbstractAsyncStateStreamOperator.java | 55 ++++++++-------
.../AbstractAsyncStateStreamOperatorV2.java | 47 +++++++------
.../runtime/state/DefaultOperatorStateBackend.java | 29 ++++++++
.../flink/runtime/state/OperatorStateBackend.java | 6 +-
...ateStoreV2.java => DefaultKeyedStateStore.java} | 6 +-
...KeyedStateStoreV2.java => KeyedStateStore.java} | 2 +-
.../runtime/state/v2}/OperatorStateStore.java | 8 ++-
.../runtime/state/v2/StateDescriptorUtils.java | 39 +++++++++++
...eAdaptor.java => OperatorListStateAdaptor.java} | 36 +++++++---
.../api/operators/StreamOperatorStateHandler.java | 9 ++-
.../api/operators/StreamingRuntimeContext.java | 39 +++++------
.../api/operators/StreamingRuntimeContextTest.java | 5 +-
.../collect/utils/MockOperatorStateStore.java | 28 +++++++-
.../api/datastream/StatefulDataStreamV2ITCase.java | 10 +--
34 files changed, 443 insertions(+), 256 deletions(-)
rename
flink-runtime/src/main/java/org/apache/flink/runtime/state/v2/{DefaultKeyedStateStoreV2.java
=> DefaultKeyedStateStore.java} (95%)
rename
flink-runtime/src/main/java/org/apache/flink/runtime/state/v2/{KeyedStateStoreV2.java
=> KeyedStateStore.java} (99%)
copy {flink-core/src/main/java/org/apache/flink/api/common/state =>
flink-runtime/src/main/java/org/apache/flink/runtime/state/v2}/OperatorStateStore.java
(96%)
copy
flink-runtime/src/main/java/org/apache/flink/runtime/state/v2/adaptor/{ListStateAdaptor.java
=> OperatorListStateAdaptor.java} (80%)