diff --git a/src/iceberg/test/merging_snapshot_update_test.cc b/src/iceberg/test/merging_snapshot_update_test.cc index 1c83e4869..00c4f6426 100644 --- a/src/iceberg/test/merging_snapshot_update_test.cc +++ b/src/iceberg/test/merging_snapshot_update_test.cc @@ -26,6 +26,7 @@ #include #include #include +#include #include #include @@ -37,6 +38,8 @@ #include "iceberg/manifest/manifest_entry.h" #include "iceberg/manifest/manifest_reader.h" #include "iceberg/manifest/manifest_writer.h" +#include "iceberg/metrics/commit_report.h" +#include "iceberg/metrics/metrics_reporter.h" #include "iceberg/partition_spec.h" #include "iceberg/puffin_dv_io.h" #include "iceberg/row/partition_values.h" @@ -59,6 +62,20 @@ namespace iceberg { +namespace { + +class MergingSnapshotCapturingReporter final : public MetricsReporter { + public: + Status Report(const MetricsReport& report) override { + reports.push_back(report); + return {}; + } + + std::vector reports; +}; + +} // namespace + class RecordingPuffinDVIO final : public PuffinDVIO { public: struct Call { @@ -564,6 +581,29 @@ TEST_F(MergingSnapshotUpdateTest, CommitNewDataFile) { EXPECT_EQ(snapshot->summary.at(SnapshotSummaryFields::kAddedRecords), "100"); } +TEST_F(MergingSnapshotUpdateTest, CommitReportsCreatedSnapshot) { + auto reporter = std::make_shared(); + ICEBERG_UNWRAP_OR_FAIL(auto op, NewMergeAppend()); + op->ReportWith(reporter); + EXPECT_THAT(op->AddFile(file_a_), IsOk()); + EXPECT_THAT(op->Commit(), IsOk()); + + ASSERT_EQ(reporter->reports.size(), 1U); + ASSERT_TRUE(std::holds_alternative(reporter->reports.front())); + const auto& report = std::get(reporter->reports.front()); + + EXPECT_THAT(table_->Refresh(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, table_->current_snapshot()); + EXPECT_EQ(report.table_name, table_->full_name()); + EXPECT_EQ(report.operation, DataOperation::kAppend); + EXPECT_EQ(report.snapshot_id, snapshot->snapshot_id); + EXPECT_EQ(report.sequence_number, snapshot->sequence_number); + ASSERT_TRUE(report.commit_metrics.added_data_files.has_value()); + EXPECT_EQ(report.commit_metrics.added_data_files->value, 1); + ASSERT_TRUE(report.commit_metrics.added_records.has_value()); + EXPECT_EQ(report.commit_metrics.added_records->value, 100); +} + TEST_F(MergingSnapshotUpdateTest, CommitV3NewDataFileAssignsRowLineage) { UpgradeTableToV3(); diff --git a/src/iceberg/update/merging_snapshot_update.h b/src/iceberg/update/merging_snapshot_update.h index 1c08a0667..6b9d90c40 100644 --- a/src/iceberg/update/merging_snapshot_update.h +++ b/src/iceberg/update/merging_snapshot_update.h @@ -60,12 +60,6 @@ namespace iceberg { /// 5. Write new delete manifests (cached for commit retry) /// 6. Merge data manifests (via data_merge_manager_) /// 7. Merge delete manifests (via delete_merge_manager_) -/// -/// TODO(Guotao): Java MergingSnapshotProducer overrides updateEvent() to return a -/// CreateSnapshotEvent(tableName, operation, snapshotId, sequenceNumber, summary) -/// for commit listeners. The C++ update framework does not yet have an event -/// notification mechanism, so this is intentionally not implemented here. Add it -/// once an equivalent CreateSnapshotEvent / listener facility exists. class ICEBERG_EXPORT MergingSnapshotUpdate : public SnapshotUpdate { public: ~MergingSnapshotUpdate() override = default;