Skip to content
16 changes: 13 additions & 3 deletions example/demo_example.cc
Original file line number Diff line number Diff line change
Expand Up @@ -79,16 +79,26 @@ int main(int argc, char** argv) {
}

auto scan = std::move(scan_result.value());
auto plan_result = scan->PlanFiles();
auto plan_result = scan->PlanFilesIterator();
if (!plan_result.has_value()) {
std::cerr << "Failed to plan files: " << plan_result.error().message << std::endl;
return 1;
}

std::cout << "Scan tasks: " << std::endl;
auto scan_tasks = std::move(plan_result.value());
for (const auto& scan_task : scan_tasks) {
std::cout << " - " << scan_task->data_file()->file_path << std::endl;
while (true) {
auto task_result = scan_tasks->Next();
if (!task_result.has_value()) {
std::cerr << "Failed to plan next file: " << task_result.error().message
<< std::endl;
return 1;
}
if (!task_result.value().has_value()) {
break;
}
std::cout << " - " << task_result.value().value()->data_file()->file_path
<< std::endl;
}

return 0;
Expand Down
329 changes: 324 additions & 5 deletions src/iceberg/manifest/manifest_group.cc
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,308 @@ ManifestGroup::~ManifestGroup() = default;
ManifestGroup::ManifestGroup(ManifestGroup&&) noexcept = default;
ManifestGroup& ManifestGroup::operator=(ManifestGroup&&) noexcept = default;

class ManifestGroup::FilePlanningIterator final
: public Iterator<std::shared_ptr<FileScanTask>> {
public:
static Result<FileScanTaskIterator> Make(
std::unique_ptr<ManifestGroup> group) {
ICEBERG_RETURN_UNEXPECTED(group->CheckErrors());

group->delete_index_builder_.WithScanMetrics(group->scan_metrics_);
ICEBERG_ASSIGN_OR_RAISE(auto delete_index, group->delete_index_builder_.Build());

const bool drop_stats =
group->PrepareStatsProjection(delete_index->has_equality_deletes());

std::unique_ptr<Evaluator> 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_));
}

return FileScanTaskIterator(
new FilePlanningIterator(std::move(group), std::move(delete_index),
std::move(data_file_evaluator), drop_stats));
}

Result<std::optional<std::shared_ptr<FileScanTask>>> NextImpl() override {
while (true) {
ICEBERG_ASSIGN_OR_RAISE(auto entry, NextEntry());
if (!entry.has_value()) {
return std::nullopt;
}

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();
continue;
}

ICEBERG_ASSIGN_OR_RAISE(auto delete_files, delete_index_->ForEntry(value));

// Equality-delete matching uses data-file statistics. Drop unrequested stats only
// after the delete index has finished matching this entry.
if (drop_stats_) {
ContentFileUtil::DropAllStats(*value.data_file);
} else if (!group_->columns_to_keep_stats_.empty()) {
ContentFileUtil::DropUnselectedStats(*value.data_file,
group_->columns_to_keep_stats_);
}

UpdateResultMetrics(*value.data_file, delete_files);

ICEBERG_ASSIGN_OR_RAISE(auto residuals, GetResidualEvaluator(spec_id));
ICEBERG_ASSIGN_OR_RAISE(auto residual,
residuals->ResidualFor(value.data_file->partition));

return std::optional<std::shared_ptr<FileScanTask>>{std::make_shared<FileScanTask>(
std::move(value.data_file), std::move(delete_files), std::move(residual))};
}
}

private:
FilePlanningIterator(std::unique_ptr<ManifestGroup> group,
std::unique_ptr<DeleteFileIndex> delete_index,
std::unique_ptr<Evaluator> data_file_evaluator, bool drop_stats)
: group_(std::move(group)),
delete_index_(std::move(delete_index)),
data_file_evaluator_(std::move(data_file_evaluator)),
drop_stats_(drop_stats) {}

using TaggedEntry = std::pair<int32_t, ManifestEntry>;
using TaggedIterator = std::pair<int32_t, std::unique_ptr<Iterator<ManifestEntry>>>;

