From f7857c5e2560e1b8ad700a27a8147996c04c4b6a Mon Sep 17 00:00:00 2001 From: kid Date: Fri, 24 Jul 2026 13:02:47 +0800 Subject: [PATCH 1/4] fix: track batch position delete references --- src/iceberg/data/position_delete_writer.cc | 34 ++++++- src/iceberg/test/data_writer_test.cc | 101 +++++++++++++++++++-- 2 files changed, 123 insertions(+), 12 deletions(-) diff --git a/src/iceberg/data/position_delete_writer.cc b/src/iceberg/data/position_delete_writer.cc index fea1ffc19..3e9181d82 100644 --- a/src/iceberg/data/position_delete_writer.cc +++ b/src/iceberg/data/position_delete_writer.cc @@ -62,10 +62,40 @@ class PositionDeleteWriter::Impl { } Status Write(ArrowArray* data) { + ICEBERG_PRECHECK(data != nullptr, "Position delete data must not be null"); ICEBERG_PRECHECK(buffered_paths_.empty(), "Cannot write batch data when there are buffered deletes."); - // TODO(anyone): Extract file paths from ArrowArray to update referenced_paths_. - return writer_->Write(data); + + ArrowSchema arrow_schema; + ICEBERG_RETURN_UNEXPECTED(ToArrowSchema(*delete_schema_, &arrow_schema)); + internal::ArrowSchemaGuard schema_guard(&arrow_schema); + + ArrowArrayView array_view; + internal::ArrowArrayViewGuard view_guard(&array_view); + ArrowError error; + ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( + ArrowArrayViewInitFromSchema(&array_view, &arrow_schema, &error), error); + ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( + ArrowArrayViewSetArray(&array_view, data, &error), error); + + const auto* path_view = array_view.children[0]; + if (ArrowArrayViewComputeNullCount(path_view) != 0) { + return InvalidArrowData("Position delete file paths must not contain null values"); + } + + std::set pending_paths; + for (int64_t i = 0; i < data->length; ++i) { + auto path = ArrowArrayViewGetStringUnsafe(path_view, i); + if (path.size_bytes == 0) { + pending_paths.emplace(); + } else { + pending_paths.emplace(path.data, static_cast(path.size_bytes)); + } + } + + ICEBERG_RETURN_UNEXPECTED(writer_->Write(data)); + referenced_paths_.merge(pending_paths); + return {}; } Status WriteDelete(std::string_view file_path, int64_t pos) { diff --git a/src/iceberg/test/data_writer_test.cc b/src/iceberg/test/data_writer_test.cc index ca1986429..4aef0c67f 100644 --- a/src/iceberg/test/data_writer_test.cc +++ b/src/iceberg/test/data_writer_test.cc @@ -26,6 +26,7 @@ #include #include "iceberg/arrow/arrow_io_internal.h" +#include "iceberg/arrow_c_data_guard_internal.h" #include "iceberg/avro/avro_register.h" #include "iceberg/data/equality_delete_writer.h" #include "iceberg/data/position_delete_writer.h" @@ -344,18 +345,12 @@ class PositionDeleteWriterTest : public DataWriterTest { }; } - std::shared_ptr<::arrow::Array> CreatePositionDeleteData() { + std::shared_ptr<::arrow::Array> CreatePositionDeleteData( + std::string_view json = + R"([["data_file_1.parquet", 0], ["data_file_1.parquet", 5], ["data_file_1.parquet", 10]])") { auto delete_schema = std::make_shared(std::vector{ MetadataColumns::kDeleteFilePath, MetadataColumns::kDeleteFilePos}); - - ArrowSchema arrow_c_schema; - ICEBERG_THROW_NOT_OK(ToArrowSchema(*delete_schema, &arrow_c_schema)); - auto arrow_type = ::arrow::ImportType(&arrow_c_schema).ValueOrDie(); - - return ::arrow::json::ArrayFromJSONString( - ::arrow::struct_(arrow_type->fields()), - R"([["data_file_1.parquet", 0], ["data_file_1.parquet", 5], ["data_file_1.parquet", 10]])") - .ValueOrDie(); + return CreateArray(*delete_schema, json); } }; @@ -458,6 +453,92 @@ TEST_F(PositionDeleteWriterTest, WriteBatchData) { const auto& data_file = metadata_result.value().data_files[0]; EXPECT_EQ(data_file->content, DataFile::Content::kPositionDeletes); EXPECT_GT(data_file->file_size_in_bytes, 0); + ASSERT_TRUE(data_file->referenced_data_file.has_value()); + EXPECT_EQ(data_file->referenced_data_file.value(), "data_file_1.parquet"); +} + +TEST_F(PositionDeleteWriterTest, WriteBatchDataForMultipleFiles) { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + auto test_data = CreatePositionDeleteData( + R"([["data_file_1.parquet", 0], ["data_file_2.parquet", 5]])"); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); + ASSERT_THAT(writer->Write(&arrow_array), IsOk()); + ASSERT_THAT(writer->Close(), IsOk()); + + auto metadata_result = writer->Metadata(); + ASSERT_THAT(metadata_result, IsOk()); + + const auto& data_file = metadata_result.value().data_files[0]; + EXPECT_FALSE(data_file->referenced_data_file.has_value()); + EXPECT_FALSE( + data_file->lower_bounds.contains(MetadataColumns::kDeleteFilePathColumnId)); + EXPECT_FALSE(data_file->lower_bounds.contains(MetadataColumns::kDeleteFilePosColumnId)); + EXPECT_FALSE( + data_file->upper_bounds.contains(MetadataColumns::kDeleteFilePathColumnId)); + EXPECT_FALSE(data_file->upper_bounds.contains(MetadataColumns::kDeleteFilePosColumnId)); +} + +TEST_F(PositionDeleteWriterTest, WriteBatchThenDeleteTracksAllReferencedFiles) { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + auto test_data = CreatePositionDeleteData(R"([["data_file_1.parquet", 0]])"); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); + ASSERT_THAT(writer->Write(&arrow_array), IsOk()); + ASSERT_THAT(writer->WriteDelete("data_file_2.parquet", 5), IsOk()); + ASSERT_THAT(writer->Close(), IsOk()); + + auto metadata_result = writer->Metadata(); + ASSERT_THAT(metadata_result, IsOk()); + EXPECT_FALSE(metadata_result.value().data_files[0]->referenced_data_file.has_value()); +} + +TEST_F(PositionDeleteWriterTest, WriteBatchRejectsNullFilePath) { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + auto test_data = CreatePositionDeleteData(R"([[null, 0]])"); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); + internal::ArrowArrayGuard array_guard(&arrow_array); + + auto result = writer->Write(&arrow_array); + ASSERT_THAT(result, IsError(ErrorKind::kInvalidArrowData)); + EXPECT_THAT(result, + HasErrorMessage("Position delete file paths must not contain null values")); +} + +TEST_F(PositionDeleteWriterTest, WriteBatchRejectsNullData) { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + auto result = writer->Write(nullptr); + ASSERT_THAT(result, IsError(ErrorKind::kInvalidArgument)); + EXPECT_THAT(result, HasErrorMessage("Position delete data must not be null")); +} + +TEST_F(PositionDeleteWriterTest, WriteEmptyBatchDoesNotAddReferencedFiles) { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + auto test_data = CreatePositionDeleteData("[]"); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); + ASSERT_THAT(writer->Write(&arrow_array), IsOk()); + ASSERT_THAT(writer->Close(), IsOk()); + + auto metadata_result = writer->Metadata(); + ASSERT_THAT(metadata_result, IsOk()); + EXPECT_FALSE(metadata_result.value().data_files[0]->referenced_data_file.has_value()); } TEST_F(PositionDeleteWriterTest, AutoFlushOnThreshold) { From 08ae1d990e53ed590ad33fb316c93c6ee5bb87cd Mon Sep 17 00:00:00 2001 From: kid Date: Sat, 25 Jul 2026 01:18:36 +0800 Subject: [PATCH 2/4] fix: validate batch offset and harden array view lifetime --- src/iceberg/data/position_delete_writer.cc | 4 +- src/iceberg/test/data_writer_test.cc | 51 ++++++++++++++++++++++ 2 files changed, 54 insertions(+), 1 deletion(-) diff --git a/src/iceberg/data/position_delete_writer.cc b/src/iceberg/data/position_delete_writer.cc index 3e9181d82..2d4682d0a 100644 --- a/src/iceberg/data/position_delete_writer.cc +++ b/src/iceberg/data/position_delete_writer.cc @@ -63,6 +63,8 @@ class PositionDeleteWriter::Impl { Status Write(ArrowArray* data) { ICEBERG_PRECHECK(data != nullptr, "Position delete data must not be null"); + ICEBERG_PRECHECK(data->offset == 0, + "Position delete data with a non-zero offset is not supported"); ICEBERG_PRECHECK(buffered_paths_.empty(), "Cannot write batch data when there are buffered deletes."); @@ -71,10 +73,10 @@ class PositionDeleteWriter::Impl { internal::ArrowSchemaGuard schema_guard(&arrow_schema); ArrowArrayView array_view; - internal::ArrowArrayViewGuard view_guard(&array_view); ArrowError error; ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( ArrowArrayViewInitFromSchema(&array_view, &arrow_schema, &error), error); + internal::ArrowArrayViewGuard view_guard(&array_view); ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( ArrowArrayViewSetArray(&array_view, data, &error), error); diff --git a/src/iceberg/test/data_writer_test.cc b/src/iceberg/test/data_writer_test.cc index 4aef0c67f..70af9c7c5 100644 --- a/src/iceberg/test/data_writer_test.cc +++ b/src/iceberg/test/data_writer_test.cc @@ -455,6 +455,57 @@ TEST_F(PositionDeleteWriterTest, WriteBatchData) { EXPECT_GT(data_file->file_size_in_bytes, 0); ASSERT_TRUE(data_file->referenced_data_file.has_value()); EXPECT_EQ(data_file->referenced_data_file.value(), "data_file_1.parquet"); + // Bounds for delete metadata columns are kept when referencing a single file. + EXPECT_TRUE(data_file->lower_bounds.contains(MetadataColumns::kDeleteFilePathColumnId)); + EXPECT_TRUE(data_file->lower_bounds.contains(MetadataColumns::kDeleteFilePosColumnId)); + EXPECT_TRUE(data_file->upper_bounds.contains(MetadataColumns::kDeleteFilePathColumnId)); + EXPECT_TRUE(data_file->upper_bounds.contains(MetadataColumns::kDeleteFilePosColumnId)); +} + +TEST_F(PositionDeleteWriterTest, WriteBatchRejectsSlicedData) { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + auto test_data = CreatePositionDeleteData( + R"([["data_file_1.parquet", 0], ["data_file_1.parquet", 5]])"); + auto sliced = test_data->Slice(1, 1); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*sliced, &arrow_array).ok()); + internal::ArrowArrayGuard array_guard(&arrow_array); + + auto result = writer->Write(&arrow_array); + ASSERT_THAT(result, IsError(ErrorKind::kInvalidArgument)); + EXPECT_THAT( + result, + HasErrorMessage("Position delete data with a non-zero offset is not supported")); +} + +TEST_F(PositionDeleteWriterTest, FailedBatchWriteDoesNotTrackReferencedFiles) { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + // A rejected batch must not contribute referenced paths. + auto bad_data = + CreatePositionDeleteData(R"([[null, 0], ["data_file_bad.parquet", 1]])"); + ArrowArray bad_array; + ASSERT_TRUE(::arrow::ExportArray(*bad_data, &bad_array).ok()); + internal::ArrowArrayGuard bad_array_guard(&bad_array); + ASSERT_THAT(writer->Write(&bad_array), IsError(ErrorKind::kInvalidArrowData)); + + auto test_data = CreatePositionDeleteData(R"([["data_file_1.parquet", 0]])"); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); + ASSERT_THAT(writer->Write(&arrow_array), IsOk()); + ASSERT_THAT(writer->Close(), IsOk()); + + auto metadata_result = writer->Metadata(); + ASSERT_THAT(metadata_result, IsOk()); + + const auto& data_file = metadata_result.value().data_files[0]; + ASSERT_TRUE(data_file->referenced_data_file.has_value()); + EXPECT_EQ(data_file->referenced_data_file.value(), "data_file_1.parquet"); } TEST_F(PositionDeleteWriterTest, WriteBatchDataForMultipleFiles) { From 7a0e768d3588cce96bd89c7f324af92124e25991 Mon Sep 17 00:00:00 2001 From: kid Date: Tue, 28 Jul 2026 13:24:27 +0800 Subject: [PATCH 3/4] fix: release rejected position delete batches Take ownership before validation so early-return paths honor the FileWriter contract and do not leak Arrow buffers. Assert that rejected sliced and null-path batches are released. --- src/iceberg/data/position_delete_writer.cc | 1 + src/iceberg/test/data_writer_test.cc | 6 ++++-- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/src/iceberg/data/position_delete_writer.cc b/src/iceberg/data/position_delete_writer.cc index 2d4682d0a..faaed504b 100644 --- a/src/iceberg/data/position_delete_writer.cc +++ b/src/iceberg/data/position_delete_writer.cc @@ -63,6 +63,7 @@ class PositionDeleteWriter::Impl { Status Write(ArrowArray* data) { ICEBERG_PRECHECK(data != nullptr, "Position delete data must not be null"); + internal::ArrowArrayGuard data_guard(data); ICEBERG_PRECHECK(data->offset == 0, "Position delete data with a non-zero offset is not supported"); ICEBERG_PRECHECK(buffered_paths_.empty(), diff --git a/src/iceberg/test/data_writer_test.cc b/src/iceberg/test/data_writer_test.cc index 70af9c7c5..d8fdfd90e 100644 --- a/src/iceberg/test/data_writer_test.cc +++ b/src/iceberg/test/data_writer_test.cc @@ -472,9 +472,10 @@ TEST_F(PositionDeleteWriterTest, WriteBatchRejectsSlicedData) { auto sliced = test_data->Slice(1, 1); ArrowArray arrow_array; ASSERT_TRUE(::arrow::ExportArray(*sliced, &arrow_array).ok()); - internal::ArrowArrayGuard array_guard(&arrow_array); auto result = writer->Write(&arrow_array); + EXPECT_EQ(arrow_array.release, nullptr); + internal::ArrowArrayGuard array_guard(&arrow_array); ASSERT_THAT(result, IsError(ErrorKind::kInvalidArgument)); EXPECT_THAT( result, @@ -558,9 +559,10 @@ TEST_F(PositionDeleteWriterTest, WriteBatchRejectsNullFilePath) { auto test_data = CreatePositionDeleteData(R"([[null, 0]])"); ArrowArray arrow_array; ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); - internal::ArrowArrayGuard array_guard(&arrow_array); auto result = writer->Write(&arrow_array); + EXPECT_EQ(arrow_array.release, nullptr); + internal::ArrowArrayGuard array_guard(&arrow_array); ASSERT_THAT(result, IsError(ErrorKind::kInvalidArrowData)); EXPECT_THAT(result, HasErrorMessage("Position delete file paths must not contain null values")); From ce7184aad9cbbda8e031af5881dcc376245c67a6 Mon Sep 17 00:00:00 2001 From: kid Date: Sat, 12 Sep 2026 18:06:41 +0800 Subject: [PATCH 4/4] fix: address review feedback on batch position delete tracking Build the Arrow delete schema and its array view once in the writer instead of rebuilding them for every batch, and validate the batch per row: null file paths and null positions are rejected, and an empty file path is now an error rather than an empty referenced path. Record referenced paths as they are seen using a transparent lookup and roll back the entries added by a batch that is rejected. A batch that keeps referencing the same file now costs a lookup instead of a scratch allocation per unique path, while a rejected batch still leaves no trace in the metadata. WriteResult also exposes the public referenced_data_files list now. Tests: union disjoint paths across two successful batches, assert referenced_data_files instead of only the per-file hint, and fold the null input cases into one invalid-input test that covers empty paths too. --- src/iceberg/data/position_delete_writer.cc | 102 ++++++++++++----- src/iceberg/test/data_writer_test.cc | 127 +++++++++++++++------ 2 files changed, 168 insertions(+), 61 deletions(-) diff --git a/src/iceberg/data/position_delete_writer.cc b/src/iceberg/data/position_delete_writer.cc index faaed504b..6238774eb 100644 --- a/src/iceberg/data/position_delete_writer.cc +++ b/src/iceberg/data/position_delete_writer.cc @@ -19,8 +19,11 @@ #include "iceberg/data/position_delete_writer.h" +#include #include #include +#include +#include #include #include @@ -57,8 +60,17 @@ class PositionDeleteWriter::Impl { ICEBERG_ASSIGN_OR_RAISE(auto writer, WriterFactoryRegistry::Open(options.format, writer_options)); - return std::unique_ptr( + auto impl = std::unique_ptr( new Impl(std::move(options), std::move(delete_schema), std::move(writer))); + ICEBERG_RETURN_UNEXPECTED(impl->InitSchema()); + return impl; + } + + ~Impl() { + ArrowArrayViewReset(&array_view_); + if (arrow_schema_.release != nullptr) { + ArrowSchemaRelease(&arrow_schema_); + } } Status Write(ArrowArray* data) { @@ -69,35 +81,46 @@ class PositionDeleteWriter::Impl { ICEBERG_PRECHECK(buffered_paths_.empty(), "Cannot write batch data when there are buffered deletes."); - ArrowSchema arrow_schema; - ICEBERG_RETURN_UNEXPECTED(ToArrowSchema(*delete_schema_, &arrow_schema)); - internal::ArrowSchemaGuard schema_guard(&arrow_schema); - - ArrowArrayView array_view; ArrowError error; ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( - ArrowArrayViewInitFromSchema(&array_view, &arrow_schema, &error), error); - internal::ArrowArrayViewGuard view_guard(&array_view); - ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( - ArrowArrayViewSetArray(&array_view, data, &error), error); + ArrowArrayViewSetArray(&array_view_, data, &error), error); - const auto* path_view = array_view.children[0]; - if (ArrowArrayViewComputeNullCount(path_view) != 0) { - return InvalidArrowData("Position delete file paths must not contain null values"); - } + const auto* path_view = array_view_.children[0]; + const auto* pos_view = array_view_.children[1]; - std::set pending_paths; + // A batch usually references files that earlier batches already referenced, so + // record paths optimistically: a known path costs a lookup instead of a scratch + // allocation. Entries added by this batch are rolled back unless the write + // succeeds, so a rejected batch still leaves no trace in the metadata. + pending_references_.clear(); for (int64_t i = 0; i < data->length; ++i) { + if (ArrowArrayViewIsNull(path_view, i)) { + RollbackPendingReferences(); + return InvalidArrowData( + "Position delete file paths must not contain null values"); + } + if (ArrowArrayViewIsNull(pos_view, i)) { + RollbackPendingReferences(); + return InvalidArrowData("Position delete positions must not contain null values"); + } auto path = ArrowArrayViewGetStringUnsafe(path_view, i); if (path.size_bytes == 0) { - pending_paths.emplace(); - } else { - pending_paths.emplace(path.data, static_cast(path.size_bytes)); + RollbackPendingReferences(); + return InvalidArrowData("Position delete file paths must not be empty"); + } + std::string_view file_path(path.data, static_cast(path.size_bytes)); + if (!referenced_paths_.contains(file_path)) { + pending_references_.push_back( + referenced_paths_.insert(std::string(file_path)).first); } } - ICEBERG_RETURN_UNEXPECTED(writer_->Write(data)); - referenced_paths_.merge(pending_paths); + Status status = writer_->Write(data); + if (!status) { + RollbackPendingReferences(); + return status; + } + pending_references_.clear(); return {}; } @@ -196,6 +219,8 @@ class PositionDeleteWriter::Impl { WriteResult result; result.data_files.push_back(std::move(data_file)); + result.referenced_data_files.assign(referenced_paths_.begin(), + referenced_paths_.end()); return result; } @@ -206,15 +231,29 @@ class PositionDeleteWriter::Impl { delete_schema_(std::move(delete_schema)), writer_(std::move(writer)) {} - Status FlushBuffer() { - ArrowSchema arrow_schema; - ICEBERG_RETURN_UNEXPECTED(ToArrowSchema(*delete_schema_, &arrow_schema)); - internal::ArrowSchemaGuard schema_guard(&arrow_schema); + Status InitSchema() { + ICEBERG_RETURN_UNEXPECTED(ToArrowSchema(*delete_schema_, &arrow_schema_)); + ArrowError error; + // The delete schema never changes, so the view is initialized once here and + // merely rebound to each incoming batch in Write. + ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( + ArrowArrayViewInitFromSchema(&array_view_, &arrow_schema_, &error), error); + return {}; + } + /// \brief Undo the referenced paths added by the batch that failed to write. + void RollbackPendingReferences() { + for (auto it : pending_references_) { + referenced_paths_.erase(it); + } + pending_references_.clear(); + } + + Status FlushBuffer() { ArrowArray array; ArrowError error; ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( - ArrowArrayInitFromSchema(&array, &arrow_schema, &error), error); + ArrowArrayInitFromSchema(&array, &arrow_schema_, &error), error); internal::ArrowArrayGuard array_guard(&array); ICEBERG_NANOARROW_RETURN_UNEXPECTED(ArrowArrayStartAppending(&array)); @@ -238,13 +277,24 @@ class PositionDeleteWriter::Impl { return {}; } + // Transparent comparator so that paths arriving as string_view can be looked up + // without first materializing a std::string. + using ReferencedPaths = std::set>; + PositionDeleteWriterOptions options_; std::shared_ptr delete_schema_; std::unique_ptr writer_; + // The immutable delete schema in Arrow form, paired with the view bound to it. + // Declared before the view and released after it in the destructor. + ArrowSchema arrow_schema_{}; + ArrowArrayView array_view_{}; bool closed_ = false; std::vector buffered_paths_; std::vector buffered_positions_; - std::set referenced_paths_; + ReferencedPaths referenced_paths_; + // Iterators of the entries added to referenced_paths_ by the batch in flight, so + // that they can be removed again if that batch is rejected. + std::vector pending_references_; }; PositionDeleteWriter::PositionDeleteWriter(std::unique_ptr impl) diff --git a/src/iceberg/test/data_writer_test.cc b/src/iceberg/test/data_writer_test.cc index d8fdfd90e..d95361675 100644 --- a/src/iceberg/test/data_writer_test.cc +++ b/src/iceberg/test/data_writer_test.cc @@ -46,7 +46,9 @@ namespace iceberg { +using ::testing::ElementsAre; using ::testing::HasSubstr; +using ::testing::UnorderedElementsAre; class DataWriterTest : public ::testing::Test { protected: @@ -487,26 +489,30 @@ TEST_F(PositionDeleteWriterTest, FailedBatchWriteDoesNotTrackReferencedFiles) { ASSERT_THAT(writer_result, IsOk()); auto writer = std::move(writer_result.value()); - // A rejected batch must not contribute referenced paths. + auto good_data = CreatePositionDeleteData(R"([["data_file_1.parquet", 0]])"); + ArrowArray good_array; + ASSERT_TRUE(::arrow::ExportArray(*good_data, &good_array).ok()); + ASSERT_THAT(writer->Write(&good_array), IsOk()); + + // The batch references a valid path before the null path rejects it, and none of + // its paths may end up in the metadata. auto bad_data = - CreatePositionDeleteData(R"([[null, 0], ["data_file_bad.parquet", 1]])"); + CreatePositionDeleteData(R"([["data_file_bad.parquet", 1], [null, 2]])"); ArrowArray bad_array; ASSERT_TRUE(::arrow::ExportArray(*bad_data, &bad_array).ok()); internal::ArrowArrayGuard bad_array_guard(&bad_array); ASSERT_THAT(writer->Write(&bad_array), IsError(ErrorKind::kInvalidArrowData)); - auto test_data = CreatePositionDeleteData(R"([["data_file_1.parquet", 0]])"); - ArrowArray arrow_array; - ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); - ASSERT_THAT(writer->Write(&arrow_array), IsOk()); ASSERT_THAT(writer->Close(), IsOk()); auto metadata_result = writer->Metadata(); ASSERT_THAT(metadata_result, IsOk()); - const auto& data_file = metadata_result.value().data_files[0]; + const auto& write_result = metadata_result.value(); + const auto& data_file = write_result.data_files[0]; ASSERT_TRUE(data_file->referenced_data_file.has_value()); EXPECT_EQ(data_file->referenced_data_file.value(), "data_file_1.parquet"); + EXPECT_THAT(write_result.referenced_data_files, ElementsAre("data_file_1.parquet")); } TEST_F(PositionDeleteWriterTest, WriteBatchDataForMultipleFiles) { @@ -514,18 +520,27 @@ TEST_F(PositionDeleteWriterTest, WriteBatchDataForMultipleFiles) { ASSERT_THAT(writer_result, IsOk()); auto writer = std::move(writer_result.value()); - auto test_data = CreatePositionDeleteData( - R"([["data_file_1.parquet", 0], ["data_file_2.parquet", 5]])"); - ArrowArray arrow_array; - ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); - ASSERT_THAT(writer->Write(&arrow_array), IsOk()); + // Disjoint paths across two successful batches must be unioned, not replaced. + auto first_data = CreatePositionDeleteData(R"([["data_file_1.parquet", 0]])"); + ArrowArray first_array; + ASSERT_TRUE(::arrow::ExportArray(*first_data, &first_array).ok()); + ASSERT_THAT(writer->Write(&first_array), IsOk()); + + auto second_data = CreatePositionDeleteData(R"([["data_file_2.parquet", 5]])"); + ArrowArray second_array; + ASSERT_TRUE(::arrow::ExportArray(*second_data, &second_array).ok()); + ASSERT_THAT(writer->Write(&second_array), IsOk()); + ASSERT_THAT(writer->Close(), IsOk()); auto metadata_result = writer->Metadata(); ASSERT_THAT(metadata_result, IsOk()); - const auto& data_file = metadata_result.value().data_files[0]; + const auto& write_result = metadata_result.value(); + const auto& data_file = write_result.data_files[0]; EXPECT_FALSE(data_file->referenced_data_file.has_value()); + EXPECT_THAT(write_result.referenced_data_files, + UnorderedElementsAre("data_file_1.parquet", "data_file_2.parquet")); EXPECT_FALSE( data_file->lower_bounds.contains(MetadataColumns::kDeleteFilePathColumnId)); EXPECT_FALSE(data_file->lower_bounds.contains(MetadataColumns::kDeleteFilePosColumnId)); @@ -548,34 +563,76 @@ TEST_F(PositionDeleteWriterTest, WriteBatchThenDeleteTracksAllReferencedFiles) { auto metadata_result = writer->Metadata(); ASSERT_THAT(metadata_result, IsOk()); - EXPECT_FALSE(metadata_result.value().data_files[0]->referenced_data_file.has_value()); + const auto& write_result = metadata_result.value(); + EXPECT_FALSE(write_result.data_files[0]->referenced_data_file.has_value()); + EXPECT_THAT(write_result.referenced_data_files, + UnorderedElementsAre("data_file_1.parquet", "data_file_2.parquet")); } -TEST_F(PositionDeleteWriterTest, WriteBatchRejectsNullFilePath) { - auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); - ASSERT_THAT(writer_result, IsOk()); - auto writer = std::move(writer_result.value()); +TEST_F(PositionDeleteWriterTest, WriteBatchRejectsInvalidInput) { + // A null array. + { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); - auto test_data = CreatePositionDeleteData(R"([[null, 0]])"); - ArrowArray arrow_array; - ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); + auto result = writer->Write(nullptr); + ASSERT_THAT(result, IsError(ErrorKind::kInvalidArgument)); + EXPECT_THAT(result, HasErrorMessage("Position delete data must not be null")); + } - auto result = writer->Write(&arrow_array); - EXPECT_EQ(arrow_array.release, nullptr); - internal::ArrowArrayGuard array_guard(&arrow_array); - ASSERT_THAT(result, IsError(ErrorKind::kInvalidArrowData)); - EXPECT_THAT(result, - HasErrorMessage("Position delete file paths must not contain null values")); -} + // A null file path. + { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); -TEST_F(PositionDeleteWriterTest, WriteBatchRejectsNullData) { - auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); - ASSERT_THAT(writer_result, IsOk()); - auto writer = std::move(writer_result.value()); + auto test_data = CreatePositionDeleteData(R"([[null, 0]])"); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); - auto result = writer->Write(nullptr); - ASSERT_THAT(result, IsError(ErrorKind::kInvalidArgument)); - EXPECT_THAT(result, HasErrorMessage("Position delete data must not be null")); + auto result = writer->Write(&arrow_array); + EXPECT_EQ(arrow_array.release, nullptr); + internal::ArrowArrayGuard array_guard(&arrow_array); + ASSERT_THAT(result, IsError(ErrorKind::kInvalidArrowData)); + EXPECT_THAT(result, HasErrorMessage( + "Position delete file paths must not contain null values")); + } + + // A null position. + { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + auto test_data = CreatePositionDeleteData(R"([["data_file_1.parquet", null]])"); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); + + auto result = writer->Write(&arrow_array); + EXPECT_EQ(arrow_array.release, nullptr); + internal::ArrowArrayGuard array_guard(&arrow_array); + ASSERT_THAT(result, IsError(ErrorKind::kInvalidArrowData)); + EXPECT_THAT(result, HasErrorMessage( + "Position delete positions must not contain null values")); + } + + // An empty file path. + { + auto writer_result = PositionDeleteWriter::Make(MakeDeleteOptions()); + ASSERT_THAT(writer_result, IsOk()); + auto writer = std::move(writer_result.value()); + + auto test_data = CreatePositionDeleteData(R"([["", 0]])"); + ArrowArray arrow_array; + ASSERT_TRUE(::arrow::ExportArray(*test_data, &arrow_array).ok()); + + auto result = writer->Write(&arrow_array); + EXPECT_EQ(arrow_array.release, nullptr); + internal::ArrowArrayGuard array_guard(&arrow_array); + ASSERT_THAT(result, IsError(ErrorKind::kInvalidArrowData)); + EXPECT_THAT(result, HasErrorMessage("Position delete file paths must not be empty")); + } } TEST_F(PositionDeleteWriterTest, WriteEmptyBatchDoesNotAddReferencedFiles) {