Skip to content

Commit 7e81fa2

Browse files
authored
test: complete v3 row lineage test coverage (#869)
- enable RewriteFiles V3 coverage and verify DV cleanup - cover V2-to-V3 assignment and incremental/changelog lineage - verify delete-aware readers preserve lineage
1 parent 4b9ff79 commit 7e81fa2

5 files changed

Lines changed: 459 additions & 28 deletions

‎src/iceberg/test/file_scan_task_reader_test.cc‎

Lines changed: 83 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -113,6 +113,14 @@ class FileScanTaskReaderTest : public TempFileTestBase {
113113
table_schema_->schema_id());
114114
}
115115

116+
std::shared_ptr<Schema> RowLineageProjection() const {
117+
return std::make_shared<Schema>(
118+
std::vector<SchemaField>{SchemaField::MakeRequired(1, "id", int32()),
119+
MetadataColumns::kRowId,
120+
MetadataColumns::kLastUpdatedSequenceNumber},
121+
table_schema_->schema_id());
122+
}
123+
116124
Result<ExportedBatch> MakeBatch(const Schema& schema,
117125
const std::string& json_data) const {
118126
ICEBERG_ASSIGN_OR_RAISE(auto arrow_schema, MakeArrowSchema(schema));
@@ -409,17 +417,12 @@ TEST_F(FileScanTaskReaderTest, ReadLastUpdatedFromDataSeq) {
409417
data_file->first_row_id = 100L;
410418
data_file->data_sequence_number = 5L;
411419
FileScanTask task(data_file);
412-
auto projected_schema = std::make_shared<Schema>(
413-
std::vector<SchemaField>{SchemaField::MakeRequired(1, "id", int32()),
414-
MetadataColumns::kRowId,
415-
MetadataColumns::kLastUpdatedSequenceNumber},
416-
table_schema_->schema_id());
417420

418421
FileScanTaskReader::Options options{
419422
.io = file_io_,
420423
.table_schema = table_schema_,
421424
.schemas = {table_schema_},
422-
.projected_schema = projected_schema,
425+
.projected_schema = RowLineageProjection(),
423426
};
424427
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
425428
auto stream_result = reader->Open(task);
@@ -504,6 +507,80 @@ TEST_F(FileScanTaskReaderTest, OpenWithEqualityDeletesAddsAndPrunesDeleteOnlyCol
504507
ASSERT_NO_FATAL_FAILURE(VerifyStream(&stream, R"([[1, "Foo"], [3, "Baz"]])"));
505508
}
506509

510+
TEST_F(FileScanTaskReaderTest, PositionDeletesPreserveRowLineage) {
511+
ICEBERG_UNWRAP_OR_FAIL(
512+
auto data_file,
513+
MakeDataFile(table_schema_,
514+
R"([[1, "Foo", "blue"], [2, "Bar", "red"], [3, "Baz", "green"]])"));
515+
data_file->first_row_id = 100;
516+
data_file->data_sequence_number = 5;
517+
ICEBERG_UNWRAP_OR_FAIL(
518+
auto pos_delete, MakePositionDeleteFile(CreateNewTempFilePathWithSuffix(".parquet"),
519+
{1}, data_file->file_path));
520+
FileScanTask task(data_file, {pos_delete});
521+
522+
FileScanTaskReader::Options options{
523+
.io = file_io_,
524+
.table_schema = table_schema_,
525+
.schemas = {table_schema_},
526+
.projected_schema = RowLineageProjection(),
527+
};
528+
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
529+
ICEBERG_UNWRAP_OR_FAIL(auto stream, reader->Open(task));
530+
531+
ASSERT_NO_FATAL_FAILURE(VerifyStream(&stream, R"([[1, 100, 5], [3, 102, 5]])"));
532+
}
533+
534+
TEST_F(FileScanTaskReaderTest, DeletionVectorDeletesPreserveRowLineage) {
535+
ICEBERG_UNWRAP_OR_FAIL(
536+
auto data_file,
537+
MakeDataFile(table_schema_,
538+
R"([[1, "Foo", "blue"], [2, "Bar", "red"], [3, "Baz", "green"]])"));
539+
data_file->first_row_id = 100;
540+
data_file->data_sequence_number = 5;
541+
ICEBERG_UNWRAP_OR_FAIL(
542+
auto deletion_vector,
543+
MakeDeletionVectorFile(CreateNewTempFilePathWithSuffix(".puffin"), {1},
544+
data_file->file_path));
545+
FileScanTask task(data_file, {deletion_vector});
546+
547+
FileScanTaskReader::Options options{
548+
.io = file_io_,
549+
.table_schema = table_schema_,
550+
.schemas = {table_schema_},
551+
.projected_schema = RowLineageProjection(),
552+
};
553+
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
554+
ICEBERG_UNWRAP_OR_FAIL(auto stream, reader->Open(task));
555+
556+
ASSERT_NO_FATAL_FAILURE(VerifyStream(&stream, R"([[1, 100, 5], [3, 102, 5]])"));
557+
}
558+
559+
TEST_F(FileScanTaskReaderTest, EqualityDeletesPreserveRowLineage) {
560+
ICEBERG_UNWRAP_OR_FAIL(
561+
auto data_file,
562+
MakeDataFile(table_schema_,
563+
R"([[1, "Foo", "blue"], [2, "Bar", "red"], [3, "Baz", "green"]])"));
564+
data_file->first_row_id = 100;
565+
data_file->data_sequence_number = 5;
566+
ICEBERG_UNWRAP_OR_FAIL(
567+
auto equality_delete,
568+
MakeEqualityDeleteFile(CreateNewTempFilePathWithSuffix(".parquet"), table_schema_,
569+
R"([[0, "unused", "red"]])", {3}));
570+
FileScanTask task(data_file, {equality_delete});
571+
572+
FileScanTaskReader::Options options{
573+
.io = file_io_,
574+
.table_schema = table_schema_,
575+
.schemas = {table_schema_},
576+
.projected_schema = RowLineageProjection(),
577+
};
578+
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
579+
ICEBERG_UNWRAP_OR_FAIL(auto stream, reader->Open(task));
580+
581+
ASSERT_NO_FATAL_FAILURE(VerifyStream(&stream, R"([[1, 100, 5], [3, 102, 5]])"));
582+
}
583+
507584
TEST_F(FileScanTaskReaderTest, OpenWithEqualityDeletesKeepsInputBatchWhenAllRowsAlive) {
508585
ICEBERG_UNWRAP_OR_FAIL(
509586
auto data_file,

‎src/iceberg/test/incremental_append_scan_test.cc‎

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
#include <memory>
2121
#include <optional>
2222
#include <string>
23+
#include <unordered_map>
2324
#include <vector>
2425

2526
#include <gmock/gmock.h>
@@ -588,6 +589,57 @@ TEST_P(IncrementalAppendScanTest, MultipleRootSnapshots) {
588589
}
589590
}
590591

592+
TEST_P(IncrementalAppendScanTest, PlanRowLineage) {
593+
if (GetParam() < 3) {
594+
GTEST_SKIP() << "Row lineage is only assigned in v3 manifests";
595+
}
596+
597+
auto snapshot_a =
598+
MakeAppendSnapshot(3, 1000L, std::nullopt, 5L, {"/path/to/file_a.parquet"});
599+
snapshot_a->first_row_id = 0;
600+
snapshot_a->added_rows = 1;
601+
602+
auto file_b = MakeDataFile("/path/to/file_b.parquet");
603+
auto entry_b = MakeEntry(ManifestStatus::kAdded, 2000L, 7L, file_b);
604+
auto manifest_b = WriteDataManifest(3, 2000L, {std::move(entry_b)});
605+
manifest_b.first_row_id = 1;
606+
auto manifest_list_b = WriteManifestList(3, 2000L, 1000L, 7L, {manifest_b});
607+
auto snapshot_b = std::make_shared<Snapshot>(Snapshot{
608+
.snapshot_id = 2000L,
609+
.parent_snapshot_id = 1000L,
610+
.sequence_number = 7L,
611+
.timestamp_ms = TimePointMsFromUnixMs(1609459200000L + 2000),
612+
.manifest_list = manifest_list_b,
613+
.summary = {{"operation", "append"}},
614+
.schema_id = schema_->schema_id(),
615+
.first_row_id = 1L,
616+
.added_rows = 1L,
617+
});
618+
auto metadata = MakeTableMetadata(
619+
{snapshot_a, snapshot_b}, 2000L,
620+
{{"main", std::make_shared<SnapshotRef>(SnapshotRef{
621+
.snapshot_id = 2000L, .retention = SnapshotRef::Branch{}})}});
622+
metadata->next_row_id = 2;
623+
624+
ICEBERG_UNWRAP_OR_FAIL(auto builder, MakeScanBuilder<IncrementalAppendScan>(metadata));
625+
builder->FromSnapshot(1000L, /*inclusive=*/true).ToSnapshot(2000L);
626+
ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build());
627+
ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles());
628+
ASSERT_EQ(tasks.size(), 2);
629+
630+
std::unordered_map<std::string, std::shared_ptr<FileScanTask>> tasks_by_path;
631+
for (const auto& task : tasks) {
632+
tasks_by_path.emplace(task->data_file()->file_path, task);
633+
}
634+
ASSERT_EQ(tasks_by_path.size(), 2);
635+
EXPECT_EQ(tasks_by_path.at("/path/to/file_a.parquet")->data_file()->first_row_id, 0);
636+
EXPECT_EQ(
637+
tasks_by_path.at("/path/to/file_a.parquet")->data_file()->data_sequence_number, 5);
638+
EXPECT_EQ(tasks_by_path.at("/path/to/file_b.parquet")->data_file()->first_row_id, 1);
639+
EXPECT_EQ(
640+
tasks_by_path.at("/path/to/file_b.parquet")->data_file()->data_sequence_number, 7);
641+
}
642+
591643
INSTANTIATE_TEST_SUITE_P(IncrementalAppendScanVersions, IncrementalAppendScanTest,
592644
testing::Values(1, 2, 3));
593645

‎src/iceberg/test/incremental_changelog_scan_test.cc‎

Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -454,6 +454,104 @@ TEST_P(IncrementalChangelogScanTest, ManifestRewritesAreIgnored) {
454454
EXPECT_EQ(insert_t3->data_file()->file_path, "/path/to/file_c.parquet");
455455
}
456456

457+
TEST_P(IncrementalChangelogScanTest, PlanAddedRowLineage) {
458+
if (GetParam() < 3) {
459+
GTEST_SKIP() << "Row lineage is only assigned in v3 manifests";
460+
}
461+
462+
auto snapshot_a =
463+
MakeAppendSnapshot(3, 1000L, std::nullopt, 5L, {"/path/to/file_a.parquet"});
464+
snapshot_a->first_row_id = 0;
465+
snapshot_a->added_rows = 1;
466+
467+
auto file_b = MakeDataFile("/path/to/file_b.parquet");
468+
auto entry_b = MakeEntry(ManifestStatus::kAdded, 2000L, 7L, file_b);
469+
auto manifest_b = WriteDataManifest(3, 2000L, {std::move(entry_b)});
470+
manifest_b.first_row_id = 1;
471+
auto manifest_list_b = WriteManifestList(3, 2000L, 1000L, 7L, {manifest_b});
472+
auto snapshot_b = std::make_shared<Snapshot>(Snapshot{
473+
.snapshot_id = 2000L,
474+
.parent_snapshot_id = 1000L,
475+
.sequence_number = 7L,
476+
.timestamp_ms = TimePointMsFromUnixMs(1609459200000L + 2000),
477+
.manifest_list = manifest_list_b,
478+
.summary = {{"operation", "append"}},
479+
.schema_id = schema_->schema_id(),
480+
.first_row_id = 1L,
481+
.added_rows = 1L,
482+
});
483+
auto metadata = MakeTableMetadata(
484+
{snapshot_a, snapshot_b}, 2000L,
485+
{{"main", std::make_shared<SnapshotRef>(SnapshotRef{
486+
.snapshot_id = 2000L, .retention = SnapshotRef::Branch{}})}});
487+
metadata->next_row_id = 2;
488+
489+
ICEBERG_UNWRAP_OR_FAIL(auto builder,
490+
MakeScanBuilder<IncrementalChangelogScan>(metadata));
491+
builder->FromSnapshot(1000L, /*inclusive=*/true).ToSnapshot(2000L);
492+
ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build());
493+
ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles());
494+
ASSERT_EQ(tasks.size(), 2);
495+
SortTasks(tasks);
496+
497+
auto added_a = std::dynamic_pointer_cast<AddedRowsScanTask>(tasks[0]);
498+
ASSERT_NE(added_a, nullptr);
499+
EXPECT_EQ(added_a->commit_snapshot_id(), 1000L);
500+
EXPECT_EQ(added_a->data_file()->first_row_id, 0);
501+
EXPECT_EQ(added_a->data_file()->data_sequence_number, 5);
502+
503+
auto added_b = std::dynamic_pointer_cast<AddedRowsScanTask>(tasks[1]);
504+
ASSERT_NE(added_b, nullptr);
505+
EXPECT_EQ(added_b->commit_snapshot_id(), 2000L);
506+
EXPECT_EQ(added_b->data_file()->first_row_id, 1);
507+
EXPECT_EQ(added_b->data_file()->data_sequence_number, 7);
508+
}
509+
510+
TEST_P(IncrementalChangelogScanTest, PlanDeletedRowLineage) {
511+
if (GetParam() < 3) {
512+
GTEST_SKIP() << "Row lineage is only assigned in v3 manifests";
513+
}
514+
515+
auto snapshot_a =
516+
MakeAppendSnapshot(3, 1000L, std::nullopt, 5L, {"/path/to/file_a.parquet"});
517+
518+
auto deleted_file = MakeDataFile("/path/to/file_a.parquet");
519+
deleted_file->first_row_id = 10;
520+
auto deleted_entry = MakeEntry(ManifestStatus::kDeleted, /*snapshot_id=*/2000L,
521+
/*sequence_number=*/5L, deleted_file);
522+
auto delete_manifest =
523+
WriteDataManifest(3, 2000L, {std::move(deleted_entry)}, unpartitioned_spec_);
524+
auto manifest_list = WriteManifestList(3, 2000L, 1000L, 7L, {delete_manifest});
525+
auto snapshot_b = std::make_shared<Snapshot>(Snapshot{
526+
.snapshot_id = 2000L,
527+
.parent_snapshot_id = 1000L,
528+
.sequence_number = 7L,
529+
.timestamp_ms = TimePointMsFromUnixMs(1609459200000L + 2000),
530+
.manifest_list = manifest_list,
531+
.summary = {{"operation", "overwrite"}},
532+
.schema_id = schema_->schema_id(),
533+
.first_row_id = 1L,
534+
.added_rows = 0L,
535+
});
536+
auto metadata = MakeTableMetadata(
537+
{snapshot_a, snapshot_b}, 2000L,
538+
{{"main", std::make_shared<SnapshotRef>(SnapshotRef{
539+
.snapshot_id = 2000L, .retention = SnapshotRef::Branch{}})}});
540+
541+
ICEBERG_UNWRAP_OR_FAIL(auto builder,
542+
MakeScanBuilder<IncrementalChangelogScan>(metadata));
543+
builder->FromSnapshot(1000L, /*inclusive=*/false).ToSnapshot(2000L);
544+
ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build());
545+
ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles());
546+
ASSERT_EQ(tasks.size(), 1);
547+
548+
auto deleted = std::dynamic_pointer_cast<DeletedDataFileScanTask>(tasks[0]);
549+
ASSERT_NE(deleted, nullptr);
550+
EXPECT_EQ(deleted->data_file()->first_row_id, 10);
551+
EXPECT_EQ(deleted->data_file()->data_sequence_number, 5);
552+
EXPECT_EQ(deleted->commit_snapshot_id(), 2000L);
553+
}
554+
457555
TEST_P(IncrementalChangelogScanTest, DeleteFilesAreNotSupported) {
458556
auto version = GetParam();
459557
if (version < 2) {

0 commit comments

Comments
 (0)