This is an automated email from the ASF dual-hosted git repository.
zhangyue19921010 pushed a change to branch stream-binary-copy
in repository https://gitbox.apache.org/repos/asf/hudi.git
from 7cfd82b0779 code review
add 1fa3f2ed6a0 [MINOR] Fix flakiness in
TestHoodieFileSystemViews.testFileSystemViewConsistency (#13375)
add 328e259d4a7 [HUDI-9417] Add validation for handling write stats during
commit (#13307)
add 56e66eb7a1f [HUDI-9501] Fix problem of schema discrepancy in partial
update flink record merger (#13398)
add 79863dbf5b8 [HUDI-9235] Refactor iterator generation logic (#13388)
add adfb952bd6a [HUDI-9235] Refactor InstantRange handling for MDT table
(#13400)
add 01dc67a7e83 [HUDI-9491] Add checkpoint id into hudi commit metedata
(#13403)
add 7c6e2e2f20f [HUDI-9481] Remove extra io.netty dependencies from
hudi-aws bundle (#13404)
add bb3f0796515 [HUDI-9511] Use local engine context for timeline server
markers (#13399)
add f9555910e5d [MINOR] Fix SparkContext conflict in HoodieSparkQuickstart
(#13406)
add a5c531f64fe [HUDI-9513] Update test profiles to just use tags, fix
current test issues (#13407)
add 5146dc81f49 [HUDI-9518] Upgrade Flink 1.15/1.16 parquet version from
1.12.2 to 1.12.3 (#13418)
add 2d99668881f [HUDI-9499] Add customized SizeEstimator and Serializer
for Avro buffered record in FileGroup reader (#13408)
add f7157954a98 [MINOR] Fix minor issues in tests, lower wait threshold
for remote markers in tests (#13419)
add 1876998d22c [HUDI-8401] Remove payload class FirstValueAvroPayload
(#13420)
add 18c098fbdb5 [HUDI-9421][Part 2] Fix more cases around rollback
scheduling, update TTL action scheduling (#13380)
add 0e527404547 [HUDI-9405] Adding new writer apis and impl to Metadata
writer to support streaming writes (#13286)
add b1af4d4c546 [HUDI-9514] Scanner resources not properly closed in
HoodieCompact (#13415)
add c74b27faf88 [MINOR] Rebalance CI on Jun 12 (#13426)
add 58191bdf5ff [HUDI-9524] Fix include in DFSPropertiesConfiguration
(#13433)
add 03757d34ec7 [MINOR] Removed unused imports in HoodieSparkQuickstart
(#13440)
add d4edee24f6e [HUDI-9334] Optimize Parallelism of show_invalid_parquet
(#13206)
add 8c7e93f0700 [HUDI-9405] Adding end to end streaming writes to metadata
table support for SPARK engine (#13402)
add cb566b9785f [HUDI-9520] Support new sink based on Flink sink V2 API
(#13423)
add 4912b8a472d [MINOR] Flakey test in TestRecordLevelIndexTableVersionSix
(#13446)
add ac9f46c1ebb [HUDI-9532] Add sort columns in plan generated by Flink
(#13455)
add f8219428f35 [HUDI-8286] Standardize compaction and log compaction to
use FileGroupReader (#13411)
add 74e58bb0d7c [HUDI-8881] Fails fast in Flink append write function
(#12666)
add 6e7a7207902 [MINOR] Removed unused imports in HoodieWriteClientExample
(#13467)
add 07d176e2d46 [HUDI-9531] Avoid generating timestamp again in
createCompleteFileInMetaPath (#13453)
add fcc03c35bc5 [HUDI-9534] Blocking instant generation for Flink COW
table (#13464)
add e017d85d76b [HUDI-9527] Update HoodieTableMetadataUtil to use
HoodieMergedLogRecordReader and FileGroupReader (#13470)
add 245ef0c2085 [MINOR] Add JVM flags for tests to avoid flakiness (#13474)
add c23b5641473 [HUDI-8286] Add special handling for ordering fields when
building FileGroupReader (#13469)
add 1f561d08d7c [HUDI-9535] Prevent Unnecessary data transfer b/w driver
and executor in Spark Partitioner (#13459)
add 45312d437a5 [HUDI-9340] Add MDT streaming write support for secondary
index (#13449)
add 35529720d5d Merge branch 'master' into stream-binary-copy
add c35fd3e4253 code review
No new revisions were added by this update.
Summary of changes:
.github/workflows/bot.yml | 171 ++++++++++--
azure-pipelines-20230430.yml | 131 ++++-----
hudi-aws/pom.xml | 10 +
.../hudi/cli/commands/ClusteringCommand.java | 10 +-
.../hudi/cli/commands/CompactionCommand.java | 10 +-
.../org/apache/hudi/cli/commands/SparkMain.java | 8 +-
hudi-client/hudi-client-common/pom.xml | 6 +
.../hudi/callback/common/WriteStatusValidator.java | 44 +++
.../org/apache/hudi/client/BaseHoodieClient.java | 19 ++
.../hudi/client/BaseHoodieTableServiceClient.java | 95 +++++--
.../apache/hudi/client/BaseHoodieWriteClient.java | 88 ++++--
.../java/org/apache/hudi/client/IndexStats.java | 67 +++++
...BaseCompactor.java => SecondaryIndexStats.java} | 40 ++-
.../org/apache/hudi/client/TableWriteStats.java | 54 ++++
.../java/org/apache/hudi/client/WriteStatus.java | 23 +-
.../client/embedded/EmbeddedTimelineService.java | 8 +-
.../client/transaction/TransactionManager.java | 3 +-
.../apache/hudi/index/bucket/BucketIdentifier.java | 14 +-
.../index/inmemory/HoodieInMemoryHashIndex.java | 2 +-
.../hudi/io/FileGroupReaderBasedAppendHandle.java | 118 ++++++++
...e.java => FileGroupReaderBasedMergeHandle.java} | 130 +++++----
.../org/apache/hudi/io/HoodieAppendHandle.java | 28 +-
.../org/apache/hudi/io/HoodieBinaryCopyHandle.java | 2 +-
.../org/apache/hudi/io/HoodieConcatHandle.java | 10 -
.../org/apache/hudi/io/HoodieCreateHandle.java | 14 +-
.../java/org/apache/hudi/io/HoodieMergeHandle.java | 106 +++++---
.../java/org/apache/hudi/io/HoodieWriteHandle.java | 37 ++-
.../hudi/io/SecondaryIndexStreamingTracker.java | 201 ++++++++++++++
.../metadata/HoodieBackedTableMetadataWriter.java | 265 +++++++++++++++---
...ieBackedTableMetadataWriterTableVersionSix.java | 5 +
.../hudi/metadata/HoodieTableMetadataWriter.java | 45 ++++
.../hudi/metadata/MetadataIndexGenerator.java | 132 +++++++++
.../apache/hudi/table/HoodieCompactionHandler.java | 10 -
.../java/org/apache/hudi/table/HoodieTable.java | 90 ++++---
.../hudi/table/action/compact/HoodieCompactor.java | 184 +++++++------
.../restore/CopyOnWriteRestoreActionExecutor.java | 16 +-
.../restore/MergeOnReadRestoreActionExecutor.java | 18 +-
.../rollback/BaseRollbackPlanActionExecutor.java | 1 -
.../hudi/table/marker/DirectWriteMarkers.java | 20 +-
.../org/apache/hudi/ClientFunctionalTestSuite.java | 32 ---
.../client/TestBaseHoodieTableServiceClient.java | 4 +-
.../hudi/client/TestBaseHoodieWriteClient.java | 3 +-
.../embedded/TestEmbeddedTimelineService.java | 14 +-
...TestPreferWriterConflictResolutionStrategy.java | 9 -
.../apache/hudi/config/TestHoodieWriteConfig.java | 1 +
.../bucket/TestPartitionBucketIndexCalculator.java | 2 +-
.../metadata/TestHoodieMetadataWriteUtils.java | 3 +-
.../hudi/metadata/TestMetadataIndexGenerator.java | 121 +++++++++
.../org/apache/hudi/table/TestHoodieTable.java | 11 +-
.../hudi/utils/HoodieWriterClientTestHarness.java | 35 +--
.../hudi/client/HoodieFlinkTableServiceClient.java | 6 +-
.../apache/hudi/client/HoodieFlinkWriteClient.java | 4 +-
.../client/common/HoodieFlinkEngineContext.java | 13 +
.../model/PartialUpdateFlinkRecordMerger.java | 63 +++--
.../hudi/execution/FlinkLazyInsertIterable.java | 3 +-
.../apache/hudi/io/FlinkMergeAndReplaceHandle.java | 8 -
.../java/org/apache/hudi/io/FlinkMergeHandle.java | 8 -
.../v2/FlinkFileGroupReaderBasedMergeHandle.java | 141 ----------
.../FlinkHoodieBackedTableMetadataWriter.java | 21 +-
.../hudi/table/HoodieFlinkCopyOnWriteTable.java | 15 --
.../hudi/table/HoodieFlinkMergeOnReadTable.java | 4 +-
.../org/apache/hudi/table/HoodieFlinkTable.java | 3 +-
.../HoodieFlinkMergeOnReadTableCompactor.java | 25 +-
.../merge/TestPartialUpdateFlinkRecordMerger.java | 45 +++-
.../hudi/client/HoodieJavaTableServiceClient.java | 5 +-
.../apache/hudi/client/HoodieJavaWriteClient.java | 7 +-
.../run/strategy/JavaExecutionStrategy.java | 2 +-
.../hudi/execution/JavaLazyInsertIterable.java | 3 +-
.../JavaHoodieBackedTableMetadataWriter.java | 19 +-
.../hudi/table/HoodieJavaMergeOnReadTable.java | 2 +-
.../org/apache/hudi/table/HoodieJavaTable.java | 4 +-
.../HoodieJavaMergeOnReadTableCompactor.java | 10 +-
.../hudi/client/TestJavaHoodieBackedMetadata.java | 2 +-
hudi-client/hudi-spark-client/pom.xml | 4 +
.../hudi/client/SparkRDDMetadataWriteClient.java | 47 +++-
.../org/apache/hudi/client/SparkRDDReadClient.java | 6 +-
.../hudi/client/SparkRDDTableServiceClient.java | 51 +++-
.../apache/hudi/client/SparkRDDWriteClient.java | 147 +++++++++-
.../hudi/client/StreamingMetadataWriteHandler.java | 126 +++++++++
.../client/common/HoodieSparkEngineContext.java | 12 +-
.../hudi/common/model/HoodieSparkRecord.java | 11 +-
.../hudi/execution/SparkLazyInsertIterable.java | 2 +
...HoodieSparkFileGroupReaderBasedMergeHandle.java | 137 ----------
.../SparkHoodieBackedTableMetadataWriter.java | 59 +++-
...ieBackedTableMetadataWriterTableVersionSix.java | 33 ++-
.../hudi/metadata/SparkMetadataWriterFactory.java | 5 +
.../hudi/table/HoodieSparkCopyOnWriteTable.java | 16 --
.../hudi/table/HoodieSparkMergeOnReadTable.java | 2 +-
.../org/apache/hudi/table/HoodieSparkTable.java | 9 +-
.../BaseSparkBucketIndexBucketInfoGetter.java | 68 +++++
.../commit/BaseSparkCommitActionExecutor.java | 33 ++-
.../commit/InsertOverwriteBucketInfoGetter.java | 31 ++-
.../commit/ListBasedSparkBucketInfoGetter.java | 17 +-
.../commit/MapBasedSparkBucketInfoGetter.java | 21 +-
.../commit/SparkBucketIndexBucketInfoGetter.java | 47 ++++
.../action/commit/SparkBucketIndexPartitioner.java | 28 +-
.../table/action/commit/SparkBucketInfoGetter.java | 21 +-
.../SparkDeletePartitionCommitActionExecutor.java | 9 +-
.../action/commit/SparkHoodiePartitioner.java | 2 +-
.../SparkInsertOverwriteCommitActionExecutor.java | 8 +-
.../commit/SparkInsertOverwritePartitioner.java | 19 +-
.../SparkMetadataTableUpsertPartitioner.java | 4 +-
.../SparkPartitionBucketIndexBucketInfoGetter.java | 47 ++++
.../SparkPartitionBucketIndexPartitioner.java | 39 +--
.../table/action/commit/UpsertPartitioner.java | 7 +-
.../HoodieSparkMergeOnReadTableCompactor.java | 5 +
.../hudi/BaseSparkInternalRowReaderContext.java | 2 +-
.../apache/hudi/client/TestHoodieReadClient.java | 4 +-
.../java/org/apache/hudi/client/TestMultiFS.java | 8 +-
.../hudi/client/TestPartitionTTLManagement.java | 13 +-
.../client/TestSparkRDDMetadataWriteClient.java | 2 +-
.../org/apache/hudi/client/TestWriteStatus.java | 22 +-
.../functional/SparkClientFunctionalTestSuite.java | 35 ---
.../java/org/apache/hudi/io/BaseTestHandle.java | 113 ++++++++
.../io/KeyGeneratorForDataGeneratorRecords.java | 22 +-
.../org/apache/hudi/io/TestHoodieMergeHandle.java | 2 +-
.../commit/TestCopyOnWriteActionExecutor.java | 1 -
.../TestSparkMetadataTableUpsertPartitioner.java | 3 +-
.../table/action/commit/TestUpsertPartitioner.java | 59 +++-
.../TestTimelineServerBasedWriteMarkers.java | 5 +-
.../hudi/testutils/HoodieClientTestBase.java | 2 +-
.../hudi/testutils/HoodieClientTestUtils.java | 4 +-
.../SparkClientFunctionalTestHarness.java | 10 +-
.../hudi/testutils/providers/SparkProvider.java | 3 +-
.../AvroRecordSerializer.java} | 37 +--
.../apache/hudi/avro/AvroRecordSizeEstimator.java | 50 ++++
.../java/org/apache/hudi/avro/AvroSchemaUtils.java | 6 +-
.../apache/hudi/avro/HoodieAvroReaderContext.java | 59 ++--
.../hudi/common/engine/HoodieEngineContext.java | 8 +-
.../hudi/common/engine/HoodieReaderContext.java | 53 +++-
.../java/org/apache/hudi/common/fs/FSUtils.java | 16 +-
.../hudi/common/model/FirstValueAvroPayload.java | 125 ---------
.../hudi/common/model/HoodieAvroIndexedRecord.java | 9 +-
.../apache/hudi/common/model/HoodieAvroRecord.java | 28 +-
.../hudi/common/model/HoodieEmptyRecord.java | 2 +-
.../hudi/common/model/HoodieIndexDefinition.java | 16 +-
.../org/apache/hudi/common/model/HoodieRecord.java | 6 +-
.../apache/hudi/common/model/HoodieWriteStat.java | 1 +
.../hudi/common/model/WriteOperationType.java | 13 +
...efaultSerializer.java => RecordSerializer.java} | 34 ++-
.../hudi/common/table/HoodieTableMetaClient.java | 2 +-
.../table/log/AbstractHoodieLogRecordScanner.java | 6 +-
.../hudi/common/table/read/BufferedRecord.java | 21 +-
.../table/read/BufferedRecordSerializer.java | 82 ++++++
.../common/table/read/FileGroupRecordBuffer.java | 19 +-
.../common/table/read/HoodieFileGroupReader.java | 51 ++--
.../read/SortedKeyBasedFileGroupRecordBuffer.java | 106 ++++++++
.../table/timeline/HoodieInstantTimeGenerator.java | 22 +-
.../timeline/versioning/v2/ActiveTimelineV2.java | 4 +-
.../table/view/AbstractTableFileSystemView.java | 12 +
.../table/view/PriorityBasedFileSystemView.java | 6 +
.../view/RemoteHoodieTableFileSystemView.java | 8 +
.../common/table/view/TableFileSystemView.java | 10 +
.../org/apache/hudi/common/util/ConfigUtils.java | 6 +-
.../org/apache/hudi/common/util/MarkerUtils.java | 17 +-
.../hudi/common/util/SerializationUtils.java | 2 +-
.../hudi/io/storage/HoodieAvroFileReader.java | 28 ++
.../io/storage/HoodieAvroHFileReaderImplBase.java | 7 -
.../io/storage/HoodieNativeAvroHFileReader.java | 51 +++-
.../hudi/metadata/HoodieBackedTableMetadata.java | 3 +-
.../hudi/metadata/HoodieMetadataPayload.java | 49 +++-
.../hudi/metadata/HoodieTableMetadataUtil.java | 296 ++++++++++++---------
.../SecondaryIndexRecordGenerationUtils.java | 7 +-
.../hudi/timeline/TimelineServiceClient.java | 2 +-
.../hudi/timeline/TimelineServiceClientBase.java | 22 +-
.../hudi/avro/TestAvroRecordSizeEstimator.java | 58 ++++
.../common/model/TestFirstValueAvroPayload.java | 81 ------
.../common/model/TestHoodieCommitMetadata.java | 4 +-
.../hudi/common/model/TestHoodieIndexMetadata.java | 49 ++++
.../hudi/common/model/TestHoodieWriteStat.java | 2 +-
.../TestBufferedRecordSerializer.java | 90 +++++++
.../table/read/TestFileGroupRecordBuffer.java | 13 +-
.../table/read/TestHoodieFileGroupReaderBase.java | 35 ++-
.../TestSortedKeyBasedFileGroupRecordBuffer.java | 163 ++++++++++++
.../timeline/TestHoodieInstantTimeGenerator.java | 24 +-
.../view/TestPriorityBasedFileSystemView.java | 37 +++
.../storage/TestHoodieNativeAvroHFileReader.java | 134 ++++++++++
.../examples/quickstart/HoodieSparkQuickstart.java | 4 +-
.../examples/spark/HoodieWriteClientExample.java | 1 -
hudi-flink-datasource/hudi-flink/pom.xml | 5 +
.../apache/hudi/configuration/FlinkOptions.java | 17 +-
.../apache/hudi/configuration/OptionsResolver.java | 7 +
.../org/apache/hudi/sink/StreamWriteFunction.java | 5 +-
.../hudi/sink/StreamWriteOperatorCoordinator.java | 41 ++-
.../hudi/sink/append/AppendWriteFunction.java | 6 +-
.../append/AppendWriteFunctionWithRateLimit.java | 3 +-
.../sink/bucket/BucketStreamWriteFunction.java | 5 +-
.../hudi/sink/bulk/BulkInsertWriteFunction.java | 2 +-
.../sink/clustering/ClusteringPlanOperator.java | 5 +-
.../sink/clustering/HoodieFlinkClusteringJob.java | 3 +-
.../hudi/sink/common/AbstractWriteFunction.java | 3 +-
.../hudi/sink/common/AbstractWriteOperator.java | 3 +-
.../hudi/sink/common/WriteOperatorFactory.java | 7 +-
.../apache/hudi/sink/compact/CompactOperator.java | 4 +-
.../hudi/sink/compact/CompactionCommitSink.java | 4 +-
.../hudi/sink/compact/CompactionPlanOperator.java | 5 +-
.../hudi/sink/compact/HoodieFlinkCompactor.java | 2 +-
.../org/apache/hudi/sink/event/Correspondent.java | 4 +-
.../org/apache/hudi/sink/utils/CommitGuard.java | 80 ++++++
...nseSeDe.java => CoordinationResponseSerDe.java} | 2 +-
.../org/apache/hudi/sink/utils/EventBuffers.java | 21 +-
.../java/org/apache/hudi/sink/utils/Pipelines.java | 40 +--
.../CleanFunctionV2.java} | 27 +-
.../java/org/apache/hudi/sink/v2/HoodieSink.java | 78 ++++++
.../clustering/ClusteringCommitSinkV2.java} | 28 +-
.../compact/CompactionCommitSinkV2.java} | 30 ++-
.../org/apache/hudi/sink/v2/utils/PipelinesV2.java | 277 +++++++++++++++++++
.../apache/hudi/streamer/HoodieFlinkStreamer.java | 2 +-
.../org/apache/hudi/table/HoodieTableSink.java | 7 +-
.../apache/hudi/table/catalog/HoodieCatalog.java | 4 +-
.../hudi/table/catalog/HoodieHiveCatalog.java | 5 +-
.../table/format/FlinkReaderContextFactory.java | 25 +-
.../java/org/apache/hudi/util/ClusteringUtil.java | 6 +-
.../java/org/apache/hudi/util/CompactionUtil.java | 15 +-
.../org/apache/hudi/util/FlinkWriteClients.java | 9 +-
.../java/org/apache/hudi/util/StreamerUtil.java | 46 ++++
.../apache/hudi/sink/ITTestDataStreamV2Write.java | 159 +++++++++++
.../apache/hudi/sink/ITTestDataStreamWrite.java | 33 ++-
.../sink/TestStreamWriteOperatorCoordinator.java | 12 +-
.../org/apache/hudi/sink/TestWriteCopyOnWrite.java | 97 ++++++-
.../hudi/sink/append/TestAppendWriteFunction.java | 74 ++++++
.../bucket/ITTestConsistentBucketStreamWrite.java | 2 +-
.../apache/hudi/sink/utils/MockCorrespondent.java | 3 +-
...dent.java => MockCorrespondentWithTimeout.java} | 18 +-
.../org/apache/hudi/sink/utils/TestWriteBase.java | 10 +
.../table/TestHoodieFileGroupReaderOnFlink.java | 4 +-
.../apache/hudi/table/format/TestInputFormat.java | 28 ++
.../org/apache/hudi/utils/TestStreamerUtil.java | 13 +
.../connector/sink2/SupportsPreWriteTopology.java} | 20 +-
.../connector/sink2/SupportsPreWriteTopology.java} | 20 +-
.../connector/sink2/SupportsPreWriteTopology.java} | 20 +-
.../connector/sink2/SupportsPreWriteTopology.java} | 20 +-
.../common/config/DFSPropertiesConfiguration.java | 2 +-
.../apache/hudi/io/hadoop/HoodieAvroOrcReader.java | 13 +
.../hudi/io/hadoop/HoodieAvroParquetReader.java | 12 +
.../lock/TestInProcessLockProvider.java | 3 +-
.../org/apache/hudi/common/fs/TestFSUtils.java | 26 --
.../hudi/common/table/read/TestCustomMerger.java | 2 +
.../common/table/read/TestEventTimeMerging.java | 2 +
.../table/view/TestHoodieTableFileSystemView.java | 12 +
.../reader/HoodieFileGroupReaderTestHarness.java | 12 +-
.../util/TestDFSPropertiesConfiguration.java | 22 ++
.../apache/hudi/common/util/TestMarkerUtils.java | 6 +-
.../hudi/metadata/TestHoodieTableMetadataUtil.java | 15 +-
.../org/apache/hudi/io/hfile/TestHFileWriter.java | 1 -
.../main/java/org/apache/hudi/DataSourceUtils.java | 50 ++++
.../internal/DataSourceInternalWriterHelper.java | 7 +-
.../org/apache/hudi/HoodieSparkSqlWriter.scala | 45 ++--
.../RunRollbackInflightTableServiceProcedure.scala | 4 +-
.../hudi/command/procedures/RunTTLProcedure.scala | 2 +-
.../ShowColumnStatsOverlapProcedure.scala | 15 +-
.../procedures/ShowInvalidParquetProcedure.scala | 62 +++--
.../apache/hudi/TestDataSourceReadWithDeletes.java | 9 +-
.../hudi/client/TestHoodieClientMultiWriter.java | 15 +-
.../TestHoodieClientOnCopyOnWriteStorage.java | 8 +-
.../TestHoodieClientOnMergeOnReadStorage.java | 5 +-
.../TestMetadataUtilRLIandSIRecordGeneration.java | 14 +-
.../TestRemoteFileSystemViewWithMetadataTable.java | 2 +-
.../HoodieSparkFunctionalTestSuiteA.java | 29 --
.../HoodieSparkFunctionalTestSuiteB.java | 29 --
.../hudi/functional/TestBootstrapReadBase.java | 3 +-
.../TestColStatsRecordWithMetadataRecord.java | 3 +-
.../hudi/functional/TestHoodieBackedMetadata.java | 4 +-
.../functional/TestHoodieFileSystemViews.java | 68 +++--
.../apache/hudi/functional/TestHoodieIndex.java | 160 ++++++++++-
.../TestSparkConsistentBucketClustering.java | 3 +-
.../apache/hudi/functional/TestWriteClient.java | 5 +-
.../java/org/apache/hudi/io/TestAppendHandle.java | 135 ++++++++++
.../java/org/apache/hudi/io/TestCreateHandle.java | 114 ++++++++
.../java/org/apache/hudi/io/TestMergeHandle.java | 131 +++++++++
.../apache/hudi/io/TestMetadataWriterCommit.java | 237 +++++++++++++++++
.../io/storage/row/TestHoodieRowCreateHandle.java | 7 +-
.../hudi/table/TestHoodieMergeOnReadTable.java | 2 +
.../table/action/compact/TestAsyncCompaction.java | 2 +-
.../TestHoodieSparkMergeOnReadTableCompaction.java | 3 +
...dieSparkMergeOnReadTableInsertUpdateDelete.java | 3 +-
.../TestHoodieSparkMergeOnReadTableRollback.java | 18 +-
.../table/functional/TestHoodieSparkRollback.java | 4 +
.../TestSparkNonBlockingConcurrencyControl.java | 16 ++
.../src/test/resources/exampleSchema.txt | 6 +-
.../TestIncrementalQueryWithArchivedInstants.scala | 16 +-
.../hudi/functional/ColumnStatIndexTestBase.scala | 4 +-
.../hudi/functional/HoodieStatsIndexTestBase.scala | 11 +-
.../hudi/functional/RecordLevelIndexTestBase.scala | 8 +-
.../hudi/functional/TestBucketIndexSupport.scala | 2 +-
.../apache/hudi/functional/TestCOWDataSource.scala | 23 +-
.../functional/TestColumnStatsIndexWithSQL.scala | 3 +-
.../hudi/functional/TestLayoutOptimization.scala | 2 +-
.../hudi/functional/TestMORDataSourceStorage.scala | 3 +-
.../hudi/functional/TestMetadataRecordIndex.scala | 5 -
.../hudi/functional/TestPartitionStatsIndex.scala | 14 +-
.../functional/TestPartitionStatsPruning.scala | 2 +
.../hudi/functional/TestRecordLevelIndex.scala | 18 +-
.../TestRecordLevelIndexTableVersionSix.scala | 3 +
.../functional/TestSecondaryIndexPruning.scala | 3 +-
.../spark/sql/avro/TestSchemaConverters.scala | 2 +-
.../hudi/dml/{ => insert}/TestInsertTable.scala | 86 ++----
.../TestInsertTableWithPartitionBucketIndex.scala | 28 +-
.../dml/{ => others}/TestDeleteFromTable.scala | 28 +-
.../hudi/dml/{ => others}/TestDeleteTable.scala | 28 +-
.../TestHoodieTableValuedFunction.scala | 28 +-
.../{ => others}/TestMergeIntoLogOnlyTable.scala | 28 +-
.../hudi/dml/{ => others}/TestMergeIntoTable.scala | 28 +-
.../dml/{ => others}/TestMergeIntoTable2.scala | 108 ++------
.../TestMergeIntoTableWithNonRecordKeyField.scala | 28 +-
.../TestMergeModeCommitTimeOrdering.scala | 28 +-
.../TestMergeModeEventTimeOrdering.scala | 28 +-
.../TestPartialUpdateForMergeInto.scala | 28 +-
.../dml/{ => others}/TestTimeTravelTable.scala | 28 +-
.../hudi/dml/{ => others}/TestUpdateTable.scala | 28 +-
.../hudi/feature/index/TestExpressionIndex.scala | 3 +-
.../sql/hudi/feature/index/TestGlobalIndex.scala | 2 +-
.../TestShowInvalidParquetProcedure.scala | 24 +-
.../sql/hudi/procedure/TestTTLProcedure.scala | 1 +
.../functional/HiveSyncFunctionalTestSuite.java | 33 ---
hudi-tests-common/pom.xml | 24 +-
.../hudi/timeline/service/RequestHandler.java | 30 ++-
.../hudi/timeline/service/TimelineService.java | 12 +-
.../service/handlers/FileSliceHandler.java | 6 +
.../timeline/service/handlers/MarkerHandler.java | 17 +-
.../service/handlers/marker/MarkerDirState.java | 26 +-
.../hudi/timeline/service/TestRequestHandler.java | 2 +-
.../hudi/timeline/service/TestTimelineService.java | 12 +-
.../service/TimelineServiceTestHarness.java | 11 +-
.../TestRemoteHoodieTableFileSystemView.java | 1 -
.../org/apache/hudi/utilities/HoodieTTLJob.java | 2 +-
.../hudi/utilities/perf/TimelineServerPerf.java | 2 -
.../hudi/utilities/streamer/HoodieStreamer.java | 7 +-
.../apache/hudi/utilities/streamer/StreamSync.java | 213 ++++++++++-----
.../deltastreamer/TestHoodieDeltaStreamer.java | 3 +-
.../functional/UtilitiesFunctionalTestSuite.java | 32 ---
.../TestHoodieMultiTableServicesMain.java | 23 +-
.../utilities/testutils/UtilitiesTestBase.java | 2 +-
pom.xml | 106 +++++---
334 files changed, 7375 insertions(+), 2845 deletions(-)
create mode 100644
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/callback/common/WriteStatusValidator.java
create mode 100644
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/IndexStats.java
copy
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/{BaseCompactor.java
=> SecondaryIndexStats.java} (50%)
create mode 100644
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/TableWriteStats.java
create mode 100644
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedAppendHandle.java
rename
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/{BaseFileGroupReaderBasedMergeHandle.java
=> FileGroupReaderBasedMergeHandle.java} (55%)
create mode 100644
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/SecondaryIndexStreamingTracker.java
create mode 100644
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/MetadataIndexGenerator.java
delete mode 100644
hudi-client/hudi-client-common/src/test/java/org/apache/hudi/ClientFunctionalTestSuite.java
create mode 100644
hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestMetadataIndexGenerator.java
delete mode 100644
hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/FlinkFileGroupReaderBasedMergeHandle.java
create mode 100644
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/StreamingMetadataWriteHandler.java
delete mode 100644
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/HoodieSparkFileGroupReaderBasedMergeHandle.java
create mode 100644
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkBucketIndexBucketInfoGetter.java
copy
hudi-hadoop-common/src/test/java/org/apache/hudi/storage/hadoop/TestHadoopStorageConfiguration.java
=>
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/InsertOverwriteBucketInfoGetter.java
(52%)
copy
hudi-sync/hudi-datahub-sync/src/test/java/org/apache/hudi/sync/datahub/DummyPartitionValueExtractor.java
=>
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/ListBasedSparkBucketInfoGetter.java
(68%)
copy hudi-common/src/main/java/org/apache/hudi/common/util/MapUtils.java =>
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/MapBasedSparkBucketInfoGetter.java
(66%)
create mode 100644
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/SparkBucketIndexBucketInfoGetter.java
copy
hudi-common/src/main/java/org/apache/hudi/common/data/HoodieAccumulator.java =>
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/SparkBucketInfoGetter.java
(62%)
create mode 100644
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/SparkPartitionBucketIndexBucketInfoGetter.java
delete mode 100644
hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/functional/SparkClientFunctionalTestSuite.java
create mode 100644
hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/BaseTestHandle.java
copy
hudi-utilities/src/main/java/org/apache/hudi/utilities/schema/SimpleSchemaProvider.java
=>
hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/KeyGeneratorForDataGeneratorRecords.java
(63%)
copy
hudi-common/src/main/java/org/apache/hudi/{common/model/EmptyHoodieRecordPayload.java
=> avro/AvroRecordSerializer.java} (51%)
create mode 100644
hudi-common/src/main/java/org/apache/hudi/avro/AvroRecordSizeEstimator.java
delete mode 100644
hudi-common/src/main/java/org/apache/hudi/common/model/FirstValueAvroPayload.java
copy
hudi-common/src/main/java/org/apache/hudi/common/serialization/{DefaultSerializer.java
=> RecordSerializer.java} (57%)
create mode 100644
hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecordSerializer.java
create mode 100644
hudi-common/src/main/java/org/apache/hudi/common/table/read/SortedKeyBasedFileGroupRecordBuffer.java
create mode 100644
hudi-common/src/test/java/org/apache/hudi/avro/TestAvroRecordSizeEstimator.java
delete mode 100644
hudi-common/src/test/java/org/apache/hudi/common/model/TestFirstValueAvroPayload.java
create mode 100644
hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieIndexMetadata.java
create mode 100644
hudi-common/src/test/java/org/apache/hudi/common/serialization/TestBufferedRecordSerializer.java
create mode 100644
hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java
copy
hudi-client/hudi-client-common/src/test/java/org/apache/hudi/optimize/TestHilbertCurveUtils.java
=>
hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieInstantTimeGenerator.java
(60%)
create mode 100644
hudi-common/src/test/java/org/apache/hudi/io/storage/TestHoodieNativeAvroHFileReader.java
create mode 100644
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CommitGuard.java
rename
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/{CoordinationResponseSeDe.java
=> CoordinationResponseSerDe.java} (99%)
copy
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/{CleanFunction.java
=> v2/CleanFunctionV2.java} (81%)
create mode 100644
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/v2/HoodieSink.java
copy
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/{clustering/ClusteringCommitSink.java
=> v2/clustering/ClusteringCommitSinkV2.java} (91%)
copy
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/{compact/CompactionCommitSink.java
=> v2/compact/CompactionCommitSinkV2.java} (86%)
create mode 100644
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/v2/utils/PipelinesV2.java
copy
hudi-common/src/main/java/org/apache/hudi/common/engine/AvroReaderContextFactory.java
=>
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkReaderContextFactory.java
(50%)
create mode 100644
hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/ITTestDataStreamV2Write.java
create mode 100644
hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/append/TestAppendWriteFunction.java
copy
hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/{MockCorrespondent.java
=> MockCorrespondentWithTimeout.java} (61%)
copy
hudi-flink-datasource/{hudi-flink/src/main/java/org/apache/hudi/sink/transform/Transformer.java
=>
hudi-flink1.15.x/src/main/java/org/apache/flink/streaming/api/connector/sink2/SupportsPreWriteTopology.java}
(56%)
copy
hudi-flink-datasource/{hudi-flink/src/main/java/org/apache/hudi/sink/transform/Transformer.java
=>
hudi-flink1.16.x/src/main/java/org/apache/flink/streaming/api/connector/sink2/SupportsPreWriteTopology.java}
(56%)
copy
hudi-flink-datasource/{hudi-flink/src/main/java/org/apache/hudi/sink/transform/Transformer.java
=>
hudi-flink1.17.x/src/main/java/org/apache/flink/streaming/api/connector/sink2/SupportsPreWriteTopology.java}
(56%)
copy
hudi-flink-datasource/{hudi-flink/src/main/java/org/apache/hudi/sink/transform/Transformer.java
=>
hudi-flink1.18.x/src/main/java/org/apache/flink/streaming/api/connector/sink2/SupportsPreWriteTopology.java}
(56%)
delete mode 100644
hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/HoodieSparkFunctionalTestSuiteA.java
delete mode 100644
hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/HoodieSparkFunctionalTestSuiteB.java
rename hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/{client
=> }/functional/TestHoodieFileSystemViews.java (86%)
create mode 100644
hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/io/TestAppendHandle.java
create mode 100644
hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/io/TestCreateHandle.java
create mode 100644
hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/io/TestMergeHandle.java
create mode 100644
hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/io/TestMetadataWriterCommit.java
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> insert}/TestInsertTable.scala (98%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> insert}/TestInsertTableWithPartitionBucketIndex.scala (97%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestDeleteFromTable.scala (78%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestDeleteTable.scala (93%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestHoodieTableValuedFunction.scala (96%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestMergeIntoLogOnlyTable.scala (76%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestMergeIntoTable.scala (98%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestMergeIntoTable2.scala (92%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestMergeIntoTableWithNonRecordKeyField.scala (93%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestMergeModeCommitTimeOrdering.scala (94%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestMergeModeEventTimeOrdering.scala (94%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestPartialUpdateForMergeInto.scala (97%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestTimeTravelTable.scala (92%)
rename
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/{
=> others}/TestUpdateTable.scala (95%)
delete mode 100644
hudi-sync/hudi-hive-sync/src/test/java/org/apache/hudi/hive/functional/HiveSyncFunctionalTestSuite.java
delete mode 100644
hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/UtilitiesFunctionalTestSuite.java