u70b3 commented on code in PR #842:
URL: https://github.com/apache/iceberg-cpp/pull/842#discussion_r3996101704
##########
src/iceberg/data/position_delete_writer.cc:
##########
@@ -62,10 +62,43 @@ 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(),
"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;
+ 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);
+
+ const auto* path_view = array_view.children[0];
+ if (ArrowArrayViewComputeNullCount(path_view) != 0) {
Review Comment:
Done in ce7184a. The null check is now per row over `i < data->length` using
`ArrowArrayViewIsNull`, applied to both the `file_path` and the `pos` child
views, so it no longer depends on the child view length.
##########
src/iceberg/data/position_delete_writer.cc:
##########
@@ -62,10 +62,43 @@ 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(),
"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;
Review Comment:
Done in ce7184a. The delete schema is immutable, so it is converted to Arrow
once in `Impl::InitSchema()` and the `ArrowArrayView` is initialized from it
there. `Write` now only rebinds the existing view with
`ArrowArrayViewSetArray`, and `FlushBuffer` reuses the same Arrow schema
instead of rebuilding it. The view and the schema are released in `Impl`'s
destructor.
##########
src/iceberg/data/position_delete_writer.cc:
##########
@@ -62,10 +62,43 @@ 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(),
"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;
+ 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);
+
+ 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<std::string> pending_paths;
Review Comment:
Done in ce7184a, with a slightly different approach. A reusable scratch set
does not actually avoid the churn, since `clear()` frees every node, so the
allocations would keep coming for paths that are already tracked.
Instead the paths are recorded optimistically: `referenced_paths_` is now
`std::set<std::string, std::less<>>`, so a `string_view` that is already
tracked costs only a lookup and no allocation at all. The iterators of the
entries inserted by the current batch are kept in a small reused vector and
erased again if the batch is rejected, which keeps the failure-safe delayed
merge: a rejected batch still leaves no trace in the metadata.
##########
src/iceberg/test/data_writer_test.cc:
##########
@@ -458,6 +453,145 @@ 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");
+ // 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());
+
+ 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,
+ 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) {
Review Comment:
Done in ce7184a. `WriteBatchDataForMultipleFiles` now performs two
successful writes with disjoint paths, so a bug that replaced
`referenced_paths_` instead of unioning across batches would fail the test. It
also asserts the public `WriteResult::referenced_data_files`.
##########
src/iceberg/test/data_writer_test.cc:
##########
@@ -458,6 +453,145 @@ 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");
+ // 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());
+
+ 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,
+ 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));
Review Comment:
Agreed, reworked in ce7184a. The test now writes a successful batch first
and a rejected batch last, so the writer is never used after a failure. The
rejected batch references a valid path and is then rejected by the null path,
and after `Close()` the test verifies that only `data_file_1.parquet` is
tracked. Because the bad path is inserted before the batch fails, this also
covers the rollback path.
##########
src/iceberg/test/data_writer_test.cc:
##########
@@ -458,6 +453,145 @@ 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");
+ // 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());
+
+ 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,
+ 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) {
+ 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());
Review Comment:
Done in ce7184a. `Metadata()` now populates
`WriteResult::referenced_data_files`, and both
`WriteBatchThenDeleteTracksAllReferencedFiles` and
`WriteBatchDataForMultipleFiles` assert it in addition to the per-file hint.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]