This is an automated email from the ASF dual-hosted git repository.
gyfora pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
from 6282650a754 [FLINK-40292][table-runtime] Add UdfMetrics helper for UDF
metrics (#28878)
new 51074d3e860 [FLINK-40176][state-processor-api] Introduce basic
StateCatalog functionality
new 5e16408ec34 [FLINK-40177][state-processor-api] State table utils and
type inference for keyed state
new 3d3fad39055 [FLINK-40177][state-processor-api] Serializer and runtime
changes for type inference
new c3b196fdcbf [FLINK-40177][state-processor-api] Keyed state table
mapping, factory and provider interfaces
new 2b20571445d [FLINK-40177][state-processor-api] Flattened keyed state
table mapping and implementation
new d28f941743b [FLINK-40178][state-processor-api] Keyed state reading
integration tests
The 6 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:
docs/content/docs/libs/state_processor_api.md | 12 +-
.../typeutils/CustomRestoreSerializerFactory.java | 118 ++++
.../api/common/typeutils/base/EnumSerializer.java | 36 +-
.../api/java/typeutils/runtime/PojoSerializer.java | 8 +-
.../typeutils/runtime/PojoSerializerSnapshot.java | 58 ++
.../runtime/PojoSerializerSnapshotData.java | 32 +-
.../EnumSerializerSnapshotMissingClassTest.java | 104 +++
.../PojoSerializerSnapshotLenientReadTest.java | 173 +++++
.../avro/typeutils/AvroSerializerSnapshot.java | 59 +-
.../avro/typeutils/AvroSerializerSnapshotTest.java | 85 +++
flink-libraries/flink-state-processing-api/pom.xml | 9 +-
.../apache/flink/state/api/StateTableUtils.java | 641 +++++++++++++++++++
.../state/api/input/KeyedStateInputFormat.java | 63 +-
.../state/api/input/MultiStateKeyIterator.java | 18 +-
.../state/api/input/OperatorStateInputFormat.java | 11 +-
.../input/deserializer/EnumNameDeserializer.java | 145 +++++
.../input/deserializer/InternalTypeConverter.java | 293 +++++++++
.../MissingClassSerializerFactory.java | 76 +++
.../PojoDeserializerCompatibilitySnapshot.java | 89 +++
.../deserializer/PojoToRowDataDeserializer.java | 314 +++++++++
.../input/operator/KeyedStateReaderOperator.java | 94 +--
.../api/input/operator/StateReaderOperator.java | 105 ++-
.../api/input/operator/WindowReaderOperator.java | 76 +--
.../flink/state/api/runtime/SavepointLoader.java | 58 +-
.../flink/state/api/schema/AvroStateUtils.java | 121 ++++
.../state/api/schema/KeyedStateSchemaInfo.java | 90 +++
.../SerializerSnapshotToLogicalTypeConverter.java | 230 +++++++
.../state/api/schema/StateSchemaExtractor.java | 141 +++++
.../flink/state/api/schema/StateSchemaInfo.java | 82 +++
.../flink/state/catalog/SnapshotDiscovery.java | 337 ++++++++++
.../apache/flink/state/catalog/StateCatalog.java | 702 +++++++++++++++++++++
.../flink/state/catalog/StateCatalogFactory.java | 86 +++
.../flink/state/catalog/StateCatalogOptions.java | 59 ++
.../table/AbstractMultiColumnScanProvider.java | 83 +++
.../AbstractSavepointDataStreamScanProvider.java | 295 +++++++++
.../table/AbstractSavepointDynamicTableSource.java | 138 ++++
.../table/AbstractSingleColumnScanProvider.java | 73 +++
.../state/table/FlattenedKeyedStateReader.java | 148 +++++
.../FlattenedSavepointDataStreamScanProvider.java | 56 ++
.../FlattenedSavepointDynamicTableSource.java | 76 +++
.../state/table/FlattenedStateTableMapping.java | 234 +++++++
.../apache/flink/state/table/KeyedStateReader.java | 319 ++--------
.../flink/state/table/MultiColumnStateMapping.java | 46 ++
.../state/table/SavepointConnectorOptions.java | 157 ++---
.../state/table/SavepointConnectorOptionsUtil.java | 16 +-
.../table/SavepointDataStreamScanProvider.java | 133 +---
.../state/table/SavepointDynamicTableSource.java | 114 ++--
.../table/SavepointDynamicTableSourceFactory.java | 349 ++++------
.../state/table/SavepointFallbackSchemaLoader.java | 115 ++++
.../state/table/SavepointFilterTranslator.java | 53 +-
...tionFactory.java => SavepointStateMapping.java} | 20 +-
.../state/table/SavepointTypeInfoResolver.java | 528 +++++++---------
.../state/table/SingleColumnStateMapping.java | 47 ++
.../flink/state/table/StateTableMapping.java | 197 ++++++
.../state/table/StateValueColumnConfiguration.java | 26 +
.../flink/state/table/StateValueConverter.java | 281 +++++++++
.../flink/state/table/TableMappingSupport.java | 349 ++++++++++
.../org.apache.flink.table.factories.Factory | 1 +
.../EmbeddedRocksDBKeyedStateReadingITCase.java} | 18 +-
.../state/api/HashMapKeyedStateReadingITCase.java} | 18 +-
.../flink/state/api/KeyedStateReadingITCase.java | 404 ++++++++++++
.../flink/state/api/StateTableUtilsTest.java | 129 ++++
.../state/api/input/MultiStateKeyIteratorTest.java | 13 +-
.../deserializer/InternalTypeConverterTest.java | 301 +++++++++
.../PojoToRowDataDeserializerTest.java | 239 +++++++
...rializerSnapshotToLogicalTypeConverterTest.java | 425 +++++++++++++
.../state/api/schema/StateSchemaExtractorTest.java | 204 ++++++
.../flink/state/catalog/SnapshotDiscoveryTest.java | 248 ++++++++
.../state/catalog/StateCatalogDiscoveryITCase.java | 210 ++++++
.../StateCatalogGeneratedSavepointITCase.java | 691 ++++++++++++++++++++
.../flink/state/catalog/StateCatalogTest.java | 143 +++++
.../apache/flink/state/catalog/TuplePojoField.java | 59 ++
.../table/SavepointDynamicTableSourceTest.java | 91 ++-
.../state/table/SavepointTypeInfoResolverTest.java | 104 +++
.../table/SavepointTypeInformationFactoryTest.java | 129 ----
.../state/table/TypeConversionDriftGuardTest.java | 273 ++++++++
.../KeyedStateCatalogSavepointGenerator.java | 356 +++++++++++
.../KeyedStatePojoAvroKeySavepointGenerator.java | 321 ++++++++++
.../KeyedStateTupleKeySavepointGenerator.java | 277 ++++++++
.../test/resources/generator/StateTestRecord.avsc | 32 +
.../savepoint-01d134-a82d2259b86b/_metadata | Bin 0 -> 17739 bytes
.../savepoint-07a18b-da0f2ab0e5e0/_metadata | Bin 0 -> 12750 bytes
.../savepoint-515de8-42f928682f3b/_metadata | Bin 0 -> 11181 bytes
.../resources/table-state-missing-avro/_metadata | Bin 0 -> 4391 bytes
.../resources/table-state-missing-class/_metadata | Bin 0 -> 4453 bytes
.../flink/runtime/state/KeyedStateBackend.java | 16 +
.../runtime/state/heap/HeapKeyedStateBackend.java | 48 ++
.../flink/runtime/state/heap/StateTable.java | 19 +
.../operators/InternalTimeServiceManagerImpl.java | 2 +-
.../state/BatchExecutionKeyedStateBackend.java | 9 +
.../operators/windowing/WindowOperator.java | 5 +-
.../flink/runtime/state/StateBackendTestBase.java | 44 ++
.../flink/runtime/state/StateBackendTestUtils.java | 6 +
.../state/ttl/mock/MockKeyedStateBackend.java | 12 +
.../streaming/runtime/tasks/TestStateBackend.java | 6 +
.../changelog/ChangelogKeyedStateBackend.java | 5 +
.../restore/ChangelogMigrationRestoreTarget.java | 6 +
.../forst/sync/ForStMultiStateKeysIterator.java} | 45 +-
.../forst/sync/ForStSyncKeyedStateBackend.java | 96 ++-
.../state/rocksdb/RocksDBKeyedStateBackend.java | 60 +-
.../iterator/RocksMultiStateKeysIterator.java | 12 +
.../table/runtime/typeutils/ExternalTypeInfo.java | 10 +
.../table/runtime/typeutils/RowDataSerializer.java | 17 +
103 files changed, 11995 insertions(+), 1487 deletions(-)
create mode 100644
flink-core/src/main/java/org/apache/flink/api/common/typeutils/CustomRestoreSerializerFactory.java
create mode 100644
flink-core/src/test/java/org/apache/flink/api/common/typeutils/base/EnumSerializerSnapshotMissingClassTest.java
create mode 100644
flink-core/src/test/java/org/apache/flink/api/java/typeutils/runtime/PojoSerializerSnapshotLenientReadTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/StateTableUtils.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/input/deserializer/EnumNameDeserializer.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/input/deserializer/InternalTypeConverter.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/input/deserializer/MissingClassSerializerFactory.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/input/deserializer/PojoDeserializerCompatibilitySnapshot.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/input/deserializer/PojoToRowDataDeserializer.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/schema/AvroStateUtils.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/schema/KeyedStateSchemaInfo.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/schema/SerializerSnapshotToLogicalTypeConverter.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/schema/StateSchemaExtractor.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/schema/StateSchemaInfo.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/catalog/SnapshotDiscovery.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/catalog/StateCatalog.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/catalog/StateCatalogFactory.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/catalog/StateCatalogOptions.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/AbstractMultiColumnScanProvider.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/AbstractSavepointDataStreamScanProvider.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/AbstractSavepointDynamicTableSource.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/AbstractSingleColumnScanProvider.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/FlattenedKeyedStateReader.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/FlattenedSavepointDataStreamScanProvider.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/FlattenedSavepointDynamicTableSource.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/FlattenedStateTableMapping.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/MultiColumnStateMapping.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/SavepointFallbackSchemaLoader.java
copy
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/{SavepointTypeInformationFactory.java
=> SavepointStateMapping.java} (64%)
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/SingleColumnStateMapping.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/StateTableMapping.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/StateValueConverter.java
create mode 100644
flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/table/TableMappingSupport.java
copy
flink-libraries/flink-state-processing-api/src/{main/java/org/apache/flink/state/table/SavepointTypeInformationFactory.java
=>
test/java/org/apache/flink/state/api/EmbeddedRocksDBKeyedStateReadingITCase.java}
(62%)
rename
flink-libraries/flink-state-processing-api/src/{main/java/org/apache/flink/state/table/SavepointTypeInformationFactory.java
=> test/java/org/apache/flink/state/api/HashMapKeyedStateReadingITCase.java}
(63%)
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/api/KeyedStateReadingITCase.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/api/StateTableUtilsTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/api/input/deserializer/InternalTypeConverterTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/api/input/deserializer/PojoToRowDataDeserializerTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/api/schema/SerializerSnapshotToLogicalTypeConverterTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/api/schema/StateSchemaExtractorTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/catalog/SnapshotDiscoveryTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/catalog/StateCatalogDiscoveryITCase.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/catalog/StateCatalogGeneratedSavepointITCase.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/catalog/StateCatalogTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/catalog/TuplePojoField.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/table/SavepointTypeInfoResolverTest.java
delete mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/table/SavepointTypeInformationFactoryTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/table/TypeConversionDriftGuardTest.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/generator/KeyedStateCatalogSavepointGenerator.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/generator/KeyedStatePojoAvroKeySavepointGenerator.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/generator/KeyedStateTupleKeySavepointGenerator.java
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/generator/StateTestRecord.avsc
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/keyed-state-catalog/savepoint-01d134-a82d2259b86b/_metadata
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/keyed-state-pojo-avro-key/savepoint-07a18b-da0f2ab0e5e0/_metadata
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/keyed-state-tuple-key/savepoint-515de8-42f928682f3b/_metadata
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/table-state-missing-avro/_metadata
create mode 100644
flink-libraries/flink-state-processing-api/src/test/resources/table-state-missing-class/_metadata
copy
flink-state-backends/{flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/iterator/RocksMultiStateKeysIterator.java
=>
flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/sync/ForStMultiStateKeysIterator.java}
(82%)