Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
89 changes: 83 additions & 6 deletions src/iceberg/test/file_scan_task_reader_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,14 @@ class FileScanTaskReaderTest : public TempFileTestBase {
table_schema_->schema_id());
}

std::shared_ptr<Schema> RowLineageProjection() const {
return std::make_shared<Schema>(
std::vector<SchemaField>{SchemaField::MakeRequired(1, "id", int32()),
MetadataColumns::kRowId,
MetadataColumns::kLastUpdatedSequenceNumber},
table_schema_->schema_id());
}

Result<ExportedBatch> MakeBatch(const Schema& schema,
const std::string& json_data) const {
ICEBERG_ASSIGN_OR_RAISE(auto arrow_schema, MakeArrowSchema(schema));
Expand Down Expand Up @@ -409,17 +417,12 @@ TEST_F(FileScanTaskReaderTest, ReadLastUpdatedFromDataSeq) {
data_file->first_row_id = 100L;
data_file->data_sequence_number = 5L;
FileScanTask task(data_file);
auto projected_schema = std::make_shared<Schema>(
std::vector<SchemaField>{SchemaField::MakeRequired(1, "id", int32()),
MetadataColumns::kRowId,
MetadataColumns::kLastUpdatedSequenceNumber},
table_schema_->schema_id());

FileScanTaskReader::Options options{
.io = file_io_,
.table_schema = table_schema_,
.schemas = {table_schema_},
.projected_schema = projected_schema,
.projected_schema = RowLineageProjection(),
};
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
auto stream_result = reader->Open(task);
Expand Down Expand Up @@ -504,6 +507,80 @@ TEST_F(FileScanTaskReaderTest, OpenWithEqualityDeletesAddsAndPrunesDeleteOnlyCol
ASSERT_NO_FATAL_FAILURE(VerifyStream(&stream, R"([[1, "Foo"], [3, "Baz"]])"));
}

TEST_F(FileScanTaskReaderTest, PositionDeletesPreserveRowLineage) {
ICEBERG_UNWRAP_OR_FAIL(
auto data_file,
MakeDataFile(table_schema_,
R"([[1, "Foo", "blue"], [2, "Bar", "red"], [3, "Baz", "green"]])"));
data_file->first_row_id = 100;
data_file->data_sequence_number = 5;
ICEBERG_UNWRAP_OR_FAIL(
auto pos_delete, MakePositionDeleteFile(CreateNewTempFilePathWithSuffix(".parquet"),
{1}, data_file->file_path));
FileScanTask task(data_file, {pos_delete});

FileScanTaskReader::Options options{
.io = file_io_,
.table_schema = table_schema_,
.schemas = {table_schema_},
.projected_schema = RowLineageProjection(),
};
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
ICEBERG_UNWRAP_OR_FAIL(auto stream, reader->Open(task));

ASSERT_NO_FATAL_FAILURE(VerifyStream(&stream, R"([[1, 100, 5], [3, 102, 5]])"));
}

TEST_F(FileScanTaskReaderTest, DeletionVectorDeletesPreserveRowLineage) {
ICEBERG_UNWRAP_OR_FAIL(
auto data_file,
MakeDataFile(table_schema_,
R"([[1, "Foo", "blue"], [2, "Bar", "red"], [3, "Baz", "green"]])"));
data_file->first_row_id = 100;
data_file->data_sequence_number = 5;
ICEBERG_UNWRAP_OR_FAIL(
auto deletion_vector,
MakeDeletionVectorFile(CreateNewTempFilePathWithSuffix(".puffin"), {1},
data_file->file_path));
FileScanTask task(data_file, {deletion_vector});

FileScanTaskReader::Options options{
.io = file_io_,
.table_schema = table_schema_,
.schemas = {table_schema_},
.projected_schema = RowLineageProjection(),
};
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
ICEBERG_UNWRAP_OR_FAIL(auto stream, reader->Open(task));

ASSERT_NO_FATAL_FAILURE(VerifyStream(&stream, R"([[1, 100, 5], [3, 102, 5]])"));
}

TEST_F(FileScanTaskReaderTest, EqualityDeletesPreserveRowLineage) {
ICEBERG_UNWRAP_OR_FAIL(
auto data_file,
MakeDataFile(table_schema_,
R"([[1, "Foo", "blue"], [2, "Bar", "red"], [3, "Baz", "green"]])"));
data_file->first_row_id = 100;
data_file->data_sequence_number = 5;
ICEBERG_UNWRAP_OR_FAIL(
auto equality_delete,
MakeEqualityDeleteFile(CreateNewTempFilePathWithSuffix(".parquet"), table_schema_,
R"([[0, "unused", "red"]])", {3}));
FileScanTask task(data_file, {equality_delete});

FileScanTaskReader::Options options{
.io = file_io_,
.table_schema = table_schema_,
.schemas = {table_schema_},
.projected_schema = RowLineageProjection(),
};
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
ICEBERG_UNWRAP_OR_FAIL(auto stream, reader->Open(task));

ASSERT_NO_FATAL_FAILURE(VerifyStream(&stream, R"([[1, 100, 5], [3, 102, 5]])"));
}

TEST_F(FileScanTaskReaderTest, OpenWithEqualityDeletesKeepsInputBatchWhenAllRowsAlive) {
ICEBERG_UNWRAP_OR_FAIL(
auto data_file,
Expand Down
52 changes: 52 additions & 0 deletions src/iceberg/test/incremental_append_scan_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
#include <memory>
#include <optional>
#include <string>
#include <unordered_map>
#include <vector>

#include <gmock/gmock.h>
Expand Down Expand Up @@ -588,6 +589,57 @@ TEST_P(IncrementalAppendScanTest, MultipleRootSnapshots) {
}
}

TEST_P(IncrementalAppendScanTest, PlanRowLineage) {
if (GetParam() < 3) {
GTEST_SKIP() << "Row lineage is only assigned in v3 manifests";
}

auto snapshot_a =
MakeAppendSnapshot(3, 1000L, std::nullopt, 5L, {"/path/to/file_a.parquet"});
snapshot_a->first_row_id = 0;
snapshot_a->added_rows = 1;

auto file_b = MakeDataFile("/path/to/file_b.parquet");
auto entry_b = MakeEntry(ManifestStatus::kAdded, 2000L, 7L, file_b);
auto manifest_b = WriteDataManifest(3, 2000L, {std::move(entry_b)});
manifest_b.first_row_id = 1;
auto manifest_list_b = WriteManifestList(3, 2000L, 1000L, 7L, {manifest_b});
auto snapshot_b = std::make_shared<Snapshot>(Snapshot{
.snapshot_id = 2000L,
.parent_snapshot_id = 1000L,
.sequence_number = 7L,
.timestamp_ms = TimePointMsFromUnixMs(1609459200000L + 2000),
.manifest_list = manifest_list_b,
.summary = {{"operation", "append"}},
.schema_id = schema_->schema_id(),
.first_row_id = 1L,
.added_rows = 1L,
});
auto metadata = MakeTableMetadata(
{snapshot_a, snapshot_b}, 2000L,
{{"main", std::make_shared<SnapshotRef>(SnapshotRef{
.snapshot_id = 2000L, .retention = SnapshotRef::Branch{}})}});
metadata->next_row_id = 2;

ICEBERG_UNWRAP_OR_FAIL(auto builder, MakeScanBuilder<IncrementalAppendScan>(metadata));
builder->FromSnapshot(1000L, /*inclusive=*/true).ToSnapshot(2000L);
ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build());
ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles());
ASSERT_EQ(tasks.size(), 2);

std::unordered_map<std::string, std::shared_ptr<FileScanTask>> tasks_by_path;
for (const auto& task : tasks) {
tasks_by_path.emplace(task->data_file()->file_path, task);
}
ASSERT_EQ(tasks_by_path.size(), 2);
EXPECT_EQ(tasks_by_path.at("/path/to/file_a.parquet")->data_file()->first_row_id, 0);
EXPECT_EQ(
tasks_by_path.at("/path/to/file_a.parquet")->data_file()->data_sequence_number, 5);
EXPECT_EQ(tasks_by_path.at("/path/to/file_b.parquet")->data_file()->first_row_id, 1);
EXPECT_EQ(
tasks_by_path.at("/path/to/file_b.parquet")->data_file()->data_sequence_number, 7);
}

INSTANTIATE_TEST_SUITE_P(IncrementalAppendScanVersions, IncrementalAppendScanTest,
testing::Values(1, 2, 3));

Expand Down
98 changes: 98 additions & 0 deletions src/iceberg/test/incremental_changelog_scan_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -454,6 +454,104 @@ TEST_P(IncrementalChangelogScanTest, ManifestRewritesAreIgnored) {
EXPECT_EQ(insert_t3->data_file()->file_path, "/path/to/file_c.parquet");
}

TEST_P(IncrementalChangelogScanTest, PlanAddedRowLineage) {
if (GetParam() < 3) {
GTEST_SKIP() << "Row lineage is only assigned in v3 manifests";
}

auto snapshot_a =
MakeAppendSnapshot(3, 1000L, std::nullopt, 5L, {"/path/to/file_a.parquet"});
snapshot_a->first_row_id = 0;
snapshot_a->added_rows = 1;

auto file_b = MakeDataFile("/path/to/file_b.parquet");
auto entry_b = MakeEntry(ManifestStatus::kAdded, 2000L, 7L, file_b);
auto manifest_b = WriteDataManifest(3, 2000L, {std::move(entry_b)});
manifest_b.first_row_id = 1;
auto manifest_list_b = WriteManifestList(3, 2000L, 1000L, 7L, {manifest_b});
auto snapshot_b = std::make_shared<Snapshot>(Snapshot{
.snapshot_id = 2000L,
.parent_snapshot_id = 1000L,
.sequence_number = 7L,
.timestamp_ms = TimePointMsFromUnixMs(1609459200000L + 2000),
.manifest_list = manifest_list_b,
.summary = {{"operation", "append"}},
.schema_id = schema_->schema_id(),
.first_row_id = 1L,
.added_rows = 1L,
});
auto metadata = MakeTableMetadata(
{snapshot_a, snapshot_b}, 2000L,
{{"main", std::make_shared<SnapshotRef>(SnapshotRef{
.snapshot_id = 2000L, .retention = SnapshotRef::Branch{}})}});
metadata->next_row_id = 2;

ICEBERG_UNWRAP_OR_FAIL(auto builder,
MakeScanBuilder<IncrementalChangelogScan>(metadata));
builder->FromSnapshot(1000L, /*inclusive=*/true).ToSnapshot(2000L);
ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build());
ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles());
ASSERT_EQ(tasks.size(), 2);
SortTasks(tasks);

