From 9ba09755e1d6f8176dae64890afd58ac0e37cf33 Mon Sep 17 00:00:00 2001 From: "wangyong.alen" Date: Sat, 5 Sep 2026 01:32:51 -0400 Subject: [PATCH 1/8] feat(scan): prune append buckets from equality predicates --- docs/source/api/scan.rst | 11 +++ .../operation/append_only_file_store_scan.cpp | 24 +++++ .../operation/append_only_file_store_scan.h | 3 + .../append_only_file_store_scan_test.cpp | 99 +++++++++++++++++++ test/inte/scan_and_read_inte_test.cpp | 91 +++++++++++++++++ 5 files changed, 228 insertions(+) diff --git a/docs/source/api/scan.rst b/docs/source/api/scan.rst index 142a51682..952e3ffb9 100644 --- a/docs/source/api/scan.rst +++ b/docs/source/api/scan.rst @@ -21,6 +21,17 @@ Scan .. _cpp-api-scan: +Bucket pruning +============== + +For fixed-bucket append tables, an equality predicate on every bucket key lets +the scan derive the target bucket using the table's bucket function. Other buckets +are excluded from the scan plan without requiring an explicit bucket ID from the +caller. An explicit bucket filter takes precedence. Queries that do not constrain +all bucket keys with equality, and bucket-unaware tables, keep the existing scan +behavior. Inferred pruning applies only to files matching the scan schema and +bucket count; older layouts retain the existing filtering behavior. + Interface ========= diff --git a/src/paimon/core/operation/append_only_file_store_scan.cpp b/src/paimon/core/operation/append_only_file_store_scan.cpp index 953bcadb8..d5e189a13 100644 --- a/src/paimon/core/operation/append_only_file_store_scan.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan.cpp @@ -26,11 +26,13 @@ #include #include #include +#include #include "arrow/type.h" #include "fmt/format.h" #include "paimon/common/predicate/predicate_filter.h" #include "paimon/common/types/data_field.h" +#include "paimon/core/bucket/bucket_select_converter.h" #include "paimon/core/core_options.h" #include "paimon/core/io/data_file_meta.h" #include "paimon/core/io/file_index_evaluator.h" @@ -42,6 +44,7 @@ #include "paimon/core/utils/field_mapping.h" #include "paimon/file_index/file_index_result.h" #include "paimon/predicate/predicate_utils.h" +#include "paimon/scan_context.h" #include "paimon/status.h" namespace paimon { @@ -66,6 +69,21 @@ Result> AppendOnlyFileStoreScan::Create table_schema, arrow_schema, core_options, executor, pool)); PAIMON_RETURN_NOT_OK( scan->SplitAndSetFilter(table_schema->PartitionKeys(), arrow_schema, scan_filters)); + const auto& bucket_keys = table_schema->BucketKeys(); + int32_t num_buckets = core_options.GetBucket(); + if (scan->predicates_ && !scan_filters->GetBucketFilter().has_value() && num_buckets > 0 && + !bucket_keys.empty()) { + std::vector> bucket_key_types; + bucket_key_types.reserve(bucket_keys.size()); + for (const auto& key : bucket_keys) { + PAIMON_ASSIGN_OR_RAISE(DataField field, table_schema->GetField(key)); + bucket_key_types.push_back(field.Type()); + } + PAIMON_ASSIGN_OR_RAISE(scan->predicate_bucket_, + BucketSelectConverter::Convert( + scan->predicates_, bucket_keys, bucket_key_types, + core_options.GetBucketFunctionType(), num_buckets, pool.get())); + } return scan; } @@ -86,6 +104,12 @@ Result AppendOnlyFileStoreScan::FilterByStats(const ManifestEntry& entry) if (!predicates_) { return true; } + // A historical file may use a different schema or bucket count after a rescale. + // Keep the inferred bucket separate from the caller's explicit bucket filter. + if (predicate_bucket_ && entry.TotalBuckets() == core_options_.GetBucket() && + entry.File()->schema_id == table_schema_->Id() && entry.Bucket() != *predicate_bucket_) { + return false; + } const auto& meta = entry.File(); std::shared_ptr data_schema = table_schema_; std::shared_ptr trimmed_predicates = predicates_; diff --git a/src/paimon/core/operation/append_only_file_store_scan.h b/src/paimon/core/operation/append_only_file_store_scan.h index 65ed66da4..d16a7720f 100644 --- a/src/paimon/core/operation/append_only_file_store_scan.h +++ b/src/paimon/core/operation/append_only_file_store_scan.h @@ -18,7 +18,9 @@ #pragma once +#include #include +#include #include #include "paimon/core/operation/file_store_scan.h" @@ -79,5 +81,6 @@ class AppendOnlyFileStoreScan : public FileStoreScan { private: std::shared_ptr evolutions_; + std::optional predicate_bucket_; }; } // namespace paimon diff --git a/src/paimon/core/operation/append_only_file_store_scan_test.cpp b/src/paimon/core/operation/append_only_file_store_scan_test.cpp index fd16d47fe..4033da7c0 100644 --- a/src/paimon/core/operation/append_only_file_store_scan_test.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan_test.cpp @@ -20,15 +20,18 @@ #include #include +#include #include #include #include #include +#include "arrow/type.h" #include "gtest/gtest.h" #include "paimon/common/data/binary_row.h" #include "paimon/common/data/binary_row_writer.h" #include "paimon/common/io/cache/lru_cache.h" +#include "paimon/core/bucket/default_bucket_function.h" #include "paimon/core/manifest/manifest_entry.h" #include "paimon/core/manifest/partition_entry.h" #include "paimon/core/schema/schema_manager.h" @@ -37,6 +40,7 @@ #include "paimon/core/table/source/abstract_table_scan.h" #include "paimon/core/table/source/snapshot/snapshot_reader.h" #include "paimon/defs.h" +#include "paimon/executor.h" #include "paimon/fs/local/local_file_system.h" #include "paimon/memory/memory_pool.h" #include "paimon/metrics.h" @@ -46,10 +50,105 @@ #include "paimon/status.h" #include "paimon/table/source/scan_metrics.h" #include "paimon/table/source/table_scan.h" +#include "paimon/testing/utils/binary_row_generator.h" #include "paimon/testing/utils/testharness.h" #include "paimon/testing/utils/timezone_guard.h" namespace paimon::test { +class AppendBucketPruningTest : public testing::Test { + public: + Result> CreateScan( + const std::shared_ptr& predicate, + const std::optional& bucket = std::nullopt) const { + auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema( + {DataField(0, arrow::field("rowkey", arrow::utf8())), + DataField(1, arrow::field("value", arrow::int32()))}); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr schema, + TableSchema::Create(schema_id_, arrow_schema, {}, {}, options_)); + PAIMON_ASSIGN_OR_RAISE(CoreOptions options, CoreOptions::FromMap(options_)); + auto filters = std::make_shared( + predicate, std::vector>(), bucket); + return AppendOnlyFileStoreScan::Create(nullptr, schema_manager_, nullptr, nullptr, schema, + arrow_schema, filters, options, + GetGlobalDefaultExecutor(), pool_); + } + + void CheckBuckets(const std::shared_ptr& predicate, + const std::optional& expected_bucket, + const std::optional& explicit_bucket = std::nullopt, + int32_t total_buckets = kNumBuckets) const { + ASSERT_OK_AND_ASSIGN(auto scan, CreateScan(predicate, explicit_bucket)); + // Every file's stats match the lookup, so only bucket pruning can discard a file. + SimpleStats stats = BinaryRowGenerator::GenerateStats( + {std::string("a"), 0}, {std::string("z"), 100}, {0, 0}, pool_.get()); + ASSERT_OK_AND_ASSIGN( + auto file, + DataFileMeta::ForAppend("data.parquet", 100, 10, stats, 0, 9, 0, std::nullopt, + std::nullopt, std::nullopt, std::nullopt, std::nullopt)); + for (int32_t bucket = 0; bucket < total_buckets; ++bucket) { + ManifestEntry entry(FileKind::Add(), BinaryRow::EmptyRow(), bucket, total_buckets, + file); + ASSERT_OK_AND_ASSIGN(auto stats_scan, CreateScan(predicate, bucket)); + ASSERT_OK_AND_ASSIGN(bool stats_match, stats_scan->FilterByStats(entry)); + ASSERT_TRUE(stats_match); + ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry)); + ASSERT_EQ(keep, !expected_bucket.has_value() || bucket == expected_bucket.value()); + } + } + + std::shared_ptr KeyEquals() const { + return PredicateBuilder::Equal(0, "rowkey", FieldType::STRING, + Literal(FieldType::STRING, "key", 3)); + } + + static constexpr int32_t kNumBuckets = 4; + int64_t schema_id_ = 0; + std::shared_ptr schema_manager_; + std::shared_ptr pool_ = GetDefaultPool(); + std::map options_ = {{Options::BUCKET, "4"}, + {Options::BUCKET_KEY, "rowkey"}}; +}; + +TEST_F(AppendBucketPruningTest, PrunesStringKeyLookup) { + BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool_.get()); + int32_t expected_bucket = DefaultBucketFunction().Bucket(key, kNumBuckets); + CheckBuckets(KeyEquals(), expected_bucket); +} + +TEST_F(AppendBucketPruningTest, PreservesExplicitBucketFilter) { + BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool_.get()); + int32_t other_bucket = (DefaultBucketFunction().Bucket(key, kNumBuckets) + 1) % kNumBuckets; + CheckBuckets(KeyEquals(), other_bucket, other_bucket); +} + +TEST_F(AppendBucketPruningTest, DoesNotPruneWithoutCompleteEqualKeys) { + CheckBuckets(nullptr, std::nullopt); + CheckBuckets(PredicateBuilder::GreaterThan(0, "rowkey", FieldType::STRING, + Literal(FieldType::STRING, "key", 3)), + std::nullopt); + options_[Options::BUCKET_KEY] = "rowkey,value"; + CheckBuckets(KeyEquals(), std::nullopt); +} + +TEST_F(AppendBucketPruningTest, DoesNotPruneBucketUnawareTable) { + options_[Options::BUCKET] = "-1"; + CheckBuckets(KeyEquals(), std::nullopt); +} + +TEST_F(AppendBucketPruningTest, DoesNotPruneDifferentBucketCount) { + CheckBuckets(KeyEquals(), std::nullopt, std::nullopt, 8); +} + +TEST_F(AppendBucketPruningTest, DoesNotPruneDifferentSchema) { + auto test_dir = UniqueTestDirectory::Create("local"); + schema_manager_ = + std::make_shared(std::make_shared(), test_dir->Str()); + ASSERT_OK_AND_ASSIGN(auto scan, CreateScan(nullptr)); + ASSERT_OK(schema_manager_->CreateTable(scan->schema_, {}, {}, options_)); + schema_id_ = 1; + CheckBuckets(KeyEquals(), std::nullopt); +} + TEST(AppendOnlyFileStoreScanTest, TestReconstructPredicateWithNonCastedFields) { std::string table_root = paimon::test::GetDataDir() + diff --git a/test/inte/scan_and_read_inte_test.cpp b/test/inte/scan_and_read_inte_test.cpp index c4c50cc43..e416f312f 100644 --- a/test/inte/scan_and_read_inte_test.cpp +++ b/test/inte/scan_and_read_inte_test.cpp @@ -29,8 +29,12 @@ #include #include "arrow/api.h" +#include "arrow/c/bridge.h" #include "arrow/ipc/json_simple.h" +#include "fmt/format.h" +#include "fmt/ranges.h" #include "gtest/gtest.h" +#include "paimon/bucket/bucket_id_calculator.h" #include "paimon/common/factories/io_hook.h" #include "paimon/common/io/cache/lru_cache.h" #include "paimon/common/table/special_fields.h" @@ -421,6 +425,93 @@ TEST_P(ScanAndReadInteTest, TestWithAppendSnapshot3) { ASSERT_TRUE(expected->Equals(read_result)) << read_result->ToString(); } +TEST_P(ScanAndReadInteTest, TestWithAppendBucketKeyPointLookup) { + for (const auto& key_type : {arrow::utf8(), arrow::binary()}) { + SCOPED_TRACE(key_type->ToString()); + auto test_dir = UniqueTestDirectory::Create("local"); + arrow::FieldVector fields = {arrow::field("rowkey", key_type), + arrow::field("value", arrow::int32())}; + std::map options = {{Options::FILE_FORMAT, FileFormat()}, + {Options::FILE_SYSTEM, "local"}, + {Options::BUCKET, "4"}, + {Options::BUCKET_KEY, "rowkey"}}; + ASSERT_OK_AND_ASSIGN( + auto helper, TestHelper::Create(test_dir->Str(), arrow::schema(fields), {}, {}, options, + /*is_streaming_mode=*/false)); + std::string table_path = test_dir->Str() + "/foo.db/bar"; + // Route writes through the public bucket calculator, as a KV client does. + constexpr int32_t kNumKeys = 64; + constexpr int32_t kNumBuckets = 4; + std::vector keys; + for (int32_t i = 0; i < kNumKeys; ++i) { + keys.push_back(fmt::format(R"(["key{:03}"])", i)); + } + auto key_array = arrow::ipc::internal::json::ArrayFromJSON( + arrow::struct_({fields[0]}), fmt::format("[{}]", fmt::join(keys, ","))) + .ValueOrDie(); + ::ArrowArray c_keys; + ::ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportArray(*key_array, &c_keys, &c_schema).ok()); + ASSERT_OK_AND_ASSIGN(auto calculator, + BucketIdCalculator::Create(false, kNumBuckets, GetDefaultPool())); + std::vector bucket_ids(kNumKeys); + ASSERT_OK(calculator->CalculateBucketIds(&c_keys, &c_schema, bucket_ids.data())); + std::vector> rows(kNumBuckets); + for (int32_t i = 0; i < kNumKeys; ++i) { + rows[bucket_ids[i]].push_back(fmt::format(R"(["key{:03}", {}])", i, i)); + } + std::vector> batches; + for (int32_t bucket = 0; bucket < kNumBuckets; ++bucket) { + ASSERT_FALSE(rows[bucket].empty()); + ASSERT_OK_AND_ASSIGN( + auto batch, TestHelper::MakeRecordBatch( + arrow::struct_(fields), + fmt::format("[{}]", fmt::join(rows[bucket], ",")), {}, bucket, {})); + batches.push_back(std::move(batch)); + } + ASSERT_OK(helper->WriteAndCommit(std::move(batches), 0, std::nullopt)); + + FieldType field_type = + key_type->id() == arrow::Type::STRING ? FieldType::STRING : FieldType::BINARY; + auto predicate = + PredicateBuilder::Equal(0, "rowkey", field_type, Literal(field_type, "key032", 6)); + // All buckets have overlapping stats for this key. An explicit filter bypasses + // inference, proving that ordinary statistics cannot account for the pruning. + for (int32_t bucket = 0; bucket < kNumBuckets; ++bucket) { + ScanContextBuilder builder(table_path); + builder.SetPredicate(predicate).SetBucketFilter(bucket); + ASSERT_OK_AND_ASSIGN(auto context, FinishScanContext(builder)); + ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(context))); + ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan()); + ASSERT_FALSE(plan->Splits().empty()); + } + ScanContextBuilder scan_builder(table_path); + scan_builder.SetPredicate(predicate); + ASSERT_OK_AND_ASSIGN(auto scan_context, FinishScanContext(scan_builder)); + ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(scan_context))); + ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan()); + ASSERT_FALSE(plan->Splits().empty()); + for (const auto& split : plan->Splits()) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + ASSERT_EQ(data_split->Bucket(), bucket_ids[32]); + } + ReadContextBuilder read_builder(table_path); + AddReadOptionsForPrefetch(&read_builder); + read_builder.SetPredicate(predicate).EnablePredicateFilter(true); + ASSERT_OK_AND_ASSIGN(auto read_context, read_builder.Finish()); + ASSERT_OK_AND_ASSIGN(auto read, TableRead::Create(std::move(read_context))); + ASSERT_OK_AND_ASSIGN(auto reader, read->CreateReader(plan->Splits())); + ASSERT_OK_AND_ASSIGN(auto result, ReadResultCollector::CollectResult(std::move(reader))); + fields.insert(fields.begin(), arrow::field("_VALUE_KIND", arrow::int8())); + auto expected_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), + R"([[0, "key032", 32]])") + .ValueOrDie(); + auto expected = std::make_shared(expected_array); + ASSERT_TRUE(expected->Equals(result)) << result->ToString(); + } +} + TEST_P(ScanAndReadInteTest, TestWithAppendSnapshot5) { auto file_format = FileFormat(); std::string table_path = GetDataDir() + "/" + file_format + "/append_09.db/append_09"; From 3fde80505bf42573e6befca48e745f2458eac8cb Mon Sep 17 00:00:00 2001 From: "wangyong.alen" Date: Sat, 5 Sep 2026 02:04:58 -0400 Subject: [PATCH 2/8] fix(scan): preserve decimal and NaN bucket matches --- docs/source/api/scan.rst | 8 +++ .../core/bucket/bucket_select_converter.cpp | 18 ++++++- .../bucket/bucket_select_converter_test.cpp | 50 +++++++++++++++++++ .../append_only_file_store_scan_test.cpp | 47 ++++++++++++++++- 4 files changed, 121 insertions(+), 2 deletions(-) diff --git a/docs/source/api/scan.rst b/docs/source/api/scan.rst index 952e3ffb9..834ccaae0 100644 --- a/docs/source/api/scan.rst +++ b/docs/source/api/scan.rst @@ -32,6 +32,14 @@ all bucket keys with equality, and bucket-unaware tables, keep the existing scan behavior. Inferred pruning applies only to files matching the scan schema and bucket count; older layouts retain the existing filtering behavior. +This inference prunes data files at the manifest-entry level. It does not enable +manifest min/max-bucket skipping or the bucket-specific live-manifest-entry cache, +which require an explicit bucket filter. + +Decimal literals are rescaled to the bucket field's type only when the conversion +is exact. NaN literals and decimals that cannot be represented exactly disable +inferred bucket pruning. + Interface ========= diff --git a/src/paimon/core/bucket/bucket_select_converter.cpp b/src/paimon/core/bucket/bucket_select_converter.cpp index 3e8aa17eb..df9def000 100644 --- a/src/paimon/core/bucket/bucket_select_converter.cpp +++ b/src/paimon/core/bucket/bucket_select_converter.cpp @@ -19,6 +19,7 @@ #include "paimon/core/bucket/bucket_select_converter.h" #include +#include #include #include #include @@ -34,6 +35,7 @@ #include "paimon/core/bucket/default_bucket_function.h" #include "paimon/core/bucket/hive_bucket_function.h" #include "paimon/core/bucket/mod_bucket_function.h" +#include "paimon/core/casting/decimal_to_decimal_cast_executor.h" #include "paimon/data/timestamp.h" #include "paimon/memory/memory_pool.h" #include "paimon/predicate/leaf_predicate.h" @@ -79,7 +81,21 @@ Result> BucketSelectConverter::Convert( for (int32_t i = 0; i < num_fields; i++) { const auto& field_name = bucket_key_names[i]; - const auto& literal = literals_map.at(field_name); + Literal literal = literals_map.at(field_name); + // Equal NaNs can have different stored bit patterns and therefore different buckets. + if ((bucket_key_types[i] == FieldType::FLOAT && std::isnan(literal.GetValue())) || + (bucket_key_types[i] == FieldType::DOUBLE && std::isnan(literal.GetValue()))) { + return std::optional(std::nullopt); + } + if (bucket_key_types[i] == FieldType::DECIMAL) { + PAIMON_ASSIGN_OR_RAISE(Literal scaled, DecimalToDecimalCastExecutor().Cast( + literal, bucket_key_arrow_types[i])); + // Hash the field's representation only when rescaling preserves the exact value. + if (scaled.IsNull() || !(scaled.GetValue() == literal.GetValue())) { + return std::optional(std::nullopt); + } + literal = std::move(scaled); + } PAIMON_RETURN_NOT_OK( WriteLiteralToRow(i, literal, bucket_key_types[i], bucket_key_arrow_types[i], &writer)); } diff --git a/src/paimon/core/bucket/bucket_select_converter_test.cpp b/src/paimon/core/bucket/bucket_select_converter_test.cpp index 2707eccb6..41fd21882 100644 --- a/src/paimon/core/bucket/bucket_select_converter_test.cpp +++ b/src/paimon/core/bucket/bucket_select_converter_test.cpp @@ -18,6 +18,7 @@ #include "paimon/core/bucket/bucket_select_converter.h" +#include #include #include #include @@ -246,6 +247,55 @@ TEST_F(BucketSelectConverterTest, HiveBucketFunctionWithDecimal) { ASSERT_EQ(function->Bucket(row, num_buckets), selected_bucket.value()); } +TEST_F(BucketSelectConverterTest, RescalesDecimalLiteralsExactly) { + for (int32_t precision : {10, 20}) { + for (BucketFunctionType function_type : + {BucketFunctionType::DEFAULT, BucketFunctionType::HIVE}) { + Decimal stored = Decimal::FromUnscaledLong(120, precision, 2); + auto stored_predicate = + PredicateBuilder::Equal(0, "amount", FieldType::DECIMAL, Literal(stored)); + ASSERT_OK_AND_ASSIGN(auto expected, + BucketSelectConverter::Convert(stored_predicate, {"amount"}, + {arrow::decimal128(precision, 2)}, + function_type, 17, pool_.get())); + ASSERT_TRUE(expected.has_value()); + for (const auto& query : {Decimal::FromUnscaledLong(12, precision, 1), + Decimal::FromUnscaledLong(1200, precision, 3)}) { + auto predicate = + PredicateBuilder::Equal(0, "amount", FieldType::DECIMAL, Literal(query)); + ASSERT_OK_AND_ASSIGN( + auto result, BucketSelectConverter::Convert(predicate, {"amount"}, + {arrow::decimal128(precision, 2)}, + function_type, 17, pool_.get())); + ASSERT_EQ(result, expected); + } + } + } +} + +TEST_F(BucketSelectConverterTest, InexactDecimalConversionReturnsNullopt) { + for (const auto& value : + {Decimal::FromUnscaledLong(123, 10, 3), Decimal::FromUnscaledLong(9999999999LL, 10, 0)}) { + auto predicate = PredicateBuilder::Equal(0, "amount", FieldType::DECIMAL, Literal(value)); + ASSERT_OK_AND_ASSIGN(auto result, BucketSelectConverter::Convert( + predicate, {"amount"}, {arrow::decimal128(10, 2)}, + BucketFunctionType::DEFAULT, 17, pool_.get())); + ASSERT_FALSE(result.has_value()); + } +} + +TEST_F(BucketSelectConverterTest, NaNReturnsNullopt) { + for (const auto& value : {Literal(std::numeric_limits::quiet_NaN()), + Literal(std::numeric_limits::quiet_NaN())}) { + auto type = value.GetType() == FieldType::FLOAT ? arrow::float32() : arrow::float64(); + auto predicate = PredicateBuilder::Equal(0, "key", value.GetType(), value); + ASSERT_OK_AND_ASSIGN(auto result, BucketSelectConverter::Convert( + predicate, {"key"}, {type}, + BucketFunctionType::DEFAULT, 17, pool_.get())); + ASSERT_FALSE(result.has_value()); + } +} + TEST_F(BucketSelectConverterTest, UnsupportedFieldTypeReturnsError) { auto predicate = PredicateBuilder::Equal(0, "items", FieldType::ARRAY, Literal(static_cast(42))); diff --git a/src/paimon/core/operation/append_only_file_store_scan_test.cpp b/src/paimon/core/operation/append_only_file_store_scan_test.cpp index 4033da7c0..3fed8f791 100644 --- a/src/paimon/core/operation/append_only_file_store_scan_test.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan_test.cpp @@ -20,6 +20,7 @@ #include #include +#include #include #include #include @@ -61,7 +62,7 @@ class AppendBucketPruningTest : public testing::Test { const std::shared_ptr& predicate, const std::optional& bucket = std::nullopt) const { auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema( - {DataField(0, arrow::field("rowkey", arrow::utf8())), + {DataField(0, arrow::field("rowkey", rowkey_type_)), DataField(1, arrow::field("value", arrow::int32()))}); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr schema, TableSchema::Create(schema_id_, arrow_schema, {}, {}, options_)); @@ -101,6 +102,33 @@ class AppendBucketPruningTest : public testing::Test { Literal(FieldType::STRING, "key", 3)); } + template + void CheckMatchingValue(FieldType field_type, const T& query_value, const T& stored_value) { + options_[Options::BUCKET] = "17"; + auto predicate = PredicateBuilder::Equal(0, "rowkey", field_type, Literal(query_value)); + ASSERT_OK_AND_ASSIGN(auto comparison, + Literal(query_value).CompareTo(Literal(stored_value))); + ASSERT_EQ(comparison, 0); + BinaryRow stored_key = BinaryRowGenerator::GenerateRow({stored_value}, pool_.get()); + BinaryRow query_key = BinaryRowGenerator::GenerateRow({query_value}, pool_.get()); + int32_t bucket = DefaultBucketFunction().Bucket(stored_key, 17); + ASSERT_NE(bucket, DefaultBucketFunction().Bucket(query_key, 17)); + SimpleStats stats = BinaryRowGenerator::GenerateStats( + {stored_value, 0}, {stored_value, 100}, {0, 0}, pool_.get()); + ASSERT_OK_AND_ASSIGN( + auto file, + DataFileMeta::ForAppend("data.parquet", 100, 10, stats, 0, 9, 0, std::nullopt, + std::nullopt, std::nullopt, std::nullopt, std::nullopt)); + ManifestEntry entry(FileKind::Add(), BinaryRow::EmptyRow(), bucket, 17, file); + ASSERT_OK_AND_ASSIGN(auto stats_scan, CreateScan(predicate, bucket)); + ASSERT_OK_AND_ASSIGN(bool stats_match, stats_scan->FilterByStats(entry)); + ASSERT_TRUE(stats_match); + ASSERT_OK_AND_ASSIGN(auto scan, CreateScan(predicate)); + ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry)); + ASSERT_TRUE(keep); + } + + std::shared_ptr rowkey_type_ = arrow::utf8(); static constexpr int32_t kNumBuckets = 4; int64_t schema_id_ = 0; std::shared_ptr schema_manager_; @@ -149,6 +177,23 @@ TEST_F(AppendBucketPruningTest, DoesNotPruneDifferentSchema) { CheckBuckets(KeyEquals(), std::nullopt); } +TEST_F(AppendBucketPruningTest, PreservesCrossScaleDecimalMatch) { + rowkey_type_ = arrow::decimal128(10, 2); + CheckMatchingValue(FieldType::DECIMAL, Decimal::FromUnscaledLong(12, 10, 1), + Decimal::FromUnscaledLong(120, 10, 2)); +} + +TEST_F(AppendBucketPruningTest, PreservesDifferentNaNPayloadMatch) { + rowkey_type_ = arrow::float64(); + uint64_t query_bits = 0x7ff8000000000000ULL; + uint64_t stored_bits = 0x7ff8000000000001ULL; + double query_value; + double stored_value; + std::memcpy(&query_value, &query_bits, sizeof(query_value)); + std::memcpy(&stored_value, &stored_bits, sizeof(stored_value)); + CheckMatchingValue(FieldType::DOUBLE, query_value, stored_value); +} + TEST(AppendOnlyFileStoreScanTest, TestReconstructPredicateWithNonCastedFields) { std::string table_root = paimon::test::GetDataDir() + From d5dc6ef27429cff35c68aa0319d77aeafc148862 Mon Sep 17 00:00:00 2001 From: "wangyong.alen" Date: Sat, 5 Sep 2026 06:21:34 -0400 Subject: [PATCH 3/8] fix(scan): share total-aware bucket pruning across scans --- docs/source/api/scan.rst | 8 +- .../core/bucket/bucket_select_converter.cpp | 31 +++-- .../core/bucket/bucket_select_converter.h | 27 ++++ .../bucket/bucket_select_converter_test.cpp | 24 ++++ .../operation/append_only_file_store_scan.cpp | 22 ---- .../operation/append_only_file_store_scan.h | 2 - .../append_only_file_store_scan_test.cpp | 8 +- src/paimon/core/operation/file_store_scan.cpp | 20 +++ src/paimon/core/operation/file_store_scan.h | 10 +- .../operation/key_value_file_store_scan.cpp | 21 ---- .../key_value_file_store_scan_test.cpp | 37 ++++++ test/inte/scan_and_read_inte_test.cpp | 118 ++++++++++-------- 12 files changed, 209 insertions(+), 119 deletions(-) diff --git a/docs/source/api/scan.rst b/docs/source/api/scan.rst index 834ccaae0..dcb13c4e6 100644 --- a/docs/source/api/scan.rst +++ b/docs/source/api/scan.rst @@ -24,13 +24,15 @@ Scan Bucket pruning ============== -For fixed-bucket append tables, an equality predicate on every bucket key lets +For fixed-bucket append and primary-key tables, an equality predicate on every bucket key lets the scan derive the target bucket using the table's bucket function. Other buckets are excluded from the scan plan without requiring an explicit bucket ID from the caller. An explicit bucket filter takes precedence. Queries that do not constrain all bucket keys with equality, and bucket-unaware tables, keep the existing scan -behavior. Inferred pruning applies only to files matching the scan schema and -bucket count; older layouts retain the existing filtering behavior. +behavior. Both scan types use a shared selector that computes the bucket with +each manifest entry's total bucket count, so rescaled files are not filtered using +the current table's bucket count. Files with an older schema ID or a nonpositive +total bucket count retain the existing filtering behavior. This inference prunes data files at the manifest-entry level. It does not enable manifest min/max-bucket skipping or the bucket-specific live-manifest-entry cache, diff --git a/src/paimon/core/bucket/bucket_select_converter.cpp b/src/paimon/core/bucket/bucket_select_converter.cpp index df9def000..5ae2b243e 100644 --- a/src/paimon/core/bucket/bucket_select_converter.cpp +++ b/src/paimon/core/bucket/bucket_select_converter.cpp @@ -48,9 +48,25 @@ Result> BucketSelectConverter::Convert( const std::shared_ptr& predicate, const std::vector& bucket_key_names, const std::vector>& bucket_key_arrow_types, BucketFunctionType bucket_function_type, int32_t num_buckets, MemoryPool* pool) { + if (num_buckets <= 0) { + return std::optional(); + } + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr selector, + ConvertToSelector(predicate, bucket_key_names, bucket_key_arrow_types, + bucket_function_type, pool)); + if (!selector) { + return std::optional(); + } + return std::optional(selector->Bucket(num_buckets)); +} + +Result> BucketSelectConverter::ConvertToSelector( + const std::shared_ptr& predicate, const std::vector& bucket_key_names, + const std::vector>& bucket_key_arrow_types, + BucketFunctionType bucket_function_type, MemoryPool* pool) { assert(pool); - if (!predicate || bucket_key_names.empty() || num_buckets <= 0) { - return std::optional(std::nullopt); + if (!predicate || bucket_key_names.empty()) { + return std::unique_ptr(); } if (bucket_key_names.size() != bucket_key_arrow_types.size()) { @@ -68,7 +84,7 @@ Result> BucketSelectConverter::Convert( auto literals_opt = ExtractEqualLiterals(predicate, bucket_key_names); if (!literals_opt.has_value()) { - return std::optional(std::nullopt); + return std::unique_ptr(); } const auto& literals_map = literals_opt.value(); @@ -85,14 +101,14 @@ Result> BucketSelectConverter::Convert( // Equal NaNs can have different stored bit patterns and therefore different buckets. if ((bucket_key_types[i] == FieldType::FLOAT && std::isnan(literal.GetValue())) || (bucket_key_types[i] == FieldType::DOUBLE && std::isnan(literal.GetValue()))) { - return std::optional(std::nullopt); + return std::unique_ptr(); } if (bucket_key_types[i] == FieldType::DECIMAL) { PAIMON_ASSIGN_OR_RAISE(Literal scaled, DecimalToDecimalCastExecutor().Cast( literal, bucket_key_arrow_types[i])); // Hash the field's representation only when rescaling preserves the exact value. if (scaled.IsNull() || !(scaled.GetValue() == literal.GetValue())) { - return std::optional(std::nullopt); + return std::unique_ptr(); } literal = std::move(scaled); } @@ -101,12 +117,11 @@ Result> BucketSelectConverter::Convert( } writer.Complete(); - // Create the bucket function and compute the bucket + // Retain the key and function; the bucket count belongs to each manifest entry. PAIMON_ASSIGN_OR_RAISE( std::unique_ptr bucket_function, CreateBucketFunction(bucket_function_type, bucket_key_types, bucket_key_arrow_types)); - int32_t bucket = bucket_function->Bucket(row, num_buckets); - return std::optional(bucket); + return std::make_unique(std::move(row), std::move(bucket_function)); } std::optional> BucketSelectConverter::ExtractEqualLiterals( diff --git a/src/paimon/core/bucket/bucket_select_converter.h b/src/paimon/core/bucket/bucket_select_converter.h index 25eb31aab..e6921c390 100644 --- a/src/paimon/core/bucket/bucket_select_converter.h +++ b/src/paimon/core/bucket/bucket_select_converter.h @@ -23,10 +23,13 @@ #include #include #include +#include #include #include "arrow/type_fwd.h" #include "paimon/bucket/bucket_function_type.h" +#include "paimon/common/data/binary_row.h" +#include "paimon/core/bucket/bucket_function.h" #include "paimon/defs.h" #include "paimon/predicate/literal.h" #include "paimon/result.h" @@ -38,6 +41,22 @@ class BucketFunction; class MemoryPool; class Predicate; +/// Selects a bucket using the bucket count recorded in each manifest entry. +class BucketSelector { + public: + BucketSelector(BinaryRow key, std::unique_ptr function) + : key_(std::move(key)), function_(std::move(function)) {} + + /// Compute the bucket for a positive total bucket count. + int32_t Bucket(int32_t num_buckets) const { + return function_->Bucket(key_, num_buckets); + } + + private: + BinaryRow key_; + std::unique_ptr function_; +}; + /// Converts predicates on bucket key fields to a target bucket ID. /// When all bucket key fields have EQUAL predicates, the converter computes /// which bucket the data must reside in, enabling bucket pruning during scan. @@ -62,6 +81,14 @@ class BucketSelectConverter { const std::vector>& bucket_key_arrow_types, BucketFunctionType bucket_function_type, int32_t num_buckets, MemoryPool* pool); + /// Build a selector once, then evaluate it with each file's bucket count. + /// Returns nullptr when the predicate cannot safely constrain all bucket keys. + static Result> ConvertToSelector( + const std::shared_ptr& predicate, + const std::vector& bucket_key_names, + const std::vector>& bucket_key_arrow_types, + BucketFunctionType bucket_function_type, MemoryPool* pool); + private: /// Extract single literal per bucket key field from EQUAL predicates. /// Splits the predicate by AND and looks for EQUAL leaf predicates on bucket key fields. diff --git a/src/paimon/core/bucket/bucket_select_converter_test.cpp b/src/paimon/core/bucket/bucket_select_converter_test.cpp index 41fd21882..85a3f3e29 100644 --- a/src/paimon/core/bucket/bucket_select_converter_test.cpp +++ b/src/paimon/core/bucket/bucket_select_converter_test.cpp @@ -61,6 +61,30 @@ class BucketSelectConverterTest : public ::testing::Test { std::shared_ptr pool_ = GetDefaultPool(); }; +TEST_F(BucketSelectConverterTest, SelectorUsesEachBucketCount) { + auto pool = GetDefaultPool(); + auto predicate = + PredicateBuilder::Equal(0, "key", FieldType::INT, Literal(static_cast(-23))); + BinaryRow row = BinaryRowGenerator::GenerateRow({static_cast(-23)}, pool.get()); + ASSERT_OK_AND_ASSIGN(auto mod_function, ModBucketFunction::Create(FieldType::INT)); + ASSERT_OK_AND_ASSIGN(auto hive_function, + HiveBucketFunction::Create({HiveFieldInfo(FieldType::INT)})); + DefaultBucketFunction default_function; + const std::map functions = { + {BucketFunctionType::DEFAULT, &default_function}, + {BucketFunctionType::MOD, mod_function.get()}, + {BucketFunctionType::HIVE, hive_function.get()}}; + for (const auto& [type, function] : functions) { + ASSERT_OK_AND_ASSIGN(auto selector, + BucketSelectConverter::ConvertToSelector( + predicate, {"key"}, {arrow::int32()}, type, pool.get())); + ASSERT_TRUE(selector); + for (int32_t total_buckets : {2, 4, 8, 17, 2}) { + ASSERT_EQ(selector->Bucket(total_buckets), function->Bucket(row, total_buckets)); + } + } +} + TEST_F(BucketSelectConverterTest, SingleStringEqualDefault) { std::string value = "hello_world"; AssertDefaultBucket(FieldType::STRING, Literal(FieldType::STRING, value.c_str(), value.size()), diff --git a/src/paimon/core/operation/append_only_file_store_scan.cpp b/src/paimon/core/operation/append_only_file_store_scan.cpp index d5e189a13..435ece577 100644 --- a/src/paimon/core/operation/append_only_file_store_scan.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan.cpp @@ -32,7 +32,6 @@ #include "fmt/format.h" #include "paimon/common/predicate/predicate_filter.h" #include "paimon/common/types/data_field.h" -#include "paimon/core/bucket/bucket_select_converter.h" #include "paimon/core/core_options.h" #include "paimon/core/io/data_file_meta.h" #include "paimon/core/io/file_index_evaluator.h" @@ -69,21 +68,6 @@ Result> AppendOnlyFileStoreScan::Create table_schema, arrow_schema, core_options, executor, pool)); PAIMON_RETURN_NOT_OK( scan->SplitAndSetFilter(table_schema->PartitionKeys(), arrow_schema, scan_filters)); - const auto& bucket_keys = table_schema->BucketKeys(); - int32_t num_buckets = core_options.GetBucket(); - if (scan->predicates_ && !scan_filters->GetBucketFilter().has_value() && num_buckets > 0 && - !bucket_keys.empty()) { - std::vector> bucket_key_types; - bucket_key_types.reserve(bucket_keys.size()); - for (const auto& key : bucket_keys) { - PAIMON_ASSIGN_OR_RAISE(DataField field, table_schema->GetField(key)); - bucket_key_types.push_back(field.Type()); - } - PAIMON_ASSIGN_OR_RAISE(scan->predicate_bucket_, - BucketSelectConverter::Convert( - scan->predicates_, bucket_keys, bucket_key_types, - core_options.GetBucketFunctionType(), num_buckets, pool.get())); - } return scan; } @@ -104,12 +88,6 @@ Result AppendOnlyFileStoreScan::FilterByStats(const ManifestEntry& entry) if (!predicates_) { return true; } - // A historical file may use a different schema or bucket count after a rescale. - // Keep the inferred bucket separate from the caller's explicit bucket filter. - if (predicate_bucket_ && entry.TotalBuckets() == core_options_.GetBucket() && - entry.File()->schema_id == table_schema_->Id() && entry.Bucket() != *predicate_bucket_) { - return false; - } const auto& meta = entry.File(); std::shared_ptr data_schema = table_schema_; std::shared_ptr trimmed_predicates = predicates_; diff --git a/src/paimon/core/operation/append_only_file_store_scan.h b/src/paimon/core/operation/append_only_file_store_scan.h index d16a7720f..ba4feb1ff 100644 --- a/src/paimon/core/operation/append_only_file_store_scan.h +++ b/src/paimon/core/operation/append_only_file_store_scan.h @@ -20,7 +20,6 @@ #include #include -#include #include #include "paimon/core/operation/file_store_scan.h" @@ -81,6 +80,5 @@ class AppendOnlyFileStoreScan : public FileStoreScan { private: std::shared_ptr evolutions_; - std::optional predicate_bucket_; }; } // namespace paimon diff --git a/src/paimon/core/operation/append_only_file_store_scan_test.cpp b/src/paimon/core/operation/append_only_file_store_scan_test.cpp index 3fed8f791..b2f16944b 100644 --- a/src/paimon/core/operation/append_only_file_store_scan_test.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan_test.cpp @@ -163,8 +163,12 @@ TEST_F(AppendBucketPruningTest, DoesNotPruneBucketUnawareTable) { CheckBuckets(KeyEquals(), std::nullopt); } -TEST_F(AppendBucketPruningTest, DoesNotPruneDifferentBucketCount) { - CheckBuckets(KeyEquals(), std::nullopt, std::nullopt, 8); +TEST_F(AppendBucketPruningTest, UsesEachEntriesBucketCount) { + BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool_.get()); + for (int32_t total_buckets : {2, 4, 8, 17}) { + CheckBuckets(KeyEquals(), DefaultBucketFunction().Bucket(key, total_buckets), std::nullopt, + total_buckets); + } } TEST_F(AppendBucketPruningTest, DoesNotPruneDifferentSchema) { diff --git a/src/paimon/core/operation/file_store_scan.cpp b/src/paimon/core/operation/file_store_scan.cpp index befe7eb97..1c5b15880 100644 --- a/src/paimon/core/operation/file_store_scan.cpp +++ b/src/paimon/core/operation/file_store_scan.cpp @@ -557,6 +557,13 @@ Result FileStoreScan::FilterManifestEntry(const ManifestEntry& entry) cons if (bucket_filter_ != std::nullopt && entry.Bucket() != bucket_filter_.value()) { return false; } + // Hash with the file's bucket count, never the current table's count. An older schema + // may encode bucket keys differently, so leave those files to the existing stats filter. + if (bucket_selector_ && entry.TotalBuckets() > 0 && + entry.File()->schema_id == table_schema_->Id() && + entry.Bucket() != bucket_selector_->Bucket(entry.TotalBuckets())) { + return false; + } if (level_filter_ != nullptr && !level_filter_(entry.Level())) { return false; } @@ -602,6 +609,19 @@ Status FileStoreScan::SplitAndSetFilter(const std::vector& partitio } } bucket_filter_ = scan_filters->GetBucketFilter(); + const auto& bucket_keys = table_schema_->BucketKeys(); + if (predicates_ && !bucket_filter_ && core_options_.GetBucket() > 0 && !bucket_keys.empty()) { + std::vector> bucket_key_types; + bucket_key_types.reserve(bucket_keys.size()); + for (const auto& key : bucket_keys) { + PAIMON_ASSIGN_OR_RAISE(DataField field, table_schema_->GetField(key)); + bucket_key_types.push_back(field.Type()); + } + PAIMON_ASSIGN_OR_RAISE(bucket_selector_, + BucketSelectConverter::ConvertToSelector( + predicates_, bucket_keys, bucket_key_types, + core_options_.GetBucketFunctionType(), pool_.get())); + } if (!scan_filters->GetPartitionFilters().empty()) { PAIMON_ASSIGN_OR_RAISE( partition_filter_, diff --git a/src/paimon/core/operation/file_store_scan.h b/src/paimon/core/operation/file_store_scan.h index ea963a716..ed941152e 100644 --- a/src/paimon/core/operation/file_store_scan.h +++ b/src/paimon/core/operation/file_store_scan.h @@ -36,6 +36,7 @@ #include "paimon/common/predicate/predicate_filter.h" #include "paimon/common/utils/field_type_utils.h" #include "paimon/common/utils/linked_hash_map.h" +#include "paimon/core/bucket/bucket_select_converter.h" #include "paimon/core/core_options.h" #include "paimon/core/manifest/manifest_entry.h" #include "paimon/core/manifest/manifest_file.h" @@ -241,14 +242,6 @@ class FileStoreScan { const std::shared_ptr& arrow_schema, const std::shared_ptr& scan_filters); - /// Set the bucket filter derived from predicate analysis (e.g., BucketSelectConverter). - /// Only sets the filter if no explicit bucket filter was already set. - void SetBucketFilterIfAbsent(int32_t bucket) { - if (!bucket_filter_.has_value()) { - bucket_filter_ = bucket; - } - } - // When schema evolves, predicates might contain fields requiring casting. To avoid false // negatives when filtering by stats, we exclude those fields from predicate. static Result> ReconstructPredicateWithNonCastedFields( @@ -321,6 +314,7 @@ class FileStoreScan { std::shared_ptr partition_filter_; std::shared_ptr executor_; std::optional bucket_filter_; + std::unique_ptr bucket_selector_; std::function level_filter_; std::optional specified_snapshot_; std::shared_ptr metrics_; diff --git a/src/paimon/core/operation/key_value_file_store_scan.cpp b/src/paimon/core/operation/key_value_file_store_scan.cpp index 54af98a0c..c8b8bacc5 100644 --- a/src/paimon/core/operation/key_value_file_store_scan.cpp +++ b/src/paimon/core/operation/key_value_file_store_scan.cpp @@ -31,7 +31,6 @@ #include "paimon/common/predicate/predicate_filter.h" #include "paimon/common/types/data_field.h" #include "paimon/common/utils/object_utils.h" -#include "paimon/core/bucket/bucket_select_converter.h" #include "paimon/core/core_options.h" #include "paimon/core/io/data_file_meta.h" #include "paimon/core/options/merge_engine.h" @@ -121,26 +120,6 @@ Status KeyValueFileStoreScan::SplitAndSetKeyValueFilter( return Status::Invalid("invalid key predicate, cannot cast to PredicateFilter"); } WithKeyFilter(key_filter); - - // Bucket select conversion: derive target bucket from EQUAL predicates on bucket keys - const auto& bucket_keys = table_schema_->BucketKeys(); - int32_t num_buckets = core_options_.GetBucket(); - if (num_buckets > 0 && !bucket_keys.empty()) { - std::vector> bucket_key_arrow_types; - bucket_key_arrow_types.reserve(bucket_keys.size()); - for (const auto& key : bucket_keys) { - PAIMON_ASSIGN_OR_RAISE(DataField field, table_schema_->GetField(key)); - bucket_key_arrow_types.push_back(field.Type()); - } - PAIMON_ASSIGN_OR_RAISE( - std::optional selected_bucket, - BucketSelectConverter::Convert(key_predicate, bucket_keys, bucket_key_arrow_types, - core_options_.GetBucketFunctionType(), num_buckets, - pool_.get())); - if (selected_bucket.has_value()) { - SetBucketFilterIfAbsent(selected_bucket.value()); - } - } } // Only set value filtering when there are predicates on non-primary-key fields. diff --git a/src/paimon/core/operation/key_value_file_store_scan_test.cpp b/src/paimon/core/operation/key_value_file_store_scan_test.cpp index f6c13cb3e..271222758 100644 --- a/src/paimon/core/operation/key_value_file_store_scan_test.cpp +++ b/src/paimon/core/operation/key_value_file_store_scan_test.cpp @@ -27,6 +27,7 @@ #include "paimon/common/data/binary_row.h" #include "paimon/common/predicate/predicate_filter.h" #include "paimon/common/types/data_field.h" +#include "paimon/core/bucket/default_bucket_function.h" #include "paimon/core/core_options.h" #include "paimon/core/io/data_file_meta.h" #include "paimon/core/manifest/file_kind.h" @@ -58,6 +59,42 @@ class Schema; } // namespace arrow namespace paimon::test { +TEST(KeyValueBucketPruningTest, UsesEachEntriesBucketCount) { + auto pool = GetDefaultPool(); + auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema( + {DataField(0, arrow::field("rowkey", arrow::utf8()))}); + std::map options = {{Options::BUCKET, "4"}}; + ASSERT_OK_AND_ASSIGN(std::shared_ptr schema, + TableSchema::Create(0, arrow_schema, {}, {"rowkey"}, options)); + ASSERT_OK_AND_ASSIGN(auto core_options, CoreOptions::FromMap(options)); + auto predicate = PredicateBuilder::Equal(0, "rowkey", FieldType::STRING, + Literal(FieldType::STRING, "key", 3)); + auto filters = std::make_shared( + predicate, std::vector>(), std::nullopt); + ASSERT_OK_AND_ASSIGN(auto scan, KeyValueFileStoreScan::Create( + nullptr, nullptr, nullptr, nullptr, schema, arrow_schema, + filters, core_options, GetGlobalDefaultExecutor(), pool)); + SimpleStats stats = + BinaryRowGenerator::GenerateStats({std::string("a")}, {std::string("z")}, {0}, pool.get()); + ASSERT_OK_AND_ASSIGN( + auto file, DataFileMeta::ForAppend("data.parquet", 100, 10, stats, 0, 9, 0, std::nullopt, + std::nullopt, std::nullopt, std::nullopt, std::nullopt)); + file->key_stats = stats; + BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool.get()); + for (int32_t total_buckets : {2, 4, 8, 17}) { + int32_t expected = DefaultBucketFunction().Bucket(key, total_buckets); + for (int32_t bucket = 0; bucket < total_buckets; ++bucket) { + ManifestEntry entry(FileKind::Add(), BinaryRow::EmptyRow(), bucket, total_buckets, + file); + ASSERT_OK_AND_ASSIGN(bool stats_match, scan->FilterByStats(entry)); + ASSERT_TRUE(stats_match); + ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry)); + ASSERT_EQ(keep, bucket == expected) + << "bucket=" << bucket << " total=" << total_buckets; + } + } +} + class KeyValueFileStoreScanTest : public testing::Test { public: void SetUp() override { diff --git a/test/inte/scan_and_read_inte_test.cpp b/test/inte/scan_and_read_inte_test.cpp index e416f312f..45be64970 100644 --- a/test/inte/scan_and_read_inte_test.cpp +++ b/test/inte/scan_and_read_inte_test.cpp @@ -485,30 +485,37 @@ TEST_P(ScanAndReadInteTest, TestWithAppendBucketKeyPointLookup) { ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan()); ASSERT_FALSE(plan->Splits().empty()); } - ScanContextBuilder scan_builder(table_path); - scan_builder.SetPredicate(predicate); - ASSERT_OK_AND_ASSIGN(auto scan_context, FinishScanContext(scan_builder)); - ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(scan_context))); - ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan()); - ASSERT_FALSE(plan->Splits().empty()); - for (const auto& split : plan->Splits()) { - auto data_split = std::dynamic_pointer_cast(split); - ASSERT_TRUE(data_split); - ASSERT_EQ(data_split->Bucket(), bucket_ids[32]); + // The data was written with four buckets. Read it with rescaled table options too. + for (const std::string& current_bucket_count : {"2", "4", "8", "17"}) { + SCOPED_TRACE(current_bucket_count); + ScanContextBuilder scan_builder(table_path); + scan_builder.SetPredicate(predicate).AddOption(Options::BUCKET, current_bucket_count); + ASSERT_OK_AND_ASSIGN(auto scan_context, FinishScanContext(scan_builder)); + ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(scan_context))); + ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan()); + ASSERT_FALSE(plan->Splits().empty()); + for (const auto& split : plan->Splits()) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + ASSERT_EQ(data_split->Bucket(), bucket_ids[32]); + } + ReadContextBuilder read_builder(table_path); + AddReadOptionsForPrefetch(&read_builder); + read_builder.SetPredicate(predicate).EnablePredicateFilter(true); + ASSERT_OK_AND_ASSIGN(auto read_context, read_builder.Finish()); + ASSERT_OK_AND_ASSIGN(auto read, TableRead::Create(std::move(read_context))); + ASSERT_OK_AND_ASSIGN(auto reader, read->CreateReader(plan->Splits())); + ASSERT_OK_AND_ASSIGN(auto result, + ReadResultCollector::CollectResult(std::move(reader))); + auto expected_fields = fields; + expected_fields.insert(expected_fields.begin(), + arrow::field("_VALUE_KIND", arrow::int8())); + auto expected_array = arrow::ipc::internal::json::ArrayFromJSON( + arrow::struct_(expected_fields), R"([[0, "key032", 32]])") + .ValueOrDie(); + auto expected = std::make_shared(expected_array); + ASSERT_TRUE(expected->Equals(result)) << result->ToString(); } - ReadContextBuilder read_builder(table_path); - AddReadOptionsForPrefetch(&read_builder); - read_builder.SetPredicate(predicate).EnablePredicateFilter(true); - ASSERT_OK_AND_ASSIGN(auto read_context, read_builder.Finish()); - ASSERT_OK_AND_ASSIGN(auto read, TableRead::Create(std::move(read_context))); - ASSERT_OK_AND_ASSIGN(auto reader, read->CreateReader(plan->Splits())); - ASSERT_OK_AND_ASSIGN(auto result, ReadResultCollector::CollectResult(std::move(reader))); - fields.insert(fields.begin(), arrow::field("_VALUE_KIND", arrow::int8())); - auto expected_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), - R"([[0, "key032", 32]])") - .ValueOrDie(); - auto expected = std::make_shared(expected_array); - ASSERT_TRUE(expected->Equals(result)) << result->ToString(); } } @@ -3116,46 +3123,51 @@ TEST_P(ScanAndReadInteTest, TestWithPKBucketSelectByPredicate) { auto predicate = PredicateBuilder::Equal(/*field_index=*/2, /*field_name=*/"f2", FieldType::INT, Literal(static_cast(0))); - ScanContextBuilder scan_context_builder(table_path); - scan_context_builder.AddOption(Options::SCAN_SNAPSHOT_ID, "6"); - scan_context_builder.SetPartitionFilter({{{"f1", "10"}}}); - scan_context_builder.SetPredicate(predicate); - ASSERT_OK_AND_ASSIGN(auto scan_context, FinishScanContext(scan_context_builder)); - ASSERT_OK_AND_ASSIGN(auto table_scan, TableScan::Create(std::move(scan_context))); + // Historical files use two buckets even when the current option has changed. + for (const std::string& current_bucket_count : {"2", "4", "8", "17"}) { + SCOPED_TRACE(current_bucket_count); + ScanContextBuilder scan_context_builder(table_path); + scan_context_builder.AddOption(Options::SCAN_SNAPSHOT_ID, "6"); + scan_context_builder.AddOption(Options::BUCKET, current_bucket_count); + scan_context_builder.SetPartitionFilter({{{"f1", "10"}}}); + scan_context_builder.SetPredicate(predicate); + ASSERT_OK_AND_ASSIGN(auto scan_context, FinishScanContext(scan_context_builder)); + ASSERT_OK_AND_ASSIGN(auto table_scan, TableScan::Create(std::move(scan_context))); - ReadContextBuilder read_context_builder(table_path); - AddReadOptionsForPrefetch(&read_context_builder); - read_context_builder.SetPredicate(predicate); - ASSERT_OK_AND_ASSIGN(auto read_context, read_context_builder.Finish()); - ASSERT_OK_AND_ASSIGN(auto table_read, TableRead::Create(std::move(read_context))); + ReadContextBuilder read_context_builder(table_path); + AddReadOptionsForPrefetch(&read_context_builder); + read_context_builder.SetPredicate(predicate); + ASSERT_OK_AND_ASSIGN(auto read_context, read_context_builder.Finish()); + ASSERT_OK_AND_ASSIGN(auto table_read, TableRead::Create(std::move(read_context))); - ASSERT_OK_AND_ASSIGN(auto result_plan, table_scan->CreatePlan()); - ASSERT_EQ(result_plan->SnapshotId().value(), 6); + ASSERT_OK_AND_ASSIGN(auto result_plan, table_scan->CreatePlan()); + ASSERT_EQ(result_plan->SnapshotId().value(), 6); - // Verify all returned splits are from bucket 1 (f2=0 hashes to bucket 1) - auto splits = result_plan->Splits(); - ASSERT_FALSE(splits.empty()); - for (const auto& split : splits) { - auto data_split = std::dynamic_pointer_cast(split); - ASSERT_TRUE(data_split); - ASSERT_EQ(data_split->Bucket(), 1); - } + // Verify all returned splits are from bucket 1 (f2=0 hashes to bucket 1) + auto splits = result_plan->Splits(); + ASSERT_FALSE(splits.empty()); + for (const auto& split : splits) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + ASSERT_EQ(data_split->Bucket(), 1); + } - ASSERT_OK_AND_ASSIGN(auto batch_reader, table_read->CreateReader(splits)); - ASSERT_OK_AND_ASSIGN(auto read_result, - ReadResultCollector::CollectResult(std::move(batch_reader))); + ASSERT_OK_AND_ASSIGN(auto batch_reader, table_read->CreateReader(splits)); + ASSERT_OK_AND_ASSIGN(auto read_result, + ReadResultCollector::CollectResult(std::move(batch_reader))); - // Only rows with f2=0 in partition f1=10 should be returned - auto expected = std::make_shared( - arrow::ipc::internal::json::ArrayFromJSON(arrow_data_type_, R"([ + // Only rows with f2=0 in partition f1=10 should be returned + auto expected = std::make_shared( + arrow::ipc::internal::json::ArrayFromJSON(arrow_data_type_, R"([ [0, "Alex", 10, 0, 16.1], [0, "Bob", 10, 0, 12.1], [0, "David", 10, 0, 17.1], [0, "Emily", 10, 0, 13.1] ])") - .ValueOrDie()); - ASSERT_TRUE(expected); - ASSERT_TRUE(expected->Equals(read_result)) << read_result->ToString(); + .ValueOrDie()); + ASSERT_TRUE(expected); + ASSERT_TRUE(expected->Equals(read_result)) << read_result->ToString(); + } } TEST_P(ScanAndReadInteTest, TestReadNullableMapKey) { From a1db16536d949fe13780db21d10e9b8db6a65e07 Mon Sep 17 00:00:00 2001 From: "wangyong.alen" Date: Sat, 5 Sep 2026 07:51:52 -0400 Subject: [PATCH 4/8] fix(scan): preserve historical pruning and inferred cache --- docs/source/api/scan.rst | 17 ++- include/paimon/cache/cache.h | 10 ++ src/paimon/common/io/cache/cache_key.cpp | 23 +++- src/paimon/common/io/cache/lru_cache_test.cpp | 10 ++ .../append_only_file_store_scan_test.cpp | 29 ++++- src/paimon/core/operation/file_store_scan.cpp | 62 ++++++++-- src/paimon/core/operation/file_store_scan.h | 3 + .../key_value_file_store_scan_test.cpp | 64 +++++----- src/paimon/core/schema/schema_manager.cpp | 8 +- src/paimon/core/schema/schema_manager.h | 3 +- .../core/schema/schema_manager_test.cpp | 19 ++- test/inte/scan_and_read_inte_test.cpp | 112 +++++++++++++----- 12 files changed, 272 insertions(+), 88 deletions(-) diff --git a/docs/source/api/scan.rst b/docs/source/api/scan.rst index dcb13c4e6..0194aa6d3 100644 --- a/docs/source/api/scan.rst +++ b/docs/source/api/scan.rst @@ -31,12 +31,17 @@ caller. An explicit bucket filter takes precedence. Queries that do not constrai all bucket keys with equality, and bucket-unaware tables, keep the existing scan behavior. Both scan types use a shared selector that computes the bucket with each manifest entry's total bucket count, so rescaled files are not filtered using -the current table's bucket count. Files with an older schema ID or a nonpositive -total bucket count retain the existing filtering behavior. - -This inference prunes data files at the manifest-entry level. It does not enable -manifest min/max-bucket skipping or the bucket-specific live-manifest-entry cache, -which require an explicit bucket filter. +the current table's bucket count. Historical schemas remain eligible when their ordered bucket-key field IDs, +types and bucket function match. Changes to unrelated columns do not disable +pruning. Incompatible bucket schemas and nonpositive total bucket counts retain +the existing filtering behavior. + +This inference prunes data files at the manifest-entry level. Manifest min/max-bucket +skipping requires an explicit bucket filter. When the snapshot live-manifest-entry +cache is enabled, inferred scans cache candidates by bucket, current bucket count +and current schema ID. Files with other bucket counts or schema IDs remain in the +cached candidates and are filtered after lookup. This permits cache reuse without +discarding files that require a different bucket calculation or schema fallback. Decimal literals are rescaled to the bucket field's type only when the conversion is exact. NaN literals and decimals that cannot be represented exactly disable diff --git a/include/paimon/cache/cache.h b/include/paimon/cache/cache.h index 5bff2a120..ea1104a3e 100644 --- a/include/paimon/cache/cache.h +++ b/include/paimon/cache/cache.h @@ -48,6 +48,16 @@ class PAIMON_EXPORT CacheKey { const std::string& branch, int32_t bucket); + /// Cache candidates for an inferred bucket, including files with other bucket counts or + /// schemas. + /// @param total_buckets Positive bucket count used to compute the inferred bucket. + /// @param schema_id Schema used to build the selector. + static std::shared_ptr ForSnapshotLiveManifestEntries(const std::string& table_path, + const std::string& branch, + int32_t bucket, + int32_t total_buckets, + int64_t schema_id); + public: virtual ~CacheKey() = default; diff --git a/src/paimon/common/io/cache/cache_key.cpp b/src/paimon/common/io/cache/cache_key.cpp index c84c8f908..615d77e00 100644 --- a/src/paimon/common/io/cache/cache_key.cpp +++ b/src/paimon/common/io/cache/cache_key.cpp @@ -24,11 +24,14 @@ namespace { class SnapshotLiveManifestEntriesCacheKey : public CacheKey { public: SnapshotLiveManifestEntriesCacheKey(const std::string& table_path, const std::string& branch, - int32_t bucket) + int32_t bucket, int32_t total_buckets = 0, + int64_t schema_id = 0) : CacheKey(CacheKind::SNAPSHOT_LIVE_MANIFEST), table_path_(table_path), branch_(branch), - bucket_(bucket) {} + bucket_(bucket), + total_buckets_(total_buckets), + schema_id_(schema_id) {} bool IsIndex() const override { return false; @@ -40,7 +43,8 @@ class SnapshotLiveManifestEntriesCacheKey : public CacheKey { return false; } return table_path_ == rhs->table_path_ && branch_ == rhs->branch_ && - bucket_ == rhs->bucket_ && GetKind() == rhs->GetKind(); + bucket_ == rhs->bucket_ && total_buckets_ == rhs->total_buckets_ && + schema_id_ == rhs->schema_id_ && GetKind() == rhs->GetKind(); } size_t HashCode() const override { @@ -50,6 +54,8 @@ class SnapshotLiveManifestEntriesCacheKey : public CacheKey { seed ^= std::hash{}(bucket_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); seed ^= std::hash{}(static_cast(GetKind())) + HASH_CONSTANT + (seed << 6) + (seed >> 2); + seed ^= std::hash{}(total_buckets_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); + seed ^= std::hash{}(schema_id_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); return seed; } @@ -59,6 +65,8 @@ class SnapshotLiveManifestEntriesCacheKey : public CacheKey { const std::string table_path_; const std::string branch_; const int32_t bucket_; + const int32_t total_buckets_; + const int64_t schema_id_; }; } // namespace @@ -82,6 +90,15 @@ std::shared_ptr CacheKey::ForSnapshotLiveManifestEntries(const std::st return std::make_shared(table_path, branch, bucket); } +std::shared_ptr CacheKey::ForSnapshotLiveManifestEntries(const std::string& table_path, + const std::string& branch, + int32_t bucket, + int32_t total_buckets, + int64_t schema_id) { + return std::make_shared(table_path, branch, bucket, + total_buckets, schema_id); +} + bool PositionCacheKey::IsIndex() const { return is_index_; } diff --git a/src/paimon/common/io/cache/lru_cache_test.cpp b/src/paimon/common/io/cache/lru_cache_test.cpp index 1d644c70e..6992cdb38 100644 --- a/src/paimon/common/io/cache/lru_cache_test.cpp +++ b/src/paimon/common/io/cache/lru_cache_test.cpp @@ -383,6 +383,16 @@ TEST_F(LruCacheTest, TestForKindSetsKeyKind) { ASSERT_EQ(CacheKind::MANIFEST, put_key->GetKind()); } +TEST_F(LruCacheTest, InferredManifestCacheKeysIncludeBucketCountAndSchema) { + auto key = CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1, 4, 0); + auto same = CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1, 4, 0); + ASSERT_TRUE(key->Equals(*same)); + ASSERT_EQ(key->HashCode(), same->HashCode()); + ASSERT_FALSE(key->Equals(*CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1))); + ASSERT_FALSE(key->Equals(*CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1, 8, 0))); + ASSERT_FALSE(key->Equals(*CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1, 4, 1))); +} + TEST_F(LruCacheTest, TestForSnapshotLiveManifestEntries) { auto main_key = CacheKey::ForSnapshotLiveManifestEntries("table_path", "main", 0); auto same_key = CacheKey::ForSnapshotLiveManifestEntries("table_path", "main", 0); diff --git a/src/paimon/core/operation/append_only_file_store_scan_test.cpp b/src/paimon/core/operation/append_only_file_store_scan_test.cpp index b2f16944b..46cfcf1e3 100644 --- a/src/paimon/core/operation/append_only_file_store_scan_test.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan_test.cpp @@ -61,9 +61,12 @@ class AppendBucketPruningTest : public testing::Test { Result> CreateScan( const std::shared_ptr& predicate, const std::optional& bucket = std::nullopt) const { - auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema( - {DataField(0, arrow::field("rowkey", rowkey_type_)), - DataField(1, arrow::field("value", arrow::int32()))}); + std::vector fields = {DataField(0, arrow::field("rowkey", rowkey_type_)), + DataField(1, arrow::field("value", arrow::int32()))}; + if (extra_field_) { + fields.emplace_back(2, arrow::field("extra", arrow::int32())); + } + auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(fields); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr schema, TableSchema::Create(schema_id_, arrow_schema, {}, {}, options_)); PAIMON_ASSIGN_OR_RAISE(CoreOptions options, CoreOptions::FromMap(options_)); @@ -131,6 +134,7 @@ class AppendBucketPruningTest : public testing::Test { std::shared_ptr rowkey_type_ = arrow::utf8(); static constexpr int32_t kNumBuckets = 4; int64_t schema_id_ = 0; + bool extra_field_ = false; std::shared_ptr schema_manager_; std::shared_ptr pool_ = GetDefaultPool(); std::map options_ = {{Options::BUCKET, "4"}, @@ -171,14 +175,29 @@ TEST_F(AppendBucketPruningTest, UsesEachEntriesBucketCount) { } } -TEST_F(AppendBucketPruningTest, DoesNotPruneDifferentSchema) { +TEST_F(AppendBucketPruningTest, PrunesCompatibleHistoricalSchema) { auto test_dir = UniqueTestDirectory::Create("local"); schema_manager_ = std::make_shared(std::make_shared(), test_dir->Str()); ASSERT_OK_AND_ASSIGN(auto scan, CreateScan(nullptr)); ASSERT_OK(schema_manager_->CreateTable(scan->schema_, {}, {}, options_)); schema_id_ = 1; - CheckBuckets(KeyEquals(), std::nullopt); + extra_field_ = true; + BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool_.get()); + CheckBuckets(KeyEquals(), DefaultBucketFunction().Bucket(key, kNumBuckets)); +} + +TEST_F(AppendBucketPruningTest, PreservesIncompatibleHistoricalBucketKeys) { + auto test_dir = UniqueTestDirectory::Create("local"); + schema_manager_ = + std::make_shared(std::make_shared(), test_dir->Str()); + ASSERT_OK_AND_ASSIGN(auto scan, CreateScan(nullptr)); + ASSERT_OK(schema_manager_->CreateTable(scan->schema_, {}, {}, options_)); + schema_id_ = 1; + rowkey_type_ = arrow::binary(); + auto predicate = PredicateBuilder::Equal(0, "rowkey", FieldType::BINARY, + Literal(FieldType::BINARY, "key", 3)); + CheckBuckets(predicate, std::nullopt); } TEST_F(AppendBucketPruningTest, PreservesCrossScaleDecimalMatch) { diff --git a/src/paimon/core/operation/file_store_scan.cpp b/src/paimon/core/operation/file_store_scan.cpp index 1c5b15880..6bab2de65 100644 --- a/src/paimon/core/operation/file_store_scan.cpp +++ b/src/paimon/core/operation/file_store_scan.cpp @@ -146,16 +146,20 @@ Result> FileStoreScan::CreatePlan() cons ReadManifests(&snapshot, &all_manifest_file_metas, &filtered_manifest_file_metas)); std::vector manifest_entries; + std::optional cache_bucket = bucket_filter_; + if (!cache_bucket && bucket_selector_) { + cache_bucket = bucket_selector_->Bucket(core_options_.GetBucket()); + } const bool use_snapshot_live_manifest_cache = snapshot.has_value() && scan_mode_ == ScanMode::ALL && core_options_.GetScanManifestEntryCacheMaxSnapshots() > 0 && core_options_.GetCache() != nullptr && !table_path_.empty() && - !row_range_index_.has_value() && bucket_filter_.has_value(); + !row_range_index_.has_value() && cache_bucket.has_value(); uint64_t lazy_decode_scanned_rows = 0; bool snapshot_cache_hit = false; if (use_snapshot_live_manifest_cache) { PAIMON_RETURN_NOT_OK(ReadManifestEntriesWithCache(snapshot.value(), all_manifest_file_metas, - bucket_filter_.value(), &manifest_entries, + cache_bucket.value(), &manifest_entries, &snapshot_cache_hit)); lazy_decode_scanned_rows = manifest_entries.size(); std::vector filtered_entries; @@ -325,7 +329,8 @@ Status FileStoreScan::ReadManifestEntries(const std::vector& m } // Cache merged live manifest entries for one bucket before applying scan filters. Each cache value -// keeps a bounded number of snapshot results for the same table/branch/bucket. Exact snapshot hits +// keeps bounded snapshot results for a table/branch/bucket. Inferred entries also include other +// bucket counts and schema IDs, with the current count and schema in the cache key. Exact hits // can be returned directly; cache misses rebuild the target snapshot bucket from the target // snapshot's data manifests. Status FileStoreScan::ReadManifestEntriesWithCache( @@ -351,7 +356,7 @@ Status FileStoreScan::ReadManifestEntriesWithCache( // cache. std::vector bucket_manifest_metas; for (const auto& meta : all_manifest_metas) { - if (MayContainBucket(meta, bucket)) { + if ((!bucket_filter_ && bucket_selector_) || MayContainBucket(meta, bucket)) { bucket_manifest_metas.push_back(meta); } } @@ -369,6 +374,11 @@ Status FileStoreScan::ReadManifestEntriesWithCache( } std::shared_ptr FileStoreScan::SnapshotLiveManifestEntriesCacheKey(int32_t bucket) const { + if (!bucket_filter_ && bucket_selector_) { + return CacheKey::ForSnapshotLiveManifestEntries( + table_path_, BranchManager::NormalizeBranch(core_options_.GetBranch()), bucket, + core_options_.GetBucket(), table_schema_->Id()); + } return CacheKey::ForSnapshotLiveManifestEntries( table_path_, BranchManager::NormalizeBranch(core_options_.GetBranch()), bucket); } @@ -409,7 +419,9 @@ Status FileStoreScan::StoreSnapshotLiveManifestEntries( Status FileStoreScan::ReadAndMergeBucketFileEntries( const std::vector& manifest_metas, int32_t bucket, std::vector* merged_entries) const { - if (core_options_.ScanManifestEntryLazyDecodeEnabled()) { + const bool inferred_bucket = !bucket_filter_ && bucket_selector_ != nullptr; + // Explicit-bucket lazy decoding cannot retain entries with a different layout. + if (!inferred_bucket && core_options_.ScanManifestEntryLazyDecodeEnabled()) { std::vector>>> futures; futures.reserve(manifest_metas.size()); for (const auto& meta : manifest_metas) { @@ -441,7 +453,9 @@ Status FileStoreScan::ReadAndMergeBucketFileEntries( PAIMON_RETURN_NOT_OK(ReadFileEntries(manifest_metas, &entries, /*apply_scan_filter=*/false)); unmerged_entries.reserve(entries.size()); for (auto& entry : entries) { - if (entry.Bucket() == bucket) { + if (entry.Bucket() == bucket || + (inferred_bucket && (entry.TotalBuckets() != core_options_.GetBucket() || + entry.File()->schema_id != table_schema_->Id()))) { unmerged_entries.emplace_back(std::move(entry)); } } @@ -557,12 +571,11 @@ Result FileStoreScan::FilterManifestEntry(const ManifestEntry& entry) cons if (bucket_filter_ != std::nullopt && entry.Bucket() != bucket_filter_.value()) { return false; } - // Hash with the file's bucket count, never the current table's count. An older schema - // may encode bucket keys differently, so leave those files to the existing stats filter. - if (bucket_selector_ && entry.TotalBuckets() > 0 && - entry.File()->schema_id == table_schema_->Id() && - entry.Bucket() != bucket_selector_->Bucket(entry.TotalBuckets())) { - return false; + if (bucket_selector_ && entry.TotalBuckets() > 0) { + PAIMON_ASSIGN_OR_RAISE(bool compatible, HasCompatibleBucketKeys(entry.File()->schema_id)); + if (compatible && entry.Bucket() != bucket_selector_->Bucket(entry.TotalBuckets())) { + return false; + } } if (level_filter_ != nullptr && !level_filter_(entry.Level())) { return false; @@ -570,6 +583,31 @@ Result FileStoreScan::FilterManifestEntry(const ManifestEntry& entry) cons return FilterByStats(entry); } +Result FileStoreScan::HasCompatibleBucketKeys(int64_t data_schema_id) const { + if (data_schema_id == table_schema_->Id()) { + return true; + } + auto cached = bucket_schema_compatibility_.Find(data_schema_id); + if (cached) { + return cached.value(); + } + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr data_schema, + schema_manager_->ReadSchema(data_schema_id)); + const auto& current_keys = table_schema_->BucketKeys(); + const auto& data_keys = data_schema->BucketKeys(); + PAIMON_ASSIGN_OR_RAISE(CoreOptions data_options, CoreOptions::FromMap(data_schema->Options())); + bool compatible = current_keys.size() == data_keys.size() && + core_options_.GetBucketFunctionType() == data_options.GetBucketFunctionType(); + for (size_t i = 0; compatible && i < current_keys.size(); ++i) { + PAIMON_ASSIGN_OR_RAISE(DataField current_field, table_schema_->GetField(current_keys[i])); + PAIMON_ASSIGN_OR_RAISE(DataField data_field, data_schema->GetField(data_keys[i])); + compatible = current_field.Id() == data_field.Id() && + current_field.Type()->Equals(data_field.Type()); + } + bucket_schema_compatibility_.Insert(data_schema_id, compatible); + return compatible; +} + Status FileStoreScan::SplitAndSetFilter(const std::vector& partition_keys, const std::shared_ptr& arrow_schema, const std::shared_ptr& scan_filters) { diff --git a/src/paimon/core/operation/file_store_scan.h b/src/paimon/core/operation/file_store_scan.h index ed941152e..6cc1e4e37 100644 --- a/src/paimon/core/operation/file_store_scan.h +++ b/src/paimon/core/operation/file_store_scan.h @@ -34,6 +34,7 @@ #include "paimon/common/predicate/leaf_predicate_impl.h" #include "paimon/common/predicate/literal_converter.h" #include "paimon/common/predicate/predicate_filter.h" +#include "paimon/common/utils/concurrent_hash_map.h" #include "paimon/common/utils/field_type_utils.h" #include "paimon/common/utils/linked_hash_map.h" #include "paimon/core/bucket/bucket_select_converter.h" @@ -293,6 +294,7 @@ class FileStoreScan { std::vector* entries) const; Result FilterManifestEntry(const ManifestEntry& entry) const; + Result HasCompatibleBucketKeys(int64_t data_schema_id) const; protected: std::shared_ptr pool_; @@ -315,6 +317,7 @@ class FileStoreScan { std::shared_ptr executor_; std::optional bucket_filter_; std::unique_ptr bucket_selector_; + mutable ConcurrentHashMap bucket_schema_compatibility_; std::function level_filter_; std::optional specified_snapshot_; std::shared_ptr metrics_; diff --git a/src/paimon/core/operation/key_value_file_store_scan_test.cpp b/src/paimon/core/operation/key_value_file_store_scan_test.cpp index 271222758..08867d4cb 100644 --- a/src/paimon/core/operation/key_value_file_store_scan_test.cpp +++ b/src/paimon/core/operation/key_value_file_store_scan_test.cpp @@ -45,6 +45,7 @@ #include "paimon/defs.h" #include "paimon/executor.h" #include "paimon/format/file_format.h" +#include "paimon/fs/local/local_file_system.h" #include "paimon/memory/memory_pool.h" #include "paimon/metrics.h" #include "paimon/predicate/literal.h" @@ -64,33 +65,42 @@ TEST(KeyValueBucketPruningTest, UsesEachEntriesBucketCount) { auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema( {DataField(0, arrow::field("rowkey", arrow::utf8()))}); std::map options = {{Options::BUCKET, "4"}}; - ASSERT_OK_AND_ASSIGN(std::shared_ptr schema, - TableSchema::Create(0, arrow_schema, {}, {"rowkey"}, options)); - ASSERT_OK_AND_ASSIGN(auto core_options, CoreOptions::FromMap(options)); - auto predicate = PredicateBuilder::Equal(0, "rowkey", FieldType::STRING, - Literal(FieldType::STRING, "key", 3)); - auto filters = std::make_shared( - predicate, std::vector>(), std::nullopt); - ASSERT_OK_AND_ASSIGN(auto scan, KeyValueFileStoreScan::Create( - nullptr, nullptr, nullptr, nullptr, schema, arrow_schema, - filters, core_options, GetGlobalDefaultExecutor(), pool)); - SimpleStats stats = - BinaryRowGenerator::GenerateStats({std::string("a")}, {std::string("z")}, {0}, pool.get()); - ASSERT_OK_AND_ASSIGN( - auto file, DataFileMeta::ForAppend("data.parquet", 100, 10, stats, 0, 9, 0, std::nullopt, - std::nullopt, std::nullopt, std::nullopt, std::nullopt)); - file->key_stats = stats; - BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool.get()); - for (int32_t total_buckets : {2, 4, 8, 17}) { - int32_t expected = DefaultBucketFunction().Bucket(key, total_buckets); - for (int32_t bucket = 0; bucket < total_buckets; ++bucket) { - ManifestEntry entry(FileKind::Add(), BinaryRow::EmptyRow(), bucket, total_buckets, - file); - ASSERT_OK_AND_ASSIGN(bool stats_match, scan->FilterByStats(entry)); - ASSERT_TRUE(stats_match); - ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry)); - ASSERT_EQ(keep, bucket == expected) - << "bucket=" << bucket << " total=" << total_buckets; + auto test_dir = UniqueTestDirectory::Create("local"); + auto manager = + std::make_shared(std::make_shared(), test_dir->Str()); + ASSERT_OK(manager->CreateTable(arrow_schema, {}, {"rowkey"}, options)); + for (int64_t current_schema_id : {0, 1}) { + ASSERT_OK_AND_ASSIGN( + std::shared_ptr schema, + TableSchema::Create(current_schema_id, arrow_schema, {}, {"rowkey"}, options)); + ASSERT_OK_AND_ASSIGN(auto core_options, CoreOptions::FromMap(options)); + auto predicate = PredicateBuilder::Equal(0, "rowkey", FieldType::STRING, + Literal(FieldType::STRING, "key", 3)); + auto filters = std::make_shared( + predicate, std::vector>(), std::nullopt); + ASSERT_OK_AND_ASSIGN( + auto scan, + KeyValueFileStoreScan::Create(nullptr, manager, nullptr, nullptr, schema, arrow_schema, + filters, core_options, GetGlobalDefaultExecutor(), pool)); + SimpleStats stats = BinaryRowGenerator::GenerateStats({std::string("a")}, + {std::string("z")}, {0}, pool.get()); + ASSERT_OK_AND_ASSIGN( + auto file, + DataFileMeta::ForAppend("data.parquet", 100, 10, stats, 0, 9, 0, std::nullopt, + std::nullopt, std::nullopt, std::nullopt, std::nullopt)); + file->key_stats = stats; + BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool.get()); + for (int32_t total_buckets : {2, 4, 8, 17}) { + int32_t expected = DefaultBucketFunction().Bucket(key, total_buckets); + for (int32_t bucket = 0; bucket < total_buckets; ++bucket) { + ManifestEntry entry(FileKind::Add(), BinaryRow::EmptyRow(), bucket, total_buckets, + file); + ASSERT_OK_AND_ASSIGN(bool stats_match, scan->FilterByStats(entry)); + ASSERT_TRUE(stats_match); + ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry)); + ASSERT_EQ(keep, bucket == expected) + << "bucket=" << bucket << " total=" << total_buckets; + } } } } diff --git a/src/paimon/core/schema/schema_manager.cpp b/src/paimon/core/schema/schema_manager.cpp index 2425cc4a6..fe44d3c67 100644 --- a/src/paimon/core/schema/schema_manager.cpp +++ b/src/paimon/core/schema/schema_manager.cpp @@ -69,16 +69,16 @@ Result>> SchemaManager::Latest() cons } Result> SchemaManager::ReadSchema(int64_t schema_id) const { - auto iter = schema_cache_.find(schema_id); - if (iter != schema_cache_.end()) { - return iter->second; + auto cached = schema_cache_.Find(schema_id); + if (cached) { + return cached.value(); } auto path = ToSchemaPath(schema_id); std::string content; PAIMON_RETURN_NOT_OK(file_system_->ReadFile(path, &content)); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr schema, TableSchema::CreateFromJson(content)); - schema_cache_[schema_id] = schema; + schema_cache_.Insert(schema_id, schema); return schema; } diff --git a/src/paimon/core/schema/schema_manager.h b/src/paimon/core/schema/schema_manager.h index 382fa3033..7480ec5da 100644 --- a/src/paimon/core/schema/schema_manager.h +++ b/src/paimon/core/schema/schema_manager.h @@ -26,6 +26,7 @@ #include #include +#include "paimon/common/utils/concurrent_hash_map.h" #include "paimon/core/schema/table_schema.h" #include "paimon/fs/file_system.h" #include "paimon/result.h" @@ -69,7 +70,7 @@ class SchemaManager { std::shared_ptr file_system_; std::string table_root_; const std::string branch_; - mutable std::map> schema_cache_; + mutable ConcurrentHashMap> schema_cache_; }; } // namespace paimon diff --git a/src/paimon/core/schema/schema_manager_test.cpp b/src/paimon/core/schema/schema_manager_test.cpp index 83f995a99..0b078462b 100644 --- a/src/paimon/core/schema/schema_manager_test.cpp +++ b/src/paimon/core/schema/schema_manager_test.cpp @@ -19,6 +19,7 @@ #include "paimon/core/schema/schema_manager.h" +#include #include #include @@ -30,6 +31,22 @@ namespace paimon::test { +TEST(SchemaManagerTest, ConcurrentHistoricalSchemaReads) { + SchemaManager manager( + std::make_shared(), + GetDataDir() + "/orc/pk_table_with_alter_table.db/pk_table_with_alter_table/"); + std::vector>>> reads; + for (int32_t i = 0; i < 16; ++i) { + reads.push_back( + std::async(std::launch::async, [&manager, i]() { return manager.ReadSchema(i % 2); })); + } + for (int32_t i = 0; i < 16; ++i) { + ASSERT_OK_AND_ASSIGN(auto schema, reads[i].get()); + ASSERT_EQ(schema->Id(), i % 2); + } + ASSERT_EQ(manager.schema_cache_.Size(), 2); +} + TEST(SchemaManagerTest, TestSimple) { auto fs = std::make_shared(); std::string table_root = @@ -91,7 +108,7 @@ TEST(SchemaManagerTest, TestSimple) { })"; ASSERT_OK_AND_ASSIGN(auto expected_schema, TableSchema::CreateFromJson(schema_json)); ASSERT_EQ(*ret, *expected_schema); - ASSERT_FALSE(manager.schema_cache_.empty()); + ASSERT_GT(manager.schema_cache_.Size(), 0); ASSERT_EQ(*manager.ReadSchema(/*schema_id=*/1).value(), *expected_schema); ASSERT_EQ(*(manager.Latest().value().value()), *expected_schema); } diff --git a/test/inte/scan_and_read_inte_test.cpp b/test/inte/scan_and_read_inte_test.cpp index 45be64970..a94d4b31b 100644 --- a/test/inte/scan_and_read_inte_test.cpp +++ b/test/inte/scan_and_read_inte_test.cpp @@ -60,6 +60,7 @@ #include "paimon/scan_context.h" #include "paimon/status.h" #include "paimon/table/source/plan.h" +#include "paimon/table/source/scan_metrics.h" #include "paimon/table/source/startup_mode.h" #include "paimon/table/source/table_read.h" #include "paimon/table/source/table_scan.h" @@ -483,38 +484,74 @@ TEST_P(ScanAndReadInteTest, TestWithAppendBucketKeyPointLookup) { ASSERT_OK_AND_ASSIGN(auto context, FinishScanContext(builder)); ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(context))); ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan()); + ASSERT_FALSE(plan->Splits().empty()); } - // The data was written with four buckets. Read it with rescaled table options too. - for (const std::string& current_bucket_count : {"2", "4", "8", "17"}) { - SCOPED_TRACE(current_bucket_count); - ScanContextBuilder scan_builder(table_path); - scan_builder.SetPredicate(predicate).AddOption(Options::BUCKET, current_bucket_count); - ASSERT_OK_AND_ASSIGN(auto scan_context, FinishScanContext(scan_builder)); - ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(scan_context))); - ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan()); - ASSERT_FALSE(plan->Splits().empty()); - for (const auto& split : plan->Splits()) { - auto data_split = std::dynamic_pointer_cast(split); - ASSERT_TRUE(data_split); - ASSERT_EQ(data_split->Bucket(), bucket_ids[32]); + int32_t other_index = -1; + for (int32_t i = 0; i < kNumKeys; ++i) { + if (bucket_ids[i] % 2 == bucket_ids[32] % 2 && bucket_ids[i] != bucket_ids[32]) { + other_index = i; + break; + } + } + ASSERT_GE(other_index, 0); + // These keys share a current two-bucket selection but need different historical buckets. + for (int32_t lookup_index : {32, other_index}) { + // The data was written with four buckets. Read it with rescaled table options too. + for (const std::string& current_bucket_count : {"2", "4", "8", "17"}) { + SCOPED_TRACE(current_bucket_count); + const std::string lookup_key = fmt::format("key{:03}", lookup_index); + auto lookup_predicate = PredicateBuilder::Equal( + 0, "rowkey", field_type, + Literal(field_type, lookup_key.data(), lookup_key.size())); + ScanContextBuilder scan_builder(table_path); + scan_builder.SetPredicate(lookup_predicate) + .AddOption(Options::BUCKET, current_bucket_count); + ASSERT_OK_AND_ASSIGN(auto scan_context, FinishScanContext(scan_builder)); + ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(scan_context))); + ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan()); + if (EnableSnapshotLiveManifestCache()) { + ASSERT_OK_AND_ASSIGN( + uint64_t enabled, + scan->GetMetrics()->GetCounter(ScanMetrics::LAST_SNAPSHOT_CACHE_ENABLED)); + ASSERT_EQ(enabled, 1); + ScanContextBuilder cached_builder(table_path); + cached_builder.SetPredicate(lookup_predicate) + .AddOption(Options::BUCKET, current_bucket_count); + ASSERT_OK_AND_ASSIGN(auto cached_context, FinishScanContext(cached_builder)); + ASSERT_OK_AND_ASSIGN(auto cached_scan, + TableScan::Create(std::move(cached_context))); + ASSERT_OK_AND_ASSIGN(plan, cached_scan->CreatePlan()); + ASSERT_OK_AND_ASSIGN(uint64_t hit, cached_scan->GetMetrics()->GetCounter( + ScanMetrics::LAST_SNAPSHOT_CACHE_HIT)); + ASSERT_EQ(hit, 1); + } + + ASSERT_FALSE(plan->Splits().empty()); + for (const auto& split : plan->Splits()) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + ASSERT_EQ(data_split->Bucket(), bucket_ids[lookup_index]); + } + ReadContextBuilder read_builder(table_path); + AddReadOptionsForPrefetch(&read_builder); + read_builder.SetPredicate(lookup_predicate).EnablePredicateFilter(true); + ASSERT_OK_AND_ASSIGN(auto read_context, read_builder.Finish()); + ASSERT_OK_AND_ASSIGN(auto read, TableRead::Create(std::move(read_context))); + ASSERT_OK_AND_ASSIGN(auto reader, read->CreateReader(plan->Splits())); + ASSERT_OK_AND_ASSIGN(auto result, + ReadResultCollector::CollectResult(std::move(reader))); + auto expected_fields = fields; + expected_fields.insert(expected_fields.begin(), + arrow::field("_VALUE_KIND", arrow::int8())); + auto expected_array = + arrow::ipc::internal::json::ArrayFromJSON( + arrow::struct_(expected_fields), + fmt::format(R"([[0, "{}", {}]])", lookup_key, lookup_index)) + .ValueOrDie(); + auto expected = std::make_shared(expected_array); + ASSERT_TRUE(expected->Equals(result)) << result->ToString(); } - ReadContextBuilder read_builder(table_path); - AddReadOptionsForPrefetch(&read_builder); - read_builder.SetPredicate(predicate).EnablePredicateFilter(true); - ASSERT_OK_AND_ASSIGN(auto read_context, read_builder.Finish()); - ASSERT_OK_AND_ASSIGN(auto read, TableRead::Create(std::move(read_context))); - ASSERT_OK_AND_ASSIGN(auto reader, read->CreateReader(plan->Splits())); - ASSERT_OK_AND_ASSIGN(auto result, - ReadResultCollector::CollectResult(std::move(reader))); - auto expected_fields = fields; - expected_fields.insert(expected_fields.begin(), - arrow::field("_VALUE_KIND", arrow::int8())); - auto expected_array = arrow::ipc::internal::json::ArrayFromJSON( - arrow::struct_(expected_fields), R"([[0, "key032", 32]])") - .ValueOrDie(); - auto expected = std::make_shared(expected_array); - ASSERT_TRUE(expected->Equals(result)) << result->ToString(); } } } @@ -3141,6 +3178,23 @@ TEST_P(ScanAndReadInteTest, TestWithPKBucketSelectByPredicate) { ASSERT_OK_AND_ASSIGN(auto table_read, TableRead::Create(std::move(read_context))); ASSERT_OK_AND_ASSIGN(auto result_plan, table_scan->CreatePlan()); + if (EnableSnapshotLiveManifestCache()) { + ASSERT_OK_AND_ASSIGN(uint64_t enabled, table_scan->GetMetrics()->GetCounter( + ScanMetrics::LAST_SNAPSHOT_CACHE_ENABLED)); + ASSERT_EQ(enabled, 1); + ScanContextBuilder cached_builder(table_path); + cached_builder.AddOption(Options::SCAN_SNAPSHOT_ID, "6") + .AddOption(Options::BUCKET, current_bucket_count) + .SetPartitionFilter({{{"f1", "10"}}}) + .SetPredicate(predicate); + ASSERT_OK_AND_ASSIGN(auto cached_context, FinishScanContext(cached_builder)); + ASSERT_OK_AND_ASSIGN(auto cached_scan, TableScan::Create(std::move(cached_context))); + ASSERT_OK_AND_ASSIGN(result_plan, cached_scan->CreatePlan()); + ASSERT_OK_AND_ASSIGN(uint64_t hit, cached_scan->GetMetrics()->GetCounter( + ScanMetrics::LAST_SNAPSHOT_CACHE_HIT)); + ASSERT_EQ(hit, 1); + } + ASSERT_EQ(result_plan->SnapshotId().value(), 6); // Verify all returned splits are from bucket 1 (f2=0 hashes to bucket 1) From 898ed7b72cb63afb3544cf69d8448bd9de3acd4e Mon Sep 17 00:00:00 2001 From: "wangyong.alen" Date: Sun, 6 Sep 2026 22:25:11 -0400 Subject: [PATCH 5/8] fix(test): avoid gcc range-loop-construct error in scan inte test --- test/inte/scan_and_read_inte_test.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/test/inte/scan_and_read_inte_test.cpp b/test/inte/scan_and_read_inte_test.cpp index a94d4b31b..ad1345405 100644 --- a/test/inte/scan_and_read_inte_test.cpp +++ b/test/inte/scan_and_read_inte_test.cpp @@ -498,7 +498,7 @@ TEST_P(ScanAndReadInteTest, TestWithAppendBucketKeyPointLookup) { // These keys share a current two-bucket selection but need different historical buckets. for (int32_t lookup_index : {32, other_index}) { // The data was written with four buckets. Read it with rescaled table options too. - for (const std::string& current_bucket_count : {"2", "4", "8", "17"}) { + for (const std::string current_bucket_count : {"2", "4", "8", "17"}) { SCOPED_TRACE(current_bucket_count); const std::string lookup_key = fmt::format("key{:03}", lookup_index); auto lookup_predicate = PredicateBuilder::Equal( @@ -3161,7 +3161,7 @@ TEST_P(ScanAndReadInteTest, TestWithPKBucketSelectByPredicate) { Literal(static_cast(0))); // Historical files use two buckets even when the current option has changed. - for (const std::string& current_bucket_count : {"2", "4", "8", "17"}) { + for (const std::string current_bucket_count : {"2", "4", "8", "17"}) { SCOPED_TRACE(current_bucket_count); ScanContextBuilder scan_context_builder(table_path); scan_context_builder.AddOption(Options::SCAN_SNAPSHOT_ID, "6"); From e6bb0bfbd7df0ded6a0be5e30c5000d4f31a3c9e Mon Sep 17 00:00:00 2001 From: "wangyong.alen" Date: Mon, 7 Sep 2026 22:21:53 -0400 Subject: [PATCH 6/8] fix(scan): keep inferred cache keys internal --- include/paimon/cache/cache.h | 10 ---------- src/paimon/common/io/cache/cache_key.cpp | 15 +++++++-------- src/paimon/common/io/cache/cache_key.h | 5 +++++ src/paimon/common/io/cache/lru_cache_test.cpp | 10 ++++++---- .../append_only_file_store_scan_test.cpp | 12 ++++-------- src/paimon/core/operation/file_store_scan.cpp | 3 ++- .../operation/key_value_file_store_scan_test.cpp | 2 +- 7 files changed, 25 insertions(+), 32 deletions(-) diff --git a/include/paimon/cache/cache.h b/include/paimon/cache/cache.h index ea1104a3e..5bff2a120 100644 --- a/include/paimon/cache/cache.h +++ b/include/paimon/cache/cache.h @@ -48,16 +48,6 @@ class PAIMON_EXPORT CacheKey { const std::string& branch, int32_t bucket); - /// Cache candidates for an inferred bucket, including files with other bucket counts or - /// schemas. - /// @param total_buckets Positive bucket count used to compute the inferred bucket. - /// @param schema_id Schema used to build the selector. - static std::shared_ptr ForSnapshotLiveManifestEntries(const std::string& table_path, - const std::string& branch, - int32_t bucket, - int32_t total_buckets, - int64_t schema_id); - public: virtual ~CacheKey() = default; diff --git a/src/paimon/common/io/cache/cache_key.cpp b/src/paimon/common/io/cache/cache_key.cpp index 615d77e00..1e3c61731 100644 --- a/src/paimon/common/io/cache/cache_key.cpp +++ b/src/paimon/common/io/cache/cache_key.cpp @@ -24,8 +24,7 @@ namespace { class SnapshotLiveManifestEntriesCacheKey : public CacheKey { public: SnapshotLiveManifestEntriesCacheKey(const std::string& table_path, const std::string& branch, - int32_t bucket, int32_t total_buckets = 0, - int64_t schema_id = 0) + int32_t bucket, int32_t total_buckets, int64_t schema_id) : CacheKey(CacheKind::SNAPSHOT_LIVE_MANIFEST), table_path_(table_path), branch_(branch), @@ -87,14 +86,14 @@ std::shared_ptr CacheKey::ForKind(const std::string& file_path, int64_ std::shared_ptr CacheKey::ForSnapshotLiveManifestEntries(const std::string& table_path, const std::string& branch, int32_t bucket) { - return std::make_shared(table_path, branch, bucket); + return std::make_shared(table_path, branch, bucket, + /*total_buckets=*/0, + /*schema_id=*/0); } -std::shared_ptr CacheKey::ForSnapshotLiveManifestEntries(const std::string& table_path, - const std::string& branch, - int32_t bucket, - int32_t total_buckets, - int64_t schema_id) { +std::shared_ptr CreateInferredSnapshotLiveManifestEntriesCacheKey( + const std::string& table_path, const std::string& branch, int32_t bucket, int32_t total_buckets, + int64_t schema_id) { return std::make_shared(table_path, branch, bucket, total_buckets, schema_id); } diff --git a/src/paimon/common/io/cache/cache_key.h b/src/paimon/common/io/cache/cache_key.h index 988735d11..796481792 100644 --- a/src/paimon/common/io/cache/cache_key.h +++ b/src/paimon/common/io/cache/cache_key.h @@ -26,6 +26,11 @@ namespace paimon { +// Cache inferred candidates separately from explicit buckets and other bucket layouts. +std::shared_ptr CreateInferredSnapshotLiveManifestEntriesCacheKey( + const std::string& table_path, const std::string& branch, int32_t bucket, int32_t total_buckets, + int64_t schema_id); + class PositionCacheKey : public CacheKey { public: PositionCacheKey(const std::string& file_path, int64_t position, int32_t length, bool is_index, diff --git a/src/paimon/common/io/cache/lru_cache_test.cpp b/src/paimon/common/io/cache/lru_cache_test.cpp index 6992cdb38..f8f6c9c7f 100644 --- a/src/paimon/common/io/cache/lru_cache_test.cpp +++ b/src/paimon/common/io/cache/lru_cache_test.cpp @@ -384,13 +384,15 @@ TEST_F(LruCacheTest, TestForKindSetsKeyKind) { } TEST_F(LruCacheTest, InferredManifestCacheKeysIncludeBucketCountAndSchema) { - auto key = CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1, 4, 0); - auto same = CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1, 4, 0); + auto key = CreateInferredSnapshotLiveManifestEntriesCacheKey("table", "main", 1, 4, 0); + auto same = CreateInferredSnapshotLiveManifestEntriesCacheKey("table", "main", 1, 4, 0); ASSERT_TRUE(key->Equals(*same)); ASSERT_EQ(key->HashCode(), same->HashCode()); ASSERT_FALSE(key->Equals(*CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1))); - ASSERT_FALSE(key->Equals(*CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1, 8, 0))); - ASSERT_FALSE(key->Equals(*CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1, 4, 1))); + ASSERT_FALSE( + key->Equals(*CreateInferredSnapshotLiveManifestEntriesCacheKey("table", "main", 1, 8, 0))); + ASSERT_FALSE( + key->Equals(*CreateInferredSnapshotLiveManifestEntriesCacheKey("table", "main", 1, 4, 1))); } TEST_F(LruCacheTest, TestForSnapshotLiveManifestEntries) { diff --git a/src/paimon/core/operation/append_only_file_store_scan_test.cpp b/src/paimon/core/operation/append_only_file_store_scan_test.cpp index 46cfcf1e3..0d0a37585 100644 --- a/src/paimon/core/operation/append_only_file_store_scan_test.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan_test.cpp @@ -20,7 +20,6 @@ #include #include -#include #include #include #include @@ -32,6 +31,7 @@ #include "paimon/common/data/binary_row.h" #include "paimon/common/data/binary_row_writer.h" #include "paimon/common/io/cache/lru_cache.h" +#include "paimon/common/utils/math.h" #include "paimon/core/bucket/default_bucket_function.h" #include "paimon/core/manifest/manifest_entry.h" #include "paimon/core/manifest/partition_entry.h" @@ -74,7 +74,7 @@ class AppendBucketPruningTest : public testing::Test { predicate, std::vector>(), bucket); return AppendOnlyFileStoreScan::Create(nullptr, schema_manager_, nullptr, nullptr, schema, arrow_schema, filters, options, - GetGlobalDefaultExecutor(), pool_); + CreateDefaultExecutor(), pool_); } void CheckBuckets(const std::shared_ptr& predicate, @@ -208,12 +208,8 @@ TEST_F(AppendBucketPruningTest, PreservesCrossScaleDecimalMatch) { TEST_F(AppendBucketPruningTest, PreservesDifferentNaNPayloadMatch) { rowkey_type_ = arrow::float64(); - uint64_t query_bits = 0x7ff8000000000000ULL; - uint64_t stored_bits = 0x7ff8000000000001ULL; - double query_value; - double stored_value; - std::memcpy(&query_value, &query_bits, sizeof(query_value)); - std::memcpy(&stored_value, &stored_bits, sizeof(stored_value)); + double query_value = FloatingPointFromBits(uint64_t{0x7ff8000000000000ULL}); + double stored_value = FloatingPointFromBits(uint64_t{0x7ff8000000000001ULL}); CheckMatchingValue(FieldType::DOUBLE, query_value, stored_value); } diff --git a/src/paimon/core/operation/file_store_scan.cpp b/src/paimon/core/operation/file_store_scan.cpp index 6bab2de65..29b1c0995 100644 --- a/src/paimon/core/operation/file_store_scan.cpp +++ b/src/paimon/core/operation/file_store_scan.cpp @@ -33,6 +33,7 @@ #include "paimon/common/data/binary_array.h" #include "paimon/common/data/blob_utils.h" #include "paimon/common/executor/future.h" +#include "paimon/common/io/cache/cache_key.h" #include "paimon/common/predicate/literal_converter.h" #include "paimon/common/types/data_field.h" #include "paimon/common/utils/field_type_utils.h" @@ -375,7 +376,7 @@ Status FileStoreScan::ReadManifestEntriesWithCache( std::shared_ptr FileStoreScan::SnapshotLiveManifestEntriesCacheKey(int32_t bucket) const { if (!bucket_filter_ && bucket_selector_) { - return CacheKey::ForSnapshotLiveManifestEntries( + return CreateInferredSnapshotLiveManifestEntriesCacheKey( table_path_, BranchManager::NormalizeBranch(core_options_.GetBranch()), bucket, core_options_.GetBucket(), table_schema_->Id()); } diff --git a/src/paimon/core/operation/key_value_file_store_scan_test.cpp b/src/paimon/core/operation/key_value_file_store_scan_test.cpp index 08867d4cb..ae53d8aa3 100644 --- a/src/paimon/core/operation/key_value_file_store_scan_test.cpp +++ b/src/paimon/core/operation/key_value_file_store_scan_test.cpp @@ -81,7 +81,7 @@ TEST(KeyValueBucketPruningTest, UsesEachEntriesBucketCount) { ASSERT_OK_AND_ASSIGN( auto scan, KeyValueFileStoreScan::Create(nullptr, manager, nullptr, nullptr, schema, arrow_schema, - filters, core_options, GetGlobalDefaultExecutor(), pool)); + filters, core_options, CreateDefaultExecutor(), pool)); SimpleStats stats = BinaryRowGenerator::GenerateStats({std::string("a")}, {std::string("z")}, {0}, pool.get()); ASSERT_OK_AND_ASSIGN( From a824461af268d3c04ed6589b9f3232aca67a5c85 Mon Sep 17 00:00:00 2001 From: "wangyong.alen" Date: Tue, 8 Sep 2026 00:01:48 -0400 Subject: [PATCH 7/8] fix(scan): distinguish explicit and inferred cache keys --- src/paimon/common/io/cache/cache_key.cpp | 66 +++--------------- src/paimon/common/io/cache/cache_key.h | 68 +++++++++++++++++-- src/paimon/common/io/cache/lru_cache_test.cpp | 16 +++-- .../append_only_file_store_scan_test.cpp | 4 +- src/paimon/core/operation/file_store_scan.cpp | 2 +- 5 files changed, 88 insertions(+), 68 deletions(-) diff --git a/src/paimon/common/io/cache/cache_key.cpp b/src/paimon/common/io/cache/cache_key.cpp index 1e3c61731..e6c44aa92 100644 --- a/src/paimon/common/io/cache/cache_key.cpp +++ b/src/paimon/common/io/cache/cache_key.cpp @@ -19,56 +19,6 @@ #include "paimon/common/io/cache/cache_key.h" namespace paimon { -namespace { - -class SnapshotLiveManifestEntriesCacheKey : public CacheKey { - public: - SnapshotLiveManifestEntriesCacheKey(const std::string& table_path, const std::string& branch, - int32_t bucket, int32_t total_buckets, int64_t schema_id) - : CacheKey(CacheKind::SNAPSHOT_LIVE_MANIFEST), - table_path_(table_path), - branch_(branch), - bucket_(bucket), - total_buckets_(total_buckets), - schema_id_(schema_id) {} - - bool IsIndex() const override { - return false; - } - - bool Equals(const CacheKey& other) const override { - const auto* rhs = dynamic_cast(&other); - if (!rhs) { - return false; - } - return table_path_ == rhs->table_path_ && branch_ == rhs->branch_ && - bucket_ == rhs->bucket_ && total_buckets_ == rhs->total_buckets_ && - schema_id_ == rhs->schema_id_ && GetKind() == rhs->GetKind(); - } - - size_t HashCode() const override { - size_t seed = 0; - seed ^= std::hash{}(table_path_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); - seed ^= std::hash{}(branch_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); - seed ^= std::hash{}(bucket_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); - seed ^= std::hash{}(static_cast(GetKind())) + HASH_CONSTANT + - (seed << 6) + (seed >> 2); - seed ^= std::hash{}(total_buckets_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); - seed ^= std::hash{}(schema_id_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); - return seed; - } - - private: - static constexpr uint64_t HASH_CONSTANT = 0x9e3779b97f4a7c15ULL; - - const std::string table_path_; - const std::string branch_; - const int32_t bucket_; - const int32_t total_buckets_; - const int64_t schema_id_; -}; - -} // namespace std::shared_ptr CacheKey::ForPosition(const std::string& file_path, int64_t position, int32_t length, bool is_index) { @@ -86,16 +36,20 @@ std::shared_ptr CacheKey::ForKind(const std::string& file_path, int64_ std::shared_ptr CacheKey::ForSnapshotLiveManifestEntries(const std::string& table_path, const std::string& branch, int32_t bucket) { - return std::make_shared(table_path, branch, bucket, - /*total_buckets=*/0, - /*schema_id=*/0); + return SnapshotLiveManifestEntriesCacheKey::ForExplicit(table_path, branch, bucket); +} + +std::shared_ptr SnapshotLiveManifestEntriesCacheKey::ForExplicit( + const std::string& table_path, const std::string& branch, int32_t bucket) { + return std::shared_ptr(new SnapshotLiveManifestEntriesCacheKey( + table_path, branch, bucket, Mode::kExplicit, std::nullopt, std::nullopt)); } -std::shared_ptr CreateInferredSnapshotLiveManifestEntriesCacheKey( +std::shared_ptr SnapshotLiveManifestEntriesCacheKey::ForInferred( const std::string& table_path, const std::string& branch, int32_t bucket, int32_t total_buckets, int64_t schema_id) { - return std::make_shared(table_path, branch, bucket, - total_buckets, schema_id); + return std::shared_ptr(new SnapshotLiveManifestEntriesCacheKey( + table_path, branch, bucket, Mode::kInferred, total_buckets, schema_id)); } bool PositionCacheKey::IsIndex() const { diff --git a/src/paimon/common/io/cache/cache_key.h b/src/paimon/common/io/cache/cache_key.h index 796481792..82fa9752c 100644 --- a/src/paimon/common/io/cache/cache_key.h +++ b/src/paimon/common/io/cache/cache_key.h @@ -20,16 +20,76 @@ #include #include +#include #include #include "paimon/cache/cache.h" namespace paimon { -// Cache inferred candidates separately from explicit buckets and other bucket layouts. -std::shared_ptr CreateInferredSnapshotLiveManifestEntriesCacheKey( - const std::string& table_path, const std::string& branch, int32_t bucket, int32_t total_buckets, - int64_t schema_id); +class SnapshotLiveManifestEntriesCacheKey : public CacheKey { + public: + static std::shared_ptr ForExplicit(const std::string& table_path, + const std::string& branch, int32_t bucket); + static std::shared_ptr ForInferred(const std::string& table_path, + const std::string& branch, int32_t bucket, + int32_t total_buckets, int64_t schema_id); + + bool IsIndex() const override { + return false; + } + + bool Equals(const CacheKey& other) const override { + const auto* rhs = dynamic_cast(&other); + if (!rhs) { + return false; + } + return table_path_ == rhs->table_path_ && branch_ == rhs->branch_ && + bucket_ == rhs->bucket_ && mode_ == rhs->mode_ && + total_buckets_ == rhs->total_buckets_ && schema_id_ == rhs->schema_id_ && + GetKind() == rhs->GetKind(); + } + + size_t HashCode() const override { + size_t seed = 0; + seed ^= std::hash{}(table_path_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); + seed ^= std::hash{}(branch_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); + seed ^= std::hash{}(bucket_) + HASH_CONSTANT + (seed << 6) + (seed >> 2); + seed ^= std::hash{}(static_cast(GetKind())) + HASH_CONSTANT + + (seed << 6) + (seed >> 2); + seed ^= std::hash>{}(total_buckets_) + HASH_CONSTANT + (seed << 6) + + (seed >> 2); + seed ^= std::hash>{}(schema_id_) + HASH_CONSTANT + (seed << 6) + + (seed >> 2); + seed ^= std::hash{}(static_cast(mode_)) + HASH_CONSTANT + (seed << 6) + + (seed >> 2); + return seed; + } + + private: + enum class Mode { kExplicit, kInferred }; + + SnapshotLiveManifestEntriesCacheKey(const std::string& table_path, const std::string& branch, + int32_t bucket, Mode mode, + std::optional total_buckets, + std::optional schema_id) + : CacheKey(CacheKind::SNAPSHOT_LIVE_MANIFEST), + table_path_(table_path), + branch_(branch), + bucket_(bucket), + mode_(mode), + total_buckets_(total_buckets), + schema_id_(schema_id) {} + + static constexpr uint64_t HASH_CONSTANT = 0x9e3779b97f4a7c15ULL; + + const std::string table_path_; + const std::string branch_; + const int32_t bucket_; + const Mode mode_; + const std::optional total_buckets_; + const std::optional schema_id_; +}; class PositionCacheKey : public CacheKey { public: diff --git a/src/paimon/common/io/cache/lru_cache_test.cpp b/src/paimon/common/io/cache/lru_cache_test.cpp index f8f6c9c7f..80133c4c3 100644 --- a/src/paimon/common/io/cache/lru_cache_test.cpp +++ b/src/paimon/common/io/cache/lru_cache_test.cpp @@ -384,15 +384,21 @@ TEST_F(LruCacheTest, TestForKindSetsKeyKind) { } TEST_F(LruCacheTest, InferredManifestCacheKeysIncludeBucketCountAndSchema) { - auto key = CreateInferredSnapshotLiveManifestEntriesCacheKey("table", "main", 1, 4, 0); - auto same = CreateInferredSnapshotLiveManifestEntriesCacheKey("table", "main", 1, 4, 0); + auto key = SnapshotLiveManifestEntriesCacheKey::ForInferred("table", "main", 1, 4, 0); + auto same = SnapshotLiveManifestEntriesCacheKey::ForInferred("table", "main", 1, 4, 0); ASSERT_TRUE(key->Equals(*same)); ASSERT_EQ(key->HashCode(), same->HashCode()); - ASSERT_FALSE(key->Equals(*CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1))); + auto explicit_key = SnapshotLiveManifestEntriesCacheKey::ForExplicit("table", "main", 1); + auto public_key = CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1); + ASSERT_TRUE(explicit_key->Equals(*public_key)); + ASSERT_EQ(explicit_key->HashCode(), public_key->HashCode()); + // Schema zero is real inferred metadata, not an explicit-mode placeholder. + ASSERT_FALSE(key->Equals(*explicit_key)); + ASSERT_FALSE(explicit_key->Equals(*key)); ASSERT_FALSE( - key->Equals(*CreateInferredSnapshotLiveManifestEntriesCacheKey("table", "main", 1, 8, 0))); + key->Equals(*SnapshotLiveManifestEntriesCacheKey::ForInferred("table", "main", 1, 8, 0))); ASSERT_FALSE( - key->Equals(*CreateInferredSnapshotLiveManifestEntriesCacheKey("table", "main", 1, 4, 1))); + key->Equals(*SnapshotLiveManifestEntriesCacheKey::ForInferred("table", "main", 1, 4, 1))); } TEST_F(LruCacheTest, TestForSnapshotLiveManifestEntries) { diff --git a/src/paimon/core/operation/append_only_file_store_scan_test.cpp b/src/paimon/core/operation/append_only_file_store_scan_test.cpp index 0d0a37585..b45a56846 100644 --- a/src/paimon/core/operation/append_only_file_store_scan_test.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan_test.cpp @@ -208,8 +208,8 @@ TEST_F(AppendBucketPruningTest, PreservesCrossScaleDecimalMatch) { TEST_F(AppendBucketPruningTest, PreservesDifferentNaNPayloadMatch) { rowkey_type_ = arrow::float64(); - double query_value = FloatingPointFromBits(uint64_t{0x7ff8000000000000ULL}); - double stored_value = FloatingPointFromBits(uint64_t{0x7ff8000000000001ULL}); + auto query_value = FloatingPointFromBits(uint64_t{0x7ff8000000000000ULL}); + auto stored_value = FloatingPointFromBits(uint64_t{0x7ff8000000000001ULL}); CheckMatchingValue(FieldType::DOUBLE, query_value, stored_value); } diff --git a/src/paimon/core/operation/file_store_scan.cpp b/src/paimon/core/operation/file_store_scan.cpp index 29b1c0995..ba0c3fef6 100644 --- a/src/paimon/core/operation/file_store_scan.cpp +++ b/src/paimon/core/operation/file_store_scan.cpp @@ -376,7 +376,7 @@ Status FileStoreScan::ReadManifestEntriesWithCache( std::shared_ptr FileStoreScan::SnapshotLiveManifestEntriesCacheKey(int32_t bucket) const { if (!bucket_filter_ && bucket_selector_) { - return CreateInferredSnapshotLiveManifestEntriesCacheKey( + return paimon::SnapshotLiveManifestEntriesCacheKey::ForInferred( table_path_, BranchManager::NormalizeBranch(core_options_.GetBranch()), bucket, core_options_.GetBucket(), table_schema_->Id()); } From e918bd482aacdc7f55e50aedf333f05297fadf0c Mon Sep 17 00:00:00 2001 From: "wangyong.alen" Date: Tue, 8 Sep 2026 03:20:13 -0400 Subject: [PATCH 8/8] refactor(cache): remove public manifest key factory --- include/paimon/cache/cache.h | 3 --- src/paimon/common/io/cache/cache_key.cpp | 6 ----- src/paimon/common/io/cache/lru_cache_test.cpp | 22 +++++++++---------- src/paimon/core/operation/file_store_scan.cpp | 13 ++++++----- src/paimon/core/operation/file_store_scan.h | 2 +- 5 files changed, 19 insertions(+), 27 deletions(-) diff --git a/include/paimon/cache/cache.h b/include/paimon/cache/cache.h index 5bff2a120..b282da74b 100644 --- a/include/paimon/cache/cache.h +++ b/include/paimon/cache/cache.h @@ -44,9 +44,6 @@ class PAIMON_EXPORT CacheKey { int32_t length, bool is_index); static std::shared_ptr ForKind(const std::string& file_path, int64_t position, int32_t length, CacheKind kind); - static std::shared_ptr ForSnapshotLiveManifestEntries(const std::string& table_path, - const std::string& branch, - int32_t bucket); public: virtual ~CacheKey() = default; diff --git a/src/paimon/common/io/cache/cache_key.cpp b/src/paimon/common/io/cache/cache_key.cpp index e6c44aa92..998d387f4 100644 --- a/src/paimon/common/io/cache/cache_key.cpp +++ b/src/paimon/common/io/cache/cache_key.cpp @@ -33,12 +33,6 @@ std::shared_ptr CacheKey::ForKind(const std::string& file_path, int64_ return key; } -std::shared_ptr CacheKey::ForSnapshotLiveManifestEntries(const std::string& table_path, - const std::string& branch, - int32_t bucket) { - return SnapshotLiveManifestEntriesCacheKey::ForExplicit(table_path, branch, bucket); -} - std::shared_ptr SnapshotLiveManifestEntriesCacheKey::ForExplicit( const std::string& table_path, const std::string& branch, int32_t bucket) { return std::shared_ptr(new SnapshotLiveManifestEntriesCacheKey( diff --git a/src/paimon/common/io/cache/lru_cache_test.cpp b/src/paimon/common/io/cache/lru_cache_test.cpp index 80133c4c3..29eb10fef 100644 --- a/src/paimon/common/io/cache/lru_cache_test.cpp +++ b/src/paimon/common/io/cache/lru_cache_test.cpp @@ -389,9 +389,6 @@ TEST_F(LruCacheTest, InferredManifestCacheKeysIncludeBucketCountAndSchema) { ASSERT_TRUE(key->Equals(*same)); ASSERT_EQ(key->HashCode(), same->HashCode()); auto explicit_key = SnapshotLiveManifestEntriesCacheKey::ForExplicit("table", "main", 1); - auto public_key = CacheKey::ForSnapshotLiveManifestEntries("table", "main", 1); - ASSERT_TRUE(explicit_key->Equals(*public_key)); - ASSERT_EQ(explicit_key->HashCode(), public_key->HashCode()); // Schema zero is real inferred metadata, not an explicit-mode placeholder. ASSERT_FALSE(key->Equals(*explicit_key)); ASSERT_FALSE(explicit_key->Equals(*key)); @@ -401,14 +398,17 @@ TEST_F(LruCacheTest, InferredManifestCacheKeysIncludeBucketCountAndSchema) { key->Equals(*SnapshotLiveManifestEntriesCacheKey::ForInferred("table", "main", 1, 4, 1))); } -TEST_F(LruCacheTest, TestForSnapshotLiveManifestEntries) { - auto main_key = CacheKey::ForSnapshotLiveManifestEntries("table_path", "main", 0); - auto same_key = CacheKey::ForSnapshotLiveManifestEntries("table_path", "main", 0); - auto branch_key = CacheKey::ForSnapshotLiveManifestEntries("table_path", "dev", 0); - auto table_key = CacheKey::ForSnapshotLiveManifestEntries("other_table_path", "main", 0); - auto bucket_key = CacheKey::ForSnapshotLiveManifestEntries("table_path", "main", 1); - auto hash_in_path_key = CacheKey::ForSnapshotLiveManifestEntries("table#path", "main", 0); - auto hash_in_branch_key = CacheKey::ForSnapshotLiveManifestEntries("table", "path#main", 0); +TEST_F(LruCacheTest, TestExplicitSnapshotLiveManifestEntriesKeys) { + auto main_key = SnapshotLiveManifestEntriesCacheKey::ForExplicit("table_path", "main", 0); + auto same_key = SnapshotLiveManifestEntriesCacheKey::ForExplicit("table_path", "main", 0); + auto branch_key = SnapshotLiveManifestEntriesCacheKey::ForExplicit("table_path", "dev", 0); + auto table_key = + SnapshotLiveManifestEntriesCacheKey::ForExplicit("other_table_path", "main", 0); + auto bucket_key = SnapshotLiveManifestEntriesCacheKey::ForExplicit("table_path", "main", 1); + auto hash_in_path_key = + SnapshotLiveManifestEntriesCacheKey::ForExplicit("table#path", "main", 0); + auto hash_in_branch_key = + SnapshotLiveManifestEntriesCacheKey::ForExplicit("table", "path#main", 0); ASSERT_EQ(CacheKind::SNAPSHOT_LIVE_MANIFEST, main_key->GetKind()); ASSERT_TRUE(CacheKeyEqual()(main_key, same_key)); diff --git a/src/paimon/core/operation/file_store_scan.cpp b/src/paimon/core/operation/file_store_scan.cpp index ba0c3fef6..2b182107c 100644 --- a/src/paimon/core/operation/file_store_scan.cpp +++ b/src/paimon/core/operation/file_store_scan.cpp @@ -374,13 +374,14 @@ Status FileStoreScan::ReadManifestEntriesWithCache( return Status::OK(); } -std::shared_ptr FileStoreScan::SnapshotLiveManifestEntriesCacheKey(int32_t bucket) const { +std::shared_ptr FileStoreScan::CreateSnapshotLiveManifestEntriesCacheKey( + int32_t bucket) const { if (!bucket_filter_ && bucket_selector_) { - return paimon::SnapshotLiveManifestEntriesCacheKey::ForInferred( + return SnapshotLiveManifestEntriesCacheKey::ForInferred( table_path_, BranchManager::NormalizeBranch(core_options_.GetBranch()), bucket, core_options_.GetBucket(), table_schema_->Id()); } - return CacheKey::ForSnapshotLiveManifestEntries( + return SnapshotLiveManifestEntriesCacheKey::ForExplicit( table_path_, BranchManager::NormalizeBranch(core_options_.GetBranch()), bucket); } @@ -389,7 +390,7 @@ Result FileStoreScan::LoadSnapshotLiveManifestEntri auto supplier = [](const std::shared_ptr&) -> Result> { return std::shared_ptr(); }; - std::shared_ptr cache_key = SnapshotLiveManifestEntriesCacheKey(bucket); + std::shared_ptr cache_key = CreateSnapshotLiveManifestEntriesCacheKey(bucket); const auto max_snapshots = core_options_.GetScanManifestEntryCacheMaxSnapshots(); Result> cache_result = core_options_.GetCache()->Get(cache_key, supplier); @@ -412,8 +413,8 @@ Status FileStoreScan::StoreSnapshotLiveManifestEntries( } auto cache_value = std::make_shared(MemorySegment::Wrap(bytes_result.value()), CacheCallback()); - Status status = - core_options_.GetCache()->Put(SnapshotLiveManifestEntriesCacheKey(bucket), cache_value); + Status status = core_options_.GetCache()->Put(CreateSnapshotLiveManifestEntriesCacheKey(bucket), + cache_value); return status.ok() ? status : Status::OK(); } diff --git a/src/paimon/core/operation/file_store_scan.h b/src/paimon/core/operation/file_store_scan.h index 6cc1e4e37..83c18d499 100644 --- a/src/paimon/core/operation/file_store_scan.h +++ b/src/paimon/core/operation/file_store_scan.h @@ -265,7 +265,7 @@ class FileStoreScan { int32_t bucket, std::vector* manifest_entries, bool* cache_hit) const; - std::shared_ptr SnapshotLiveManifestEntriesCacheKey(int32_t bucket) const; + std::shared_ptr CreateSnapshotLiveManifestEntriesCacheKey(int32_t bucket) const; Result LoadSnapshotLiveManifestEntries(int32_t bucket) const; Status StoreSnapshotLiveManifestEntries(int32_t bucket, const SnapshotLiveManifestEntries& entries) const;