diff --git a/src/paimon/core/manifest/manifest_file.cpp b/src/paimon/core/manifest/manifest_file.cpp index ecbcb6bd..d525f62a 100644 --- a/src/paimon/core/manifest/manifest_file.cpp +++ b/src/paimon/core/manifest/manifest_file.cpp @@ -26,7 +26,10 @@ #include "arrow/c/abi.h" #include "arrow/c/bridge.h" #include "paimon/common/data/columnar/columnar_row.h" +#include "paimon/common/predicate/predicate_validator.h" +#include "paimon/common/reader/late_materializing_file_batch_reader.h" #include "paimon/common/utils/arrow/status_utils.h" +#include "paimon/common/utils/checked_cast.h" #include "paimon/core/io/rolling_file_writer.h" #include "paimon/core/manifest/manifest_entry.h" #include "paimon/core/manifest/manifest_entry_serializer.h" @@ -39,6 +42,7 @@ #include "paimon/format/file_format.h" #include "paimon/format/reader_builder.h" #include "paimon/format/writer_builder.h" +#include "paimon/predicate/predicate_builder.h" #include "paimon/status.h" namespace arrow { @@ -49,6 +53,12 @@ class Schema; namespace paimon { class MemoryPool; +namespace { +constexpr int32_t kVersionFieldIndex = 0; +constexpr int32_t kBucketFieldIndex = 3; +constexpr int32_t kTotalBucketsFieldIndex = 4; +} // namespace + ManifestFile::ManifestFile(const std::shared_ptr& file_system, const std::shared_ptr& reader_builder, const std::shared_ptr& writer_builder, @@ -91,17 +101,26 @@ Result> ManifestFile::Create( } Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t bucket, + const std::optional& expected_total_buckets, std::optional file_size, std::vector* entries) const { + // Readers without precise bitmap selection still filter aligned Arrow columns + // before constructing ManifestEntry and DataFileMeta objects. return ReadArrowBatches( - file_name, - [this, bucket, entries](const std::shared_ptr& batch) -> Status { - const arrow::ArrayVector& fields = batch->fields(); - ColumnarRow row(fields, pool_, /*row_id=*/0); - for (int64_t i = 0; i < batch->length(); i++) { + file_name, file_size, + [this, bucket, expected_total_buckets, + entries](const std::shared_ptr& batch) -> Status { + ColumnarRow row(batch->fields(), pool_, /*row_id=*/0); + for (int64_t i = 0; i < batch->length(); ++i) { row.SetRowId(i); - PAIMON_RETURN_NOT_OK(ManifestEntrySerializer::ValidateVersion(row.GetInt(0))); - if (ManifestEntrySerializer::GetBucket(row) != bucket) { + PAIMON_RETURN_NOT_OK( + ManifestEntrySerializer::ValidateVersion(row.GetInt(kVersionFieldIndex))); + // Different or unknown bucket counts must reach per-entry bucket filtering. + const bool historical_layout = + expected_total_buckets && + (row.IsNullAt(kBucketFieldIndex) || row.IsNullAt(kTotalBucketsFieldIndex) || + row.GetInt(kTotalBucketsFieldIndex) != expected_total_buckets.value()); + if (!historical_layout && ManifestEntrySerializer::GetBucket(row) != bucket) { continue; } PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry, serializer_->FromRow(row)); @@ -109,7 +128,60 @@ Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t buc } return Status::OK(); }, - file_size); + [this, bucket, expected_total_buckets](std::unique_ptr* reader) { + return PrepareBucketRead(bucket, expected_total_buckets, reader); + }); +} + +Status ManifestFile::PrepareBucketRead(int32_t bucket, + const std::optional& expected_total_buckets, + std::unique_ptr* reader) const { + if (!(*reader)->SupportPreciseBitmapSelection()) { + return Status::OK(); + } + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr c_schema, (*reader)->GetFileSchema()); + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr file_schema, + arrow::ImportSchema(c_schema.get())); + const auto& target_type = serializer_->GetDataType(); + const std::string& bucket_name = target_type->field(kBucketFieldIndex)->name(); + std::shared_ptr selector = + PredicateBuilder::Equal(kBucketFieldIndex, bucket_name, FieldType::INT, Literal(bucket)); + if (expected_total_buckets) { + const std::string& total_name = target_type->field(kTotalBucketsFieldIndex)->name(); + PAIMON_ASSIGN_OR_RAISE( + selector, + PredicateBuilder::Or( + {selector, + PredicateBuilder::IsNull(kTotalBucketsFieldIndex, total_name, FieldType::INT), + PredicateBuilder::NotEqual(kTotalBucketsFieldIndex, total_name, FieldType::INT, + Literal(expected_total_buckets.value()))})); + } + // Retain unsupported versions regardless of bucket so the consumer validates every version + // before bucket filtering, including when the probe would otherwise select no entries. + const std::string& version_name = target_type->field(kVersionFieldIndex)->name(); + PAIMON_ASSIGN_OR_RAISE( + selector, + PredicateBuilder::Or( + {selector, PredicateBuilder::IsNull(kVersionFieldIndex, version_name, FieldType::INT), + PredicateBuilder::NotEqual( + kVersionFieldIndex, version_name, FieldType::INT, + Literal( + checked_cast(serializer_.get())->GetVersion()))})); + if (!PredicateValidator::ValidatePredicateWithSchema(*file_schema, selector, + /*validate_field_idx=*/true) + .ok()) { + // An incompatible projection must fall back to ordinary schema evolution. + return Status::OK(); + } + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr selective_reader, + LateMaterializingFileBatchReader::Create(std::move(*reader), arrow_pool_)); + // Keep the on-disk schema; ManifestMetaReader still performs schema evolution afterwards. + ArrowSchema full_c_schema; + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*file_schema, &full_c_schema)); + PAIMON_RETURN_NOT_OK(selective_reader->SetReadSchema(&full_c_schema, selector, std::nullopt)); + *reader = std::move(selective_reader); + return Status::OK(); } Result> ManifestFile::Write( diff --git a/src/paimon/core/manifest/manifest_file.h b/src/paimon/core/manifest/manifest_file.h index 2f3a60cb..3840368f 100644 --- a/src/paimon/core/manifest/manifest_file.h +++ b/src/paimon/core/manifest/manifest_file.h @@ -63,16 +63,23 @@ class ManifestFile : public ObjectsFile { /// @note This method is atomic. Result> Write(const std::vector& entries); - /// Read a manifest file and deserialize only entries for the specified bucket. + /// Read entries for a bucket. An inferred total bucket count also retains entries with + /// different or unknown bucket counts for per-entry bucket filtering. + /// Bucket pruning assumes stable bucket-key hashing across schema versions. + /// Serialization versions are validated for all entries before bucket filtering. /// /// @param file_size Length of the manifest when the caller already has it from the manifest /// list, which saves the read a metadata request on a remote store. Pass /// std::nullopt when the length is not known. Status ReadBucketEntries(const std::string& file_name, int32_t bucket, + const std::optional& expected_total_buckets, std::optional file_size, std::vector* entries) const; private: + Status PrepareBucketRead(int32_t bucket, const std::optional& expected_total_buckets, + std::unique_ptr* reader) const; + ManifestFile(const std::shared_ptr& file_system, const std::shared_ptr& reader_builder, const std::shared_ptr& writer_builder, diff --git a/src/paimon/core/manifest/manifest_file_test.cpp b/src/paimon/core/manifest/manifest_file_test.cpp index 0ce9cc72..df5bcbd8 100644 --- a/src/paimon/core/manifest/manifest_file_test.cpp +++ b/src/paimon/core/manifest/manifest_file_test.cpp @@ -25,13 +25,18 @@ #include #include "arrow/api.h" +#include "arrow/ipc/json_simple.h" +#include "fmt/format.h" #include "gtest/gtest.h" #include "paimon/common/data/binary_row.h" #include "paimon/common/data/data_define.h" +#include "paimon/common/utils/checked_cast.h" #include "paimon/core/io/data_file_meta.h" +#include "paimon/core/io/meta_to_arrow_array_converter.h" #include "paimon/core/manifest/file_kind.h" #include "paimon/core/manifest/file_source.h" #include "paimon/core/manifest/manifest_entry.h" +#include "paimon/core/manifest/manifest_entry_serializer.h" #include "paimon/core/manifest/manifest_file_meta.h" #include "paimon/core/stats/simple_stats.h" #include "paimon/core/utils/file_store_path_factory.h" @@ -40,6 +45,8 @@ #include "paimon/defs.h" #include "paimon/format/file_format.h" #include "paimon/format/file_format_factory.h" +#include "paimon/format/format_writer.h" +#include "paimon/format/writer_builder.h" #include "paimon/fs/local/local_file_system.h" #include "paimon/memory/memory_pool.h" #include "paimon/testing/utils/binary_row_generator.h" @@ -112,32 +119,51 @@ class ManifestFileTest : public testing::Test { const std::string& file_format_str, const std::string& root_path, const std::string& file_name, const std::shared_ptr& pool, const std::optional& bucket = std::nullopt) const { + EXPECT_OK_AND_ASSIGN( + std::vector entries, + TryReadManifestEntry(file_format_str, root_path, file_name, pool, bucket, + /*cache_enabled=*/false, /*inferred_bucket=*/false)); + return entries; + } + + Result> TryReadManifestEntry( + const std::string& file_format_str, const std::string& root_path, + const std::string& file_name, const std::shared_ptr& pool, + const std::optional& bucket, bool cache_enabled, bool inferred_bucket) const { std::shared_ptr file_system = std::make_shared(); - EXPECT_OK_AND_ASSIGN(std::shared_ptr file_format, - FileFormatFactory::Get(file_format_str, {})); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr file_format, + FileFormatFactory::Get(file_format_str, {})); auto unused_schema = arrow::schema(arrow::FieldVector({arrow::field("f0", arrow::utf8())})); - EXPECT_OK_AND_ASSIGN(std::shared_ptr path_factory, - FileStorePathFactory::Create( - root_path, unused_schema, /*partition_keys=*/{}, - /*default_part_value=*/"", file_format->Identifier(), - /*data_file_prefix=*/"data-", - /*legacy_partition_name_enabled=*/true, /*external_paths=*/{}, - /*global_index_external_path=*/std::nullopt, - /*index_file_in_data_file_dir=*/false, pool)); - EXPECT_OK_AND_ASSIGN(CoreOptions options, - CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"}})); - EXPECT_OK_AND_ASSIGN( + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr path_factory, + FileStorePathFactory::Create( + root_path, unused_schema, /*partition_keys=*/{}, + /*default_part_value=*/"", file_format->Identifier(), + /*data_file_prefix=*/"data-", + /*legacy_partition_name_enabled=*/true, /*external_paths=*/{}, + /*global_index_external_path=*/std::nullopt, + /*index_file_in_data_file_dir=*/false, pool)); + PAIMON_ASSIGN_OR_RAISE(CoreOptions options, + CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"}})); + if (cache_enabled) { + options.WithCache( + std::make_shared(CacheKind::MANIFEST, 64 * 1024 * 1024)); + } + PAIMON_ASSIGN_OR_RAISE( std::unique_ptr manifest_file, ManifestFile::Create(file_system, file_format, "zstd", path_factory, /*target_file_size=*/1024, pool, options, unused_schema)); std::vector manifest_entries; - if (bucket) { - EXPECT_OK(manifest_file->ReadBucketEntries(file_name, bucket.value(), - /*file_size=*/std::nullopt, - &manifest_entries)); + if (bucket && inferred_bucket) { + PAIMON_RETURN_NOT_OK(manifest_file->ReadBucketEntries( + file_name, bucket.value(), /*expected_total_buckets=*/2, /*file_size=*/std::nullopt, + &manifest_entries)); + } else if (bucket) { + PAIMON_RETURN_NOT_OK( + manifest_file->ReadBucketEntries(file_name, bucket.value(), std::nullopt, + /*file_size=*/std::nullopt, &manifest_entries)); } else { - EXPECT_OK(manifest_file->Read(file_name, /*filter=*/nullptr, - /*file_size=*/std::nullopt, &manifest_entries)); + PAIMON_RETURN_NOT_OK(manifest_file->Read( + file_name, /*filter=*/nullptr, /*file_size=*/std::nullopt, &manifest_entries)); } return manifest_entries; @@ -360,17 +386,17 @@ TEST_F(ManifestFileTest, TestReadBucketEntriesMaterializesOnlySelectedBucket) { ASSERT_EQ(2, all_entries.size()); std::vector bucket_one_entries; - ASSERT_OK(manifest_file->ReadBucketEntries(manifest_name, /*bucket=*/1, + ASSERT_OK(manifest_file->ReadBucketEntries(manifest_name, /*bucket=*/1, std::nullopt, /*file_size=*/std::nullopt, &bucket_one_entries)); ASSERT_EQ(std::vector({all_entries[0]}), bucket_one_entries); std::vector bucket_zero_entries; - ASSERT_OK(manifest_file->ReadBucketEntries(manifest_name, /*bucket=*/0, + ASSERT_OK(manifest_file->ReadBucketEntries(manifest_name, /*bucket=*/0, std::nullopt, /*file_size=*/std::nullopt, &bucket_zero_entries)); ASSERT_EQ(std::vector({all_entries[1]}), bucket_zero_entries); std::vector missing_bucket_entries; - ASSERT_OK(manifest_file->ReadBucketEntries(manifest_name, /*bucket=*/2, + ASSERT_OK(manifest_file->ReadBucketEntries(manifest_name, /*bucket=*/2, std::nullopt, /*file_size=*/std::nullopt, &missing_bucket_entries)); ASSERT_TRUE(missing_bucket_entries.empty()); @@ -417,14 +443,87 @@ TEST_F(ManifestFileTest, TestReadPassesKnownSizeToOpen) { ASSERT_EQ(0, counting_file_system->open_count); std::vector bucket_one_entries; - ASSERT_OK(manifest_file->ReadBucketEntries(manifest_name, /*bucket=*/1, kRecordedSize, - &bucket_one_entries)); + ASSERT_OK(manifest_file->ReadBucketEntries(manifest_name, /*bucket=*/1, std::nullopt, + kRecordedSize, &bucket_one_entries)); ASSERT_EQ(std::vector({all_entries[0]}), bucket_one_entries); ASSERT_EQ(std::vector({kRecordedSize, kRecordedSize}), counting_file_system->opened_lengths); ASSERT_EQ(0, counting_file_system->open_count); } +TEST_F(ManifestFileTest, TestInferredBucketProbeSkipsArrowMaterialization) { + auto pool = GetDefaultPool(); + std::vector source_entries = + ReadManifestEntry("orc", paimon::test::GetDataDir() + "/orc/append_09.db/append_09", + "manifest-3a44a0da-1008-463c-914e-28d271375e24-0", pool); + ASSERT_EQ(2, source_entries.size()); + + for (int32_t cache_mode : {0, 1, 2}) { + SCOPED_TRACE(cache_mode); + auto test_dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(test_dir); + auto file_system = std::make_shared(); + ASSERT_OK(file_system->Mkdirs(FileStorePathFactory::ManifestPath(test_dir->Str()))); + ASSERT_OK_AND_ASSIGN(std::shared_ptr file_format, + FileFormatFactory::Get("avro", {})); + auto unused_schema = arrow::schema(arrow::FieldVector({arrow::field("f0", arrow::utf8())})); + ASSERT_OK_AND_ASSIGN(std::shared_ptr path_factory, + FileStorePathFactory::Create( + test_dir->Str(), unused_schema, /*partition_keys=*/{}, + /*default_part_value=*/"", file_format->Identifier(), + /*data_file_prefix=*/"data-", + /*legacy_partition_name_enabled=*/true, /*external_paths=*/{}, + /*global_index_external_path=*/std::nullopt, + /*index_file_in_data_file_dir=*/false, pool)); + ASSERT_OK_AND_ASSIGN(CoreOptions options, + CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"}})); + auto cache = std::make_shared( + cache_mode == 1 ? CacheKind::MANIFEST : CacheKind::DEFAULT, 64 * 1024 * 1024); + if (cache_mode != 0) { + options.WithCache(cache); + } + ASSERT_OK_AND_ASSIGN( + std::unique_ptr manifest_file, + ManifestFile::Create(file_system, file_format, "null", path_factory, + /*target_file_size=*/1024, pool, options, unused_schema)); + + constexpr size_t payload_size = 8 * 1024 * 1024; + auto large_file = std::make_shared(*source_entries[0].File()); + large_file->file_name.assign(payload_size, 'x'); + // A schema change alone must not force unrelated buckets to be materialized. + large_file->schema_id = source_entries[1].File()->schema_id + 1; + ManifestEntry excluded(FileKind::Add(), source_entries[0].Partition(), 1, 2, large_file); + ManifestEntry selected(FileKind::Add(), source_entries[1].Partition(), 0, 2, + source_entries[1].File()); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, + manifest_file->WriteWithoutRolling({excluded, selected})); + std::vector warm; + ASSERT_OK(manifest_file->Read(written.first, nullptr, /*file_size=*/std::nullopt, &warm)); + // Cached bytes, if available, belong to a separate pool from reader allocations. + for (bool inferred : {false, true}) { + std::shared_ptr read_pool = GetMemoryPool(); + ASSERT_OK_AND_ASSIGN( + auto reader, ManifestFile::Create(file_system, file_format, "null", path_factory, + 1024, read_pool, options, unused_schema)); + std::vector entries; + if (inferred) { + ASSERT_OK(reader->ReadBucketEntries(written.first, 0, /*expected_total_buckets=*/2, + /*file_size=*/std::nullopt, &entries)); + ASSERT_EQ(std::vector({selected}), entries); + ASSERT_LT(read_pool->MaxMemoryUsage(), payload_size / 2); + } else { + ASSERT_OK( + reader->Read(written.first, nullptr, /*file_size=*/std::nullopt, &entries)); + ASSERT_EQ(warm, entries); + ASSERT_GE(read_pool->MaxMemoryUsage(), payload_size); + } + } + ASSERT_EQ(cache_mode == 1 ? 1 : 3, file_system->open_count); + ASSERT_EQ(cache_mode == 1 ? 1 : 0, cache->SupplierCallCount()); + } +} + TEST_F(ManifestFileTest, TestReadBucketEntriesSkipsDeserializingOtherBuckets) { auto pool = GetDefaultPool(); std::vector source_entries = @@ -432,47 +531,194 @@ TEST_F(ManifestFileTest, TestReadBucketEntriesSkipsDeserializingOtherBuckets) { "manifest-3a44a0da-1008-463c-914e-28d271375e24-0", pool); ASSERT_EQ(2, source_entries.size()); - auto test_dir = UniqueTestDirectory::Create(); - ASSERT_TRUE(test_dir); - std::shared_ptr file_system = test_dir->GetFileSystem(); - ASSERT_OK(file_system->Mkdirs(FileStorePathFactory::ManifestPath(test_dir->Str()))); - ASSERT_OK_AND_ASSIGN(std::shared_ptr file_format, - FileFormatFactory::Get("avro", {})); - auto unused_schema = arrow::schema(arrow::FieldVector({arrow::field("f0", arrow::utf8())})); - ASSERT_OK_AND_ASSIGN( - std::shared_ptr path_factory, - FileStorePathFactory::Create(test_dir->Str(), unused_schema, /*partition_keys=*/{}, - /*default_part_value=*/"", file_format->Identifier(), - /*data_file_prefix=*/"data-", - /*legacy_partition_name_enabled=*/true, /*external_paths=*/{}, - /*global_index_external_path=*/std::nullopt, - /*index_file_in_data_file_dir=*/false, pool)); - ASSERT_OK_AND_ASSIGN(CoreOptions options, - CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"}})); - ASSERT_OK_AND_ASSIGN( - std::unique_ptr manifest_file, - ManifestFile::Create(file_system, file_format, "zstd", path_factory, - /*target_file_size=*/1024, pool, options, unused_schema)); + // Selective decoding also works without a cache or when the cache rejects manifest keys. + for (int32_t cache_mode : {0, 1, 2}) { + SCOPED_TRACE(cache_mode); + auto test_dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(test_dir); + auto file_system = std::make_shared(); + ASSERT_OK(file_system->Mkdirs(FileStorePathFactory::ManifestPath(test_dir->Str()))); + ASSERT_OK_AND_ASSIGN(std::shared_ptr file_format, + FileFormatFactory::Get("avro", {})); + auto unused_schema = arrow::schema(arrow::FieldVector({arrow::field("f0", arrow::utf8())})); + ASSERT_OK_AND_ASSIGN(std::shared_ptr path_factory, + FileStorePathFactory::Create( + test_dir->Str(), unused_schema, /*partition_keys=*/{}, + /*default_part_value=*/"", file_format->Identifier(), + /*data_file_prefix=*/"data-", + /*legacy_partition_name_enabled=*/true, /*external_paths=*/{}, + /*global_index_external_path=*/std::nullopt, + /*index_file_in_data_file_dir=*/false, pool)); + ASSERT_OK_AND_ASSIGN(CoreOptions options, + CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"}})); + if (cache_mode != 0) { + options.WithCache(std::make_shared( + cache_mode == 1 ? CacheKind::MANIFEST : CacheKind::DEFAULT, 64 * 1024 * 1024)); + } + ASSERT_OK_AND_ASSIGN( + std::unique_ptr manifest_file, + ManifestFile::Create(file_system, file_format, "null", path_factory, + /*target_file_size=*/1024, pool, options, unused_schema)); - ManifestEntry invalid_other_bucket(FileKind(static_cast(2)), - source_entries[0].Partition(), /*bucket=*/1, - /*total_buckets=*/2, source_entries[0].File()); - ManifestEntry valid_target_bucket(FileKind::Add(), source_entries[1].Partition(), /*bucket=*/0, - /*total_buckets=*/2, source_entries[1].File()); - using WrittenFile = std::pair; - ASSERT_OK_AND_ASSIGN( - WrittenFile written_file, - manifest_file->WriteWithoutRolling({invalid_other_bucket, valid_target_bucket})); + ManifestEntry invalid_other_bucket(FileKind(static_cast(2)), + source_entries[0].Partition(), /*bucket=*/1, + /*total_buckets=*/2, source_entries[0].File()); + ManifestEntry valid_target_bucket(FileKind::Add(), source_entries[1].Partition(), + /*bucket=*/0, + /*total_buckets=*/2, source_entries[1].File()); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN( + WrittenFile written_file, + manifest_file->WriteWithoutRolling({invalid_other_bucket, valid_target_bucket})); + + // Exercise both a cold cache and reuse of the retained manifest bytes. + for (int32_t read = 0; read < 2; ++read) { + std::vector bucket_entries; + ASSERT_OK(manifest_file->ReadBucketEntries(written_file.first, /*bucket=*/0, + std::nullopt, /*file_size=*/std::nullopt, + &bucket_entries)); + ASSERT_EQ(std::vector({valid_target_bucket}), bucket_entries); + ASSERT_EQ(cache_mode == 1 ? 1 : read + 1, file_system->open_count); + } + std::vector missing_bucket_entries; + ASSERT_OK(manifest_file->ReadBucketEntries(written_file.first, /*bucket=*/2, std::nullopt, + /*file_size=*/std::nullopt, + &missing_bucket_entries)); + ASSERT_TRUE(missing_bucket_entries.empty()); + std::vector all_entries; + ASSERT_NOK_WITH_MSG(manifest_file->Read(written_file.first, /*filter=*/nullptr, + /*file_size=*/std::nullopt, &all_entries), + "Unsupported byte value 2 for file kind."); + } +} - std::vector all_entries; - ASSERT_NOK_WITH_MSG(manifest_file->Read(written_file.first, /*filter=*/nullptr, - /*file_size=*/std::nullopt, &all_entries), - "Unsupported byte value 2 for file kind."); - - std::vector bucket_entries; - ASSERT_OK(manifest_file->ReadBucketEntries(written_file.first, /*bucket=*/0, - /*file_size=*/std::nullopt, &bucket_entries)); - ASSERT_EQ(std::vector({valid_target_bucket}), bucket_entries); +TEST_F(ManifestFileTest, TestReadBucketEntriesWithReorderedFields) { + auto pool = GetDefaultPool(); + auto source = + ReadManifestEntry("orc", paimon::test::GetDataDir() + "/orc/append_09.db/append_09", + "manifest-3a44a0da-1008-463c-914e-28d271375e24-0", pool); + ASSERT_EQ(source.size(), 2); + std::vector entries; + std::vector rows; + ManifestEntrySerializer serializer(pool); + for (int32_t bucket = 0; bucket < 2; ++bucket) { + entries.emplace_back(FileKind::Add(), source[bucket].Partition(), bucket, 2, + source[bucket].File()); + ASSERT_OK_AND_ASSIGN(BinaryRow row, serializer.ToRow(entries.back())); + rows.push_back(std::move(row)); + } + ASSERT_OK_AND_ASSIGN(std::unique_ptr converter, + MetaToArrowArrayConverter::Create(serializer.GetDataType(), pool)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr array, converter->NextBatch(rows)); + auto struct_array = checked_pointer_cast(array); + ASSERT_OK_AND_ASSIGN(std::shared_ptr format, FileFormatFactory::Get("avro", {})); + // Move each selector field away from its standard position. Ordinary schema alignment + // must still restore the fields before version validation and bucket filtering. + for (int32_t field_index : {0, 3, 4}) { + SCOPED_TRACE(field_index); + auto fields = serializer.GetDataType()->fields(); + auto columns = struct_array->fields(); + std::swap(fields[1], fields[field_index]); + std::swap(columns[1], columns[field_index]); + auto reordered_result = arrow::StructArray::Make(columns, fields); + ASSERT_TRUE(reordered_result.ok()) << reordered_result.status().ToString(); + std::shared_ptr reordered = reordered_result.ValueOrDie(); + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto fs = dir->GetFileSystem(); + const std::string manifest_dir = FileStorePathFactory::ManifestPath(dir->Str()); + ASSERT_OK(fs->Mkdirs(manifest_dir)); + const std::string file_name = "manifest-reordered"; + ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportSchema(*arrow::schema(fields), &c_schema).ok()); + ASSERT_OK_AND_ASSIGN(std::shared_ptr builder, + format->CreateWriterBuilder(&c_schema, 2)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr output, + fs->Create(manifest_dir + "/" + file_name, false)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr writer, builder->Build(output, "null")); + ArrowArray c_array; + ASSERT_TRUE(arrow::ExportArray(*reordered, &c_array).ok()); + ASSERT_OK(writer->AddBatch(&c_array)); + ASSERT_OK(writer->Flush()); + ASSERT_OK(writer->Finish()); + ASSERT_OK(output->Close()); + for (bool cache_enabled : {false, true}) { + for (bool inferred_bucket : {false, true}) { + for (int32_t bucket : {0, 1, 2}) { + ASSERT_OK_AND_ASSIGN( + std::vector actual, + TryReadManifestEntry("avro", dir->Str(), file_name, pool, bucket, + cache_enabled, inferred_bucket)); + const std::vector expected = + bucket < 2 ? std::vector{entries[bucket]} + : std::vector{}; + ASSERT_EQ(actual, expected); + } + } + } + } +} + +TEST_F(ManifestFileTest, TestReadBucketEntriesValidatesAllVersions) { + auto pool = GetDefaultPool(); + ManifestEntrySerializer serializer(pool); + arrow::FieldVector fields; + for (const auto& field : serializer.GetDataType()->fields()) { + fields.push_back(field->WithNullable(true)); + } + auto type = arrow::struct_(fields); + const std::vector> cases = { + {"999", "Unsupported version: 999"}, {"1", "not compatible"}}; + for (const auto& item : cases) { + SCOPED_TRACE(item.first); + for (int32_t total_buckets : {2, 3}) { + SCOPED_TRACE(total_buckets); + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto fs = dir->GetFileSystem(); + const std::string manifest_dir = FileStorePathFactory::ManifestPath(dir->Str()); + ASSERT_OK(fs->Mkdirs(manifest_dir)); + const std::string file_name = "manifest-invalid-version"; + ASSERT_OK_AND_ASSIGN(std::shared_ptr format, + FileFormatFactory::Get("avro", {})); + ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportSchema(*arrow::schema(type->fields()), &c_schema).ok()); + ASSERT_OK_AND_ASSIGN(std::shared_ptr builder, + format->CreateWriterBuilder(&c_schema, 2)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr output, + fs->Create(manifest_dir + "/" + file_name, false)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr writer, + builder->Build(output, "null")); + // A uniform unsupported version must fail before decoding the null metadata. + auto array = + arrow::ipc::internal::json::ArrayFromJSON( + type, fmt::format("[[{},0,null,1,{},null]]", item.first, total_buckets)) + .ValueOrDie(); + ArrowArray c_array; + ASSERT_TRUE(arrow::ExportArray(*array, &c_array).ok()); + ASSERT_OK(writer->AddBatch(&c_array)); + ASSERT_OK(writer->Flush()); + ASSERT_OK(writer->Finish()); + ASSERT_OK(output->Close()); + + for (bool cache_enabled : {false, true}) { + SCOPED_TRACE(cache_enabled); + for (bool inferred_bucket : {false, true}) { + SCOPED_TRACE(inferred_bucket); + // Matching entries always undergo version validation. + ASSERT_NOK_WITH_MSG( + TryReadManifestEntry("avro", dir->Str(), file_name, pool, /*bucket=*/1, + cache_enabled, inferred_bucket), + item.second); + // An unsupported version must fail even when its bucket does not match. + ASSERT_NOK_WITH_MSG( + TryReadManifestEntry("avro", dir->Str(), file_name, pool, /*bucket=*/0, + cache_enabled, inferred_bucket), + item.second); + } + } + } + } } TEST_F(ManifestFileTest, TestLegacyManifestFormatIsReadOnly) { @@ -600,6 +846,16 @@ TEST_F(ManifestFileTest, TestManifestFileCompatibleWithJavaPaimon09) { ASSERT_EQ(expected_manifest_entries, ReadManifestEntry("avro", paimon::test::GetDataDir() + "/avro", "avro_manifest_09", pool, /*bucket=*/0)); + ASSERT_OK_AND_ASSIGN(std::vector cached_entries, + TryReadManifestEntry("avro", paimon::test::GetDataDir() + "/avro", + "avro_manifest_09", pool, /*bucket=*/0, + /*cache_enabled=*/true, false)); + ASSERT_EQ(expected_manifest_entries, cached_entries); + ASSERT_OK_AND_ASSIGN(std::vector inferred_entries, + TryReadManifestEntry("avro", paimon::test::GetDataDir() + "/avro", + "avro_manifest_09", pool, /*bucket=*/0, + /*cache_enabled=*/true, /*inferred_bucket=*/true)); + ASSERT_EQ(expected_manifest_entries, inferred_entries); } TEST_F(ManifestFileTest, TestManifestFileCompatibleWithJavaPaimon11) { @@ -639,6 +895,16 @@ TEST_F(ManifestFileTest, TestManifestFileCompatibleWithJavaPaimon11) { ASSERT_EQ(expected_manifest_entries, ReadManifestEntry("avro", paimon::test::GetDataDir() + "/avro", "avro_manifest_11", pool, /*bucket=*/0)); + ASSERT_OK_AND_ASSIGN(std::vector cached_entries, + TryReadManifestEntry("avro", paimon::test::GetDataDir() + "/avro", + "avro_manifest_11", pool, /*bucket=*/0, + /*cache_enabled=*/true, false)); + ASSERT_EQ(expected_manifest_entries, cached_entries); + ASSERT_OK_AND_ASSIGN(std::vector inferred_entries, + TryReadManifestEntry("avro", paimon::test::GetDataDir() + "/avro", + "avro_manifest_11", pool, /*bucket=*/0, + /*cache_enabled=*/true, /*inferred_bucket=*/true)); + ASSERT_EQ(expected_manifest_entries, inferred_entries); } } // namespace paimon::test 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 b45a5684..908f9bbb 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 @@ -19,6 +19,7 @@ #include "paimon/core/operation/append_only_file_store_scan.h" #include +#include #include #include #include @@ -33,15 +34,21 @@ #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/io/meta_to_arrow_array_converter.h" #include "paimon/core/manifest/manifest_entry.h" +#include "paimon/core/manifest/manifest_entry_serializer.h" +#include "paimon/core/manifest/manifest_file.h" +#include "paimon/core/manifest/manifest_file_meta.h" #include "paimon/core/manifest/partition_entry.h" #include "paimon/core/schema/schema_manager.h" #include "paimon/core/schema/table_schema.h" #include "paimon/core/stats/simple_stats_evolution.h" #include "paimon/core/table/source/abstract_table_scan.h" #include "paimon/core/table/source/snapshot/snapshot_reader.h" +#include "paimon/core/utils/file_store_path_factory.h" #include "paimon/defs.h" #include "paimon/executor.h" +#include "paimon/format/file_format_factory.h" #include "paimon/fs/local/local_file_system.h" #include "paimon/memory/memory_pool.h" #include "paimon/metrics.h" @@ -51,7 +58,10 @@ #include "paimon/status.h" #include "paimon/table/source/scan_metrics.h" #include "paimon/table/source/table_scan.h" +#include "paimon/testing/mock/mock_file_batch_reader.h" +#include "paimon/testing/mock/mock_format_reader_builder.h" #include "paimon/testing/utils/binary_row_generator.h" +#include "paimon/testing/utils/test_helper.h" #include "paimon/testing/utils/testharness.h" #include "paimon/testing/utils/timezone_guard.h" namespace paimon::test { @@ -141,12 +151,315 @@ class AppendBucketPruningTest : public testing::Test { {Options::BUCKET_KEY, "rowkey"}}; }; +namespace { + +class CountingSchemaFileSystem : public LocalFileSystem { + public: + Status ReadFile(const std::string& path, std::string* content) override { + ++read_count; + return LocalFileSystem::ReadFile(path, content); + } + + int32_t read_count = 0; +}; + +class CountingManifestEntrySerializer : public ManifestEntrySerializer { + public: + explicit CountingManifestEntrySerializer(const std::shared_ptr& pool) + : ManifestEntrySerializer(pool) {} + + Result ConvertFrom(int32_t version, const InternalRow& row) const override { + ++decoded_entries; + return ManifestEntrySerializer::ConvertFrom(version, row); + } + + mutable std::atomic decoded_entries{0}; +}; + +class CountingManifestReaderBuilder : public MockFormatReaderBuilder { + public: + explicit CountingManifestReaderBuilder(const std::shared_ptr& data) + : MockFormatReaderBuilder(data, data->type(), 2), data_(data) {} + + Result> Build( + const std::shared_ptr& input) const override { + return std::make_unique(data_, &rows_read); + } + + mutable std::atomic rows_read{0}; + + private: + class Reader : public MockFileBatchReader { + public: + Reader(const std::shared_ptr& data, std::atomic* rows_read) + : MockFileBatchReader(data, data->type(), 2), rows_read_(rows_read) { + EnableRandomizeBatchSize(false); + } + + bool SupportPreciseBitmapSelection() const override { + return true; + } + + Result NextBatchWithBitmap() override { + PAIMON_ASSIGN_OR_RAISE(ReadBatchWithBitmap batch, + MockFileBatchReader::NextBatchWithBitmap()); + if (!BatchReader::IsEofBatch(batch)) { + *rows_read_ += batch.first.first->length; + } + return batch; + } + + private: + std::atomic* rows_read_; + }; + + std::shared_ptr data_; +}; + +} // namespace + +TEST_F(AppendBucketPruningTest, ReadsSingleBucketManifestOnce) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto fs = dir->GetFileSystem(); + // The converter owns the Arrow memory pool used by the mock reader's shared data. + ManifestEntrySerializer serializer(pool_); + ASSERT_OK_AND_ASSIGN(std::unique_ptr converter, + MetaToArrowArrayConverter::Create(serializer.GetDataType(), pool_)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr initial, CreateScan(KeyEquals())); + const int32_t bucket = initial->bucket_selector_->Bucket(kNumBuckets); + ASSERT_OK_AND_ASSIGN(std::shared_ptr format, FileFormatFactory::Get("avro", {})); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr paths, + FileStorePathFactory::Create(dir->Str(), initial->schema_, {}, "", "avro", "data-", true, + {}, std::nullopt, false, pool_)); + ASSERT_OK(fs->Mkdirs(FileStorePathFactory::ManifestPath(dir->Str()))); + ASSERT_OK_AND_ASSIGN(CoreOptions options, CoreOptions::FromMap({})); + options.WithCache(std::make_shared(16 * 1024 * 1024)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr manifest, + ManifestFile::Create(fs, format, "null", paths, 1024 * 1024, pool_, + options, arrow::schema({}))); + SimpleStats stats = BinaryRowGenerator::GenerateStats( + {std::string("a"), 0}, {std::string("z"), 100}, {0, 0}, pool_.get()); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr file, + DataFileMeta::ForAppend("data.avro", 100, 10, stats, 0, 9, schema_id_, std::nullopt, + std::nullopt, std::nullopt, std::nullopt, std::nullopt)); + std::vector entries = { + ManifestEntry(FileKind::Add(), BinaryRow::EmptyRow(), bucket, kNumBuckets, file)}; + ASSERT_OK_AND_ASSIGN(std::vector metas, manifest->Write(entries)); + ASSERT_EQ(metas.size(), 1); + ASSERT_OK_AND_ASSIGN(BinaryRow row, serializer.ToRow(entries[0])); + ASSERT_OK_AND_ASSIGN(std::shared_ptr data, converter->NextBatch({row})); + auto reader_builder = std::make_shared(data); + manifest->reader_builder_ = reader_builder; + options_[Options::SCAN_MANIFEST_ENTRY_LAZY_DECODE_ENABLED] = "true"; + for (bool inferred : {false, true}) { + SCOPED_TRACE(inferred); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr scan, + CreateScan(KeyEquals(), inferred ? std::nullopt : std::optional(bucket))); + scan->manifest_file_ = manifest; + for (bool known_bounds : {false, true}) { + SCOPED_TRACE(known_bounds); + auto meta = metas[0]; + meta.min_bucket_ = known_bounds ? std::optional(bucket) : std::nullopt; + meta.max_bucket_ = meta.min_bucket_; + reader_builder->rows_read = 0; + std::vector actual; + ASSERT_OK(scan->ReadAndMergeBucketFileEntries({meta}, bucket, &actual)); + ASSERT_EQ(actual, entries); + ASSERT_EQ(reader_builder->rows_read.load(), known_bounds ? 1 : 2); + } + } +} + +TEST_F(AppendBucketPruningTest, SelectivelyDecodesInferredBucketCandidates) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto fs = dir->GetFileSystem(); + schema_manager_ = std::make_shared(fs, dir->Str()); + schema_id_ = 0; + ASSERT_OK_AND_ASSIGN(std::unique_ptr initial_scan, + CreateScan(nullptr)); + ASSERT_OK(schema_manager_->CreateTable(initial_scan->schema_, {}, {}, options_)); + schema_id_ = 1; + ASSERT_OK_AND_ASSIGN(std::unique_ptr inferred_scan, + CreateScan(KeyEquals())); + ASSERT_FALSE(inferred_scan->bucket_filter_); + ASSERT_TRUE(inferred_scan->bucket_selector_); + const int32_t bucket = inferred_scan->bucket_selector_->Bucket(kNumBuckets); + ASSERT_OK_AND_ASSIGN(std::shared_ptr format, FileFormatFactory::Get("avro", {})); + ASSERT_OK_AND_ASSIGN(std::shared_ptr paths, + FileStorePathFactory::Create( + dir->Str(), initial_scan->schema_, /*partition_keys=*/{}, + /*default_part_value=*/"", "avro", /*data_file_prefix=*/"data-", + /*legacy_partition_name_enabled=*/true, /*external_paths=*/{}, + /*global_index_external_path=*/std::nullopt, + /*index_file_in_data_file_dir=*/false, pool_)); + ASSERT_OK(fs->Mkdirs(FileStorePathFactory::ManifestPath(dir->Str()))); + SimpleStats stats = BinaryRowGenerator::GenerateStats( + {std::string("a"), 0}, {std::string("z"), 100}, {0, 0}, pool_.get()); + std::vector entries; + for (const auto& layout : std::vector>{ + {kNumBuckets, schema_id_}, {8, schema_id_}, {kNumBuckets, 0}, {-1, schema_id_}}) { + for (int32_t id = 0; id < std::max(1, layout.first); ++id) { + ASSERT_OK_AND_ASSIGN( + std::shared_ptr file, + DataFileMeta::ForAppend(fmt::format("data-{}.avro", entries.size()), 100, 10, stats, + 0, 9, layout.second, std::nullopt, std::nullopt, + std::nullopt, std::nullopt, std::nullopt)); + entries.emplace_back(FileKind::Add(), BinaryRow::EmptyRow(), id, layout.first, file); + } + } + ASSERT_OK_AND_ASSIGN( + std::shared_ptr deleted_file, + DataFileMeta::ForAppend("deleted.avro", 100, 10, stats, 0, 9, schema_id_, std::nullopt, + std::nullopt, std::nullopt, std::nullopt, std::nullopt)); + entries.emplace_back(FileKind::Add(), BinaryRow::EmptyRow(), bucket, kNumBuckets, deleted_file); + ManifestEntry deletion(FileKind::Delete(), BinaryRow::EmptyRow(), bucket, kNumBuckets, + deleted_file); + for (bool cache_enabled : {false, true}) { + SCOPED_TRACE(cache_enabled); + ASSERT_OK_AND_ASSIGN(CoreOptions manifest_options, + CoreOptions::FromMap({{Options::READ_BATCH_SIZE, "2"}})); + if (cache_enabled) { + manifest_options.WithCache(std::make_shared(16 * 1024 * 1024)); + } + ASSERT_OK_AND_ASSIGN(std::shared_ptr manifest, + ManifestFile::Create(fs, format, "null", paths, 1024 * 1024, pool_, + manifest_options, arrow::schema({}))); + auto serializer = std::make_unique(pool_); + auto* counter = serializer.get(); + manifest->serializer_ = std::move(serializer); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile additions, manifest->WriteWithoutRolling(entries)); + ASSERT_OK_AND_ASSIGN(WrittenFile deletions, manifest->WriteWithoutRolling({deletion})); + std::vector metas; + for (const auto& file : {additions, deletions}) { + metas.emplace_back(file.first, file.second, 0, 0, SimpleStats::EmptyStats(), schema_id_, + std::nullopt, std::nullopt, std::nullopt, std::nullopt, std::nullopt, + std::nullopt); + } + std::vector baseline_files; + for (bool lazy_decode : {false, true}) { + SCOPED_TRACE(lazy_decode); + options_[Options::SCAN_MANIFEST_ENTRY_LAZY_DECODE_ENABLED] = + lazy_decode ? "true" : "false"; + ASSERT_OK_AND_ASSIGN(std::unique_ptr scan, + CreateScan(KeyEquals())); + scan->manifest_file_ = manifest; + counter->decoded_entries = 0; + std::vector candidates; + ASSERT_OK(scan->ReadAndMergeBucketFileEntries(metas, bucket, &candidates)); + // Both schema versions have the same bucket layout: skip their unrelated + // buckets, but retain different bucket counts and both sides of the deletion. + ASSERT_EQ(counter->decoded_entries.load(), + entries.size() + 1 - (lazy_decode ? 2 * (kNumBuckets - 1) : 0)); + std::vector files; + for (const auto& entry : candidates) { + ASSERT_NE(entry.FileName(), "deleted.avro"); + ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry)); + if (keep) { + files.push_back(entry.FileName()); + } + } + std::sort(files.begin(), files.end()); + ASSERT_EQ(files.size(), 4); + if (!lazy_decode) { + baseline_files = files; + } else { + ASSERT_EQ(files, baseline_files); + } + } + } +} + 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, SelectivelyDecodesWithoutReadingHistoricalSchemas) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto fs = std::make_shared(); + schema_manager_ = std::make_shared(fs, dir->Str()); + schema_id_ = 0; + rowkey_type_ = arrow::int32(); + ASSERT_OK_AND_ASSIGN(std::unique_ptr initial, CreateScan(nullptr)); + ASSERT_OK(schema_manager_->CreateTable(initial->schema_, {}, {}, options_)); + schema_id_ = 1000; + auto predicate = PredicateBuilder::Equal(0, "rowkey", FieldType::INT, Literal(int32_t{-1})); + ASSERT_OK_AND_ASSIGN(std::unique_ptr current, CreateScan(predicate)); + // Only schema 0 has a schema file; unused historical versions must not be read. + const int32_t bucket = current->bucket_selector_->Bucket(kNumBuckets); + SimpleStats old_stats = BinaryRowGenerator::GenerateStats({int32_t{-100}, 0}, {int32_t{0}, 100}, + {0, 0}, pool_.get()); + std::vector entries; + for (int32_t id = 0; id <= kNumBuckets; ++id) { + const bool historical = id < kNumBuckets; + ASSERT_OK_AND_ASSIGN( + std::shared_ptr file, + DataFileMeta::ForAppend(fmt::format("data-{}.avro", id), 100, 10, old_stats, 0, 9, + historical ? 0 : schema_id_, std::nullopt, std::nullopt, + std::nullopt, std::nullopt, std::nullopt)); + entries.emplace_back(FileKind::Add(), BinaryRow::EmptyRow(), historical ? id : bucket, + kNumBuckets, file); + } + ASSERT_OK_AND_ASSIGN(std::shared_ptr format, FileFormatFactory::Get("avro", {})); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr paths, + FileStorePathFactory::Create(dir->Str(), initial->schema_, {}, "", "avro", "data-", true, + {}, std::nullopt, false, pool_)); + ASSERT_OK(fs->Mkdirs(FileStorePathFactory::ManifestPath(dir->Str()))); + for (bool cache_enabled : {false, true}) { + SCOPED_TRACE(cache_enabled); + ASSERT_OK_AND_ASSIGN(CoreOptions manifest_options, CoreOptions::FromMap({})); + if (cache_enabled) { + manifest_options.WithCache(std::make_shared(16 * 1024 * 1024)); + } + ASSERT_OK_AND_ASSIGN(std::shared_ptr manifest, + ManifestFile::Create(fs, format, "null", paths, 1024 * 1024, pool_, + manifest_options, arrow::schema({}))); + ASSERT_OK_AND_ASSIGN(std::vector metas, manifest->Write(entries)); + ASSERT_EQ(metas.size(), 1); + ASSERT_EQ(metas[0].SchemaId(), schema_id_); + auto serializer = std::make_unique(pool_); + auto* counter = serializer.get(); + manifest->serializer_ = std::move(serializer); + for (bool lazy_decode : {false, true}) { + SCOPED_TRACE(lazy_decode); + options_[Options::SCAN_MANIFEST_ENTRY_LAZY_DECODE_ENABLED] = + lazy_decode ? "true" : "false"; + ASSERT_OK_AND_ASSIGN(std::unique_ptr scan, + CreateScan(predicate)); + scan->manifest_file_ = manifest; + std::vector baseline; + ASSERT_OK(scan->ReadAndMergeFileEntries(metas, &baseline)); + ASSERT_EQ(baseline.size(), 2); + // Keep schema reads used by stats evolution out of the bucket-read measurement. + scan->schema_manager_ = std::make_shared(fs, dir->Str()); + counter->decoded_entries = 0; + fs->read_count = 0; + std::vector candidates; + ASSERT_OK(scan->ReadAndMergeBucketFileEntries(metas, bucket, &candidates)); + ASSERT_EQ(fs->read_count, 0); + ASSERT_EQ(counter->decoded_entries.load(), + entries.size() - (lazy_decode ? kNumBuckets - 1 : 0)); + std::vector filtered; + for (const auto& entry : candidates) { + ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry)); + if (keep) { + filtered.push_back(entry); + } + } + ASSERT_EQ(filtered, baseline); + } + } +} + TEST_F(AppendBucketPruningTest, PreservesExplicitBucketFilter) { BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool_.get()); int32_t other_bucket = (DefaultBucketFunction().Bucket(key, kNumBuckets) + 1) % kNumBuckets; @@ -187,17 +500,26 @@ TEST_F(AppendBucketPruningTest, PrunesCompatibleHistoricalSchema) { 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, PrunesHistoricalBucketsWithoutReadingSchemas) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto fs = std::make_shared(); + schema_manager_ = std::make_shared(fs, dir->Str()); + schema_id_ = 1000; + ASSERT_OK_AND_ASSIGN(std::unique_ptr scan, CreateScan(KeyEquals())); + const int32_t other_bucket = (scan->bucket_selector_->Bucket(kNumBuckets) + 1) % kNumBuckets; + for (int64_t data_schema_id : {0, 500, 999, 1000}) { + ASSERT_OK_AND_ASSIGN( + std::shared_ptr file, + DataFileMeta::ForAppend("data.parquet", 100, 10, SimpleStats::EmptyStats(), 0, 9, + data_schema_id, std::nullopt, std::nullopt, std::nullopt, + std::nullopt, std::nullopt)); + ManifestEntry entry(FileKind::Add(), BinaryRow::EmptyRow(), other_bucket, kNumBuckets, + file); + ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry)); + ASSERT_FALSE(keep); + } + ASSERT_EQ(fs->read_count, 0); } 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 89c137f3..45376b22 100644 --- a/src/paimon/core/operation/file_store_scan.cpp +++ b/src/paimon/core/operation/file_store_scan.cpp @@ -356,7 +356,7 @@ Status FileStoreScan::ReadManifestEntries(const std::vector& m // Cache merged live manifest entries for one bucket before applying scan filters. Each cache value // 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 +// bucket counts, 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( @@ -447,15 +447,26 @@ Status FileStoreScan::ReadAndMergeBucketFileEntries( const std::vector& manifest_metas, int32_t bucket, std::vector* merged_entries) const { 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()) { + if (core_options_.ScanManifestEntryLazyDecodeEnabled()) { std::vector>>> futures; futures.reserve(manifest_metas.size()); for (const auto& meta : manifest_metas) { - auto read_meta_task = [this, meta, bucket]() -> Result> { + auto read_meta_task = [this, meta, bucket, + inferred_bucket]() -> Result> { std::vector bucket_entries; - PAIMON_RETURN_NOT_OK(manifest_file_->ReadBucketEntries( - meta.FileName(), bucket, meta.FileSize(), &bucket_entries)); + if (meta.MinBucket() && meta.MaxBucket() && meta.MinBucket().value() == bucket && + meta.MaxBucket().value() == bucket) { + // Every entry belongs to this bucket; a projection pass cannot prune rows. + PAIMON_RETURN_NOT_OK(manifest_file_->Read(meta.FileName(), /*filter=*/nullptr, + meta.FileSize(), &bucket_entries)); + } else if (inferred_bucket) { + PAIMON_RETURN_NOT_OK(manifest_file_->ReadBucketEntries( + meta.FileName(), bucket, core_options_.GetBucket(), meta.FileSize(), + &bucket_entries)); + } else { + PAIMON_RETURN_NOT_OK(manifest_file_->ReadBucketEntries( + meta.FileName(), bucket, std::nullopt, meta.FileSize(), &bucket_entries)); + } return bucket_entries; }; futures.push_back(Via(executor_.get(), read_meta_task)); @@ -481,8 +492,7 @@ Status FileStoreScan::ReadAndMergeBucketFileEntries( unmerged_entries.reserve(entries.size()); for (auto& entry : entries) { if (entry.Bucket() == bucket || - (inferred_bucket && (entry.TotalBuckets() != core_options_.GetBucket() || - entry.File()->schema_id != table_schema_->Id()))) { + (inferred_bucket && entry.TotalBuckets() != core_options_.GetBucket())) { unmerged_entries.emplace_back(std::move(entry)); } } @@ -598,11 +608,9 @@ Result FileStoreScan::FilterManifestEntry(const ManifestEntry& entry) cons if (bucket_filter_ != std::nullopt && entry.Bucket() != bucket_filter_.value()) { 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 (bucket_selector_ && entry.TotalBuckets() > 0 && + entry.Bucket() != bucket_selector_->Bucket(entry.TotalBuckets())) { + return false; } if (level_filter_ != nullptr && !level_filter_(entry.Level())) { return false; @@ -610,31 +618,6 @@ 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 d5ee9215..0eef2909 100644 --- a/src/paimon/core/operation/file_store_scan.h +++ b/src/paimon/core/operation/file_store_scan.h @@ -34,7 +34,6 @@ #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" @@ -303,7 +302,6 @@ 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_; @@ -326,7 +324,6 @@ 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/snapshot_file_scan.cpp b/src/paimon/core/snapshot_file_scan.cpp index 9d80ec86..c7f854b8 100644 --- a/src/paimon/core/snapshot_file_scan.cpp +++ b/src/paimon/core/snapshot_file_scan.cpp @@ -182,9 +182,12 @@ class SnapshotFileCollector { } futures.push_back(Via(executor_.get(), [this, meta]() -> Result { std::vector entries; - if (bucket_id_) { - PAIMON_RETURN_NOT_OK(manifest_file_->ReadBucketEntries( - meta.FileName(), bucket_id_.value(), meta.FileSize(), &entries)); + const bool all_in_bucket = + bucket_id_ && meta.MinBucket() == bucket_id_ && meta.MaxBucket() == bucket_id_; + if (bucket_id_ && !all_in_bucket) { + PAIMON_RETURN_NOT_OK( + manifest_file_->ReadBucketEntries(meta.FileName(), bucket_id_.value(), + std::nullopt, meta.FileSize(), &entries)); } else { PAIMON_RETURN_NOT_OK(manifest_file_->Read(meta.FileName(), /*filter=*/nullptr, meta.FileSize(), &entries)); diff --git a/src/paimon/core/snapshot_file_scan_test.cpp b/src/paimon/core/snapshot_file_scan_test.cpp index e6629fbd..b0363376 100644 --- a/src/paimon/core/snapshot_file_scan_test.cpp +++ b/src/paimon/core/snapshot_file_scan_test.cpp @@ -20,6 +20,7 @@ #include "paimon/snapshot/snapshot_file_scan.h" #include +#include #include #include #include @@ -31,12 +32,15 @@ #include #include "gtest/gtest.h" +#include "paimon/common/factories/io_hook.h" #include "paimon/common/utils/path_util.h" +#include "paimon/common/utils/scope_guard.h" #include "paimon/common/utils/string_utils.h" #include "paimon/defs.h" #include "paimon/fs/local/local_file_system.h" #include "paimon/predicate/predicate_builder.h" #include "paimon/scan_context.h" +#include "paimon/testing/utils/test_helper.h" #include "paimon/testing/utils/testharness.h" namespace paimon::test { @@ -151,6 +155,37 @@ TEST(SnapshotFileScanTest, TestLatestSnapshotAndBucketFilter) { bucket_one_files); } +TEST(SnapshotFileScanTest, TestSingleBucketManifestAvoidsExtraIO) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto schema = arrow::schema({arrow::field("value", arrow::utf8())}); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr helper, + TestHelper::Create( + dir->Str(), schema, {}, {}, + {{Options::FILE_FORMAT, "orc"}, {Options::BUCKET, "1"}, {Options::BUCKET_KEY, "value"}}, + false)); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr batch, + TestHelper::MakeRecordBatch(arrow::struct_(schema->fields()), R"([["value"]])", {}, 0, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(batch), 0, std::nullopt)); + helper.reset(); + const std::string table_path = dir->Str() + "/foo.db/bar"; + + auto* io_hook = IOHook::GetInstance(); + ScopeGuard guard([io_hook]() { io_hook->Clear(); }); + io_hook->Reset(-1, IOHook::Mode::SILENT); + ASSERT_OK_AND_ASSIGN(std::set all_files, ListFiles(table_path)); + const int64_t full_read_io_count = io_hook->IOCount(); + ASSERT_GT(full_read_io_count, 0); + + io_hook->Reset(-1, IOHook::Mode::SILENT); + ASSERT_OK_AND_ASSIGN(std::set bucket_files, + ListFiles(table_path, std::nullopt, CreateFilter({}, 0))); + ASSERT_EQ(bucket_files, all_files); + ASSERT_EQ(io_hook->IOCount(), full_read_io_count); +} + TEST(SnapshotFileScanTest, TestExplicitSnapshot) { std::string table_path = GetDataDir() + "/orc/append_09.db/append_09"; diff --git a/src/paimon/core/utils/objects_file.h b/src/paimon/core/utils/objects_file.h index 845efd41..e65ee75a 100644 --- a/src/paimon/core/utils/objects_file.h +++ b/src/paimon/core/utils/objects_file.h @@ -92,10 +92,12 @@ class ObjectsFile { return Status::OK(); } + // Optional preparation may wrap the reader and reread the file, using either cached bytes + // or the underlying file stream. Status ReadArrowBatches( - const std::string& file_name, + const std::string& file_name, std::optional file_size, const std::function&)>& consumer, - std::optional file_size) const; + const std::function*)>& prepare_reader) const; std::shared_ptr path_factory_; std::shared_ptr pool_; @@ -158,7 +160,7 @@ Status ObjectsFile::Read(const std::string& file_name, const std::function(const T&)>& filter, std::optional file_size, std::vector* result) const { return ReadArrowBatches( - file_name, + file_name, file_size, [this, &filter, result](const std::shared_ptr& struct_array) -> Status { result->reserve(result->size() + struct_array->length()); const arrow::ArrayVector& fields = struct_array->fields(); @@ -177,14 +179,14 @@ Status ObjectsFile::Read(const std::string& file_name, } return Status::OK(); }, - file_size); + /*prepare_reader=*/nullptr); } template Status ObjectsFile::ReadArrowBatches( - const std::string& file_name, + const std::string& file_name, std::optional file_size, const std::function&)>& consumer, - std::optional file_size) const { + const std::function*)>& prepare_reader) const { std::string file_path = path_factory_->ToPath(file_name); std::shared_ptr file_input_stream; std::shared_ptr cached_bytes; @@ -215,6 +217,9 @@ Status ObjectsFile::ReadArrowBatches( PAIMON_ASSIGN_OR_RAISE(std::unique_ptr batch_reader, reader_builder_->Build(file_input_stream)); + if (prepare_reader) { + PAIMON_RETURN_NOT_OK(prepare_reader(&batch_reader)); + } auto reader = std::make_unique(std::move(batch_reader), serializer_->GetDataType(), arrow_pool_); while (true) { diff --git a/src/paimon/format/avro/avro_direct_decoder.cpp b/src/paimon/format/avro/avro_direct_decoder.cpp index 9f9e3642..9f2d5763 100644 --- a/src/paimon/format/avro/avro_direct_decoder.cpp +++ b/src/paimon/format/avro/avro_direct_decoder.cpp @@ -466,6 +466,10 @@ Status AvroDirectDecoder::DecodeAvroToBuilder(const ::avro::NodePtr& avro_node, return DecodeFieldToBuilder(avro_node, projection, decoder, array_builder, ctx); } +Status AvroDirectDecoder::SkipValue(const ::avro::NodePtr& avro_node, ::avro::Decoder* decoder) { + return SkipAvroValue(avro_node, decoder); +} + Status AvroDirectDecoder::ReserveBuilderCapacity(int64_t capacity, arrow::ArrayBuilder* array_builder) { return ReserveBuilderCapacityImpl(capacity, array_builder); diff --git a/src/paimon/format/avro/avro_direct_decoder.h b/src/paimon/format/avro/avro_direct_decoder.h index 3b316f54..ad5cdd9e 100644 --- a/src/paimon/format/avro/avro_direct_decoder.h +++ b/src/paimon/format/avro/avro_direct_decoder.h @@ -82,6 +82,9 @@ class AvroDirectDecoder { ::avro::Decoder* decoder, arrow::ArrayBuilder* array_builder, DecodeContext* ctx); + /// Skip one value without materializing it in an Arrow builder. + static Status SkipValue(const ::avro::NodePtr& avro_node, ::avro::Decoder* decoder); + /// Reserve slots for a builder and any struct children with the same cardinality. /// @param capacity Number of additional values to append. /// @param array_builder Builder to reserve. diff --git a/src/paimon/format/avro/avro_file_batch_reader.cpp b/src/paimon/format/avro/avro_file_batch_reader.cpp index 252f563c..de3c2837 100644 --- a/src/paimon/format/avro/avro_file_batch_reader.cpp +++ b/src/paimon/format/avro/avro_file_batch_reader.cpp @@ -18,7 +18,8 @@ #include "paimon/format/avro/avro_file_batch_reader.h" -#include +#include +#include #include #include @@ -101,25 +102,26 @@ Result> AvroFileBatchReader::CreateD } Result AvroFileBatchReader::NextBatch() { - if (next_row_to_read_ == std::numeric_limits::max()) { - next_row_to_read_ = 0; - } + previous_first_row_ = next_row_to_read_; + previous_row_ids_.clear(); + previous_batch_row_count_ = 0; try { while (array_builder_->length() < batch_size_) { - if (!reader_->hasMore()) { + const std::optional target = + selection_ ? selection_->NextRow() : std::optional(next_row_to_read_); + if (!target) { + break; + } + PAIMON_ASSIGN_OR_RAISE(bool available, AdvanceToRow(target.value())); + if (!available) { break; } - if (array_builder_->length() == 0) { - PAIMON_RETURN_NOT_OK( - AvroDirectDecoder::ReserveBuilderCapacity(batch_size_, array_builder_.get())); + PAIMON_RETURN_NOT_OK(ReadCurrentRow(/*materialize=*/true)); + if (selection_) { + previous_row_ids_.push_back(target.value()); + selection_->Advance(); } - reader_->decr(); - PAIMON_RETURN_NOT_OK(AvroDirectDecoder::DecodeAvroToBuilder( - reader_->dataSchema().root(), read_fields_projection_, &reader_->decoder(), - array_builder_.get(), &decode_context_)); } - previous_first_row_ = next_row_to_read_; - next_row_to_read_ += array_builder_->length(); if (array_builder_->length() == 0) { previous_batch_row_count_ = 0; return BatchReader::MakeEofBatch(); @@ -145,6 +147,108 @@ Result AvroFileBatchReader::NextBatch() { } } +std::optional AvroFileBatchReader::SelectionCursor::NextRow() const { + if (next_ == end_) { + return std::nullopt; + } + return static_cast(*next_); +} + +void AvroFileBatchReader::BlockIndex::Observe(uint64_t row, int64_t file_offset) { + if (state_ != State::kBuilding || + (!blocks_.empty() && blocks_.back().file_offset == file_offset)) { + return; + } + if (blocks_.size() == kMaxBlocks) { + blocks_.clear(); + state_ = State::kDisabled; + return; + } + blocks_.push_back({row, file_offset}); +} + +void AvroFileBatchReader::BlockIndex::Finish(uint64_t row_count) { + if (state_ == State::kBuilding) { + row_count_ = row_count; + state_ = State::kReady; + } +} + +void AvroFileBatchReader::BlockIndex::Reset() { + if (state_ == State::kBuilding) { + blocks_.clear(); + } +} + +std::optional AvroFileBatchReader::BlockIndex::RowCount() const { + return state_ == State::kReady ? std::optional(row_count_) : std::nullopt; +} + +std::optional AvroFileBatchReader::BlockIndex::Locate( + uint64_t row) const { + if (state_ != State::kReady || row >= row_count_ || blocks_.empty()) { + return std::nullopt; + } + auto block = std::upper_bound( + blocks_.begin(), blocks_.end(), row, + [](uint64_t target, const Position& position) { return target < position.first_row; }); + return *(block - 1); +} + +Result AvroFileBatchReader::AdvanceToRow(uint64_t row) { + assert(row >= next_row_to_read_); + const std::optional row_count = block_index_.RowCount(); + if (row_count && row >= row_count.value()) { + return false; + } + if (row > next_row_to_read_) { + const auto block = block_index_.Locate(row); + if (block && block->first_row > next_row_to_read_) { + reader_->seek(block->file_offset); + next_row_to_read_ = block->first_row; + } + } + while (next_row_to_read_ < row) { + if (!PrepareNextRow()) { + return false; + } + PAIMON_RETURN_NOT_OK(ReadCurrentRow(/*materialize=*/false)); + } + return PrepareNextRow(); +} + +bool AvroFileBatchReader::PrepareNextRow() { + if (!reader_->hasMore()) { + if (!selection_) { + block_index_.Finish(next_row_to_read_); + } + return false; + } + // Preserve index construction during full reads; bitmap reads only consume an existing index. + if (!selection_) { + block_index_.Observe(next_row_to_read_, reader_->previousSync()); + } + return true; +} + +Status AvroFileBatchReader::ReadCurrentRow(bool materialize) { + reader_->decr(); + if (materialize) { + if (array_builder_->length() == 0) { + PAIMON_RETURN_NOT_OK( + AvroDirectDecoder::ReserveBuilderCapacity(batch_size_, array_builder_.get())); + } + PAIMON_RETURN_NOT_OK(AvroDirectDecoder::DecodeAvroToBuilder( + reader_->dataSchema().root(), read_fields_projection_, &reader_->decoder(), + array_builder_.get(), &decode_context_)); + } else { + PAIMON_RETURN_NOT_OK( + AvroDirectDecoder::SkipValue(reader_->dataSchema().root(), &reader_->decoder())); + } + ++next_row_to_read_; + return Status::OK(); +} + Status AvroFileBatchReader::SetReadSchema(::ArrowSchema* read_schema, const std::shared_ptr& predicate, const std::optional& selection_bitmap) { @@ -152,9 +256,6 @@ Status AvroFileBatchReader::SetReadSchema(::ArrowSchema* read_schema, return Status::Invalid("SetReadSchema failed: read schema cannot be nullptr"); } // TODO(menglingda.mld): support predicate - if (selection_bitmap) { - // TODO(menglingda.mld): support bitmap - } PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr arrow_read_schema, arrow::ImportSchema(read_schema)); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr file_schema, @@ -180,9 +281,15 @@ Status AvroFileBatchReader::SetReadSchema(::ArrowSchema* read_schema, reader_ = std::move(reader); array_builder_ = std::move(array_builder); decode_context_.ClearBuilderMetadata(); - previous_first_row_ = std::numeric_limits::max(); + selection_.reset(); + if (selection_bitmap) { + selection_.emplace(selection_bitmap.value()); + } + block_index_.Reset(); + previous_row_ids_.clear(); + previous_first_row_.reset(); previous_batch_row_count_ = 0; - next_row_to_read_ = std::numeric_limits::max(); + next_row_to_read_ = 0; close_ = false; return Status::OK(); } @@ -212,6 +319,9 @@ Result> AvroFileBatchReader::GetFileSchema() cons } Result AvroFileBatchReader::GetNumberOfRows() const { + if (const auto row_count = block_index_.RowCount()) { + return row_count.value(); + } if (!total_rows_) { PAIMON_ASSIGN_OR_RAISE(int64_t current_pos, input_stream_->GetPos()); ScopeGuard stream_guard([this, current_pos]() -> void { diff --git a/src/paimon/format/avro/avro_file_batch_reader.h b/src/paimon/format/avro/avro_file_batch_reader.h index 05bdf912..3c0fb6d7 100644 --- a/src/paimon/format/avro/avro_file_batch_reader.h +++ b/src/paimon/format/avro/avro_file_batch_reader.h @@ -18,8 +18,9 @@ #pragma once -#include +#include #include +#include #include #include #include @@ -52,7 +53,7 @@ class AvroFileBatchReader : public FileBatchReader { Result GetPreviousBatchFileRowId(uint64_t batch_row_id) const override { if (previous_batch_row_count_ == 0) { - if (previous_first_row_ == std::numeric_limits::max()) { + if (!previous_first_row_) { return Status::Invalid("No batch has been read yet."); } else { return Status::Invalid("Last batch was EOF."); @@ -63,7 +64,8 @@ class AvroFileBatchReader : public FileBatchReader { fmt::format("batch_row_id {} is out of range, last batch row count is {}", batch_row_id, previous_batch_row_count_)); } - return previous_first_row_ + batch_row_id; + return selection_ ? previous_row_ids_[batch_row_id] + : previous_first_row_.value() + batch_row_id; } Result GetNumberOfRows() const override; @@ -77,10 +79,55 @@ class AvroFileBatchReader : public FileBatchReader { } bool SupportPreciseBitmapSelection() const override { - return false; + return true; } private: + class SelectionCursor { + public: + explicit SelectionCursor(const RoaringBitmap32& bitmap) + : bitmap_(bitmap), next_(bitmap_.Begin()), end_(bitmap_.End()) {} + SelectionCursor(const SelectionCursor&) = delete; + SelectionCursor& operator=(const SelectionCursor&) = delete; + + std::optional NextRow() const; + + void Advance() { + ++next_; + } + + private: + // Iterators must be destroyed before the bitmap they reference. + RoaringBitmap32 bitmap_; + RoaringBitmap32::Iterator next_; + RoaringBitmap32::Iterator end_; + }; + + class BlockIndex { + public: + struct Position { + uint64_t first_row; + int64_t file_offset; + }; + + void Observe(uint64_t row, int64_t file_offset); + void Finish(uint64_t row_count); + void Reset(); + std::optional Locate(uint64_t row) const; + std::optional RowCount() const; + + private: + enum class State { kBuilding, kReady, kDisabled }; + static constexpr size_t kMaxBlocks = 64 * 1024; + State state_ = State::kBuilding; + std::vector blocks_; + uint64_t row_count_ = 0; + }; + + Result AdvanceToRow(uint64_t row); + bool PrepareNextRow(); + Status ReadCurrentRow(bool materialize); + void DoClose(); static Result> CreateDataFileReader( @@ -105,8 +152,12 @@ class AvroFileBatchReader : public FileBatchReader { std::unique_ptr<::avro::DataFileReaderBase> reader_; std::unique_ptr array_builder_; std::optional> read_fields_projection_; - uint64_t previous_first_row_ = std::numeric_limits::max(); - uint64_t next_row_to_read_ = std::numeric_limits::max(); + std::optional selection_; + // File-level acceleration, independent of the current projection and bitmap. + BlockIndex block_index_; + std::vector previous_row_ids_; + std::optional previous_first_row_; + uint64_t next_row_to_read_ = 0; uint64_t previous_batch_row_count_ = 0; mutable std::optional total_rows_ = std::nullopt; const int32_t batch_size_; diff --git a/src/paimon/format/avro/avro_file_batch_reader_test.cpp b/src/paimon/format/avro/avro_file_batch_reader_test.cpp index bd0d63eb..6a61c4c9 100644 --- a/src/paimon/format/avro/avro_file_batch_reader_test.cpp +++ b/src/paimon/format/avro/avro_file_batch_reader_test.cpp @@ -18,6 +18,8 @@ #include "paimon/format/avro/avro_file_batch_reader.h" +#include +#include #include #include #include @@ -25,19 +27,44 @@ #include "arrow/api.h" #include "arrow/c/bridge.h" #include "arrow/ipc/api.h" +#include "avro/Compiler.hh" #include "gtest/gtest.h" +#include "paimon/common/utils/checked_cast.h" #include "paimon/common/utils/path_util.h" #include "paimon/core/manifest/manifest_file.h" #include "paimon/core/manifest/manifest_list.h" #include "paimon/format/file_format.h" #include "paimon/format/file_format_factory.h" #include "paimon/fs/local/local_file_system.h" +#include "paimon/io/byte_array_input_stream.h" #include "paimon/testing/utils/read_result_collector.h" #include "paimon/testing/utils/testharness.h" #include "paimon/testing/utils/timezone_guard.h" namespace paimon::avro::test { +class CountingInputStream : public ByteArrayInputStream { + public: + explicit CountingInputStream(const std::string& bytes) + : ByteArrayInputStream(bytes.data(), bytes.size()) {} + + using ByteArrayInputStream::Read; + + Result Read(char* buffer, int64_t size) override { + PAIMON_ASSIGN_OR_RAISE(int64_t count, ByteArrayInputStream::Read(buffer, size)); + bytes_read += count; + return count; + } + + Status Seek(int64_t offset, SeekOrigin origin) override { + ++seek_count; + return ByteArrayInputStream::Seek(offset, origin); + } + + int64_t bytes_read = 0; + int32_t seek_count = 0; +}; + class AvroFileBatchReaderTest : public ::testing::Test, public ::testing::WithParamInterface { public: void SetUp() override { @@ -328,6 +355,346 @@ TEST_F(AvroFileBatchReaderTest, TestSetReadSchemaRejectNestedSubFieldProjection) "does not support nested sub-field projection"); } +TEST_F(AvroFileBatchReaderTest, TestPreciseBitmapSelectionAndReset) { + auto data_type = arrow::struct_( + {arrow::field("id", arrow::int32()), + arrow::field("payload", + arrow::struct_({arrow::field("text", arrow::utf8()), + arrow::field("values", arrow::list(arrow::int32()))}))}); + auto source = + arrow::ipc::internal::json::ArrayFromJSON( + data_type, R"([[0,["a",[1,2]]],[1,null],[2,["c",[]]],[3,["d",[4]]],[4,["e",null]]])") + .ValueOrDie(); + const std::vector> selections = { + {}, {0}, {1, 3}, {0, 1, 2, 3, 4}, {2, 99}}; + for (const std::string& compression : {std::string("null"), std::string("deflate")}) { + const std::string path = PathUtil::JoinPath(dir_->Str(), compression + ".avro"); + WriteData(source, path, compression); + for (int32_t batch_size : {1, 2, 8}) { + ASSERT_OK_AND_ASSIGN(std::shared_ptr builder, + file_format_->CreateReaderBuilder(batch_size)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr input, fs_->Open(path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, builder->Build(input)); + ASSERT_TRUE(reader->SupportPreciseBitmapSelection()); + for (const auto& ids : selections) { + RoaringBitmap32 selection; + for (uint32_t id : ids) { + selection.Add(id); + } + ArrowSchema c_schema; + ASSERT_TRUE( + arrow::ExportSchema(*arrow::schema(data_type->fields()), &c_schema).ok()); + ASSERT_OK(reader->SetReadSchema(&c_schema, nullptr, selection)); + ASSERT_NOK(reader->GetPreviousBatchFileRowId(0)); + size_t cursor = 0; + while (true) { + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch batch, reader->NextBatch()); + if (BatchReader::IsEofBatch(batch)) { + break; + } + auto array = + arrow::ImportArray(batch.first.get(), batch.second.get()).ValueOrDie(); + ASSERT_GT(array->length(), 0); + ASSERT_LE(array->length(), batch_size); + for (int64_t i = 0; i < array->length(); ++i) { + ASSERT_LT(cursor, ids.size()); + ASSERT_OK_AND_ASSIGN(uint64_t file_row, + reader->GetPreviousBatchFileRowId(i)); + ASSERT_EQ(file_row, ids[cursor++]); + ASSERT_TRUE(array->Slice(i, 1)->Equals(source->Slice(file_row, 1))); + } + ASSERT_NOK(reader->GetPreviousBatchFileRowId(array->length())); + } + ASSERT_EQ(cursor, static_cast(std::count_if( + ids.begin(), ids.end(), [](uint32_t id) { return id < 5; }))); + ASSERT_NOK(reader->GetPreviousBatchFileRowId(0)); + } + // Removing the bitmap must restore ordinary contiguous file row IDs. + ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportSchema(*arrow::schema(data_type->fields()), &c_schema).ok()); + ASSERT_OK(reader->SetReadSchema(&c_schema, nullptr, std::nullopt)); + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch batch, reader->NextBatch()); + auto array = arrow::ImportArray(batch.first.get(), batch.second.get()).ValueOrDie(); + ASSERT_TRUE(array->Equals(source->Slice(0, std::min(batch_size, 5)))); + ASSERT_OK_AND_ASSIGN(uint64_t first_row, reader->GetPreviousBatchFileRowId(0)); + ASSERT_EQ(first_row, 0); + } + } +} + +TEST_F(AvroFileBatchReaderTest, TestBitmapSelectionAcrossBlocksAfterProjection) { + auto data_type = arrow::struct_( + {arrow::field("id", arrow::int32()), arrow::field("payload", arrow::utf8())}); + std::string json = "["; + for (int32_t i = 0; i < 128; ++i) { + if (i != 0) { + json += ","; + } + json += fmt::format(R"([{},"{}"])", i, std::string(16384, 'a' + i % 26)); + } + json += "]"; + auto source = arrow::ipc::internal::json::ArrayFromJSON(data_type, json).ValueOrDie(); + for (const std::string& compression : {std::string("null"), std::string("deflate")}) { + const std::string path = PathUtil::JoinPath(dir_->Str(), compression + ".avro"); + WriteData(source, path, compression); + // Fresh reader, partial projection, full projection, and an exhausted bitmap that + // stopped before physical EOF must all support the same subsequent selections. + for (int32_t preparation : {0, 1, 2, 3}) { + SCOPED_TRACE(preparation); + for (int32_t batch_size : {1, 7, 256}) { + ASSERT_OK_AND_ASSIGN(std::shared_ptr builder, + file_format_->CreateReaderBuilder(batch_size)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr input, fs_->Open(path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + builder->Build(input)); + if (preparation != 0) { + ArrowSchema probe_schema; + ASSERT_TRUE( + arrow::ExportSchema(*arrow::schema({data_type->field(0)}), &probe_schema) + .ok()); + const std::optional selection = + preparation == 3 + ? std::optional(RoaringBitmap32::From({0, 64})) + : std::nullopt; + ASSERT_OK(reader->SetReadSchema(&probe_schema, nullptr, selection)); + do { + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch batch, reader->NextBatch()); + if (BatchReader::IsEofBatch(batch)) { + break; + } + ASSERT_TRUE(arrow::ImportArray(batch.first.get(), batch.second.get()).ok()); + } while (preparation != 1); + } + for (const std::vector& ids : + std::vector>{{}, + {127}, + {0, 1, 64, 65, 127}, + {31, 32, 33}, + {0, 128}, + {999}, + {0, std::numeric_limits::min()}}) { + ArrowSchema full_schema; + ASSERT_TRUE( + arrow::ExportSchema(*arrow::schema(data_type->fields()), &full_schema) + .ok()); + ASSERT_OK( + reader->SetReadSchema(&full_schema, nullptr, RoaringBitmap32::From(ids))); + size_t cursor = 0; + while (true) { + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch batch, reader->NextBatch()); + if (BatchReader::IsEofBatch(batch)) { + break; + } + auto array = + arrow::ImportArray(batch.first.get(), batch.second.get()).ValueOrDie(); + for (int64_t i = 0; i < array->length(); ++i) { + ASSERT_LT(cursor, ids.size()); + ASSERT_OK_AND_ASSIGN(uint64_t row, + reader->GetPreviousBatchFileRowId(i)); + ASSERT_EQ(row, ids[cursor++]); + ASSERT_TRUE(array->Slice(i, 1)->Equals(source->Slice(row, 1))); + } + } + ASSERT_EQ(cursor, std::count_if(ids.begin(), ids.end(), [](int32_t id) { + return static_cast(id) < 128; + })); + } + ArrowSchema full_schema; + ASSERT_TRUE( + arrow::ExportSchema(*arrow::schema(data_type->fields()), &full_schema).ok()); + ASSERT_OK(reader->SetReadSchema(&full_schema, nullptr, std::nullopt)); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr restored, + paimon::test::ReadResultCollector::CollectResult(reader.get())); + ASSERT_TRUE(restored->Equals(std::make_shared(source))); + } + } + } +} + +TEST_F(AvroFileBatchReaderTest, TestBitmapSelectionAvoidsUnnecessaryReads) { + const std::string path = PathUtil::JoinPath(dir_->Str(), "indexed.avro"); + auto data_type = arrow::struct_( + {arrow::field("id", arrow::int32()), arrow::field("payload", arrow::utf8())}); + std::string json = "["; + for (int32_t i = 0; i < 256; ++i) { + if (i != 0) { + json += ","; + } + json += fmt::format(R"([{},"{}"])", i, std::string(16384, 'a' + i % 26)); + } + json += "]"; + auto source = arrow::ipc::internal::json::ArrayFromJSON(data_type, json).ValueOrDie(); + WriteData(source, path, "null"); + std::string bytes; + ASSERT_OK(fs_->ReadFile(path, &bytes)); + auto input = std::make_shared(bytes); + ASSERT_OK_AND_ASSIGN(std::shared_ptr builder, + file_format_->CreateReaderBuilder(/*batch_size=*/8)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, builder->Build(input)); + auto set_selection = [&](const std::optional& selection) -> Status { + ArrowSchema schema; + PAIMON_RETURN_NOT_OK_FROM_ARROW( + arrow::ExportSchema(*arrow::schema(data_type->fields()), &schema)); + return reader->SetReadSchema(&schema, nullptr, selection); + }; + ASSERT_OK(set_selection(RoaringBitmap32::From({0}))); + input->bytes_read = 0; + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch first, reader->NextBatch()); + auto first_array = arrow::ImportArray(first.first.get(), first.second.get()).ValueOrDie(); + ASSERT_TRUE(first_array->Equals(source->Slice(0, 1))); + ASSERT_LT(input->bytes_read, static_cast(bytes.size()) / 2); + const int64_t prefix_reads = input->bytes_read; + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch eof, reader->NextBatch()); + ASSERT_TRUE(BatchReader::IsEofBatch(eof)); + ASSERT_EQ(input->bytes_read, prefix_reads); + ASSERT_NOK(reader->GetPreviousBatchFileRowId(0)); + + // Exhausting a selection must not be treated as the end of the file. + ASSERT_OK_AND_ASSIGN(uint64_t row_count, reader->GetNumberOfRows()); + ASSERT_EQ(row_count, 256); + ASSERT_OK(set_selection(RoaringBitmap32::From({255}))); + input->bytes_read = 0; + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch last, reader->NextBatch()); + auto last_array = arrow::ImportArray(last.first.get(), last.second.get()).ValueOrDie(); + ASSERT_TRUE(last_array->Equals(source->Slice(255, 1))); + ASSERT_OK_AND_ASSIGN(uint64_t last_row, reader->GetPreviousBatchFileRowId(0)); + ASSERT_EQ(last_row, 255); + const int64_t sequential_reads = input->bytes_read; + + ASSERT_OK(set_selection(std::nullopt)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr all, + paimon::test::ReadResultCollector::CollectResult(reader.get())); + ASSERT_TRUE(all->Equals(std::make_shared(source))); + ASSERT_OK(set_selection(RoaringBitmap32::From({255}))); + input->bytes_read = 0; + input->seek_count = 0; + ASSERT_OK_AND_ASSIGN(last, reader->NextBatch()); + last_array = arrow::ImportArray(last.first.get(), last.second.get()).ValueOrDie(); + ASSERT_TRUE(last_array->Equals(source->Slice(255, 1))); + ASSERT_OK_AND_ASSIGN(last_row, reader->GetPreviousBatchFileRowId(0)); + ASSERT_EQ(last_row, 255); + ASSERT_GT(input->seek_count, 0); + ASSERT_LT(input->bytes_read, sequential_reads); +} + +TEST_F(AvroFileBatchReaderTest, TestBitmapSelectionOnEmptyFileAndReset) { + const std::string path = PathUtil::JoinPath(dir_->Str(), "empty.avro"); + auto schema = ::avro::compileJsonSchemaFromString( + R"({"type":"record","name":"row","fields":[{"name":"id","type":"int"}]})"); + ::avro::DataFileWriterBase writer(path.c_str(), schema, /*syncInterval=*/1024); + writer.close(); + ASSERT_OK_AND_ASSIGN(std::shared_ptr builder, + file_format_->CreateReaderBuilder(/*batch_size=*/2)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr input, fs_->Open(path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, builder->Build(input)); + // Reach EOF first so the reader has a complete, empty block index. + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch initial, reader->NextBatch()); + ASSERT_TRUE(BatchReader::IsEofBatch(initial)); + for (const std::optional& selection : + std::vector>{ + RoaringBitmap32(), RoaringBitmap32::From({0, 99}), std::nullopt}) { + ASSERT_OK_AND_ASSIGN(std::unique_ptr read_schema, reader->GetFileSchema()); + ASSERT_OK(reader->SetReadSchema(read_schema.get(), nullptr, selection)); + ASSERT_NOK(reader->GetPreviousBatchFileRowId(0)); + for (int32_t i = 0; i < 2; ++i) { + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch batch, reader->NextBatch()); + ASSERT_TRUE(BatchReader::IsEofBatch(batch)); + ASSERT_NOK(reader->GetPreviousBatchFileRowId(0)); + } + ASSERT_OK_AND_ASSIGN(uint64_t row_count, reader->GetNumberOfRows()); + ASSERT_EQ(row_count, 0); + } +} + +TEST_F(AvroFileBatchReaderTest, TestBlockIndexReset) { + AvroFileBatchReader::BlockIndex index; + index.Observe(/*row=*/0, /*file_offset=*/100); + index.Reset(); + // A partial index must be discarded before rebuilding from the first row. + index.Observe(/*row=*/0, /*file_offset=*/200); + index.Finish(/*row_count=*/2); + ASSERT_EQ(index.RowCount(), 2); + ASSERT_TRUE(index.Locate(1)); + ASSERT_EQ(index.Locate(1)->file_offset, 200); + + // A completed index remains usable after a projection or selection reset. + index.Reset(); + ASSERT_EQ(index.RowCount(), 2); + ASSERT_TRUE(index.Locate(1)); + ASSERT_EQ(index.Locate(1)->file_offset, 200); + + AvroFileBatchReader::BlockIndex oversized; + for (uint64_t row = 0; row <= 64 * 1024; ++row) { + oversized.Observe(row, static_cast(row)); + } + oversized.Finish(/*row_count=*/64 * 1024 + 1); + ASSERT_FALSE(oversized.RowCount()); + for (int32_t reset = 0; reset < 2; ++reset) { + oversized.Reset(); + // An immutable file already known to exceed the limit must not start a new index. + oversized.Observe(/*row=*/0, /*file_offset=*/100); + oversized.Finish(/*row_count=*/1); + ASSERT_FALSE(oversized.RowCount()); + ASSERT_FALSE(oversized.Locate(0)); + } +} + +TEST_F(AvroFileBatchReaderTest, TestBitmapSelectionBeyondBlockIndexLimit) { + const std::string path = PathUtil::JoinPath(dir_->Str(), "many-blocks.avro"); + auto schema = ::avro::compileJsonSchemaFromString( + R"({"type":"record","name":"row","fields":[{"name":"id","type":"int"}]})"); + // One row per block exceeds the bounded index without requiring a large payload. + constexpr int32_t kRowCount = 64 * 1024 + 2; + ::avro::DataFileWriterBase writer(path.c_str(), schema, /*syncInterval=*/1024); + for (int32_t i = 0; i < kRowCount; ++i) { + writer.encoder().encodeInt(i); + writer.incr(); + writer.flush(); + } + writer.close(); + ASSERT_OK_AND_ASSIGN(std::shared_ptr builder, + file_format_->CreateReaderBuilder(/*batch_size=*/1024)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr input, fs_->Open(path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, builder->Build(input)); + int64_t scanned_rows = 0; + while (true) { + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch batch, reader->NextBatch()); + if (BatchReader::IsEofBatch(batch)) { + break; + } + auto array = arrow::ImportArray(batch.first.get(), batch.second.get()).ValueOrDie(); + scanned_rows += array->length(); + } + ASSERT_EQ(scanned_rows, kRowCount); + // Both sides of the index limit must remain readable after falling back to sequential skips. + const std::vector ids = {0, 64 * 1024 - 1, 64 * 1024, kRowCount - 1}; + ASSERT_OK_AND_ASSIGN(std::unique_ptr read_schema, reader->GetFileSchema()); + ASSERT_OK(reader->SetReadSchema(read_schema.get(), nullptr, RoaringBitmap32::From(ids))); + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch selected, reader->NextBatch()); + ASSERT_FALSE(BatchReader::IsEofBatch(selected)); + auto array = arrow::ImportArray(selected.first.get(), selected.second.get()).ValueOrDie(); + auto rows = checked_pointer_cast(array); + auto values = checked_pointer_cast(rows->field(0)); + ASSERT_EQ(values->length(), ids.size()); + for (size_t i = 0; i < ids.size(); ++i) { + ASSERT_EQ(values->Value(i), ids[i]); + ASSERT_OK_AND_ASSIGN(uint64_t file_row, reader->GetPreviousBatchFileRowId(i)); + ASSERT_EQ(file_row, ids[i]); + } + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch eof, reader->NextBatch()); + ASSERT_TRUE(BatchReader::IsEofBatch(eof)); + ASSERT_NOK(reader->GetPreviousBatchFileRowId(0)); + ASSERT_OK_AND_ASSIGN(read_schema, reader->GetFileSchema()); + ASSERT_OK(reader->SetReadSchema(read_schema.get(), nullptr, std::nullopt)); + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch restored, reader->NextBatch()); + ASSERT_FALSE(BatchReader::IsEofBatch(restored)); + ASSERT_EQ(restored.first->length, 1024); + ASSERT_OK_AND_ASSIGN(uint64_t first_row, reader->GetPreviousBatchFileRowId(0)); + ASSERT_EQ(first_row, 0); + ASSERT_TRUE(arrow::ImportArray(restored.first.get(), restored.second.get()).ok()); +} + TEST_F(AvroFileBatchReaderTest, TestGetPreviousBatchFileRowId) { std::string path = paimon::test::GetDataDir() + "/avro/append_simple.db/" diff --git a/test/inte/scan_and_read_inte_test.cpp b/test/inte/scan_and_read_inte_test.cpp index 093d1385..5e5672e9 100644 --- a/test/inte/scan_and_read_inte_test.cpp +++ b/test/inte/scan_and_read_inte_test.cpp @@ -23,6 +23,7 @@ #include #include #include +#include #include #include #include @@ -48,6 +49,7 @@ #include "paimon/core/table/source/deletion_file.h" #include "paimon/core/table/source/key_value_table_read.h" #include "paimon/defs.h" +#include "paimon/executor.h" #include "paimon/fs/file_system.h" #include "paimon/fs/local/local_file_system.h" #include "paimon/memory/memory_pool.h" @@ -64,6 +66,7 @@ #include "paimon/table/source/startup_mode.h" #include "paimon/table/source/table_read.h" #include "paimon/table/source/table_scan.h" +#include "paimon/testing/utils/counting_cache_test_utils.h" #include "paimon/testing/utils/io_exception_helper.h" #include "paimon/testing/utils/read_result_collector.h" #include "paimon/testing/utils/test_helper.h" @@ -426,6 +429,302 @@ TEST_P(ScanAndReadInteTest, TestWithAppendSnapshot3) { ASSERT_TRUE(expected->Equals(read_result)) << read_result->ToString(); } +TEST(SelectiveManifestDecodeInteTest, TestPartitionedPointLookup) { + // Scale the many-partition/many-bucket point-lookup workload down for CI. + constexpr int32_t kPartitions = 4; + constexpr int32_t kBuckets = 32; + constexpr int32_t kCandidateKeys = 1024; + const std::string payload(16 * 1024, 'x'); + auto dir = UniqueTestDirectory::Create("local"); + ASSERT_TRUE(dir); + arrow::FieldVector fields = {arrow::field("p", arrow::int32()), + arrow::field("rowkey", arrow::utf8()), + arrow::field("payload", arrow::utf8())}; + auto schema = arrow::schema(fields); + auto data_type = arrow::struct_(fields); + std::map options = {{Options::FILE_FORMAT, "orc"}, + {Options::BUCKET, fmt::format("{}", kBuckets)}, + {Options::BUCKET_KEY, "rowkey"}, + {Options::MANIFEST_COMPRESSION, "zstd"}, + {Options::MANIFEST_TARGET_FILE_SIZE, "64 mb"}}; + ASSERT_OK_AND_ASSIGN(std::unique_ptr helper, + TestHelper::Create(dir->Str(), schema, {"p"}, {}, options, false)); + const std::string table_path = PathUtil::JoinPath(dir->Str(), "foo.db/bar"); + + std::vector candidate_rows; + for (int32_t i = 0; i < kCandidateKeys; ++i) { + candidate_rows.push_back(fmt::format(R"(["key{:04}"])", i)); + } + auto key_array = + arrow::ipc::internal::json::ArrayFromJSON( + arrow::struct_({fields[1]}), fmt::format("[{}]", fmt::join(candidate_rows, ","))) + .ValueOrDie(); + ArrowArray c_keys; + ArrowSchema c_key_schema; + ASSERT_TRUE(arrow::ExportArray(*key_array, &c_keys, &c_key_schema).ok()); + ASSERT_OK_AND_ASSIGN(std::unique_ptr calculator, + BucketIdCalculator::Create(false, kBuckets, GetDefaultPool())); + std::vector bucket_ids(kCandidateKeys); + ASSERT_OK(calculator->CalculateBucketIds(&c_keys, &c_key_schema, bucket_ids.data())); + std::vector keys(kBuckets); + for (int32_t i = 0; i < kCandidateKeys; ++i) { + if (keys[bucket_ids[i]].empty()) { + keys[bucket_ids[i]] = fmt::format("key{:04}", i); + } + } + std::vector> batches; + for (int32_t p = 0; p < kPartitions; ++p) { + for (int32_t bucket = 0; bucket < kBuckets; ++bucket) { + ASSERT_FALSE(keys[bucket].empty()); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr batch, + TestHelper::MakeRecordBatch( + data_type, fmt::format(R"([[{},"{}","{}"]])", p, keys[bucket], payload), + {{"p", fmt::format("{}", p)}}, bucket, {})); + batches.push_back(std::move(batch)); + } + } + ASSERT_OK(helper->WriteAndCommit(std::move(batches), 0, std::nullopt)); + const std::string& lookup_key = keys.back(); + auto predicate = + PredicateBuilder::Equal(1, "rowkey", FieldType::STRING, + Literal(FieldType::STRING, lookup_key.data(), lookup_key.size())); + ASSERT_OK_AND_ASSIGN(std::shared_ptr executor, CreateDefaultExecutor(1)); + // Plans and cache values can outlive a scan and reference its allocator. + // Keep every measured pool alive until all plans and caches have been destroyed. + std::vector> scan_pools; + + struct ScanResult { + std::shared_ptr plan; + uint64_t cache_enabled; + uint64_t cache_hit; + int64_t peak_bytes; + }; + auto scan = [&](bool lazy_decode, const std::shared_ptr& cache, + bool entry_cache) -> Result { + std::shared_ptr pool = GetMemoryPool(); + scan_pools.push_back(pool); + ScanContextBuilder builder(table_path); + builder.WithMemoryPool(pool) + .WithExecutor(executor) + .WithCache(cache) + .AddOption(Options::SCAN_MANIFEST_ENTRY_CACHE_MAX_SNAPSHOTS, entry_cache ? "3" : "0") + .AddOption(Options::SCAN_MANIFEST_ENTRY_LAZY_DECODE_ENABLED, + lazy_decode ? "true" : "false") + .AddOption(Options::READ_BATCH_SIZE, "128"); + if (entry_cache) { + // No explicit bucket or partition filter: the predicate must select the bucket. + builder.SetPredicate(predicate); + } + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr context, builder.Finish()); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr table_scan, + TableScan::Create(std::move(context))); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr plan, table_scan->CreatePlan()); + PAIMON_ASSIGN_OR_RAISE(uint64_t enabled, table_scan->GetMetrics()->GetCounter( + ScanMetrics::LAST_SNAPSHOT_CACHE_ENABLED)); + PAIMON_ASSIGN_OR_RAISE(uint64_t hit, table_scan->GetMetrics()->GetCounter( + ScanMetrics::LAST_SNAPSHOT_CACHE_HIT)); + return ScanResult{plan, enabled, hit, pool->MaxMemoryUsage()}; + }; + auto check_read = [&](const ScanResult& result, int32_t partition_count, bool historical) { + ASSERT_EQ(result.cache_enabled, 1); + ASSERT_EQ(result.plan->Splits().size(), partition_count); + for (const auto& split : result.plan->Splits()) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + ASSERT_TRUE(data_split->Bucket() == kBuckets - 1 || + (historical && data_split->Bucket() == kBuckets / 2 - 1)); + ASSERT_EQ(data_split->GetFileList().size(), 1); + } + ReadContextBuilder builder(table_path); + builder.SetPredicate(predicate).EnablePredicateFilter(true); + ASSERT_OK_AND_ASSIGN(std::unique_ptr context, builder.Finish()); + ASSERT_OK_AND_ASSIGN(std::unique_ptr read, + TableRead::Create(std::move(context))); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + read->CreateReader(result.plan->Splits())); + ASSERT_OK_AND_ASSIGN(std::shared_ptr rows, + ReadResultCollector::CollectResult(std::move(reader))); + ASSERT_TRUE(rows); + ASSERT_EQ(rows->length(), partition_count); + std::set partitions; + for (const auto& chunk : rows->chunks()) { + auto values = std::dynamic_pointer_cast(chunk); + ASSERT_TRUE(values); + auto p = std::dynamic_pointer_cast(values->GetFieldByName("p")); + auto key = + std::dynamic_pointer_cast(values->GetFieldByName("rowkey")); + auto data = + std::dynamic_pointer_cast(values->GetFieldByName("payload")); + ASSERT_TRUE(p && key && data); + for (int64_t i = 0; i < values->length(); ++i) { + ASSERT_TRUE(partitions.insert(p->Value(i)).second); + ASSERT_EQ(key->GetString(i), lookup_key); + ASSERT_EQ(data->GetString(i), payload); + } + } + for (int32_t p = 0; p < partition_count; ++p) { + ASSERT_EQ(partitions.count(p), 1); + } + }; + + std::vector> caches; + std::vector peaks; + for (bool lazy_decode : {false, true}) { + auto cache = std::make_shared( + std::map{{CacheKind::MANIFEST, 64 * 1024 * 1024}, + {CacheKind::SNAPSHOT_LIVE_MANIFEST, 64 * 1024 * 1024}}); + caches.push_back(cache); + // Warm only raw manifest bytes, with a separate pool, before measuring either mode. + ASSERT_OK_AND_ASSIGN(ScanResult warm, scan(false, cache, false)); + ASSERT_EQ(warm.cache_enabled, 0); + size_t file_count = 0; + for (const auto& split : warm.plan->Splits()) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + file_count += data_split->GetFileList().size(); + } + ASSERT_EQ(file_count, kPartitions * kBuckets); + const int64_t supplier_calls = cache->SupplierCallCount(CacheKind::MANIFEST); + ASSERT_GT(supplier_calls, 0); + ASSERT_OK_AND_ASSIGN(ScanResult measured, scan(lazy_decode, cache, true)); + ASSERT_EQ(measured.cache_hit, 0); + ASSERT_NO_FATAL_FAILURE(check_read(measured, kPartitions, false)); + ASSERT_EQ(cache->SupplierCallCount(CacheKind::MANIFEST), supplier_calls); + peaks.push_back(measured.peak_bytes); + RecordProperty(lazy_decode ? "selective_peak_bytes" : "baseline_peak_bytes", + fmt::format("{}", measured.peak_bytes)); + ASSERT_OK_AND_ASSIGN(ScanResult hit, scan(lazy_decode, cache, true)); + ASSERT_EQ(hit.cache_hit, 1); + ASSERT_NO_FATAL_FAILURE(check_read(hit, kPartitions, false)); + // Manifest lists are still consulted, but their bytes must also come from the cache. + ASSERT_EQ(cache->SupplierCallCount(CacheKind::MANIFEST), supplier_calls); + } + // Selective decoding also applies when only manifest entries, not raw bytes, are cached. + auto entry_cache_only = + std::make_shared(CacheKind::SNAPSHOT_LIVE_MANIFEST, 64 * 1024 * 1024); + ASSERT_OK_AND_ASSIGN(ScanResult uncached_manifest, scan(true, entry_cache_only, true)); + ASSERT_EQ(uncached_manifest.cache_hit, 0); + ASSERT_NO_FATAL_FAILURE(check_read(uncached_manifest, kPartitions, false)); + ASSERT_GT(entry_cache_only->GetCount(CacheKind::MANIFEST), 0); + ASSERT_EQ(entry_cache_only->SupplierCallCount(CacheKind::MANIFEST), 0); + RecordProperty("uncached_manifest_peak_bytes", fmt::format("{}", uncached_manifest.peak_bytes)); + // Wide, real column statistics make full Arrow materialization observable without timing gates. + ASSERT_LT(peaks[1], peaks[0] / 2); + ASSERT_LT(uncached_manifest.peak_bytes, peaks[0] / 2); + + // Change both schema ID and bucket count, then append the same key in a new partition. + // Old files must still be read from bucket 31 while the new file belongs to bucket 15. + helper.reset(); + options[Options::BUCKET] = fmt::format("{}", kBuckets / 2); + ASSERT_OK(TestHelper::WriteNextSchema( + dir->GetFileSystem(), table_path, + {DataField(0, fields[0]), DataField(1, fields[1]), DataField(2, fields[2])}, 2, options)); + ASSERT_OK_AND_ASSIGN(helper, TestHelper::Create(table_path, options, false)); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr next_batch, + TestHelper::MakeRecordBatch( + data_type, fmt::format(R"([[{},"{}","{}"]])", kPartitions, lookup_key, payload), + {{"p", fmt::format("{}", kPartitions)}}, kBuckets / 2 - 1, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(next_batch), 1, std::nullopt)); + for (size_t i = 0; i < caches.size(); ++i) { + ASSERT_OK_AND_ASSIGN(ScanResult next, scan(i != 0, caches[i], true)); + ASSERT_EQ(next.cache_hit, 0); + ASSERT_EQ(next.plan->SnapshotId(), 2); + ASSERT_NO_FATAL_FAILURE(check_read(next, kPartitions + 1, true)); + ASSERT_OK_AND_ASSIGN(ScanResult hit, scan(i != 0, caches[i], true)); + ASSERT_EQ(hit.cache_hit, 1); + ASSERT_NO_FATAL_FAILURE(check_read(hit, kPartitions + 1, true)); + } +} + +TEST(SelectiveManifestDecodeInteTest, TestSchemaEvolutionPreservesBucketPruning) { + constexpr int32_t kBuckets = 4; + auto key_field = arrow::field("rowkey", arrow::int32()); + auto added_field = arrow::field("value", arrow::int32()); + constexpr int32_t kKey = -1; + auto array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({key_field}), + fmt::format("[[{}]]", kKey)) + .ValueOrDie(); + ArrowArray c_array; + ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportArray(*array, &c_array, &c_schema).ok()); + ASSERT_OK_AND_ASSIGN(std::unique_ptr calculator, + BucketIdCalculator::Create(false, kBuckets, GetDefaultPool())); + int32_t bucket = 0; + ASSERT_OK(calculator->CalculateBucketIds(&c_array, &c_schema, &bucket)); + auto dir = UniqueTestDirectory::Create("local"); + ASSERT_TRUE(dir); + std::map options = {{Options::FILE_FORMAT, "orc"}, + {Options::BUCKET, fmt::format("{}", kBuckets)}, + {Options::BUCKET_KEY, "rowkey"}}; + ASSERT_OK_AND_ASSIGN( + std::unique_ptr helper, + TestHelper::Create(dir->Str(), arrow::schema({key_field}), {}, {}, options, false)); + const std::string table_path = PathUtil::JoinPath(dir->Str(), "foo.db/bar"); + ASSERT_OK_AND_ASSIGN(std::unique_ptr batch, + TestHelper::MakeRecordBatch(arrow::struct_({key_field}), + fmt::format("[[{}]]", kKey), {}, bucket, {})); + std::vector> batches; + batches.push_back(std::move(batch)); + ASSERT_OK(helper->WriteAndCommit(std::move(batches), 0, std::nullopt)); + helper.reset(); + ASSERT_OK(TestHelper::WriteNextSchema(dir->GetFileSystem(), table_path, + {DataField(0, key_field), DataField(1, added_field)}, 1, + options)); + auto predicate = PredicateBuilder::Equal(0, "rowkey", FieldType::INT, Literal(kKey)); + auto expected = + arrow::ipc::internal::json::ArrayFromJSON( + arrow::struct_({arrow::field("_VALUE_KIND", arrow::int8()), key_field, added_field}), + fmt::format("[[0, {}, null]]", kKey)) + .ValueOrDie(); + for (bool cache_manifest_bytes : {false, true}) { + SCOPED_TRACE(cache_manifest_bytes); + for (bool lazy_decode : {false, true}) { + SCOPED_TRACE(lazy_decode); + std::map capacities = { + {CacheKind::SNAPSHOT_LIVE_MANIFEST, 16 * 1024 * 1024}}; + if (cache_manifest_bytes) { + capacities.emplace(CacheKind::MANIFEST, 16 * 1024 * 1024); + } + auto cache = std::make_shared(capacities); + for (int32_t attempt = 0; attempt < 2; ++attempt) { + ScanContextBuilder builder(table_path); + builder.SetPredicate(predicate) + .WithCache(cache) + .AddOption(Options::SCAN_MANIFEST_ENTRY_CACHE_MAX_SNAPSHOTS, "3") + .AddOption(Options::SCAN_MANIFEST_ENTRY_LAZY_DECODE_ENABLED, + lazy_decode ? "true" : "false"); + ASSERT_OK_AND_ASSIGN(std::unique_ptr context, builder.Finish()); + ASSERT_OK_AND_ASSIGN(std::unique_ptr scan, + TableScan::Create(std::move(context))); + ASSERT_OK_AND_ASSIGN(std::shared_ptr plan, scan->CreatePlan()); + ASSERT_EQ(plan->Splits().size(), 1); + auto split = std::dynamic_pointer_cast(plan->Splits()[0]); + ASSERT_TRUE(split); + ASSERT_EQ(split->Bucket(), bucket); + ASSERT_OK_AND_ASSIGN(uint64_t hit, scan->GetMetrics()->GetCounter( + ScanMetrics::LAST_SNAPSHOT_CACHE_HIT)); + ASSERT_EQ(hit, attempt); + ReadContextBuilder read_builder(table_path); + read_builder.SetPredicate(predicate).EnablePredicateFilter(true); + ASSERT_OK_AND_ASSIGN(std::unique_ptr read_context, + read_builder.Finish()); + ASSERT_OK_AND_ASSIGN(std::unique_ptr read, + TableRead::Create(std::move(read_context))); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + read->CreateReader(plan->Splits())); + ASSERT_OK_AND_ASSIGN(std::shared_ptr rows, + ReadResultCollector::CollectResult(std::move(reader))); + ASSERT_TRUE(rows); + ASSERT_TRUE(rows->Equals(std::make_shared(expected))) + << rows->type()->ToString() << "\n" + << rows->ToString(); + } + } + } +} + TEST_P(ScanAndReadInteTest, TestWithAppendBucketKeyPointLookup) { for (const auto& key_type : {arrow::utf8(), arrow::binary()}) { SCOPED_TRACE(key_type->ToString());