| /* |
| * Licensed to the Apache Software Foundation (ASF) under one |
| * or more contributor license agreements. See the NOTICE file |
| * distributed with this work for additional information |
| * regarding copyright ownership. The ASF licenses this file |
| * to you under the Apache License, Version 2.0 (the |
| * "License"); you may not use this file except in compliance |
| * with the License. You may obtain a copy of the License at |
| * |
| * http://www.apache.org/licenses/LICENSE-2.0 |
| * |
| * Unless required by applicable law or agreed to in writing, software |
| * distributed under the License is distributed on an "AS IS" BASIS, |
| * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
| * See the License for the specific language governing permissions and |
| * limitations under the License. |
| */ |
| |
| #include "paimon/core/core_options.h" |
| |
| #include <limits> |
| |
| #include "gtest/gtest.h" |
| #include "paimon/bucket/bucket_function_type.h" |
| #include "paimon/common/fs/resolving_file_system.h" |
| #include "paimon/core/options/expire_config.h" |
| #include "paimon/defs.h" |
| #include "paimon/fs/local/local_file_system.h" |
| #include "paimon/testing/mock/mock_file_system.h" |
| #include "paimon/testing/utils/testharness.h" |
| #include "paimon/testing/utils/timezone_guard.h" |
| namespace paimon::test { |
| |
| TEST(CoreOptionsTest, TestDefaultValue) { |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap({})); |
| ASSERT_EQ(core_options.GetManifestFormat()->Identifier(), "avro"); |
| ASSERT_EQ(core_options.GetFileFormat()->Identifier(), "parquet"); |
| ASSERT_EQ(core_options.GetWriteFileFormat(0)->Identifier(), "parquet"); |
| ASSERT_EQ(core_options.GetWriteFileFormat(3)->Identifier(), "parquet"); |
| ASSERT_TRUE(core_options.GetFileSystem()); |
| ASSERT_EQ(-1, core_options.GetBucket()); |
| ASSERT_EQ(64 * 1024L, core_options.GetPageSize()); |
| ASSERT_EQ(256 * 1024 * 1024L, core_options.GetTargetFileSize(/*has_primary_key=*/false)); |
| ASSERT_EQ(128 * 1024 * 1024L, core_options.GetTargetFileSize(/*has_primary_key=*/true)); |
| ASSERT_EQ(256 * 1024 * 1024L, core_options.GetBlobTargetFileSize()); |
| ASSERT_EQ(187904815, core_options.GetCompactionFileSize(/*has_primary_key=*/false)); |
| ASSERT_EQ(93952404, core_options.GetCompactionFileSize(/*has_primary_key=*/true)); |
| |
| ASSERT_EQ("__DEFAULT_PARTITION__", core_options.GetPartitionDefaultName()); |
| ASSERT_EQ(std::nullopt, core_options.GetScanSnapshotId()); |
| ASSERT_EQ("zstd", core_options.GetFileCompression()); |
| ASSERT_EQ("zstd", core_options.GetWriteFileCompression(0)); |
| ASSERT_EQ("zstd", core_options.GetWriteFileCompression(3)); |
| ASSERT_EQ("zstd", core_options.GetManifestCompression()); |
| ASSERT_EQ(1, core_options.GetFileCompressionZstdLevel()); |
| ASSERT_EQ(StartupMode::LatestFull(), core_options.GetStartupMode()); |
| ASSERT_EQ(8 * 1024 * 1024L, core_options.GetManifestTargetFileSize()); |
| ASSERT_EQ(16 * 1024 * 1024L, core_options.GetManifestFullCompactionThresholdSize()); |
| ASSERT_EQ(30, core_options.GetManifestMergeMinCount()); |
| ASSERT_EQ(128 * 1024 * 1024L, core_options.GetSourceSplitTargetSize()); |
| ASSERT_EQ(4 * 1024 * 1024L, core_options.GetSourceSplitOpenFileCost()); |
| ASSERT_EQ(1024, core_options.GetReadBatchSize()); |
| ASSERT_EQ(1024, core_options.GetWriteBatchSize()); |
| ASSERT_EQ(256 * 1024 * 1024, core_options.GetWriteBufferSize()); |
| ASSERT_TRUE(core_options.GetWriteBufferSpillable()); |
| ASSERT_EQ(std::numeric_limits<int64_t>::max(), core_options.GetWriteBufferSpillMaxDiskSize()); |
| ASSERT_EQ(128, core_options.GetLocalSortMaxNumFileHandles()); |
| ASSERT_EQ("zstd", core_options.GetSpillCompressOptions().compress); |
| ASSERT_EQ(1, core_options.GetSpillCompressOptions().zstd_level); |
| ASSERT_FALSE(core_options.CommitForceCompact()); |
| ASSERT_EQ(std::numeric_limits<int64_t>::max(), core_options.GetCommitTimeout()); |
| ASSERT_EQ(10, core_options.GetCommitMaxRetries()); |
| ExpireConfig expire_config = core_options.GetExpireConfig(); |
| ASSERT_EQ(10, expire_config.GetSnapshotRetainMin()); |
| ASSERT_EQ(std::numeric_limits<int32_t>::max(), expire_config.GetSnapshotRetainMax()); |
| ASSERT_EQ(50, expire_config.GetSnapshotMaxDeletes()); |
| ASSERT_FALSE(expire_config.CleanEmptyDirectories()); |
| ASSERT_EQ(1 * 3600 * 1000L, expire_config.GetSnapshotTimeRetainMs()); |
| ASSERT_EQ(std::vector<std::string>(), core_options.GetSequenceField()); |
| ASSERT_TRUE(core_options.SequenceFieldSortOrderIsAscending()); |
| ASSERT_EQ(MergeEngine::DEDUPLICATE, core_options.GetMergeEngine()); |
| ASSERT_EQ(SortEngine::LOSER_TREE, core_options.GetSortEngine()); |
| ASSERT_FALSE(core_options.IgnoreDelete()); |
| ASSERT_FALSE(core_options.WriteOnly()); |
| ASSERT_EQ(5, core_options.GetCompactionMinFileNum()); |
| ASSERT_FALSE(core_options.CompactionForceRewriteAllFiles()); |
| ASSERT_FALSE(core_options.CompactionForceUpLevel0()); |
| ASSERT_EQ(std::nullopt, core_options.GetFieldsDefaultFunc()); |
| ASSERT_EQ(std::nullopt, core_options.GetFieldAggFunc("f0").value()); |
| ASSERT_FALSE(core_options.FieldAggIgnoreRetract("f1").value()); |
| ASSERT_EQ(",", core_options.FieldListAggDelimiter("f1").value()); |
| ASSERT_FALSE(core_options.FieldCollectAggDistinct("f1").value()); |
| ASSERT_FALSE(core_options.DeletionVectorsEnabled()); |
| ASSERT_FALSE(core_options.DeletionVectorsBitmap64()); |
| ASSERT_EQ(2 * 1024 * 1024, core_options.DeletionVectorTargetFileSize()); |
| ASSERT_EQ(ChangelogProducer::NONE, core_options.GetChangelogProducer()); |
| ASSERT_FALSE(core_options.NeedLookup()); |
| ASSERT_FALSE(core_options.PrepareCommitWaitCompaction()); |
| LookupStrategy expected_lookup_strategy = {/*is_first_row=*/false, |
| /*produce_changelog=*/false, |
| /*deletion_vector=*/false, /*force_lookup=*/false}; |
| ASSERT_EQ(expected_lookup_strategy, core_options.GetLookupStrategy()); |
| ASSERT_TRUE(core_options.GetFieldsSequenceGroups().empty()); |
| ASSERT_FALSE(core_options.AggregationRemoveRecordOnDelete()); |
| ASSERT_FALSE(core_options.PartialUpdateRemoveRecordOnDelete()); |
| ASSERT_TRUE(core_options.GetPartialUpdateRemoveRecordOnSequenceGroup().empty()); |
| ASSERT_EQ(std::nullopt, core_options.GetScanFallbackBranch()); |
| ASSERT_EQ("main", core_options.GetBranch()); |
| ASSERT_TRUE(core_options.FileIndexReadEnabled()); |
| ASSERT_EQ(std::nullopt, core_options.GetDataFileExternalPaths()); |
| ASSERT_EQ(ExternalPathStrategy::NONE, core_options.GetExternalPathStrategy()); |
| ASSERT_TRUE(core_options.EnableAdaptivePrefetchStrategy()); |
| ASSERT_EQ(core_options.DataFilePrefix(), "data-"); |
| ASSERT_FALSE(core_options.IndexFileInDataFileDir()); |
| ASSERT_FALSE(core_options.RowTrackingEnabled()); |
| ASSERT_TRUE(core_options.RowTrackingPartitionGroupOnCommit()); |
| ASSERT_FALSE(core_options.DataEvolutionEnabled()); |
| ASSERT_TRUE(core_options.GetBlobFields().empty()); |
| ASSERT_TRUE(core_options.GetBlobDescriptorFields().empty()); |
| ASSERT_TRUE(core_options.GetBlobViewFields().empty()); |
| ASSERT_TRUE(core_options.GetBlobInlineFields().empty()); |
| ASSERT_TRUE(core_options.GetBlobExternalStorageFields().empty()); |
| ASSERT_EQ(std::nullopt, core_options.GetBlobExternalStoragePath()); |
| ASSERT_TRUE(core_options.LegacyPartitionNameEnabled()); |
| ASSERT_TRUE(core_options.GlobalIndexEnabled()); |
| ASSERT_EQ(std::nullopt, core_options.GetGlobalIndexExternalPath()); |
| ASSERT_EQ(std::nullopt, core_options.GetGlobalIndexThreadNum()); |
| ASSERT_EQ(std::nullopt, core_options.GetScanTagName()); |
| ASSERT_EQ(std::nullopt, core_options.GetOptimizedCompactionInterval()); |
| ASSERT_EQ(std::nullopt, core_options.GetCompactionTotalSizeThreshold()); |
| ASSERT_EQ(std::nullopt, core_options.GetCompactionIncrementalSizeThreshold()); |
| ASSERT_EQ(-1, core_options.GetCompactOffPeakStartHour()); |
| ASSERT_EQ(-1, core_options.GetCompactOffPeakEndHour()); |
| ASSERT_EQ(0, core_options.GetCompactOffPeakRatio()); |
| ASSERT_TRUE(core_options.LookupCacheBloomFilterEnabled()); |
| ASSERT_EQ(0.05, core_options.GetLookupCacheBloomFilterFpp()); |
| ASSERT_EQ("zstd", core_options.GetLookupCompressOptions().compress); |
| ASSERT_EQ(1, core_options.GetLookupCompressOptions().zstd_level); |
| ASSERT_EQ(64 * 1024, core_options.GetCachePageSize()); |
| ASSERT_EQ(200, core_options.GetCompactionMaxSizeAmplificationPercent()); |
| ASSERT_EQ(1, core_options.GetCompactionSizeRatio()); |
| ASSERT_EQ(5, core_options.GetNumSortedRunsCompactionTrigger()); |
| ASSERT_EQ(8, core_options.GetNumSortedRunsStopTrigger()); |
| ASSERT_EQ(6, core_options.GetNumLevels()); |
| ASSERT_EQ(LookupCompactMode::RADICAL, core_options.GetLookupCompactMode()); |
| ASSERT_EQ(10, core_options.GetLookupCompactMaxInterval()); |
| ASSERT_EQ(256 * 1024 * 1024, core_options.GetLookupCacheMaxMemory()); |
| ASSERT_EQ(0.25, core_options.GetLookupCacheHighPrioPoolRatio()); |
| ASSERT_EQ(1 * 3600 * 1000, core_options.GetLookupCacheFileRetentionMs()); |
| ASSERT_FALSE(core_options.TableReadSequenceNumberEnabled()); |
| ASSERT_FALSE(core_options.KeyValueSequenceNumberEnabled()); |
| ASSERT_EQ(INT64_MAX, core_options.GetLookupCacheMaxDiskSize()); |
| ASSERT_FALSE(core_options.LookupRemoteFileEnabled()); |
| ASSERT_EQ(core_options.GetLookupRemoteLevelThreshold(), INT32_MIN); |
| ASSERT_EQ(BucketFunctionType::DEFAULT, core_options.GetBucketFunctionType()); |
| } |
| |
| TEST(CoreOptionsTest, TestFromMap) { |
| std::map<std::string, std::string> options = { |
| {Options::FILE_SYSTEM, "Local"}, |
| {Options::FILE_FORMAT, "ORC"}, |
| {Options::MANIFEST_FORMAT, "avRo"}, |
| {Options::BUCKET, "3"}, |
| {Options::PAGE_SIZE, "128 kb"}, |
| {Options::TARGET_FILE_SIZE, "512MB"}, |
| {Options::BLOB_TARGET_FILE_SIZE, "1G"}, |
| {Options::PARTITION_DEFAULT_NAME, "foo"}, |
| {Options::MANIFEST_TARGET_FILE_SIZE, "16MB"}, |
| {Options::MANIFEST_FULL_COMPACTION_FILE_SIZE, "32MB"}, |
| {Options::MANIFEST_MERGE_MIN_COUNT, "2"}, |
| {Options::SOURCE_SPLIT_TARGET_SIZE, "24MB"}, |
| {Options::SOURCE_SPLIT_OPEN_FILE_COST, "32MB"}, |
| {Options::READ_BATCH_SIZE, "2048"}, |
| {Options::WRITE_BUFFER_SIZE, "16MB"}, |
| {Options::WRITE_BATCH_SIZE, "1234"}, |
| {Options::WRITE_BUFFER_SPILLABLE, "false"}, |
| {Options::WRITE_BUFFER_SPILL_MAX_DISK_SIZE, "7GB"}, |
| {Options::LOCAL_SORT_MAX_NUM_FILE_HANDLES, "64"}, |
| {Options::SPILL_COMPRESSION, "lz4"}, |
| {Options::COMMIT_FORCE_COMPACT, "true"}, |
| {Options::COMMIT_TIMEOUT, "120s"}, |
| {Options::COMMIT_MAX_RETRIES, "20"}, |
| {Options::SCAN_SNAPSHOT_ID, "5"}, |
| {Options::SCAN_MODE, "from-snapshot-full"}, |
| {Options::SNAPSHOT_NUM_RETAINED_MIN, "15"}, |
| {Options::SNAPSHOT_NUM_RETAINED_MAX, "30"}, |
| {Options::SNAPSHOT_EXPIRE_LIMIT, "20"}, |
| {Options::SNAPSHOT_TIME_RETAINED, "2h"}, |
| {Options::SNAPSHOT_CLEAN_EMPTY_DIRECTORIES, "true"}, |
| {Options::SEQUENCE_FIELD, "f1,f2,f3"}, |
| {Options::SEQUENCE_FIELD_SORT_ORDER, "descending"}, |
| {Options::MERGE_ENGINE, "partial-update"}, |
| {Options::SORT_ENGINE, "min-heap"}, |
| {Options::IGNORE_DELETE, "true"}, |
| {Options::FIELDS_DEFAULT_AGG_FUNC, "sum"}, |
| {"fields.f0.aggregate-function", "min"}, |
| {"fields.f1.ignore-retract", "true"}, |
| {"fields.f2.list-agg-delimiter", " | "}, |
| {"fields.f2.distinct", "true"}, |
| {Options::DELETION_VECTORS_ENABLED, "true"}, |
| {Options::DELETION_VECTOR_BITMAP64, "true"}, |
| {Options::DELETION_VECTOR_INDEX_FILE_TARGET_SIZE, "4MB"}, |
| {Options::CHANGELOG_PRODUCER, "full-compaction"}, |
| {Options::FORCE_LOOKUP, "true"}, |
| {"fields.g_1,g_3.sequence-group", "c,d"}, |
| {Options::AGGREGATION_REMOVE_RECORD_ON_DELETE, "true"}, |
| {Options::PARTIAL_UPDATE_REMOVE_RECORD_ON_DELETE, "true"}, |
| {Options::PARTIAL_UPDATE_REMOVE_RECORD_ON_SEQUENCE_GROUP, "a,b"}, |
| {Options::SCAN_FALLBACK_BRANCH, "fallback"}, |
| {Options::BRANCH, "rt"}, |
| {Options::FILE_INDEX_READ_ENABLED, "false"}, |
| {Options::DATA_FILE_EXTERNAL_PATHS, "FILE:///tmp/index"}, |
| {Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "round-robin"}, |
| {Options::FILE_COMPRESSION, "snappy"}, |
| {Options::MANIFEST_COMPRESSION, "zlib"}, |
| {Options::FILE_COMPRESSION_ZSTD_LEVEL, "2"}, |
| {"test.enable-adaptive-prefetch-strategy", "false"}, |
| {Options::DATA_FILE_PREFIX, "test-data-"}, |
| {Options::INDEX_FILE_IN_DATA_FILE_DIR, "true"}, |
| {Options::ROW_TRACKING_ENABLED, "true"}, |
| {Options::ROW_TRACKING_PARTITION_GROUP_ON_COMMIT, "false"}, |
| {Options::DATA_EVOLUTION_ENABLED, "true"}, |
| {Options::BLOB_FIELD, "blob1,blob2"}, |
| {Options::BLOB_DESCRIPTOR_FIELD, "blob3,blob4"}, |
| {Options::BLOB_VIEW_FIELD, "blob5"}, |
| {Options::BLOB_EXTERNAL_STORAGE_FIELD, "blob3,blob4"}, |
| {Options::BLOB_EXTERNAL_STORAGE_PATH, "FILE:///tmp/blob_external_storage/"}, |
| {Options::PARTITION_GENERATE_LEGACY_NAME, "false"}, |
| {Options::GLOBAL_INDEX_ENABLED, "false"}, |
| {Options::GLOBAL_INDEX_THREAD_NUM, "4"}, |
| {Options::GLOBAL_INDEX_EXTERNAL_PATH, "FILE:///tmp/global_index/"}, |
| {Options::SCAN_TAG_NAME, "test-tag"}, |
| {Options::WRITE_ONLY, "true"}, |
| {Options::COMPACTION_MIN_FILE_NUM, "10"}, |
| {Options::COMPACTION_FORCE_REWRITE_ALL_FILES, "true"}, |
| {Options::COMPACTION_FORCE_UP_LEVEL_0, "true"}, |
| {Options::COMPACTION_MAX_SIZE_AMPLIFICATION_PERCENT, "123"}, |
| {Options::COMPACTION_SIZE_RATIO, "9"}, |
| {Options::NUM_SORTED_RUNS_COMPACTION_TRIGGER, "11"}, |
| {Options::NUM_SORTED_RUNS_STOP_TRIGGER, "17"}, |
| {Options::NUM_LEVELS, "9"}, |
| {Options::LOOKUP_COMPACT, "gentle"}, |
| {Options::LOOKUP_COMPACT_MAX_INTERVAL, "7"}, |
| {Options::COMPACTION_OPTIMIZATION_INTERVAL, "2s"}, |
| {Options::COMPACTION_TOTAL_SIZE_THRESHOLD, "5 GB"}, |
| {Options::COMPACTION_INCREMENTAL_SIZE_THRESHOLD, "12 kB"}, |
| {Options::COMPACT_OFFPEAK_START_HOUR, "3"}, |
| {Options::COMPACT_OFFPEAK_END_HOUR, "16"}, |
| {Options::COMPACTION_OFFPEAK_RATIO, "8"}, |
| {Options::LOOKUP_CACHE_BLOOM_FILTER_ENABLED, "false"}, |
| {Options::LOOKUP_CACHE_BLOOM_FILTER_FPP, "0.5"}, |
| {Options::LOOKUP_CACHE_SPILL_COMPRESSION, "lz4"}, |
| {Options::SPILL_COMPRESSION_ZSTD_LEVEL, "2"}, |
| {Options::CACHE_PAGE_SIZE, "6MB"}, |
| {Options::FILE_FORMAT_PER_LEVEL, "0:AVRO,3:parquet"}, |
| {Options::FILE_COMPRESSION_PER_LEVEL, "0:lz4,3:none"}, |
| {Options::LOOKUP_CACHE_MAX_MEMORY_SIZE, "1MB"}, |
| {Options::LOOKUP_CACHE_HIGH_PRIO_POOL_RATIO, "0.35"}, |
| {Options::LOOKUP_CACHE_FILE_RETENTION, "30min"}, |
| {Options::LOOKUP_CACHE_MAX_DISK_SIZE, "10GB"}, |
| {Options::LOOKUP_REMOTE_FILE_ENABLED, "True"}, |
| {Options::LOOKUP_REMOTE_LEVEL_THRESHOLD, "2"}, |
| {Options::TABLE_READ_SEQUENCE_NUMBER_ENABLED, "true"}, |
| {Options::KEY_VALUE_SEQUENCE_NUMBER_ENABLED, "true"}, |
| {Options::BUCKET_FUNCTION_TYPE, "mod"}}; |
| |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| auto fs = core_options.GetFileSystem(); |
| ASSERT_TRUE(fs); |
| |
| ASSERT_EQ(core_options.GetFileFormat()->Identifier(), "orc"); |
| ASSERT_EQ(core_options.GetWriteFileFormat(0)->Identifier(), "avro"); |
| ASSERT_EQ(core_options.GetWriteFileFormat(1)->Identifier(), "orc"); |
| ASSERT_EQ(core_options.GetWriteFileFormat(3)->Identifier(), "parquet"); |
| |
| auto manifest_format = core_options.GetManifestFormat(); |
| ASSERT_EQ(manifest_format->Identifier(), "avro"); |
| |
| ASSERT_EQ(3, core_options.GetBucket()); |
| ASSERT_EQ(128 * 1024L, core_options.GetPageSize()); |
| ASSERT_EQ(512 * 1024 * 1024L, core_options.GetTargetFileSize(/*has_primary_key=*/true)); |
| ASSERT_EQ(512 * 1024 * 1024L, core_options.GetTargetFileSize(/*has_primary_key=*/false)); |
| ASSERT_EQ(1024 * 1024 * 1024L, core_options.GetBlobTargetFileSize()); |
| ASSERT_EQ("foo", core_options.GetPartitionDefaultName()); |
| ASSERT_EQ(16 * 1024 * 1024L, core_options.GetManifestTargetFileSize()); |
| ASSERT_EQ(32 * 1024 * 1024L, core_options.GetManifestFullCompactionThresholdSize()); |
| ASSERT_EQ(2, core_options.GetManifestMergeMinCount()); |
| ASSERT_EQ(24 * 1024 * 1024L, core_options.GetSourceSplitTargetSize()); |
| ASSERT_EQ(32 * 1024 * 1024L, core_options.GetSourceSplitOpenFileCost()); |
| ASSERT_EQ(2048, core_options.GetReadBatchSize()); |
| ASSERT_EQ(1234, core_options.GetWriteBatchSize()); |
| ASSERT_EQ(16 * 1024 * 1024, core_options.GetWriteBufferSize()); |
| ASSERT_FALSE(core_options.GetWriteBufferSpillable()); |
| ASSERT_EQ(7L * 1024 * 1024 * 1024, core_options.GetWriteBufferSpillMaxDiskSize()); |
| ASSERT_EQ(64, core_options.GetLocalSortMaxNumFileHandles()); |
| ASSERT_EQ("lz4", core_options.GetSpillCompressOptions().compress); |
| ASSERT_EQ(2, core_options.GetSpillCompressOptions().zstd_level); |
| ASSERT_TRUE(core_options.CommitForceCompact()); |
| ASSERT_EQ(120 * 1000, core_options.GetCommitTimeout()); |
| ASSERT_EQ(20, core_options.GetCommitMaxRetries()); |
| ASSERT_EQ(5, core_options.GetScanSnapshotId().value_or(-1)); |
| ExpireConfig expire_config = core_options.GetExpireConfig(); |
| ASSERT_EQ(15, expire_config.GetSnapshotRetainMin()); |
| ASSERT_EQ(30, expire_config.GetSnapshotRetainMax()); |
| ASSERT_EQ(20, expire_config.GetSnapshotMaxDeletes()); |
| ASSERT_EQ(2 * 3600 * 1000L, expire_config.GetSnapshotTimeRetainMs()); |
| ASSERT_TRUE(expire_config.CleanEmptyDirectories()); |
| ASSERT_EQ(std::vector<std::string>({"f1", "f2", "f3"}), core_options.GetSequenceField()); |
| ASSERT_FALSE(core_options.SequenceFieldSortOrderIsAscending()); |
| ASSERT_EQ(MergeEngine::PARTIAL_UPDATE, core_options.GetMergeEngine()); |
| ASSERT_EQ(SortEngine::MIN_HEAP, core_options.GetSortEngine()); |
| ASSERT_TRUE(core_options.IgnoreDelete()); |
| ASSERT_EQ("sum", core_options.GetFieldsDefaultFunc().value()); |
| ASSERT_EQ("min", core_options.GetFieldAggFunc("f0").value().value()); |
| ASSERT_TRUE(core_options.FieldAggIgnoreRetract("f1").value()); |
| ASSERT_TRUE(core_options.FieldAggIgnoreRetract("f1").value()); |
| ASSERT_EQ(" | ", core_options.FieldListAggDelimiter("f2").value()); |
| ASSERT_TRUE(core_options.FieldCollectAggDistinct("f2").value()); |
| ASSERT_TRUE(core_options.DeletionVectorsEnabled()); |
| ASSERT_TRUE(core_options.DeletionVectorsBitmap64()); |
| ASSERT_EQ(4 * 1024 * 1024, core_options.DeletionVectorTargetFileSize()); |
| ASSERT_EQ(ChangelogProducer::FULL_COMPACTION, core_options.GetChangelogProducer()); |
| ASSERT_TRUE(core_options.NeedLookup()); |
| ASSERT_TRUE(core_options.PrepareCommitWaitCompaction()); |
| LookupStrategy expected_lookup_strategy = {/*is_first_row=*/false, |
| /*produce_changelog=*/false, |
| /*deletion_vector=*/true, /*force_lookup=*/true}; |
| ASSERT_EQ(expected_lookup_strategy, core_options.GetLookupStrategy()); |
| |
| std::map<std::string, std::string> seq_grp; |
| seq_grp["g_1,g_3"] = "c,d"; |
| ASSERT_EQ(core_options.GetFieldsSequenceGroups(), seq_grp); |
| ASSERT_TRUE(core_options.AggregationRemoveRecordOnDelete()); |
| ASSERT_TRUE(core_options.PartialUpdateRemoveRecordOnDelete()); |
| ASSERT_EQ(core_options.GetPartialUpdateRemoveRecordOnSequenceGroup(), |
| std::vector<std::string>({"a", "b"})); |
| ASSERT_EQ(core_options.GetScanFallbackBranch(), std::optional<std::string>("fallback")); |
| ASSERT_EQ(core_options.GetBranch(), "rt"); |
| ASSERT_FALSE(core_options.FileIndexReadEnabled()); |
| ASSERT_EQ(core_options.GetDataFileExternalPaths(), |
| std::optional<std::string>("FILE:///tmp/index")); |
| ASSERT_EQ(core_options.GetExternalPathStrategy(), ExternalPathStrategy::ROUND_ROBIN); |
| ASSERT_EQ("snappy", core_options.GetFileCompression()); |
| ASSERT_EQ("lz4", core_options.GetWriteFileCompression(0)); |
| ASSERT_EQ("snappy", core_options.GetWriteFileCompression(1)); |
| ASSERT_EQ("none", core_options.GetWriteFileCompression(3)); |
| ASSERT_EQ("snappy", core_options.GetWriteFileCompression(5)); |
| ASSERT_EQ("zlib", core_options.GetManifestCompression()); |
| ASSERT_EQ(2, core_options.GetFileCompressionZstdLevel()); |
| ASSERT_FALSE(core_options.EnableAdaptivePrefetchStrategy()); |
| ASSERT_EQ(core_options.DataFilePrefix(), "test-data-"); |
| ASSERT_TRUE(core_options.IndexFileInDataFileDir()); |
| ASSERT_TRUE(core_options.RowTrackingEnabled()); |
| ASSERT_FALSE(core_options.RowTrackingPartitionGroupOnCommit()); |
| ASSERT_TRUE(core_options.DataEvolutionEnabled()); |
| ASSERT_EQ(core_options.GetBlobFields(), std::vector<std::string>({"blob1", "blob2"})); |
| ASSERT_EQ(core_options.GetBlobDescriptorFields(), std::vector<std::string>({"blob3", "blob4"})); |
| ASSERT_EQ(core_options.GetBlobViewFields(), std::vector<std::string>({"blob5"})); |
| ASSERT_EQ(core_options.GetBlobInlineFields(), |
| std::vector<std::string>({"blob3", "blob4", "blob5"})); |
| ASSERT_EQ(core_options.GetBlobExternalStorageFields(), |
| std::vector<std::string>({"blob3", "blob4"})); |
| ASSERT_EQ(core_options.GetBlobExternalStoragePath(), |
| std::optional<std::string>("FILE:///tmp/blob_external_storage/")); |
| ASSERT_FALSE(core_options.LegacyPartitionNameEnabled()); |
| ASSERT_FALSE(core_options.GlobalIndexEnabled()); |
| ASSERT_EQ(core_options.GetGlobalIndexThreadNum(), 4); |
| ASSERT_TRUE(core_options.GetGlobalIndexExternalPath()); |
| ASSERT_EQ(core_options.GetGlobalIndexExternalPath().value(), "FILE:///tmp/global_index/"); |
| ASSERT_EQ("test-tag", core_options.GetScanTagName().value()); |
| ASSERT_EQ(StartupMode::FromSnapshotFull(), core_options.GetStartupMode()); |
| ASSERT_EQ(375809637, core_options.GetCompactionFileSize(/*has_primary_key=*/true)); |
| ASSERT_EQ(375809637, core_options.GetCompactionFileSize(/*has_primary_key=*/false)); |
| ASSERT_TRUE(core_options.WriteOnly()); |
| ASSERT_EQ(10, core_options.GetCompactionMinFileNum()); |
| ASSERT_EQ(123, core_options.GetCompactionMaxSizeAmplificationPercent()); |
| ASSERT_EQ(9, core_options.GetCompactionSizeRatio()); |
| ASSERT_EQ(11, core_options.GetNumSortedRunsCompactionTrigger()); |
| ASSERT_EQ(17, core_options.GetNumSortedRunsStopTrigger()); |
| ASSERT_EQ(9, core_options.GetNumLevels()); |
| ASSERT_EQ(LookupCompactMode::GENTLE, core_options.GetLookupCompactMode()); |
| ASSERT_EQ(11, core_options.GetLookupCompactMaxInterval()); |
| ASSERT_TRUE(core_options.CompactionForceRewriteAllFiles()); |
| ASSERT_TRUE(core_options.CompactionForceUpLevel0()); |
| ASSERT_EQ(2000, core_options.GetOptimizedCompactionInterval().value()); |
| ASSERT_EQ(5l * 1024 * 1024 * 1024, core_options.GetCompactionTotalSizeThreshold().value()); |
| ASSERT_EQ(12l * 1024, core_options.GetCompactionIncrementalSizeThreshold().value()); |
| ASSERT_EQ(3, core_options.GetCompactOffPeakStartHour()); |
| ASSERT_EQ(16, core_options.GetCompactOffPeakEndHour()); |
| ASSERT_EQ(8, core_options.GetCompactOffPeakRatio()); |
| ASSERT_FALSE(core_options.LookupCacheBloomFilterEnabled()); |
| ASSERT_EQ(0.5, core_options.GetLookupCacheBloomFilterFpp()); |
| ASSERT_EQ("lz4", core_options.GetLookupCompressOptions().compress); |
| ASSERT_EQ(2, core_options.GetLookupCompressOptions().zstd_level); |
| ASSERT_EQ(6 * 1024 * 1024, core_options.GetCachePageSize()); |
| ASSERT_EQ(1024 * 1024, core_options.GetLookupCacheMaxMemory()); |
| ASSERT_EQ(0.35, core_options.GetLookupCacheHighPrioPoolRatio()); |
| ASSERT_EQ(30 * 60 * 1000, core_options.GetLookupCacheFileRetentionMs()); |
| ASSERT_EQ(10L * 1024 * 1024 * 1024, core_options.GetLookupCacheMaxDiskSize()); |
| ASSERT_TRUE(core_options.TableReadSequenceNumberEnabled()); |
| ASSERT_TRUE(core_options.KeyValueSequenceNumberEnabled()); |
| ASSERT_TRUE(core_options.LookupRemoteFileEnabled()); |
| ASSERT_EQ(core_options.GetLookupRemoteLevelThreshold(), 2); |
| ASSERT_EQ(BucketFunctionType::MOD, core_options.GetBucketFunctionType()); |
| } |
| |
| TEST(CoreOptionsTest, TestInvalidCase) { |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::BUCKET, "3.5"}}), |
| "Invalid Config [bucket: 3.5]"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::SCAN_SNAPSHOT_ID, "3.5"}}), |
| "Invalid Config [scan.snapshot-id: 3.5]"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::SEQUENCE_FIELD_SORT_ORDER, "invalid"}}), |
| "invalid sort order: invalid"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::SORT_ENGINE, "invalid"}}), |
| "invalid sort engine: invalid"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::MERGE_ENGINE, "invalid"}}), |
| "invalid merge engine: invalid"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::CHANGELOG_PRODUCER, "invalid"}}), |
| "invalid changelog producer: invalid"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::LOOKUP_COMPACT, "invalid"}}), |
| "invalid lookup mode: invalid"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::LOOKUP_COMPACT_MAX_INTERVAL, "invalid"}}), |
| "Invalid Config [lookup-compact.max-interval: invalid]"); |
| ASSERT_NOK_WITH_MSG( |
| CoreOptions::FromMap({{Options::LOOKUP_CACHE_HIGH_PRIO_POOL_RATIO, "1.1"}}), |
| "The high priority pool ratio should in the range [0, 1), while input is 1.1"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::BUCKET_FUNCTION_TYPE, "invalid"}}), |
| "invalid bucket function type: invalid"); |
| } |
| |
| TEST(CoreOptionsTest, TestLookupCompactMaxIntervalComputedValue) { |
| std::map<std::string, std::string> options = { |
| {Options::NUM_SORTED_RUNS_COMPACTION_TRIGGER, "11"}, |
| {Options::LOOKUP_COMPACT_MAX_INTERVAL, "13"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_EQ(13, core_options.GetLookupCompactMaxInterval()); |
| } |
| |
| TEST(CoreOptionsTest, TestNumSortedRunsStopTriggerFloorAndDefault) { |
| { |
| std::map<std::string, std::string> options = { |
| {Options::NUM_SORTED_RUNS_COMPACTION_TRIGGER, "11"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_EQ(14, core_options.GetNumSortedRunsStopTrigger()); |
| } |
| |
| { |
| std::map<std::string, std::string> options = { |
| {Options::NUM_SORTED_RUNS_COMPACTION_TRIGGER, "11"}, |
| {Options::NUM_SORTED_RUNS_STOP_TRIGGER, "7"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_EQ(11, core_options.GetNumSortedRunsStopTrigger()); |
| } |
| } |
| |
| TEST(CoreOptionsTest, TestLookupStrategy) { |
| { |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap({})); |
| auto strategy = core_options.GetLookupStrategy(); |
| ASSERT_FALSE(strategy.is_first_row); |
| ASSERT_FALSE(strategy.produce_changelog); |
| ASSERT_FALSE(strategy.deletion_vector); |
| ASSERT_FALSE(strategy.need_lookup); |
| } |
| { |
| std::map<std::string, std::string> options = { |
| {Options::MERGE_ENGINE, "first-row"}, |
| {Options::CHANGELOG_PRODUCER, "lookup"}, |
| {Options::DELETION_VECTORS_ENABLED, "true"}, |
| {Options::FORCE_LOOKUP, "true"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| auto strategy = core_options.GetLookupStrategy(); |
| ASSERT_TRUE(strategy.is_first_row); |
| ASSERT_TRUE(strategy.produce_changelog); |
| ASSERT_TRUE(strategy.deletion_vector); |
| ASSERT_TRUE(strategy.need_lookup); |
| } |
| } |
| |
| TEST(CoreOptionsTest, TestPrepareCommitWaitCompaction) { |
| { |
| std::map<std::string, std::string> options = { |
| {Options::FORCE_LOOKUP, "true"}, |
| {Options::LOOKUP_WAIT, "false"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_TRUE(core_options.NeedLookup()); |
| ASSERT_FALSE(core_options.PrepareCommitWaitCompaction()); |
| } |
| |
| { |
| std::map<std::string, std::string> options = { |
| {Options::FORCE_LOOKUP, "true"}, |
| {Options::LOOKUP_WAIT, "true"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_TRUE(core_options.NeedLookup()); |
| ASSERT_TRUE(core_options.PrepareCommitWaitCompaction()); |
| } |
| |
| { |
| std::map<std::string, std::string> options = { |
| {Options::LOOKUP_WAIT, "true"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_FALSE(core_options.NeedLookup()); |
| ASSERT_FALSE(core_options.PrepareCommitWaitCompaction()); |
| } |
| } |
| |
| TEST(CoreOptionsTest, TestInvalidFileFormatPerLevel) { |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::FILE_FORMAT_PER_LEVEL, "0:AVRO:parquet"}}), |
| "fail to parse key file.format.per.level, value 0:AVRO:parquet"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::FILE_FORMAT_PER_LEVEL, "aaa:avro"}}), |
| "fail to parse level aaa from string to int in file.format.per.level"); |
| } |
| |
| TEST(CoreOptionsTest, TestInvalidFileCompressionPerLevel) { |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::FILE_COMPRESSION_PER_LEVEL, "0:lz4:zstd"}}), |
| "fail to parse key file.compression.per.level, value 0:lz4:zstd"); |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::FILE_COMPRESSION_PER_LEVEL, "abc:lz4"}}), |
| "fail to parse level abc from string to int in file.compression.per.level"); |
| } |
| |
| TEST(CoreOptionsTest, TestCreateExternalPath) { |
| std::map<std::string, std::string> options = { |
| {Options::DATA_FILE_EXTERNAL_PATHS, |
| " FILE:///tmp/index1 ,FILE:///tmp/index2,FILE:///tmp/index3,,"}, |
| {Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "round-robin"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_OK_AND_ASSIGN(std::vector<std::string> external_paths, |
| core_options.CreateExternalPaths()); |
| ASSERT_EQ("FILE:/tmp/index1", external_paths[0]); |
| ASSERT_EQ("FILE:/tmp/index2", external_paths[1]); |
| ASSERT_EQ("FILE:/tmp/index3", external_paths[2]); |
| } |
| |
| TEST(CoreOptionsTest, TestInvalidCreateExternalPath) { |
| { |
| std::map<std::string, std::string> options = { |
| {Options::DATA_FILE_EXTERNAL_PATHS, |
| "/tmp/index1,FILE:///tmp/index2,FILE:///tmp/index3, "}, |
| {Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "round-robin"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_NOK_WITH_MSG(core_options.CreateExternalPaths(), |
| "scheme is null, path is /tmp/index1"); |
| } |
| { |
| std::map<std::string, std::string> options = { |
| {Options::DATA_FILE_EXTERNAL_PATHS, "FILE:///tmp/index"}, |
| {Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "specific-fs"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_NOK_WITH_MSG(core_options.CreateExternalPaths(), |
| "do not support specific-fs external path strategy for now"); |
| } |
| { |
| std::map<std::string, std::string> options = { |
| {Options::DATA_FILE_EXTERNAL_PATHS, ","}, |
| {Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "round-robin"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_NOK_WITH_MSG(core_options.CreateExternalPaths(), "external paths is empty"); |
| } |
| } |
| |
| TEST(CoreOptionsTest, TestCreateGlobalIndexExternalPath) { |
| std::map<std::string, std::string> options = { |
| {Options::GLOBAL_INDEX_EXTERNAL_PATH, " FILE:///tmp/index1"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_OK_AND_ASSIGN(std::optional<std::string> external_path, |
| core_options.CreateGlobalIndexExternalPath()); |
| ASSERT_EQ("FILE:/tmp/index1", external_path.value()); |
| } |
| |
| TEST(CoreOptionsTest, TestInvalidCreateGlobalIndexExternalPath) { |
| std::map<std::string, std::string> options = { |
| {Options::GLOBAL_INDEX_EXTERNAL_PATH, "/tmp/index1"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| ASSERT_NOK_WITH_MSG(core_options.CreateGlobalIndexExternalPath(), |
| "scheme is null, path is /tmp/index1"); |
| } |
| |
| TEST(CoreOptionsTest, TestFileSystem) { |
| { |
| auto mock_fs = std::make_shared<MockFileSystem>(); |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, |
| CoreOptions::FromMap({}, |
| /*specified_file_system=*/mock_fs)); |
| auto fs = core_options.GetFileSystem(); |
| ASSERT_TRUE(std::dynamic_pointer_cast<MockFileSystem>(fs)); |
| } |
| { |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, |
| CoreOptions::FromMap({}, |
| /*specified_file_system=*/nullptr)); |
| auto fs = core_options.GetFileSystem(); |
| auto typed_fs = std::dynamic_pointer_cast<ResolvingFileSystem>(fs); |
| ASSERT_TRUE(typed_fs); |
| ASSERT_TRUE(std::dynamic_pointer_cast<LocalFileSystem>( |
| typed_fs->GetRealFileSystem("/tmp").value_or(nullptr))); |
| } |
| { |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions core_options, |
| CoreOptions::FromMap( |
| {}, /*specified_file_system=*/nullptr, |
| /*fs_scheme_to_identifier_map=*/{{"hdfs", "mock_fs"}, {"oss", "local"}})); |
| auto fs = core_options.GetFileSystem(); |
| auto typed_fs = std::dynamic_pointer_cast<ResolvingFileSystem>(fs); |
| ASSERT_TRUE(typed_fs); |
| ASSERT_TRUE(std::dynamic_pointer_cast<LocalFileSystem>( |
| typed_fs->GetRealFileSystem("/tmp").value_or(nullptr))); |
| ASSERT_TRUE(std::dynamic_pointer_cast<MockFileSystem>( |
| typed_fs->GetRealFileSystem("hdfs:///tmp/").value_or(nullptr))); |
| ASSERT_TRUE(std::dynamic_pointer_cast<LocalFileSystem>( |
| typed_fs->GetRealFileSystem("oss:///tmp/").value_or(nullptr))); |
| } |
| } |
| |
| TEST(CoreOptionsTest, TestNormalizeValueInCoreOption) { |
| std::map<std::string, std::string> options = { |
| {Options::SEQUENCE_FIELD_SORT_ORDER, "ASCENDING"}, |
| {Options::SORT_ENGINE, "MIN-heap"}, |
| {Options::MERGE_ENGINE, "first-ROW"}, |
| {Options::CHANGELOG_PRODUCER, "LOOKUP"}, |
| {Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "ROUND-ROBIN"}, |
| {Options::LOOKUP_COMPACT, "GENTLE"}, |
| {Options::SCAN_MODE, "DEFAULT"}, |
| {Options::BUCKET_FUNCTION_TYPE, "MOD"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); |
| |
| ASSERT_EQ(StartupMode::LatestFull(), core_options.GetStartupMode()); |
| ASSERT_EQ(ExternalPathStrategy::ROUND_ROBIN, core_options.GetExternalPathStrategy()); |
| ASSERT_EQ(ChangelogProducer::LOOKUP, core_options.GetChangelogProducer()); |
| ASSERT_EQ(MergeEngine::FIRST_ROW, core_options.GetMergeEngine()); |
| ASSERT_EQ(SortEngine::MIN_HEAP, core_options.GetSortEngine()); |
| ASSERT_EQ(LookupCompactMode::GENTLE, core_options.GetLookupCompactMode()); |
| ASSERT_TRUE(core_options.SequenceFieldSortOrderIsAscending()); |
| ASSERT_EQ(BucketFunctionType::MOD, core_options.GetBucketFunctionType()); |
| } |
| |
| TEST(CoreOptionsTest, TestScanTimestampMillis) { |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, |
| CoreOptions::FromMap({{Options::SCAN_TIMESTAMP_MILLIS, "1721614515032"}})); |
| ASSERT_EQ(1721614515032, core_options.GetScanTimestampMillis().value()); |
| ASSERT_EQ(StartupMode::FromTimestamp(), core_options.GetStartupMode()); |
| } |
| |
| TEST(CoreOptionsTest, TestScanTimestampMillisExplicitMode) { |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, |
| CoreOptions::FromMap({{Options::SCAN_MODE, "from-timestamp"}, |
| {Options::SCAN_TIMESTAMP_MILLIS, "1721614515032"}})); |
| ASSERT_EQ(StartupMode::FromTimestamp(), core_options.GetStartupMode()); |
| ASSERT_EQ(1721614515032, core_options.GetScanTimestampMillis().value()); |
| } |
| |
| TEST(CoreOptionsTest, TestScanTimestampMillisNotSet) { |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap({})); |
| ASSERT_EQ(std::nullopt, core_options.GetScanTimestampMillis()); |
| ASSERT_EQ(StartupMode::LatestFull(), core_options.GetStartupMode()); |
| } |
| |
| TEST(CoreOptionsTest, TestScanTimestampString) { |
| TimezoneGuard tz_guard("Asia/Shanghai"); |
| ASSERT_OK_AND_ASSIGN(CoreOptions core_options, |
| CoreOptions::FromMap({{Options::SCAN_TIMESTAMP, "2023-06-01 00:00:00"}})); |
| ASSERT_EQ(core_options.GetScanTimestampMillis().value(), 1685548800000); |
| ASSERT_EQ(StartupMode::FromTimestamp(), core_options.GetStartupMode()); |
| } |
| |
| TEST(CoreOptionsTest, TestScanTimestampStringDateOnly) { |
| ASSERT_OK_AND_ASSIGN(CoreOptions opts1, |
| CoreOptions::FromMap({{Options::SCAN_TIMESTAMP, "2023-06-01"}})); |
| ASSERT_OK_AND_ASSIGN(CoreOptions opts2, |
| CoreOptions::FromMap({{Options::SCAN_TIMESTAMP, "2023-06-01 00:00:00"}})); |
| ASSERT_EQ(opts1.GetScanTimestampMillis().value(), opts2.GetScanTimestampMillis().value()); |
| } |
| |
| TEST(CoreOptionsTest, TestScanTimestampMillisAndStringMutuallyExclusive) { |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::SCAN_TIMESTAMP_MILLIS, "1721614515032"}, |
| {Options::SCAN_TIMESTAMP, "2023-06-01 00:00:00"}}), |
| "scan.timestamp-millis and scan.timestamp cannot be set at the same time"); |
| } |
| |
| TEST(CoreOptionsTest, TestScanTimestampInvalidString) { |
| ASSERT_NOK(CoreOptions::FromMap({{Options::SCAN_TIMESTAMP, "not-a-date"}})); |
| } |
| |
| TEST(CoreOptionsTest, TestOverflowProtection) { |
| std::string max_val = std::to_string(std::numeric_limits<int32_t>::max()); |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions options, |
| CoreOptions::FromMap({{Options::NUM_SORTED_RUNS_COMPACTION_TRIGGER, max_val}})); |
| |
| ASSERT_EQ(options.GetNumSortedRunsStopTrigger(), std::numeric_limits<int32_t>::max()); |
| ASSERT_EQ(options.GetNumLevels(), std::numeric_limits<int32_t>::max()); |
| ASSERT_EQ(options.GetLookupCompactMaxInterval(), std::numeric_limits<int32_t>::max()); |
| } |
| |
| TEST(CoreOptionsTest, TestExplicitNumLevels) { |
| ASSERT_OK_AND_ASSIGN(CoreOptions options, CoreOptions::FromMap({{Options::NUM_LEVELS, "10"}})); |
| ASSERT_EQ(options.GetNumLevels(), 10); |
| } |
| |
| TEST(CoreOptionsTest, TestParseChangelogProducer) { |
| { |
| ASSERT_OK_AND_ASSIGN(CoreOptions options, |
| CoreOptions::FromMap({{Options::CHANGELOG_PRODUCER, "none"}})); |
| ASSERT_EQ(options.GetChangelogProducer(), ChangelogProducer::NONE); |
| } |
| { |
| ASSERT_OK_AND_ASSIGN(CoreOptions options, |
| CoreOptions::FromMap({{Options::CHANGELOG_PRODUCER, "input"}})); |
| ASSERT_EQ(options.GetChangelogProducer(), ChangelogProducer::INPUT); |
| } |
| { |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions options, |
| CoreOptions::FromMap({{Options::CHANGELOG_PRODUCER, "full-compaction"}})); |
| ASSERT_EQ(options.GetChangelogProducer(), ChangelogProducer::FULL_COMPACTION); |
| } |
| { |
| ASSERT_OK_AND_ASSIGN(CoreOptions options, |
| CoreOptions::FromMap({{Options::CHANGELOG_PRODUCER, "lookup"}})); |
| ASSERT_EQ(options.GetChangelogProducer(), ChangelogProducer::LOOKUP); |
| } |
| { |
| // case insensitive |
| ASSERT_OK_AND_ASSIGN(CoreOptions options, |
| CoreOptions::FromMap({{Options::CHANGELOG_PRODUCER, "LOOKUP"}})); |
| ASSERT_EQ(options.GetChangelogProducer(), ChangelogProducer::LOOKUP); |
| } |
| ASSERT_NOK_WITH_MSG(CoreOptions::FromMap({{Options::CHANGELOG_PRODUCER, "invalid"}}), |
| "invalid changelog producer: invalid"); |
| } |
| |
| TEST(CoreOptionsTest, TestParseExternalPathStrategy) { |
| { |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions options, |
| CoreOptions::FromMap({{Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "none"}})); |
| ASSERT_EQ(options.GetExternalPathStrategy(), ExternalPathStrategy::NONE); |
| } |
| { |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions options, |
| CoreOptions::FromMap({{Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "specific-fs"}})); |
| ASSERT_EQ(options.GetExternalPathStrategy(), ExternalPathStrategy::SPECIFIC_FS); |
| } |
| { |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions options, |
| CoreOptions::FromMap({{Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "round-robin"}})); |
| ASSERT_EQ(options.GetExternalPathStrategy(), ExternalPathStrategy::ROUND_ROBIN); |
| } |
| { |
| // case insensitive |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions options, |
| CoreOptions::FromMap({{Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "ROUND-ROBIN"}})); |
| ASSERT_EQ(options.GetExternalPathStrategy(), ExternalPathStrategy::ROUND_ROBIN); |
| } |
| ASSERT_NOK_WITH_MSG( |
| CoreOptions::FromMap({{Options::DATA_FILE_EXTERNAL_PATHS_STRATEGY, "invalid"}}), |
| "invalid external path strategy: invalid"); |
| } |
| |
| TEST(CoreOptionsTest, TestCopyAssignmentOperator) { |
| // Build a CoreOptions with non-default values |
| std::map<std::string, std::string> options = { |
| {Options::BUCKET, "3"}, |
| {Options::PAGE_SIZE, "128 kb"}, |
| {Options::TARGET_FILE_SIZE, "512MB"}, |
| {Options::FILE_FORMAT, "ORC"}, |
| {Options::FILE_COMPRESSION, "lz4"}, |
| {Options::FILE_COMPRESSION_ZSTD_LEVEL, "5"}, |
| {Options::PARTITION_DEFAULT_NAME, "foo"}, |
| {Options::MANIFEST_MERGE_MIN_COUNT, "2"}, |
| {Options::READ_BATCH_SIZE, "2048"}, |
| {Options::WRITE_BATCH_SIZE, "1234"}, |
| {Options::WRITE_BUFFER_SIZE, "16MB"}, |
| {Options::WRITE_BUFFER_SPILLABLE, "false"}, |
| {Options::COMMIT_FORCE_COMPACT, "true"}, |
| {Options::COMMIT_MAX_RETRIES, "20"}, |
| {Options::SEQUENCE_FIELD, "f1,f2"}, |
| {Options::MERGE_ENGINE, "first-row"}, |
| {Options::SORT_ENGINE, "min-heap"}, |
| {Options::CHANGELOG_PRODUCER, "lookup"}, |
| {Options::DELETION_VECTORS_ENABLED, "true"}, |
| {Options::FORCE_LOOKUP, "true"}, |
| {Options::IGNORE_DELETE, "true"}, |
| {Options::WRITE_ONLY, "true"}, |
| {Options::COMPACTION_MIN_FILE_NUM, "10"}, |
| {Options::COMPACTION_FORCE_REWRITE_ALL_FILES, "true"}, |
| {Options::NUM_SORTED_RUNS_COMPACTION_TRIGGER, "11"}, |
| {Options::NUM_SORTED_RUNS_STOP_TRIGGER, "17"}, |
| {Options::NUM_LEVELS, "9"}, |
| {Options::LOOKUP_COMPACT, "gentle"}, |
| {Options::DATA_FILE_PREFIX, "test-data-"}, |
| {Options::ROW_TRACKING_ENABLED, "true"}, |
| {Options::DATA_EVOLUTION_ENABLED, "true"}, |
| {Options::BUCKET_FUNCTION_TYPE, "mod"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions source, CoreOptions::FromMap(options)); |
| |
| // Default-constructed target with different values |
| CoreOptions target; |
| |
| // Perform copy assignment |
| target = source; |
| |
| // Verify all fields are correctly copied |
| ASSERT_EQ(3, target.GetBucket()); |
| ASSERT_EQ(128 * 1024L, target.GetPageSize()); |
| ASSERT_EQ("orc", target.GetFileFormat()->Identifier()); |
| ASSERT_EQ("lz4", target.GetFileCompression()); |
| ASSERT_EQ(5, target.GetFileCompressionZstdLevel()); |
| ASSERT_EQ("foo", target.GetPartitionDefaultName()); |
| ASSERT_EQ(2, target.GetManifestMergeMinCount()); |
| ASSERT_EQ(2048, target.GetReadBatchSize()); |
| ASSERT_EQ(1234, target.GetWriteBatchSize()); |
| ASSERT_EQ(16 * 1024 * 1024L, target.GetWriteBufferSize()); |
| ASSERT_FALSE(target.GetWriteBufferSpillable()); |
| ASSERT_TRUE(target.CommitForceCompact()); |
| ASSERT_EQ(20, target.GetCommitMaxRetries()); |
| ASSERT_EQ(std::vector<std::string>({"f1", "f2"}), target.GetSequenceField()); |
| ASSERT_EQ(MergeEngine::FIRST_ROW, target.GetMergeEngine()); |
| ASSERT_EQ(SortEngine::MIN_HEAP, target.GetSortEngine()); |
| ASSERT_EQ(ChangelogProducer::LOOKUP, target.GetChangelogProducer()); |
| ASSERT_TRUE(target.DeletionVectorsEnabled()); |
| ASSERT_TRUE(target.NeedLookup()); |
| ASSERT_TRUE(target.IgnoreDelete()); |
| ASSERT_TRUE(target.WriteOnly()); |
| ASSERT_EQ(10, target.GetCompactionMinFileNum()); |
| ASSERT_TRUE(target.CompactionForceRewriteAllFiles()); |
| ASSERT_EQ(11, target.GetNumSortedRunsCompactionTrigger()); |
| ASSERT_EQ(17, target.GetNumSortedRunsStopTrigger()); |
| ASSERT_EQ(9, target.GetNumLevels()); |
| ASSERT_EQ(LookupCompactMode::GENTLE, target.GetLookupCompactMode()); |
| ASSERT_EQ("test-data-", target.DataFilePrefix()); |
| ASSERT_TRUE(target.RowTrackingEnabled()); |
| ASSERT_TRUE(target.DataEvolutionEnabled()); |
| ASSERT_EQ(BucketFunctionType::MOD, target.GetBucketFunctionType()); |
| |
| // Verify the target's ToMap matches the source's ToMap |
| ASSERT_EQ(source.ToMap(), target.ToMap()); |
| |
| CoreOptions target2 = source; |
| ASSERT_EQ(source.ToMap(), target2.ToMap()); |
| } |
| |
| TEST(CoreOptionsTest, TestAssignmentIndependence) { |
| std::map<std::string, std::string> options = { |
| {Options::BUCKET, "5"}, |
| {Options::MERGE_ENGINE, "first-row"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(CoreOptions source, CoreOptions::FromMap(options)); |
| |
| CoreOptions target; |
| target = source; |
| |
| // Verify target matches source |
| ASSERT_EQ(5, target.GetBucket()); |
| ASSERT_EQ(MergeEngine::FIRST_ROW, target.GetMergeEngine()); |
| |
| // Modify source by reassigning a different config |
| std::map<std::string, std::string> new_options = { |
| {Options::BUCKET, "99"}, |
| {Options::MERGE_ENGINE, "deduplicate"}, |
| }; |
| ASSERT_OK_AND_ASSIGN(source, CoreOptions::FromMap(new_options)); |
| |
| // Target should be unaffected (deep copy) |
| ASSERT_EQ(5, target.GetBucket()); |
| ASSERT_EQ(MergeEngine::FIRST_ROW, target.GetMergeEngine()); |
| |
| // Source should have new values |
| ASSERT_EQ(99, source.GetBucket()); |
| ASSERT_EQ(MergeEngine::DEDUPLICATE, source.GetMergeEngine()); |
| } |
| |
| TEST(CoreOptionsTest, TestFallback) { |
| { |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions options, |
| CoreOptions::FromMap({{Options::FALLBACK_BLOB_DESCRIPTOR_FIELD, "b1,b2"}})); |
| ASSERT_EQ(options.GetBlobDescriptorFields(), std::vector<std::string>({"b1", "b2"})); |
| } |
| { |
| ASSERT_OK_AND_ASSIGN( |
| CoreOptions options, |
| CoreOptions::FromMap({{Options::FALLBACK_BLOB_DESCRIPTOR_FIELD, "b1,b2"}, |
| {Options::BLOB_DESCRIPTOR_FIELD, "new_b1 , new_b2"}})); |
| ASSERT_EQ(options.GetBlobDescriptorFields(), |
| std::vector<std::string>({"new_b1", "new_b2"})); |
| } |
| } |
| } // namespace paimon::test |