From da852b5c56ff5e4905ae65b88b1366dc084e65d8 Mon Sep 17 00:00:00 2001 From: Zehua Zou Date: Mon, 21 Sep 2026 18:01:53 +0800 Subject: [PATCH 1/2] refactor code and simplify src/iceberg/manifest/manifest_group.cc --- src/iceberg/manifest/manifest_group.cc | 455 ++++++++++-------------- src/iceberg/manifest/manifest_group.h | 7 +- src/iceberg/test/manifest_group_test.cc | 136 +++++++ 3 files changed, 331 insertions(+), 267 deletions(-) diff --git a/src/iceberg/manifest/manifest_group.cc b/src/iceberg/manifest/manifest_group.cc index 87ee857e5..8f2970464 100644 --- a/src/iceberg/manifest/manifest_group.cc +++ b/src/iceberg/manifest/manifest_group.cc @@ -131,6 +131,148 @@ ManifestGroup::~ManifestGroup() = default; ManifestGroup::ManifestGroup(ManifestGroup&&) noexcept = default; ManifestGroup& ManifestGroup::operator=(ManifestGroup&&) noexcept = default; +class ManifestGroup::PlanningContext { + public: + // `group` is borrowed; its planning configuration must remain unchanged for the + // lifetime of this context. + static Result Make(const ManifestGroup& group, + std::vector columns) { + std::unique_ptr file_evaluator; + if (!group.file_filter_ || group.file_filter_->op() == Expression::Operation::kTrue) { + return PlanningContext(group, std::move(file_evaluator), std::move(columns)); + } + + auto data_file_schema = DataFileFilterSchema(); + ICEBERG_ASSIGN_OR_RAISE( + file_evaluator, + Evaluator::Make(*data_file_schema, group.file_filter_, group.case_sensitive_)); + if (std::ranges::contains(columns, Schema::kAllColumns)) { + return PlanningContext(group, std::move(file_evaluator), std::move(columns)); + } + + ICEBERG_ASSIGN_OR_RAISE( + auto bound_file_filter, + Binder::Bind(*data_file_schema, group.file_filter_, group.case_sensitive_)); + ICEBERG_ASSIGN_OR_RAISE(auto referenced_field_ids, + ReferenceVisitor::GetReferencedFieldIds(bound_file_filter)); + + std::unordered_set selected_columns(columns.cbegin(), columns.cend()); + for (const auto field_id : referenced_field_ids) { + if (field_id == DataFile::kSpecIdFieldId) { + continue; + } + ICEBERG_ASSIGN_OR_RAISE(auto column_name, + data_file_schema->FindColumnNameById(field_id)); + if (!column_name.has_value()) { + continue; + } + + std::string column_name_str(column_name.value()); + if (selected_columns.insert(column_name_str).second) { + columns.push_back(std::move(column_name_str)); + } + } + return PlanningContext(group, std::move(file_evaluator), std::move(columns)); + } + + Result> MakeManifestEvaluator( + int32_t spec_id) const { + auto spec_iter = group_.specs_by_id_.find(spec_id); + ICEBERG_CHECK(spec_iter != group_.specs_by_id_.cend(), + "Cannot find partition spec for ID {}", spec_id); + + auto projector = Projections::Inclusive(*spec_iter->second, *group_.schema_, + group_.case_sensitive_); + ICEBERG_ASSIGN_OR_RAISE(auto partition_filter, + projector->Project(group_.data_filter_)); + ICEBERG_ASSIGN_OR_RAISE(partition_filter, And::Make(std::move(partition_filter), + group_.partition_filter_)); + return ManifestEvaluator::MakePartitionFilter(std::move(partition_filter), + spec_iter->second, *group_.schema_, + group_.case_sensitive_); + } + + Result> MakeResidualEvaluator( + int32_t spec_id) const { + auto spec_iter = group_.specs_by_id_.find(spec_id); + ICEBERG_CHECK(spec_iter != group_.specs_by_id_.cend(), + "Cannot find partition spec for ID {}", spec_id); + + ICEBERG_ASSIGN_OR_RAISE( + auto evaluator, + ResidualEvaluator::Make( + (group_.ignore_residuals_ ? True::Instance() : group_.data_filter_), + *spec_iter->second, *group_.schema_, group_.case_sensitive_)); + return std::shared_ptr(std::move(evaluator)); + } + + Result ShouldReadManifest(const ManifestFile& manifest, + const ManifestEvaluator& evaluator) const { + ICEBERG_ASSIGN_OR_RAISE(bool should_match, evaluator.Evaluate(manifest)); + const bool has_non_deleted_files = + manifest.has_added_files() || manifest.has_existing_files(); + const bool has_non_existing_files = + manifest.has_added_files() || manifest.has_deleted_files(); + const bool has_only_ignored_files = + (group_.ignore_deleted_ && !has_non_deleted_files) || + (group_.ignore_existing_ && !has_non_existing_files); + if (!should_match || has_only_ignored_files) { + if (group_.scan_metrics_) { + group_.scan_metrics_->skipped_data_manifests->Increment(1); + } + return false; + } + + if (group_.scan_metrics_) { + group_.scan_metrics_->scanned_data_manifests->Increment(1); + } + return true; + } + + Result ShouldKeepEntry(const ManifestEntry& entry) const { + if (group_.ignore_existing_ && entry.status == ManifestStatus::kExisting) { + if (group_.scan_metrics_) { + group_.scan_metrics_->skipped_data_files->Increment(1); + } + return false; + } + + ICEBERG_DCHECK(entry.data_file != nullptr, "Data file cannot be null"); + if (file_evaluator_ != nullptr) { + DataFileStructLike data_file(*entry.data_file); + ICEBERG_ASSIGN_OR_RAISE(bool should_match, file_evaluator_->Evaluate(data_file)); + if (!should_match) { + if (group_.scan_metrics_) { + group_.scan_metrics_->skipped_data_files->Increment(1); + } + return false; + } + } + + if (!group_.manifest_entry_predicate_(entry)) { + if (group_.scan_metrics_) { + group_.scan_metrics_->skipped_data_files->Increment(1); + } + return false; + } + + return true; + } + + const std::vector& columns() const { return columns_; } + + private: + PlanningContext(const ManifestGroup& group, std::unique_ptr file_evaluator, + std::vector columns) + : group_(group), + file_evaluator_(std::move(file_evaluator)), + columns_(std::move(columns)) {} + + const ManifestGroup& group_; + std::unique_ptr file_evaluator_; + std::vector columns_; +}; + class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { public: static Result Make(std::unique_ptr group) { @@ -141,20 +283,12 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { auto stats_projection = group->PrepareStatsProjection(delete_index->has_equality_deletes()); + ICEBERG_ASSIGN_OR_RAISE( + auto context, PlanningContext::Make(*group, std::move(stats_projection.columns))); - std::unique_ptr data_file_evaluator; - if (group->file_filter_ && - group->file_filter_->op() != Expression::Operation::kTrue) { - ICEBERG_ASSIGN_OR_RAISE( - data_file_evaluator, - Evaluator::Make(*DataFileFilterSchema(), group->file_filter_, - group->case_sensitive_)); - } - const bool drop_stats = stats_projection.drop_stats; - - return FileScanTaskStreamPtr(new FilePlanningStream( - std::move(group), std::move(delete_index), std::move(data_file_evaluator), - std::move(stats_projection.columns), drop_stats)); + return FileScanTaskStreamPtr( + new FilePlanningStream(std::move(group), std::move(delete_index), + std::move(context), stats_projection.drop_stats)); } Result>> NextImpl() override { @@ -165,24 +299,8 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { } auto [spec_id, value] = std::move(entry).value(); - if (group_->ignore_existing_ && value.status == ManifestStatus::kExisting) { - IncrementSkippedDataFiles(); - continue; - } - - ICEBERG_DCHECK(value.data_file != nullptr, "Data file cannot be null"); - if (data_file_evaluator_) { - DataFileStructLike data_file(*value.data_file); - ICEBERG_ASSIGN_OR_RAISE(bool should_match, - data_file_evaluator_->Evaluate(data_file)); - if (!should_match) { - IncrementSkippedDataFiles(); - continue; - } - } - - if (!group_->manifest_entry_predicate_(value)) { - IncrementSkippedDataFiles(); + ICEBERG_ASSIGN_OR_RAISE(bool should_keep, context_.ShouldKeepEntry(value)); + if (!should_keep) { continue; } @@ -211,40 +329,16 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { private: FilePlanningStream(std::unique_ptr group, std::unique_ptr delete_index, - std::unique_ptr data_file_evaluator, - std::vector columns, bool drop_stats) + PlanningContext context, bool drop_stats) : group_(std::move(group)), delete_index_(std::move(delete_index)), - data_file_evaluator_(std::move(data_file_evaluator)), - columns_(std::move(columns)), + context_(std::move(context)), drop_stats_(drop_stats) {} using TaggedEntry = std::pair; using TaggedStream = std::pair; - // FIXME: Perhaps refactor this concurrent/sequential stream state machine into a - // generic reusable ParallelStream utility, similar to Iceberg Java's - // ParallelIterable. Result> NextEntry() { - if (!group_->executor_.has_value()) { - while (true) { - if (!entry_stream_) { - ICEBERG_ASSIGN_OR_RAISE(bool opened, OpenNextManifest()); - if (!opened) { - return std::nullopt; - } - } - - ICEBERG_ASSIGN_OR_RAISE(auto entry, entry_stream_->Next()); - if (!entry.has_value()) { - entry_stream_.reset(); - continue; - } - return std::optional{std::in_place, current_spec_id_, - std::move(entry).value()}; - } - } - while (true) { if (next_batch_stream_ == batch_streams_.size()) { ICEBERG_ASSIGN_OR_RAISE(bool loaded, LoadNextManifestBatch()); @@ -270,21 +364,7 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { return cached->second.get(); } - auto spec_iter = group_->specs_by_id_.find(spec_id); - ICEBERG_CHECK(spec_iter != group_->specs_by_id_.cend(), - "Cannot find partition spec for ID {}", spec_id); - - const auto& spec = spec_iter->second; - auto projector = - Projections::Inclusive(*spec, *group_->schema_, group_->case_sensitive_); - ICEBERG_ASSIGN_OR_RAISE(auto partition_filter, - projector->Project(group_->data_filter_)); - ICEBERG_ASSIGN_OR_RAISE(partition_filter, And::Make(std::move(partition_filter), - group_->partition_filter_)); - ICEBERG_ASSIGN_OR_RAISE(auto evaluator, - ManifestEvaluator::MakePartitionFilter( - std::move(partition_filter), spec, *group_->schema_, - group_->case_sensitive_)); + ICEBERG_ASSIGN_OR_RAISE(auto evaluator, context_.MakeManifestEvaluator(spec_id)); auto* result = evaluator.get(); manifest_evaluators_.emplace(spec_id, std::move(evaluator)); return result; @@ -296,67 +376,23 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { return cached->second.get(); } - auto spec_iter = group_->specs_by_id_.find(spec_id); - ICEBERG_CHECK(spec_iter != group_->specs_by_id_.cend(), - "Cannot find partition spec for ID {}", spec_id); - - ICEBERG_ASSIGN_OR_RAISE( - auto evaluator, - ResidualEvaluator::Make( - (group_->ignore_residuals_ ? True::Instance() : group_->data_filter_), - *spec_iter->second, *group_->schema_, group_->case_sensitive_)); + ICEBERG_ASSIGN_OR_RAISE(auto evaluator, context_.MakeResidualEvaluator(spec_id)); auto* result = evaluator.get(); residual_evaluators_.emplace(spec_id, std::move(evaluator)); return result; } - Result ShouldReadManifest(const ManifestFile& manifest) { - ICEBERG_ASSIGN_OR_RAISE(auto evaluator, - GetManifestEvaluator(manifest.partition_spec_id)); - ICEBERG_ASSIGN_OR_RAISE(bool should_match, evaluator->Evaluate(manifest)); - const bool has_non_deleted_files = - manifest.has_added_files() || manifest.has_existing_files(); - const bool has_non_existing_files = - manifest.has_added_files() || manifest.has_deleted_files(); - const bool has_only_ignored_files = - (group_->ignore_deleted_ && !has_non_deleted_files) || - (group_->ignore_existing_ && !has_non_existing_files); - if (!should_match || has_only_ignored_files) { - IncrementSkippedDataManifests(); - return false; - } - - if (group_->scan_metrics_) { - group_->scan_metrics_->scanned_data_manifests->Increment(1); - } - return true; - } - - Result OpenNextManifest() { - while (next_manifest_ < group_->data_manifests_.size()) { - const auto& manifest = group_->data_manifests_[next_manifest_++]; - ICEBERG_ASSIGN_OR_RAISE(bool should_read, ShouldReadManifest(manifest)); - if (!should_read) { - continue; - } - - ICEBERG_ASSIGN_OR_RAISE(auto reader, group_->MakeReader(manifest, columns_)); - ICEBERG_ASSIGN_OR_RAISE(entry_stream_, group_->ignore_deleted_ - ? reader->LiveEntriesStream() - : reader->EntriesStream()); - current_spec_id_ = manifest.partition_spec_id; - return true; - } - return false; - } - Result LoadNextManifestBatch() { + const size_t batch_size = group_->executor_.has_value() ? kManifestReadBatchSize : 1; std::vector manifests; - manifests.reserve(kManifestReadBatchSize); + manifests.reserve(batch_size); while (next_manifest_ < group_->data_manifests_.size() && - manifests.size() < kManifestReadBatchSize) { + manifests.size() < batch_size) { const auto& manifest = group_->data_manifests_[next_manifest_++]; - ICEBERG_ASSIGN_OR_RAISE(bool should_read, ShouldReadManifest(manifest)); + ICEBERG_ASSIGN_OR_RAISE(auto evaluator, + GetManifestEvaluator(manifest.partition_spec_id)); + ICEBERG_ASSIGN_OR_RAISE(bool should_read, + context_.ShouldReadManifest(manifest, *evaluator)); if (should_read) { manifests.push_back(&manifest); } @@ -375,7 +411,7 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { group_->executor_, manifests, [this](const ManifestFile* manifest) -> Result> { ICEBERG_ASSIGN_OR_RAISE(auto reader, - group_->MakeReader(*manifest, columns_)); + group_->MakeReader(*manifest, context_.columns())); ICEBERG_ASSIGN_OR_RAISE(auto stream, group_->ignore_deleted_ ? reader->LiveEntriesStream() : reader->EntriesStream()); @@ -388,18 +424,6 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { return true; } - void IncrementSkippedDataManifests() { - if (group_->scan_metrics_) { - group_->scan_metrics_->skipped_data_manifests->Increment(1); - } - } - - void IncrementSkippedDataFiles() { - if (group_->scan_metrics_) { - group_->scan_metrics_->skipped_data_files->Increment(1); - } - } - void UpdateResultMetrics(const DataFile& data_file, const std::vector>& delete_files) { if (!group_->scan_metrics_) { @@ -420,15 +444,12 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { std::unique_ptr group_; std::unique_ptr delete_index_; - std::unique_ptr data_file_evaluator_; - std::vector columns_; + PlanningContext context_; std::unordered_map> manifest_evaluators_; std::unordered_map> residual_evaluators_; - ManifestEntryStreamPtr entry_stream_; std::vector batch_streams_; size_t next_manifest_ = 0; size_t next_batch_stream_ = 0; - int32_t current_spec_id_ = 0; bool drop_stats_; // Limit the number of manifest readers and streams retained by executor-backed @@ -522,31 +543,26 @@ Result ManifestGroup::PlanFilesStream() && { Result>> ManifestGroup::Plan( const CreateTasksFunction& create_tasks) { + delete_index_builder_.WithScanMetrics(scan_metrics_); + ICEBERG_ASSIGN_OR_RAISE(auto delete_index, delete_index_builder_.Build()); + + auto stats_projection = PrepareStatsProjection(delete_index->has_equality_deletes()); + ICEBERG_ASSIGN_OR_RAISE( + auto context, PlanningContext::Make(*this, std::move(stats_projection.columns))); + std::unordered_map> residual_cache; auto get_residual_evaluator = [&](int32_t spec_id) -> Result { - if (residual_cache.contains(spec_id)) { - return residual_cache[spec_id].get(); + auto cached = residual_cache.find(spec_id); + if (cached != residual_cache.end()) { + return cached->second.get(); } - auto spec_iter = specs_by_id_.find(spec_id); - ICEBERG_CHECK(spec_iter != specs_by_id_.cend(), - "Cannot find partition spec for ID {}", spec_id); - - const auto& spec = spec_iter->second; - ICEBERG_ASSIGN_OR_RAISE( - auto residual_evaluator, - ResidualEvaluator::Make((ignore_residuals_ ? True::Instance() : data_filter_), - *spec, *schema_, case_sensitive_)); - residual_cache[spec_id] = std::move(residual_evaluator); - - return residual_cache[spec_id].get(); + ICEBERG_ASSIGN_OR_RAISE(auto evaluator, context.MakeResidualEvaluator(spec_id)); + auto* result = evaluator.get(); + residual_cache.emplace(spec_id, std::move(evaluator)); + return result; }; - delete_index_builder_.WithScanMetrics(scan_metrics_); - ICEBERG_ASSIGN_OR_RAISE(auto delete_index, delete_index_builder_.Build()); - - auto stats_projection = PrepareStatsProjection(delete_index->has_equality_deletes()); - std::unordered_map> task_context_cache; auto get_task_context = [&](int32_t spec_id) -> Result { if (task_context_cache.contains(spec_id)) { @@ -569,7 +585,7 @@ Result>> ManifestGroup::Plan( return task_context_cache[spec_id].get(); }; - ICEBERG_ASSIGN_OR_RAISE(auto entry_groups, ReadEntries(stats_projection.columns)); + ICEBERG_ASSIGN_OR_RAISE(auto entry_groups, ReadEntries(context)); std::vector> all_tasks; for (auto& [spec_id, entries] : entry_groups) { @@ -583,7 +599,8 @@ Result>> ManifestGroup::Plan( } Result> ManifestGroup::Entries() { - ICEBERG_ASSIGN_OR_RAISE(auto entry_groups, ReadEntries(columns_)); + ICEBERG_ASSIGN_OR_RAISE(auto context, PlanningContext::Make(*this, columns_)); + ICEBERG_ASSIGN_OR_RAISE(auto entry_groups, ReadEntries(context)); std::vector all_entries; for (auto& [_, entries] : entry_groups) { @@ -595,37 +612,10 @@ Result> ManifestGroup::Entries() { } Result> ManifestGroup::MakeReader( - const ManifestFile& manifest, std::vector columns) { + const ManifestFile& manifest, const std::vector& columns) { ICEBERG_ASSIGN_OR_RAISE(auto reader, ManifestReader::Make(manifest, io_, schema_, specs_by_id_)); - if (file_filter_ && file_filter_->op() != Expression::Operation::kTrue && - !std::ranges::contains(columns, Schema::kAllColumns)) { - auto data_file_schema = DataFileFilterSchema(); - ICEBERG_ASSIGN_OR_RAISE( - auto bound_file_filter, - Binder::Bind(*data_file_schema, file_filter_, case_sensitive_)); - ICEBERG_ASSIGN_OR_RAISE(auto referenced_field_ids, - ReferenceVisitor::GetReferencedFieldIds(bound_file_filter)); - - std::unordered_set selected_columns(columns.cbegin(), columns.cend()); - for (const auto field_id : referenced_field_ids) { - if (field_id == DataFile::kSpecIdFieldId) { - continue; - } - ICEBERG_ASSIGN_OR_RAISE(auto column_name, - data_file_schema->FindColumnNameById(field_id)); - if (column_name.has_value()) { - std::string column_name_str(column_name.value()); - if (selected_columns.contains(column_name_str)) { - continue; - } - columns.push_back(std::move(column_name_str)); - selected_columns.insert(columns.back()); - } - } - } - reader->FilterRows(data_filter_) .FilterPartitions(partition_filter_) .CaseSensitive(case_sensitive_) @@ -659,101 +649,38 @@ ManifestGroup::StatsProjection ManifestGroup::PrepareStatsProjection( } Result>> -ManifestGroup::ReadEntries(const std::vector& columns) { +ManifestGroup::ReadEntries(PlanningContext& context) { const auto cache_capacity = static_cast(specs_by_id_.size()); auto get_manifest_evaluator = internal::MemoizeLru( - [this](int32_t spec_id) -> Result> { - auto spec_iter = specs_by_id_.find(spec_id); - ICEBERG_CHECK(spec_iter != specs_by_id_.cend(), - "Cannot find partition spec for ID {}", spec_id); - - auto projector = - Projections::Inclusive(*spec_iter->second, *schema_, case_sensitive_); - ICEBERG_ASSIGN_OR_RAISE(auto partition_filter, projector->Project(data_filter_)); - ICEBERG_ASSIGN_OR_RAISE(partition_filter, - And::Make(partition_filter, partition_filter_)); - ICEBERG_ASSIGN_OR_RAISE( - auto evaluator, ManifestEvaluator::MakePartitionFilter( - std::move(partition_filter), spec_iter->second, *schema_, - case_sensitive_)); + [&context](int32_t spec_id) -> Result> { + ICEBERG_ASSIGN_OR_RAISE(auto evaluator, context.MakeManifestEvaluator(spec_id)); return std::shared_ptr(std::move(evaluator)); }, cache_capacity); - const bool has_file_filter = - file_filter_ && file_filter_->op() != Expression::Operation::kTrue; - std::unique_ptr data_file_evaluator; - if (has_file_filter) { - ICEBERG_ASSIGN_OR_RAISE( - data_file_evaluator, - Evaluator::Make(*DataFileFilterSchema(), file_filter_, case_sensitive_)); - } - return ParallelCollect( executor_, data_manifests_, - [&](const ManifestFile& manifest) + [this, &context, &get_manifest_evaluator](const ManifestFile& manifest) -> Result>> { const int32_t spec_id = manifest.partition_spec_id; ICEBERG_ASSIGN_OR_RAISE(auto manifest_evaluator, get_manifest_evaluator(spec_id)); - ICEBERG_ASSIGN_OR_RAISE(bool should_match, - manifest_evaluator->Evaluate(manifest)); - if (!should_match) { - // Skip this manifest because it doesn't match partition filter - if (scan_metrics_) { - scan_metrics_->skipped_data_manifests->Increment(1); - } + ICEBERG_ASSIGN_OR_RAISE( + bool should_read, context.ShouldReadManifest(manifest, *manifest_evaluator)); + if (!should_read) { return {}; } - if (ignore_deleted_) { - // only scan manifests that have entries other than deletes - if (!manifest.has_added_files() && !manifest.has_existing_files()) { - if (scan_metrics_) scan_metrics_->skipped_data_manifests->Increment(1); - return {}; - } - } - if (ignore_existing_) { - // only scan manifests that have entries other than existing - if (!manifest.has_added_files() && !manifest.has_deleted_files()) { - if (scan_metrics_) scan_metrics_->skipped_data_manifests->Increment(1); - return {}; - } - } - - if (scan_metrics_) { - scan_metrics_->scanned_data_manifests->Increment(1); - } - - // Read manifest entries - ICEBERG_ASSIGN_OR_RAISE(auto reader, MakeReader(manifest, columns)); + ICEBERG_ASSIGN_OR_RAISE(auto reader, MakeReader(manifest, context.columns())); ICEBERG_ASSIGN_OR_RAISE( auto entries, ignore_deleted_ ? reader->LiveEntries() : reader->Entries()); std::unordered_map> manifest_result; - for (auto& entry : entries) { - if (ignore_existing_ && entry.status == ManifestStatus::kExisting) { - if (scan_metrics_) scan_metrics_->skipped_data_files->Increment(1); - continue; + ICEBERG_ASSIGN_OR_RAISE(bool should_keep, context.ShouldKeepEntry(entry)); + if (should_keep) { + manifest_result[spec_id].push_back(std::move(entry)); } - - if (data_file_evaluator != nullptr) { - DataFileStructLike data_file(*entry.data_file); - ICEBERG_ASSIGN_OR_RAISE(bool should_match, - data_file_evaluator->Evaluate(data_file)); - if (!should_match) { - if (scan_metrics_) scan_metrics_->skipped_data_files->Increment(1); - continue; - } - } - - if (!manifest_entry_predicate_(entry)) { - if (scan_metrics_) scan_metrics_->skipped_data_files->Increment(1); - continue; - } - - manifest_result[spec_id].push_back(std::move(entry)); } return manifest_result; }); diff --git a/src/iceberg/manifest/manifest_group.h b/src/iceberg/manifest/manifest_group.h index 50448aa61..8c5231553 100644 --- a/src/iceberg/manifest/manifest_group.h +++ b/src/iceberg/manifest/manifest_group.h @@ -178,6 +178,7 @@ class ICEBERG_EXPORT ManifestGroup : public ErrorCollector { const CreateTasksFunction& create_tasks); private: + class PlanningContext; class FilePlanningStream; struct StatsProjection { @@ -191,10 +192,10 @@ class ICEBERG_EXPORT ManifestGroup : public ErrorCollector { DeleteFileIndex::Builder&& delete_index_builder); Result>> ReadEntries( - const std::vector& columns); + PlanningContext& context); - Result> MakeReader(const ManifestFile& manifest, - std::vector columns); + Result> MakeReader( + const ManifestFile& manifest, const std::vector& columns); StatsProjection PrepareStatsProjection(bool has_equality_deletes) const; diff --git a/src/iceberg/test/manifest_group_test.cc b/src/iceberg/test/manifest_group_test.cc index f5eaa8311..7534c0999 100644 --- a/src/iceberg/test/manifest_group_test.cc +++ b/src/iceberg/test/manifest_group_test.cc @@ -37,6 +37,8 @@ #include "iceberg/manifest/manifest_list.h" #include "iceberg/manifest/manifest_reader.h" #include "iceberg/manifest/manifest_writer.h" +#include "iceberg/metrics/metrics_context.h" +#include "iceberg/metrics/scan_report.h" #include "iceberg/partition_spec.h" #include "iceberg/schema.h" #include "iceberg/table_scan.h" @@ -812,6 +814,140 @@ TEST_P(ManifestGroupTest, MultipleDataManifests) { EXPECT_EQ(executor.submit_count(), 2); } +TEST_P(ManifestGroupTest, StreamBatchBoundary) { + auto version = GetParam(); + + // Use one more manifest than the current executor batch size to exercise + // loading the next batch. + constexpr size_t kManifestCount = 33; + std::vector manifests; + std::vector expected_paths; + manifests.reserve(kManifestCount); + expected_paths.reserve(kManifestCount); + const auto partition = PartitionValues(std::vector{}); + for (size_t i = 0; i < kManifestCount; ++i) { + auto path = std::format("/path/to/data-{}.parquet", i); + manifests.push_back(WriteDataManifest( + version, /*snapshot_id=*/1000L + static_cast(i), + {MakeEntry(ManifestStatus::kAdded, + /*snapshot_id=*/1000L + static_cast(i), + /*sequence_number=*/static_cast(i) + 1, + MakeDataFile(path, partition, unpartitioned_spec_->spec_id()))}, + unpartitioned_spec_)); + expected_paths.emplace_back(std::move(path)); + } + + test::ThreadExecutor executor; + for (bool use_executor : {false, true}) { + SCOPED_TRACE(std::format("use_executor={}", use_executor)); + ICEBERG_UNWRAP_OR_FAIL( + auto group, ManifestGroup::Make(file_io_, schema_, GetSpecsById(), manifests)); + + if (use_executor) { + group->PlanWith(std::ref(executor)); + } + + ICEBERG_UNWRAP_OR_FAIL(auto stream, std::move(*group).PlanFilesStream()); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, stream->ToVector()); + EXPECT_EQ(GetPaths(tasks), expected_paths); + } +} + +TEST_P(ManifestGroupTest, FilterMetricsParity) { + auto version = GetParam(); + + constexpr int64_t kSnapshotId = 1000L; + const auto partition_0 = PartitionValues({Literal::Int(0)}); + const auto partition_1 = PartitionValues({Literal::Int(1)}); + auto selected_manifest = WriteDataManifest( + version, kSnapshotId, + {MakeEntry(ManifestStatus::kAdded, kSnapshotId, /*sequence_number=*/1, + MakeDataFile("/path/to/keep.parquet", partition_0, + partitioned_spec_->spec_id(), /*record_count=*/20)), + MakeEntry(ManifestStatus::kExisting, kSnapshotId, /*sequence_number=*/1, + MakeDataFile("/path/to/existing.parquet", partition_0, + partitioned_spec_->spec_id(), /*record_count=*/20)), + MakeEntry(ManifestStatus::kAdded, kSnapshotId, /*sequence_number=*/1, + MakeDataFile("/path/to/small.parquet", partition_0, + partitioned_spec_->spec_id(), /*record_count=*/5)), + MakeEntry(ManifestStatus::kAdded, kSnapshotId, /*sequence_number=*/1, + MakeDataFile("/path/to/predicate.parquet", partition_0, + partitioned_spec_->spec_id(), /*record_count=*/20))}, + partitioned_spec_); + auto skipped_manifest = WriteDataManifest( + version, kSnapshotId, + {MakeEntry(ManifestStatus::kAdded, kSnapshotId, /*sequence_number=*/1, + MakeDataFile("/path/to/pruned.parquet", partition_1, + partitioned_spec_->spec_id(), /*record_count=*/20))}, + partitioned_spec_); + const std::vector manifests = {selected_manifest, skipped_manifest}; + + auto configure = [](ManifestGroup& group, const std::shared_ptr& metrics) { + group.FilterPartitions(Expressions::Equal("data_bucket_16_2", Literal::Int(0))) + .IgnoreExisting() + .FilterFiles(Expressions::GreaterThanOrEqual("record_count", Literal::Long(10))) + .FilterManifestEntries([](const ManifestEntry& entry) { + return entry.data_file->file_path != "/path/to/predicate.parquet"; + }) + .WithScanMetrics(metrics); + }; + auto make_metrics = [] { + auto context = MetricsContext::Default(); + return std::shared_ptr(ScanMetrics::Make(*context)); + }; + + auto entries_metrics = make_metrics(); + ICEBERG_UNWRAP_OR_FAIL( + auto entries_group, + ManifestGroup::Make(file_io_, schema_, GetSpecsById(), manifests)); + configure(*entries_group, entries_metrics); + ICEBERG_UNWRAP_OR_FAIL(auto entries, entries_group->Entries()); + EXPECT_THAT(GetEntryPaths(entries), testing::ElementsAre("/path/to/keep.parquet")); + + struct StreamResult { + std::vector paths; + std::shared_ptr metrics; + }; + auto run_stream = [&](bool use_executor) -> Result { + auto metrics = make_metrics(); + ICEBERG_ASSIGN_OR_RAISE( + auto group, ManifestGroup::Make(file_io_, schema_, GetSpecsById(), manifests)); + configure(*group, metrics); + test::ThreadExecutor executor; + if (use_executor) { + group->PlanWith(std::ref(executor)); + } + ICEBERG_ASSIGN_OR_RAISE(auto stream, std::move(*group).PlanFilesStream()); + ICEBERG_ASSIGN_OR_RAISE(auto tasks, stream->ToVector()); + return StreamResult{.paths = GetPaths(tasks), .metrics = std::move(metrics)}; + }; + + ICEBERG_UNWRAP_OR_FAIL(auto serial, run_stream(false)); + ICEBERG_UNWRAP_OR_FAIL(auto parallel, run_stream(true)); + EXPECT_THAT(serial.paths, testing::ElementsAre("/path/to/keep.parquet")); + EXPECT_EQ(parallel.paths, serial.paths); + + auto check_filter_metrics = [](const std::shared_ptr& metrics) { + EXPECT_EQ(metrics->skipped_data_manifests->value(), 1); + EXPECT_EQ(metrics->scanned_data_manifests->value(), 1); + EXPECT_EQ(metrics->skipped_data_files->value(), 3); + }; + check_filter_metrics(entries_metrics); + check_filter_metrics(serial.metrics); + check_filter_metrics(parallel.metrics); + + EXPECT_EQ(entries_metrics->result_data_files->value(), 0); + EXPECT_EQ(entries_metrics->result_delete_files->value(), 0); + EXPECT_EQ(entries_metrics->total_file_size_in_bytes->value(), 0); + EXPECT_EQ(entries_metrics->total_delete_file_size_in_bytes->value(), 0); + + EXPECT_EQ(serial.metrics->result_data_files->value(), 1); + EXPECT_EQ(serial.metrics->result_delete_files->value(), 0); + EXPECT_EQ(serial.metrics->total_file_size_in_bytes->value(), 10); + EXPECT_EQ(serial.metrics->total_delete_file_size_in_bytes->value(), 0); + EXPECT_EQ(parallel.metrics->ToResult(), serial.metrics->ToResult()); +} + TEST_P(ManifestGroupTest, PartitionFilter) { auto version = GetParam(); From 6bbfb32372a9bfc68c1a041ef96476d040cd3977 Mon Sep 17 00:00:00 2001 From: Zehua Zou Date: Mon, 21 Sep 2026 18:14:56 +0800 Subject: [PATCH 2/2] add one more comment --- src/iceberg/manifest/manifest_group.cc | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/iceberg/manifest/manifest_group.cc b/src/iceberg/manifest/manifest_group.cc index 8f2970464..d97d66071 100644 --- a/src/iceberg/manifest/manifest_group.cc +++ b/src/iceberg/manifest/manifest_group.cc @@ -386,6 +386,8 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream { const size_t batch_size = group_->executor_.has_value() ? kManifestReadBatchSize : 1; std::vector manifests; manifests.reserve(batch_size); + // Skipped manifests do not count toward the batch size. Keep scanning until enough + // readable manifests are collected or all input is exhausted. while (next_manifest_ < group_->data_manifests_.size() && manifests.size() < batch_size) { const auto& manifest = group_->data_manifests_[next_manifest_++];