Result<std::optional<TaggedEntry>> NextEntry() {
if (!group_->executor_.has_value()) {
while (true) {
if (!entry_iterator_) {
ICEBERG_ASSIGN_OR_RAISE(bool opened, OpenNextManifest());
if (!opened) {
return std::nullopt;
}
}

ICEBERG_ASSIGN_OR_RAISE(auto entry, entry_iterator_->Next());
if (!entry.has_value()) {
entry_iterator_.reset();
continue;
}
return std::optional<TaggedEntry>{std::in_place, current_spec_id_,
std::move(entry).value()};
}
}

while (true) {
if (next_batch_iterator_ == batch_iterators_.size()) {
ICEBERG_ASSIGN_OR_RAISE(bool loaded, LoadNextManifestBatch());
if (!loaded) {
return std::nullopt;
}
}

auto& [spec_id, iterator] = batch_iterators_[next_batch_iterator_];
ICEBERG_ASSIGN_OR_RAISE(auto entry, iterator->Next());
if (!entry.has_value()) {
iterator.reset();
++next_batch_iterator_;
continue;
}
return std::optional<TaggedEntry>{std::in_place, spec_id, std::move(entry).value()};
}
}

Result<ManifestEvaluator*> GetManifestEvaluator(int32_t spec_id) {
auto cached = manifest_evaluators_.find(spec_id);
if (cached != manifest_evaluators_.end()) {
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_));
auto* result = evaluator.get();
manifest_evaluators_.emplace(spec_id, std::move(evaluator));
return result;
}

Result<ResidualEvaluator*> GetResidualEvaluator(int32_t spec_id) {
auto cached = residual_evaluators_.find(spec_id);
if (cached != residual_evaluators_.end()) {
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_));
auto* result = evaluator.get();
residual_evaluators_.emplace(spec_id, std::move(evaluator));
return result;
}

Result<bool> 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<bool> 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));
ICEBERG_ASSIGN_OR_RAISE(entry_iterator_, group_->ignore_deleted_
? reader->LiveEntriesIterator()
: reader->EntriesIterator());
current_spec_id_ = manifest.partition_spec_id;
return true;
}
return false;
}

Result<bool> LoadNextManifestBatch() {
std::vector<const ManifestFile*> manifests;
manifests.reserve(kManifestReadBatchSize);
while (next_manifest_ < group_->data_manifests_.size() &&
manifests.size() < kManifestReadBatchSize) {
const auto& manifest = group_->data_manifests_[next_manifest_++];
ICEBERG_ASSIGN_OR_RAISE(bool should_read, ShouldReadManifest(manifest));
if (should_read) {
manifests.push_back(&manifest);
}
}

if (manifests.empty()) {
return false;
}

// Open the readers concurrently, but keep their iterators instead of collecting
// entries here. This preserves bounded memory for large manifests while retaining
// parallel manifest initialization when an executor is configured.
ICEBERG_ASSIGN_OR_RAISE(
batch_iterators_,
ParallelCollect(
group_->executor_, manifests,
[this](const ManifestFile* manifest) -> Result<std::vector<TaggedIterator>> {
ICEBERG_ASSIGN_OR_RAISE(auto reader, group_->MakeReader(*manifest));
ICEBERG_ASSIGN_OR_RAISE(auto iterator, group_->ignore_deleted_
? reader->LiveEntriesIterator()
: reader->EntriesIterator());

std::vector<TaggedIterator> tagged_iterators;
tagged_iterators.emplace_back(manifest->partition_spec_id,
std::move(iterator));
return tagged_iterators;
}));
next_batch_iterator_ = 0;
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<std::shared_ptr<DataFile>>& delete_files) {
if (!group_->scan_metrics_) {
return;
}

group_->scan_metrics_->total_file_size_in_bytes->Increment(
ContentFileUtil::ContentSizeInBytes(data_file));
group_->scan_metrics_->result_data_files->Increment(1);
group_->scan_metrics_->result_delete_files->Increment(
static_cast<int64_t>(delete_files.size()));
int64_t deletes_size = 0;
for (const auto& delete_file : delete_files) {
deletes_size += ContentFileUtil::ContentSizeInBytes(*delete_file);
}
group_->scan_metrics_->total_delete_file_size_in_bytes->Increment(deletes_size);
}

