Skip to content

Commit 0113cfb

Browse files
manuzhangcodex
andcommitted
fix: address stream API review feedback
Co-authored-by: Codex <codex@openai.com>
1 parent d82d9e2 commit 0113cfb

16 files changed

Lines changed: 174 additions & 105 deletions

‎example/demo_example.cc‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -79,7 +79,7 @@ int main(int argc, char** argv) {
7979
}
8080

8181
auto scan = std::move(scan_result.value());
82-
auto plan_result = scan->PlanFilesIterator();
82+
auto plan_result = scan->PlanFilesStream();
8383
if (!plan_result.has_value()) {
8484
std::cerr << "Failed to plan files: " << plan_result.error().message << std::endl;
8585
return 1;
Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
* or more contributor license agreements. See the NOTICE file
44
* distributed with this work for additional information
55
* regarding copyright ownership. The ASF licenses this file
6-
* to you under the Apache License, Version 2.0 (the
6+
* under the Apache License, Version 2.0 (the
77
* "License"); you may not use this file except in compliance
88
* with the License. You may obtain a copy of the License at
99
*
@@ -19,8 +19,8 @@
1919

2020
#pragma once
2121

22-
/// \file iceberg/file_scan_task_iterator.h
23-
/// \brief Define the owning iterator type for file scan tasks.
22+
/// \file iceberg/file_scan_task_stream.h
23+
/// \brief Define the owning stream type for file scan tasks.
2424

2525
#include <memory>
2626

@@ -30,6 +30,6 @@ namespace iceberg {
3030

3131
class FileScanTask;
3232

33-
using FileScanTaskIterator = std::unique_ptr<Iterator<std::shared_ptr<FileScanTask>>>;
33+
using FileScanTaskStream = std::unique_ptr<Iterator<std::shared_ptr<FileScanTask>>>;
3434

3535
} // namespace iceberg

‎src/iceberg/manifest/manifest_group.cc‎

Lines changed: 39 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -134,13 +134,13 @@ ManifestGroup& ManifestGroup::operator=(ManifestGroup&&) noexcept = default;
134134
class ManifestGroup::FilePlanningIterator final
135135
: public Iterator<std::shared_ptr<FileScanTask>> {
136136
public:
137-
static Result<FileScanTaskIterator> Make(std::unique_ptr<ManifestGroup> group) {
137+
static Result<FileScanTaskStream> Make(std::unique_ptr<ManifestGroup> group) {
138138
ICEBERG_RETURN_UNEXPECTED(group->CheckErrors());
139139

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

143-
const bool drop_stats =
143+
auto stats_projection =
144144
group->PrepareStatsProjection(delete_index->has_equality_deletes());
145145

146146
std::unique_ptr<Evaluator> data_file_evaluator;
@@ -151,10 +151,11 @@ class ManifestGroup::FilePlanningIterator final
151151
Evaluator::Make(*DataFileFilterSchema(), group->file_filter_,
152152
group->case_sensitive_));
153153
}
154+
const bool drop_stats = stats_projection.drop_stats;
154155

155-
return FileScanTaskIterator(
156-
new FilePlanningIterator(std::move(group), std::move(delete_index),
157-
std::move(data_file_evaluator), drop_stats));
156+
return FileScanTaskStream(new FilePlanningIterator(
157+
std::move(group), std::move(delete_index), std::move(data_file_evaluator),
158+
std::move(stats_projection.columns), drop_stats));
158159
}
159160

160161
Result<std::optional<std::shared_ptr<FileScanTask>>> NextImpl() override {
@@ -211,10 +212,12 @@ class ManifestGroup::FilePlanningIterator final
211212
private:
212213
FilePlanningIterator(std::unique_ptr<ManifestGroup> group,
213214
std::unique_ptr<DeleteFileIndex> delete_index,
214-
std::unique_ptr<Evaluator> data_file_evaluator, bool drop_stats)
215+
std::unique_ptr<Evaluator> data_file_evaluator,
216+
std::vector<std::string> columns, bool drop_stats)
215217
: group_(std::move(group)),
216218
delete_index_(std::move(delete_index)),
217219
data_file_evaluator_(std::move(data_file_evaluator)),
220+
columns_(std::move(columns)),
218221
drop_stats_(drop_stats) {}
219222

220223
using TaggedEntry = std::pair<int32_t, ManifestEntry>;
@@ -335,10 +338,10 @@ class ManifestGroup::FilePlanningIterator final
335338
continue;
336339
}
337340

338-
ICEBERG_ASSIGN_OR_RAISE(auto reader, group_->MakeReader(manifest));
341+
ICEBERG_ASSIGN_OR_RAISE(auto reader, group_->MakeReader(manifest, columns_));
339342
ICEBERG_ASSIGN_OR_RAISE(entry_iterator_, group_->ignore_deleted_
340-
? reader->LiveEntriesIterator()
341-
: reader->EntriesIterator());
343+
? reader->LiveEntriesStream()
344+
: reader->EntriesStream());
342345
current_spec_id_ = manifest.partition_spec_id;
343346
return true;
344347
}
@@ -369,10 +372,11 @@ class ManifestGroup::FilePlanningIterator final
369372
ParallelCollect(
370373
group_->executor_, manifests,
371374
[this](const ManifestFile* manifest) -> Result<std::vector<TaggedIterator>> {
372-
ICEBERG_ASSIGN_OR_RAISE(auto reader, group_->MakeReader(*manifest));
375+
ICEBERG_ASSIGN_OR_RAISE(auto reader,
376+
group_->MakeReader(*manifest, columns_));
373377
ICEBERG_ASSIGN_OR_RAISE(auto iterator, group_->ignore_deleted_
374-
? reader->LiveEntriesIterator()
375-
: reader->EntriesIterator());
378+
? reader->LiveEntriesStream()
379+
: reader->EntriesStream());
376380

377381
std::vector<TaggedIterator> tagged_iterators;
378382
tagged_iterators.emplace_back(manifest->partition_spec_id,
@@ -416,6 +420,7 @@ class ManifestGroup::FilePlanningIterator final
416420
std::unique_ptr<ManifestGroup> group_;
417421
std::unique_ptr<DeleteFileIndex> delete_index_;
418422
std::unique_ptr<Evaluator> data_file_evaluator_;
423+
std::vector<std::string> columns_;
419424
std::unordered_map<int32_t, std::unique_ptr<ManifestEvaluator>> manifest_evaluators_;
420425
std::unordered_map<int32_t, std::shared_ptr<ResidualEvaluator>> residual_evaluators_;
421426
std::unique_ptr<Iterator<ManifestEntry>> entry_iterator_;
@@ -555,7 +560,7 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> ManifestGroup::PlanFiles() {
555560
return file_tasks;
556561
}
557562

558-
Result<FileScanTaskIterator> ManifestGroup::PlanFilesIterator() && {
563+
Result<FileScanTaskStream> ManifestGroup::PlanFilesStream() && {
559564
auto group = std::make_unique<ManifestGroup>(std::move(*this));
560565
return FilePlanningIterator::Make(std::move(group));
561566
}
@@ -585,7 +590,7 @@ Result<std::vector<std::shared_ptr<ScanTask>>> ManifestGroup::Plan(
585590
delete_index_builder_.WithScanMetrics(scan_metrics_);
586591
ICEBERG_ASSIGN_OR_RAISE(auto delete_index, delete_index_builder_.Build());
587592

588-
const bool drop_stats = PrepareStatsProjection(delete_index->has_equality_deletes());
593+
auto stats_projection = PrepareStatsProjection(delete_index->has_equality_deletes());
589594

590595
std::unordered_map<int32_t, std::unique_ptr<TaskContext>> task_context_cache;
591596
auto get_task_context = [&](int32_t spec_id) -> Result<TaskContext*> {
@@ -603,13 +608,13 @@ Result<std::vector<std::shared_ptr<ScanTask>>> ManifestGroup::Plan(
603608
TaskContext{.spec = spec,
604609
.deletes = delete_index.get(),
605610
.residuals = residuals,
606-
.drop_stats = drop_stats,
611+
.drop_stats = stats_projection.drop_stats,
607612
.columns_to_keep_stats = columns_to_keep_stats_});
608613

609614
return task_context_cache[spec_id].get();
610615
};
611616

612-
ICEBERG_ASSIGN_OR_RAISE(auto entry_groups, ReadEntries());
617+
ICEBERG_ASSIGN_OR_RAISE(auto entry_groups, ReadEntries(stats_projection.columns));
613618

614619
std::vector<std::shared_ptr<ScanTask>> all_tasks;
615620
for (auto& [spec_id, entries] : entry_groups) {
@@ -623,7 +628,7 @@ Result<std::vector<std::shared_ptr<ScanTask>>> ManifestGroup::Plan(
623628
}
624629

625630
Result<std::vector<ManifestEntry>> ManifestGroup::Entries() {
626-
ICEBERG_ASSIGN_OR_RAISE(auto entry_groups, ReadEntries());
631+
ICEBERG_ASSIGN_OR_RAISE(auto entry_groups, ReadEntries(columns_));
627632

628633
std::vector<ManifestEntry> all_entries;
629634
for (auto& [_, entries] : entry_groups) {
@@ -635,21 +640,23 @@ Result<std::vector<ManifestEntry>> ManifestGroup::Entries() {
635640
}
636641

637642
Result<std::unique_ptr<ManifestReader>> ManifestGroup::MakeReader(
638-
const ManifestFile& manifest) {
643+
const ManifestFile& manifest, const std::vector<std::string>& columns) {
639644
ICEBERG_ASSIGN_OR_RAISE(auto reader,
640645
ManifestReader::Make(manifest, io_, schema_, specs_by_id_));
641646

642-
auto columns = columns_;
647+
auto reader_columns = columns;
643648
if (file_filter_ && file_filter_->op() != Expression::Operation::kTrue &&
644-
!columns.empty() && !std::ranges::contains(columns, Schema::kAllColumns)) {
649+
!reader_columns.empty() &&
650+
!std::ranges::contains(reader_columns, Schema::kAllColumns)) {
645651
auto data_file_schema = DataFileFilterSchema();
646652
ICEBERG_ASSIGN_OR_RAISE(
647653
auto bound_file_filter,
648654
Binder::Bind(*data_file_schema, file_filter_, case_sensitive_));
649655
ICEBERG_ASSIGN_OR_RAISE(auto referenced_field_ids,
650656
ReferenceVisitor::GetReferencedFieldIds(bound_file_filter));
651657

652-
std::unordered_set<std::string> selected_columns(columns.cbegin(), columns.cend());
658+
std::unordered_set<std::string> selected_columns(reader_columns.cbegin(),
659+
reader_columns.cend());
653660
for (const auto field_id : referenced_field_ids) {
654661
if (field_id == DataFile::kSpecIdFieldId) {
655662
continue;
@@ -661,16 +668,16 @@ Result<std::unique_ptr<ManifestReader>> ManifestGroup::MakeReader(
661668
if (selected_columns.contains(column_name_str)) {
662669
continue;
663670
}
664-
columns.push_back(std::move(column_name_str));
665-
selected_columns.insert(columns.back());
671+
reader_columns.push_back(std::move(column_name_str));
672+
selected_columns.insert(reader_columns.back());
666673
}
667674
}
668675
}
669676

670677
reader->FilterRows(data_filter_)
671678
.FilterPartitions(partition_filter_)
672679
.CaseSensitive(case_sensitive_)
673-
.Select(std::move(columns));
680+
.Select(std::move(reader_columns));
674681

675682
if (scan_metrics_) {
676683
reader->SkipCounter(scan_metrics_->skipped_data_files);
@@ -679,20 +686,22 @@ Result<std::unique_ptr<ManifestReader>> ManifestGroup::MakeReader(
679686
return reader;
680687
}
681688

682-
bool ManifestGroup::PrepareStatsProjection(bool has_equality_deletes) {
689+
ManifestGroup::StatsProjection ManifestGroup::PrepareStatsProjection(
690+
bool has_equality_deletes) const {
683691
// The caller's projection records whether stats were requested. Equality-delete
684692
// matching may add stats temporarily, but they should still be dropped from the
685693
// result when the original projection did not request them. Keeping this decision
686694
// here ensures eager and iterator planning use identical semantics.
687-
const bool drop_stats = ManifestReader::ShouldDropStats(columns_);
695+
StatsProjection result{.columns = columns_,
696+
.drop_stats = ManifestReader::ShouldDropStats(columns_)};
688697
if (has_equality_deletes) {
689-
columns_ = ManifestReader::WithStatsColumns(columns_);
698+
result.columns = ManifestReader::WithStatsColumns(result.columns);
690699
}
691-
return drop_stats;
700+
return result;
692701
}
693702

694703
Result<std::unordered_map<int32_t, std::vector<ManifestEntry>>>
695-
ManifestGroup::ReadEntries() {
704+
ManifestGroup::ReadEntries(const std::vector<std::string>& columns) {
696705
const auto cache_capacity = static_cast<int32_t>(specs_by_id_.size());
697706
auto get_manifest_evaluator = internal::MemoizeLru(
698707
[this](int32_t spec_id) -> Result<std::shared_ptr<ManifestEvaluator>> {
@@ -759,7 +768,7 @@ ManifestGroup::ReadEntries() {
759768
}
760769

761770
// Read manifest entries
762-
ICEBERG_ASSIGN_OR_RAISE(auto reader, MakeReader(manifest));
771+
ICEBERG_ASSIGN_OR_RAISE(auto reader, MakeReader(manifest, columns));
763772
ICEBERG_ASSIGN_OR_RAISE(
764773
auto entries, ignore_deleted_ ? reader->LiveEntries() : reader->Entries());
765774

‎src/iceberg/manifest/manifest_group.h‎

Lines changed: 18 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@
3030
#include <vector>
3131

3232
#include "iceberg/delete_file_index.h"
33-
#include "iceberg/file_scan_task_iterator.h"
33+
#include "iceberg/file_scan_task_stream.h"
3434
#include "iceberg/iceberg_export.h"
3535
#include "iceberg/manifest/manifest_entry.h"
3636
#include "iceberg/manifest/manifest_list.h"
@@ -139,13 +139,15 @@ class ICEBERG_EXPORT ManifestGroup : public ErrorCollector {
139139

140140
/// \brief Lazily plan scan tasks for matching data files.
141141
///
142-
/// The returned iterator owns the planning state and may outlive this ManifestGroup.
142+
/// The returned stream owns the planning state and may outlive this ManifestGroup.
143143
/// It reads one bounded manifest batch at a time instead of materializing all manifest
144-
/// entries and scan tasks. When PlanWith() configures an executor, entry iterators for
144+
/// entries and scan tasks. When PlanWith() configures an executor, entry streams for
145145
/// manifests in each batch are opened in parallel, while entries are consumed one
146-
/// manifest at a time. Creating the iterator consumes this group's configuration, so
147-
/// this method may only be called on an rvalue.
148-
Result<FileScanTaskIterator> PlanFilesIterator() &&;
146+
/// manifest at a time. Delete manifests are still read eagerly when creating the
147+
/// stream because delete files must be indexed before data-file planning can begin.
148+
/// Creating the stream consumes this group's configuration, so this method may only
149+
/// be called on an rvalue.
150+
Result<FileScanTaskStream> PlanFilesStream() &&;
149151

150152
/// \brief Get all matching manifest entries.
151153
Result<std::vector<ManifestEntry>> Entries();
@@ -164,16 +166,23 @@ class ICEBERG_EXPORT ManifestGroup : public ErrorCollector {
164166
private:
165167
class FilePlanningIterator;
166168

169+
struct StatsProjection {
170+
std::vector<std::string> columns;
171+
bool drop_stats;
172+
};
173+
167174
ManifestGroup(std::shared_ptr<FileIO> io, std::shared_ptr<Schema> schema,
168175
std::unordered_map<int32_t, std::shared_ptr<PartitionSpec>> specs_by_id,
169176
std::vector<ManifestFile> data_manifests,
170177
DeleteFileIndex::Builder&& delete_index_builder);
171178

172-
Result<std::unordered_map<int32_t, std::vector<ManifestEntry>>> ReadEntries();
179+
Result<std::unordered_map<int32_t, std::vector<ManifestEntry>>> ReadEntries(
180+
const std::vector<std::string>& columns);
173181

174-
Result<std::unique_ptr<ManifestReader>> MakeReader(const ManifestFile& manifest);
182+
Result<std::unique_ptr<ManifestReader>> MakeReader(
183+
const ManifestFile& manifest, const std::vector<std::string>& columns);
175184

176-
bool PrepareStatsProjection(bool has_equality_deletes);
185+
StatsProjection PrepareStatsProjection(bool has_equality_deletes) const;
177186

178187
std::shared_ptr<FileIO> io_;
179188
std::shared_ptr<Schema> schema_;

‎src/iceberg/manifest/manifest_reader.cc‎

Lines changed: 13 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -815,17 +815,17 @@ class ManifestEntryIteratorImpl final : public Iterator<ManifestEntry> {
815815

816816
} // namespace
817817

818-
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReader::EntriesIterator() {
819-
if (auto* iterable = dynamic_cast<SupportsManifestEntryIteration*>(this)) {
820-
return iterable->EntriesIterator();
818+
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReader::EntriesStream() {
819+
if (auto* iterable = dynamic_cast<SupportsManifestEntryStreaming*>(this)) {
820+
return iterable->EntriesStream();
821821
}
822822
ICEBERG_ASSIGN_OR_RAISE(auto entries, Entries());
823823
return std::make_unique<VectorIterator<ManifestEntry>>(std::move(entries));
824824
}
825825

826-
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReader::LiveEntriesIterator() {
827-
if (auto* iterable = dynamic_cast<SupportsManifestEntryIteration*>(this)) {
828-
return iterable->LiveEntriesIterator();
826+
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReader::LiveEntriesStream() {
827+
if (auto* iterable = dynamic_cast<SupportsManifestEntryStreaming*>(this)) {
828+
return iterable->LiveEntriesStream();
829829
}
830830
ICEBERG_ASSIGN_OR_RAISE(auto entries, LiveEntries());
831831
return std::make_unique<VectorIterator<ManifestEntry>>(std::move(entries));
@@ -1002,25 +1002,24 @@ Status ManifestReaderImpl::OpenReader(std::shared_ptr<Schema> projection) {
10021002
}
10031003

10041004
Result<std::vector<ManifestEntry>> ManifestReaderImpl::Entries() {
1005-
ICEBERG_ASSIGN_OR_RAISE(auto entries, EntriesIterator());
1005+
ICEBERG_ASSIGN_OR_RAISE(auto entries, EntriesStream());
10061006
return entries->ToVector();
10071007
}
10081008

10091009
Result<std::vector<ManifestEntry>> ManifestReaderImpl::LiveEntries() {
1010-
ICEBERG_ASSIGN_OR_RAISE(auto entries, LiveEntriesIterator());
1010+
ICEBERG_ASSIGN_OR_RAISE(auto entries, LiveEntriesStream());
10111011
return entries->ToVector();
10121012
}
10131013

1014-
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReaderImpl::EntriesIterator() {
1015-
return MakeEntriesIterator(/*only_live=*/false);
1014+
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReaderImpl::EntriesStream() {
1015+
return MakeEntriesStream(/*only_live=*/false);
10161016
}
10171017

1018-
Result<std::unique_ptr<Iterator<ManifestEntry>>>
1019-
ManifestReaderImpl::LiveEntriesIterator() {
1020-
return MakeEntriesIterator(/*only_live=*/true);
1018+
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReaderImpl::LiveEntriesStream() {
1019+
return MakeEntriesStream(/*only_live=*/true);
10211020
}
10221021

1023-
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReaderImpl::MakeEntriesIterator(
1022+
Result<std::unique_ptr<Iterator<ManifestEntry>>> ManifestReaderImpl::MakeEntriesStream(
10241023
bool only_live) {
10251024
ICEBERG_ASSIGN_OR_RAISE(auto partition_type, spec_->RawPartitionType(*schema_));
10261025
auto data_file_schema = DataFile::Type(std::move(partition_type))->ToSchema();

0 commit comments

Comments
 (0)