auto added_a = std::dynamic_pointer_cast<AddedRowsScanTask>(tasks[0]);
ASSERT_NE(added_a, nullptr);
EXPECT_EQ(added_a->commit_snapshot_id(), 1000L);
EXPECT_EQ(added_a->data_file()->first_row_id, 0);
EXPECT_EQ(added_a->data_file()->data_sequence_number, 5);

auto added_b = std::dynamic_pointer_cast<AddedRowsScanTask>(tasks[1]);
ASSERT_NE(added_b, nullptr);
EXPECT_EQ(added_b->commit_snapshot_id(), 2000L);
EXPECT_EQ(added_b->data_file()->first_row_id, 1);
EXPECT_EQ(added_b->data_file()->data_sequence_number, 7);
}

TEST_P(IncrementalChangelogScanTest, PlanDeletedRowLineage) {
if (GetParam() < 3) {
GTEST_SKIP() << "Row lineage is only assigned in v3 manifests";
}

auto snapshot_a =
MakeAppendSnapshot(3, 1000L, std::nullopt, 5L, {"/path/to/file_a.parquet"});

auto deleted_file = MakeDataFile("/path/to/file_a.parquet");
deleted_file->first_row_id = 10;
auto deleted_entry = MakeEntry(ManifestStatus::kDeleted, /*snapshot_id=*/2000L,
/*sequence_number=*/5L, deleted_file);
auto delete_manifest =
WriteDataManifest(3, 2000L, {std::move(deleted_entry)}, unpartitioned_spec_);
auto manifest_list = WriteManifestList(3, 2000L, 1000L, 7L, {delete_manifest});
auto snapshot_b = std::make_shared<Snapshot>(Snapshot{
.snapshot_id = 2000L,
.parent_snapshot_id = 1000L,
.sequence_number = 7L,
.timestamp_ms = TimePointMsFromUnixMs(1609459200000L + 2000),
.manifest_list = manifest_list,
.summary = {{"operation", "overwrite"}},
.schema_id = schema_->schema_id(),
.first_row_id = 1L,
.added_rows = 0L,
});
auto metadata = MakeTableMetadata(
{snapshot_a, snapshot_b}, 2000L,
{{"main", std::make_shared<SnapshotRef>(SnapshotRef{
.snapshot_id = 2000L, .retention = SnapshotRef::Branch{}})}});

ICEBERG_UNWRAP_OR_FAIL(auto builder,
MakeScanBuilder<IncrementalChangelogScan>(metadata));
builder->FromSnapshot(1000L, /*inclusive=*/false).ToSnapshot(2000L);
ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build());
ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles());
ASSERT_EQ(tasks.size(), 1);

auto deleted = std::dynamic_pointer_cast<DeletedDataFileScanTask>(tasks[0]);
ASSERT_NE(deleted, nullptr);
EXPECT_EQ(deleted->data_file()->first_row_id, 10);
EXPECT_EQ(deleted->data_file()->data_sequence_number, 5);
EXPECT_EQ(deleted->commit_snapshot_id(), 2000L);
}

TEST_P(IncrementalChangelogScanTest, DeleteFilesAreNotSupported) {
auto version = GetParam();
if (version < 2) {
Expand Down
Loading
Loading