std::unique_ptr<ManifestGroup> group_;
std::unique_ptr<DeleteFileIndex> delete_index_;
std::unique_ptr<Evaluator> data_file_evaluator_;
std::unordered_map<int32_t, std::unique_ptr<ManifestEvaluator>> manifest_evaluators_;
std::unordered_map<int32_t, std::shared_ptr<ResidualEvaluator>> residual_evaluators_;
std::unique_ptr<Iterator<ManifestEntry>> entry_iterator_;
std::vector<TaggedIterator> batch_iterators_;
size_t next_manifest_ = 0;
size_t next_batch_iterator_ = 0;
int32_t current_spec_id_ = 0;
bool drop_stats_;

// Limit the number of manifest readers and iterators retained by executor-backed
// planning. The executor still controls actual task concurrency, while this fixed
// cap prevents resource use from scaling with the total manifest count. Entries
// within each manifest remain streamed, so this does not cap manifest size.
static constexpr size_t kManifestReadBatchSize = 32;
Comment thread
manuzhang marked this conversation as resolved.
};

ManifestGroup& ManifestGroup::FilterData(std::shared_ptr<Expression> filter) {
ICEBERG_BUILDER_ASSIGN_OR_RETURN(data_filter_, And::Make(data_filter_, filter));
delete_index_builder_.DataFilter(std::move(filter));
Expand Down Expand Up @@ -211,12 +513,15 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> ManifestGroup::PlanFiles() {
tasks.reserve(entries.size());

for (auto& entry : entries) {
ICEBERG_ASSIGN_OR_RAISE(auto delete_files, ctx.deletes->ForEntry(entry));

// Equality-delete matching uses data-file statistics. Drop unrequested stats only
// after the delete index has finished matching this entry.
if (ctx.drop_stats) {
ContentFileUtil::DropAllStats(*entry.data_file);
} else if (!ctx.columns_to_keep_stats.empty()) {
ContentFileUtil::DropUnselectedStats(*entry.data_file, ctx.columns_to_keep_stats);
}
ICEBERG_ASSIGN_OR_RAISE(auto delete_files, ctx.deletes->ForEntry(entry));
// Count result metrics once per data file task. A delete file shared by
// multiple data files contributes once to each task, unlike indexed delete files.
if (scan_metrics_) {
Expand Down Expand Up @@ -251,6 +556,11 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> ManifestGroup::PlanFiles() {
return file_tasks;
}

Result<FileScanTaskIterator> ManifestGroup::PlanFilesIterator() && {
auto group = std::make_unique<ManifestGroup>(std::move(*this));
Comment thread
manuzhang marked this conversation as resolved.
return FilePlanningIterator::Make(std::move(group));
}

Result<std::vector<std::shared_ptr<ScanTask>>> ManifestGroup::Plan(
const CreateTasksFunction& create_tasks) {
std::unordered_map<int32_t, std::shared_ptr<ResidualEvaluator>> residual_cache;
Expand All @@ -276,10 +586,7 @@ Result<std::vector<std::shared_ptr<ScanTask>>> ManifestGroup::Plan(
delete_index_builder_.WithScanMetrics(scan_metrics_);
ICEBERG_ASSIGN_OR_RAISE(auto delete_index, delete_index_builder_.Build());

bool drop_stats = ManifestReader::ShouldDropStats(columns_);
if (delete_index->has_equality_deletes()) {
columns_ = ManifestReader::WithStatsColumns(columns_);
}
const bool drop_stats = PrepareStatsProjection(delete_index->has_equality_deletes());

std::unordered_map<int32_t, std::unique_ptr<TaskContext>> task_context_cache;
auto get_task_context = [&](int32_t spec_id) -> Result<TaskContext*> {
Expand Down Expand Up @@ -373,6 +680,18 @@ Result<std::unique_ptr<ManifestReader>> ManifestGroup::MakeReader(
return reader;
}

bool ManifestGroup::PrepareStatsProjection(bool has_equality_deletes) {
// The caller's projection records whether stats were requested. Equality-delete
// matching may add stats temporarily, but they should still be dropped from the
// result when the original projection did not request them. Keeping this decision
// here ensures eager and iterator planning use identical semantics.
const bool drop_stats = ManifestReader::ShouldDropStats(columns_);
if (has_equality_deletes) {
columns_ = ManifestReader::WithStatsColumns(columns_);
}
return drop_stats;
}

Result<std::unordered_map<int32_t, std::vector<ManifestEntry>>>
ManifestGroup::ReadEntries() {
const auto cache_capacity = static_cast<int32_t>(specs_by_id_.size());
Expand Down
Loading
Loading