Skip to content

Commit c65a806

Browse files
authored
feat: write snapshot v3 row lineage fields at top level (#791)
1 parent 68ee0f4 commit c65a806

4 files changed

Lines changed: 94 additions & 32 deletions

File tree

src/iceberg/json_serde.cc

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,8 @@ constexpr std::string_view kSequenceNumber = "sequence-number";
105105
constexpr std::string_view kTimestampMs = "timestamp-ms";
106106
constexpr std::string_view kManifestList = "manifest-list";
107107
constexpr std::string_view kSummary = "summary";
108+
constexpr std::string_view kFirstRowId = "first-row-id";
109+
constexpr std::string_view kAddedRows = "added-rows";
108110
constexpr std::string_view kMinSnapshotsToKeep = "min-snapshots-to-keep";
109111
constexpr std::string_view kMaxSnapshotAgeMs = "max-snapshot-age-ms";
110112
constexpr std::string_view kMaxRefAgeMs = "max-ref-age-ms";
@@ -470,6 +472,10 @@ nlohmann::json ToJson(const Snapshot& snapshot) {
470472
json[kSummary] = snapshot.summary;
471473
}
472474
SetOptionalField(json, kSchemaId, snapshot.schema_id);
475+
SetOptionalField(json, kFirstRowId, snapshot.first_row_id);
476+
if (snapshot.first_row_id.has_value()) {
477+
SetOptionalField(json, kAddedRows, snapshot.added_rows);
478+
}
473479
return json;
474480
}
475481

@@ -866,12 +872,32 @@ Result<std::unique_ptr<Snapshot>> SnapshotFromJson(const nlohmann::json& json) {
866872
}
867873
}
868874

875+
ICEBERG_ASSIGN_OR_RAISE(auto first_row_id,
876+
GetJsonValueOptional<int64_t>(json, kFirstRowId));
877+
ICEBERG_ASSIGN_OR_RAISE(auto added_rows,
878+
GetJsonValueOptional<int64_t>(json, kAddedRows));
879+
880+
if (first_row_id.has_value() && first_row_id.value() < 0) {
881+
return JsonParseError("Invalid first-row-id (cannot be negative): {}",
882+
first_row_id.value());
883+
}
884+
if (added_rows.has_value() && added_rows.value() < 0) {
885+
return JsonParseError("Invalid added-rows (cannot be negative): {}",
886+
added_rows.value());
887+
}
888+
if (first_row_id.has_value() && !added_rows.has_value()) {
889+
return JsonParseError("Invalid added-rows (required when first-row-id is set): null");
890+
}
891+
if (!first_row_id.has_value()) {
892+
added_rows = std::nullopt;
893+
}
894+
869895
ICEBERG_ASSIGN_OR_RAISE(auto schema_id, GetJsonValueOptional<int32_t>(json, kSchemaId));
870896

871897
return std::make_unique<Snapshot>(
872898
snapshot_id, parent_snapshot_id,
873899
sequence_number.value_or(TableMetadata::kInitialSequenceNumber), timestamp_ms,
874-
manifest_list, std::move(summary), schema_id);
900+
manifest_list, std::move(summary), schema_id, first_row_id, added_rows);
875901
}
876902

877903
nlohmann::json ToJson(const BlobMetadata& blob_metadata) {

src/iceberg/snapshot.cc

Lines changed: 6 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -161,23 +161,9 @@ std::optional<std::string_view> Snapshot::Operation() const {
161161
return std::nullopt;
162162
}
163163

164-
Result<std::optional<int64_t>> Snapshot::FirstRowId() const {
165-
auto it = summary.find(SnapshotSummaryFields::kFirstRowId);
166-
if (it == summary.end()) {
167-
return std::nullopt;
168-
}
169-
170-
return StringUtils::ParseNumber<int64_t>(it->second);
171-
}
172-
173-
Result<std::optional<int64_t>> Snapshot::AddedRows() const {
174-
auto it = summary.find(SnapshotSummaryFields::kAddedRows);
175-
if (it == summary.end()) {
176-
return std::nullopt;
177-
}
164+
Result<std::optional<int64_t>> Snapshot::FirstRowId() const { return first_row_id; }
178165

179-
return StringUtils::ParseNumber<int64_t>(it->second);
180-
}
166+
Result<std::optional<int64_t>> Snapshot::AddedRows() const { return added_rows; }
181167

182168
bool Snapshot::Equals(const Snapshot& other) const {
183169
if (this == &other) {
@@ -186,7 +172,8 @@ bool Snapshot::Equals(const Snapshot& other) const {
186172
return snapshot_id == other.snapshot_id &&
187173
parent_snapshot_id == other.parent_snapshot_id &&
188174
sequence_number == other.sequence_number && timestamp_ms == other.timestamp_ms &&
189-
schema_id == other.schema_id;
175+
schema_id == other.schema_id && first_row_id == other.first_row_id &&
176+
added_rows == other.added_rows;
190177
}
191178

192179
Result<std::unique_ptr<Snapshot>> Snapshot::Make(
@@ -203,12 +190,6 @@ Result<std::unique_ptr<Snapshot>> Snapshot::Make(
203190
ICEBERG_PRECHECK(!first_row_id.has_value() || added_rows.has_value(),
204191
"Missing added-rows when first-row-id is set");
205192
summary[SnapshotSummaryFields::kOperation] = operation;
206-
if (first_row_id.has_value()) {
207-
summary[SnapshotSummaryFields::kFirstRowId] = std::to_string(first_row_id.value());
208-
}
209-
if (added_rows.has_value()) {
210-
summary[SnapshotSummaryFields::kAddedRows] = std::to_string(added_rows.value());
211-
}
212193
return std::make_unique<Snapshot>(Snapshot{
213194
.snapshot_id = snapshot_id,
214195
.parent_snapshot_id = parent_snapshot_id,
@@ -217,6 +198,8 @@ Result<std::unique_ptr<Snapshot>> Snapshot::Make(
217198
.manifest_list = std::move(manifest_list),
218199
.summary = std::move(summary),
219200
.schema_id = schema_id,
201+
.first_row_id = first_row_id,
202+
.added_rows = first_row_id.has_value() ? added_rows : std::nullopt,
220203
});
221204
}
222205

src/iceberg/snapshot.h

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -403,6 +403,10 @@ struct ICEBERG_EXPORT Snapshot {
403403
std::unordered_map<std::string, std::string> summary;
404404
/// ID of the table's current schema when the snapshot was created.
405405
std::optional<int32_t> schema_id;
406+
/// The row-id of the first newly added row in this snapshot.
407+
std::optional<int64_t> first_row_id;
408+
/// The upper bound of rows with assigned row IDs in this snapshot.
409+
std::optional<int64_t> added_rows;
406410

407411
/// \brief Create a new Snapshot instance with validation on the inputs.
408412
static Result<std::unique_ptr<Snapshot>> Make(

src/iceberg/test/json_serde_test.cc

Lines changed: 57 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -359,6 +359,53 @@ TEST(JsonInternalTest, Snapshot) {
359359
TestJsonConversion(snapshot, expected_json);
360360
}
361361

362+
TEST(JsonInternalTest, SnapshotRowLineageSerializesTopLevelFields) {
363+
ICEBERG_UNWRAP_OR_FAIL(
364+
auto snapshot,
365+
Snapshot::Make(/*sequence_number=*/99, /*snapshot_id=*/1234567890,
366+
/*parent_snapshot_id=*/9876543210,
367+
TimePointMsFromUnixMs(1234567890123), DataOperation::kAppend,
368+
{{SnapshotSummaryFields::kAddedDataFiles, "50"}},
369+
/*schema_id=*/42, "/path/to/manifest_list",
370+
/*first_row_id=*/100, /*added_rows=*/25));
371+
372+
auto json = ToJson(*snapshot);
373+
EXPECT_EQ(json["first-row-id"], 100);
374+
EXPECT_EQ(json["added-rows"], 25);
375+
EXPECT_FALSE(json["summary"].contains("first-row-id"));
376+
EXPECT_FALSE(json["summary"].contains("added-rows"));
377+
}
378+
379+
TEST(JsonInternalTest, SnapshotFromJsonReadsTopLevelRowLineageFields) {
380+
nlohmann::json snapshot_json =
381+
R"({"snapshot-id":1234567890,
382+
"parent-snapshot-id":9876543210,
383+
"sequence-number":99,
384+
"timestamp-ms":1234567890123,
385+
"manifest-list":"/path/to/manifest_list",
386+
"summary":{
387+
"operation":"append",
388+
"added-data-files":"50",
389+
"first-row-id":"101",
390+
"added-rows":"26"
391+
},
392+
"schema-id":42,
393+
"first-row-id":100,
394+
"added-rows":25})"_json;
395+
396+
ICEBERG_UNWRAP_OR_FAIL(auto snapshot, SnapshotFromJson(snapshot_json));
397+
ICEBERG_UNWRAP_OR_FAIL(auto first_row_id, snapshot->FirstRowId());
398+
ICEBERG_UNWRAP_OR_FAIL(auto added_rows, snapshot->AddedRows());
399+
EXPECT_EQ(first_row_id, 100);
400+
EXPECT_EQ(added_rows, 25);
401+
402+
auto json = ToJson(*snapshot);
403+
EXPECT_EQ(json["first-row-id"], 100);
404+
EXPECT_EQ(json["added-rows"], 25);
405+
EXPECT_EQ(json["summary"]["first-row-id"], "101");
406+
EXPECT_EQ(json["summary"]["added-rows"], "26");
407+
}
408+
362409
// FIXME: disable it for now since Iceberg Spark plugin generates
363410
// custom summary keys.
364411
TEST(JsonInternalTest, DISABLED_SnapshotFromJsonWithInvalidSummary) {
@@ -583,19 +630,21 @@ TEST(JsonInternalTest, TableUpdateSetDefaultSortOrder) {
583630
}
584631

585632
TEST(JsonInternalTest, TableUpdateAddSnapshot) {
586-
auto snapshot = std::make_shared<Snapshot>(
587-
Snapshot{.snapshot_id = 123456789,
588-
.parent_snapshot_id = 987654321,
589-
.sequence_number = 5,
590-
.timestamp_ms = TimePointMsFromUnixMs(1234567890000),
591-
.manifest_list = "/path/to/manifest-list.avro",
592-
.summary = {{SnapshotSummaryFields::kOperation, DataOperation::kAppend}},
593-
.schema_id = 1});
633+
ICEBERG_UNWRAP_OR_FAIL(
634+
auto snapshot_unique,
635+
Snapshot::Make(/*sequence_number=*/5, /*snapshot_id=*/123456789,
636+
/*parent_snapshot_id=*/987654321,
637+
TimePointMsFromUnixMs(1234567890000), DataOperation::kAppend,
638+
/*summary=*/{}, /*schema_id=*/1, "/path/to/manifest-list.avro",
639+
/*first_row_id=*/100, /*added_rows=*/25));
640+
std::shared_ptr<Snapshot> snapshot(std::move(snapshot_unique));
594641
table::AddSnapshot update(snapshot);
595642

596643
ICEBERG_UNWRAP_OR_FAIL(auto json, ToJson(update));
597644
EXPECT_EQ(json["action"], "add-snapshot");
598645
EXPECT_TRUE(json.contains("snapshot"));
646+
EXPECT_EQ(json["snapshot"]["first-row-id"], 100);
647+
EXPECT_EQ(json["snapshot"]["added-rows"], 25);
599648

600649
auto parsed = TableUpdateFromJson(json);
601650
ASSERT_THAT(parsed, IsOk());

0 commit comments

Comments
 (0)