lxy-9602 commented on code in PR #392:
URL: https://github.com/apache/paimon-cpp/pull/392#discussion_r4122554341


##########
src/paimon/core/append/append_compact_coordinator.cpp:
##########
@@ -203,6 +203,15 @@ Status ValidateTable(const std::shared_ptr<TableSchema>& 
table_schema,
                      const std::shared_ptr<arrow::Schema>& arrow_schema,
                      const CoreOptions& core_options) {
     
PAIMON_RETURN_NOT_OK(BlobUtils::ValidateContainerBlobWriteSchema(arrow_schema));
+    // The rewrite reads and writes plain append files, which can neither 
merge data-evolution
+    // blob layers nor write blob files; Java compacts these through its 
data-evolution
+    // compaction instead.
+    for (const auto& field : arrow_schema->fields()) {
+        if (BlobUtils::IsArrayBlobField(field)) {
+            return Status::NotImplemented(
+                "Compacting a table with ARRAY<BLOB> is not supported by the 
C++ writer.");
+        }
+    }

Review Comment:
   How about checking here that data evolution is not enabled? It seems blob 
requires data evolution to be enabled. The current special-casing logic for 
`ContainerBlob` and `ArrayBlob` feels a bit odd.



##########
docs/source/user_guide/write.rst:
##########
@@ -71,6 +71,71 @@ RecordBatch Construction
   - Prefer batch sizes tuned for I/O throughput (e.g., tens to hundreds of MB 
per flush, depending on filesystem and cluster configuration).
   - Maintain stable sort orders within a batch only if required by downstream 
merge or compaction logic; otherwise avoid unnecessary ordering costs.
 
+Writing BLOB Columns
+~~~~~~~~~~~~~~~~~~~~
+
+A ``BLOB`` column is a ``LargeBinary`` field carrying Paimon's BLOB field
+metadata, and an ``ARRAY<BLOB>`` column is a top-level ``List`` field whose
+element field carries it. Build the BLOB field with 
``paimon::Blob::ArrowField``
+and import it into Arrow; the element field of an ``ARRAY<BLOB>`` column must
+keep that metadata:
+
+.. code-block:: cpp
+
+   PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<::ArrowSchema> c_element,
+                          paimon::Blob::ArrowField("element", 
/*nullable=*/true));
+   arrow::Result<std::shared_ptr<arrow::Field>> element = 
arrow::ImportField(c_element.get());
+   if (!element.ok()) {
+       return paimon::Status::Invalid(element.status().ToString());
+   }

Review Comment:
   PAIMON_ASSIGN_OR_RAISE_FROM_ARROW?



##########
src/paimon/format/blob/blob_format_writer.cpp:
##########
@@ -163,64 +190,128 @@ Status BlobFormatWriter::Finish() {
 Status BlobFormatWriter::WriteBlob(std::string_view blob_data) {
     // Open the blob input stream before writing any bytes, so that a failed 
fetch can be
     // converted to a NULL element without leaving partial data in the output 
stream.
+    PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<InputStream> in,
+                           OpenBlobInputStream(blob_data, 
/*element_index=*/std::nullopt));
+    // A null stream means a write-null option already converted the failure.
+    if (in == nullptr) {
+        bin_lengths_.push_back(BlobDefs::kNullBinLength);
+        return Status::OK();
+    }
+    PAIMON_ASSIGN_OR_RAISE(int64_t entry_pos, BeginEntry());
+    Result<int64_t> copied = CopyBlobData(in.get());
+    if (!copied.ok()) {
+        return AddFailureContext(copied.status(), "failed to copy",
+                                 /*element_index=*/std::nullopt);
+    }
+    return FinishEntry(entry_pos);
+}
+
+Status BlobFormatWriter::WriteArrayBlob(const arrow::ListArray& list_array) {
+    const std::shared_ptr<arrow::Array>& values = list_array.values();
+    if (values->type_id() != arrow::Type::type::LARGE_BINARY) {
+        return Status::Invalid("BlobFormatWriter only support large binary 
ARRAY<BLOB> elements.");
+    }
+    const auto& blob_values = checked_cast<const 
arrow::LargeBinaryArray&>(*values);
+    const int64_t element_offset = list_array.value_offset(0);
+    const int32_t element_count = list_array.value_length(0);
+
+    PAIMON_ASSIGN_OR_RAISE(int64_t entry_pos, BeginEntry());
+    PAIMON_RETURN_NOT_OK(
+        WriteWithCrc32(array_magic_number_bytes_->data(), 
array_magic_number_bytes_->size()));
+    PAIMON_RETURN_NOT_OK(WriteWithCrc32(reinterpret_cast<const 
char*>(&BlobDefs::kArrayBlobVersion),
+                                        sizeof(BlobDefs::kArrayBlobVersion)));
+    PAIMON_UNIQUE_PTR<Bytes> element_count_bytes =
+        IntegerToLittleEndian<int32_t>(element_count, pool_);
+    PAIMON_RETURN_NOT_OK(WriteWithCrc32(element_count_bytes->data(), 
element_count_bytes->size()));
+
+    // Each element is opened before any of its bytes are written, so an 
element converted to NULL
+    // leaves no partial data. As in Java, any other failure leaves a partial 
entry and fails the
+    // write.
+    std::vector<int64_t> element_lengths(element_count, 
BlobDefs::kNullBinLength);
+    for (int32_t i = 0; i < element_count; ++i) {
+        const int64_t value_index = element_offset + i;
+        if (blob_values.IsNull(value_index)) {
+            continue;
+        }
+        PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<InputStream> in,
+                               
OpenBlobInputStream(blob_values.GetView(value_index), i));
+        if (in == nullptr) {
+            continue;
+        }
+        Result<int64_t> copied = CopyBlobData(in.get());
+        if (!copied.ok()) {
+            return AddFailureContext(copied.status(), "failed to copy", i);
+        }
+        element_lengths[i] = copied.value();
+    }
+
+    const std::vector<char> index_bytes = 
DeltaVarintCompressor::Compress(element_lengths);
+    PAIMON_ASSIGN_OR_RAISE(int32_t index_length, 
ToIndexLength(index_bytes.size()));
+    PAIMON_RETURN_NOT_OK(WriteWithCrc32(index_bytes.data(), 
index_bytes.size()));
+    PAIMON_UNIQUE_PTR<Bytes> index_length_bytes =
+        IntegerToLittleEndian<int32_t>(index_length, pool_);
+    PAIMON_RETURN_NOT_OK(WriteWithCrc32(index_length_bytes->data(), 
index_length_bytes->size()));
+    return FinishEntry(entry_pos);
+}
+
+Result<std::unique_ptr<InputStream>> BlobFormatWriter::OpenBlobInputStream(
+    std::string_view blob_data, std::optional<int32_t> element_index) {
     // Whether blob_data is a serialized BlobDescriptor is detected by its 
magic header rather
-    // than taken from a blob_as_descriptor option, so each row may hold 
either form.
-    std::unique_ptr<InputStream> in;
+    // than taken from a blob_as_descriptor option, so each value may hold 
either form.
     PAIMON_ASSIGN_OR_RAISE(bool is_descriptor,
                            BlobDescriptor::IsBlobDescriptor(blob_data.data(), 
blob_data.size()));
     if (is_descriptor) {
-        PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<InputStream> descriptor_in,
-                               OpenDescriptorInputStream(blob_data));
-        // A null stream means a write-null option already converted the 
failure.
-        if (descriptor_in == nullptr) {
-            bin_lengths_.push_back(BlobDefs::kNullBinLength);
-            return Status::OK();
+        return OpenDescriptorInputStream(blob_data, element_index);
+    }
+    return std::unique_ptr<InputStream>(
+        std::make_unique<ByteArrayInputStream>(blob_data.data(), 
blob_data.size()));
+}
+
+Result<int64_t> BlobFormatWriter::CopyBlobData(InputStream* in) {
+    PAIMON_ASSIGN_OR_RAISE(int64_t length, in->Length());
+    int64_t copied = 0;
+    while (copied < length) {
+        const int64_t request =
+            std::min(length - copied, 
static_cast<int64_t>(tmp_buffer_->size()));
+        PAIMON_ASSIGN_OR_RAISE(int64_t read_len, in->Read(tmp_buffer_->data(), 
request));
+        if (read_len <= 0) {
+            return Status::IOError(
+                fmt::format("unexpected end of blob data after {} of {} bytes: 
read returned {}",
+                            copied, length, read_len));
+        }
+        if (read_len > request) {
+            return Status::Invalid(fmt::format("read returned {} bytes, more 
than the {} requested",
+                                               read_len, request));
         }
-        in = std::move(descriptor_in);
-    } else {
-        in = std::make_unique<ByteArrayInputStream>(blob_data.data(), 
blob_data.size());
+        PAIMON_RETURN_NOT_OK(WriteWithCrc32(tmp_buffer_->data(), read_len));
+        copied += read_len;
     }
-    PAIMON_ASSIGN_OR_RAISE(int64_t file_length, in->Length());
+    return copied;
+}
 
+Result<int64_t> BlobFormatWriter::BeginEntry() {
     crc32_ = 0;
-    PAIMON_ASSIGN_OR_RAISE(int64_t previous_pos, out_->GetPos());
-
-    // write magic number
+    PAIMON_ASSIGN_OR_RAISE(int64_t entry_pos, out_->GetPos());
     PAIMON_RETURN_NOT_OK(WriteWithCrc32(magic_number_bytes_->data(), 
magic_number_bytes_->size()));
-    int64_t total_read_length = 0;
-    int64_t read_len = std::min(file_length, 
static_cast<int64_t>(tmp_buffer_->size()));
-    while (read_len > 0) {
-        PAIMON_ASSIGN_OR_RAISE(int64_t actual_read_len, 
in->Read(tmp_buffer_->data(), read_len));
-        if (actual_read_len != read_len) {
-            return Status::Invalid(
-                fmt::format("actual read length {}, not match with expect 
length {}",
-                            actual_read_len, read_len));
-        }
-        PAIMON_RETURN_NOT_OK(WriteWithCrc32(tmp_buffer_->data(), 
actual_read_len));
-        total_read_length += actual_read_len;
-        read_len =
-            std::min(file_length - total_read_length, 
static_cast<int64_t>(tmp_buffer_->size()));
-    }
+    return entry_pos;
+}
 
-    // write bin length
+Status BlobFormatWriter::FinishEntry(int64_t entry_pos) {
     PAIMON_ASSIGN_OR_RAISE(int64_t current_pos, out_->GetPos());
     /// magic number(4) + blob content(bin length - 16) + bin length(8) + 
crc32(4)
     /// ↑                                             ↑
-    /// previous_pos                               current_pos
-    int64_t bin_length = current_pos - previous_pos + 8 + 4;
+    /// entry_pos                                  current_pos
+    int64_t bin_length = current_pos - entry_pos + 8 + 4;
     bin_lengths_.push_back(bin_length);
     PAIMON_UNIQUE_PTR<Bytes> bin_length_bytes = 
IntegerToLittleEndian<int64_t>(bin_length, pool_);
     PAIMON_RETURN_NOT_OK(WriteWithCrc32(bin_length_bytes->data(), 
bin_length_bytes->size()));
 
-    // write crc32
     PAIMON_UNIQUE_PTR<Bytes> crc32_bytes = 
IntegerToLittleEndian<int32_t>(crc32_, pool_);
-    PAIMON_RETURN_NOT_OK(WriteBytes(crc32_bytes->data(), crc32_bytes->size()));
-
-    return Status::OK();
+    return WriteBytes(crc32_bytes->data(), crc32_bytes->size());
 }
 
 Result<std::unique_ptr<InputStream>> 
BlobFormatWriter::OpenDescriptorInputStream(
-    std::string_view blob_data) {
+    std::string_view blob_data, std::optional<int32_t> element_index) {
     // A descriptor that cannot be deserialized is a fetch failure: the 
referenced data cannot be

Review Comment:
   Java seems to have `ReusingBlobRefStreamProvider` to reuse the input. Please 
add a TODO here to describe this difference and the potential performance 
optimization we may do later.



##########
src/paimon/format/blob/blob_format_writer.cpp:
##########
@@ -163,64 +190,128 @@ Status BlobFormatWriter::Finish() {
 Status BlobFormatWriter::WriteBlob(std::string_view blob_data) {
     // Open the blob input stream before writing any bytes, so that a failed 
fetch can be
     // converted to a NULL element without leaving partial data in the output 
stream.
+    PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<InputStream> in,
+                           OpenBlobInputStream(blob_data, 
/*element_index=*/std::nullopt));
+    // A null stream means a write-null option already converted the failure.
+    if (in == nullptr) {
+        bin_lengths_.push_back(BlobDefs::kNullBinLength);
+        return Status::OK();
+    }
+    PAIMON_ASSIGN_OR_RAISE(int64_t entry_pos, BeginEntry());
+    Result<int64_t> copied = CopyBlobData(in.get());
+    if (!copied.ok()) {
+        return AddFailureContext(copied.status(), "failed to copy",
+                                 /*element_index=*/std::nullopt);
+    }
+    return FinishEntry(entry_pos);
+}
+
+Status BlobFormatWriter::WriteArrayBlob(const arrow::ListArray& list_array) {
+    const std::shared_ptr<arrow::Array>& values = list_array.values();
+    if (values->type_id() != arrow::Type::type::LARGE_BINARY) {
+        return Status::Invalid("BlobFormatWriter only support large binary 
ARRAY<BLOB> elements.");
+    }
+    const auto& blob_values = checked_cast<const 
arrow::LargeBinaryArray&>(*values);
+    const int64_t element_offset = list_array.value_offset(0);
+    const int32_t element_count = list_array.value_length(0);
+
+    PAIMON_ASSIGN_OR_RAISE(int64_t entry_pos, BeginEntry());
+    PAIMON_RETURN_NOT_OK(
+        WriteWithCrc32(array_magic_number_bytes_->data(), 
array_magic_number_bytes_->size()));
+    PAIMON_RETURN_NOT_OK(WriteWithCrc32(reinterpret_cast<const 
char*>(&BlobDefs::kArrayBlobVersion),
+                                        sizeof(BlobDefs::kArrayBlobVersion)));
+    PAIMON_UNIQUE_PTR<Bytes> element_count_bytes =
+        IntegerToLittleEndian<int32_t>(element_count, pool_);
+    PAIMON_RETURN_NOT_OK(WriteWithCrc32(element_count_bytes->data(), 
element_count_bytes->size()));
+
+    // Each element is opened before any of its bytes are written, so an 
element converted to NULL
+    // leaves no partial data. As in Java, any other failure leaves a partial 
entry and fails the
+    // write.
+    std::vector<int64_t> element_lengths(element_count, 
BlobDefs::kNullBinLength);
+    for (int32_t i = 0; i < element_count; ++i) {
+        const int64_t value_index = element_offset + i;
+        if (blob_values.IsNull(value_index)) {
+            continue;
+        }
+        PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<InputStream> in,
+                               
OpenBlobInputStream(blob_values.GetView(value_index), i));
+        if (in == nullptr) {
+            continue;
+        }
+        Result<int64_t> copied = CopyBlobData(in.get());
+        if (!copied.ok()) {
+            return AddFailureContext(copied.status(), "failed to copy", i);
+        }
+        element_lengths[i] = copied.value();
+    }
+
+    const std::vector<char> index_bytes = 
DeltaVarintCompressor::Compress(element_lengths);
+    PAIMON_ASSIGN_OR_RAISE(int32_t index_length, 
ToIndexLength(index_bytes.size()));
+    PAIMON_RETURN_NOT_OK(WriteWithCrc32(index_bytes.data(), 
index_bytes.size()));
+    PAIMON_UNIQUE_PTR<Bytes> index_length_bytes =
+        IntegerToLittleEndian<int32_t>(index_length, pool_);
+    PAIMON_RETURN_NOT_OK(WriteWithCrc32(index_length_bytes->data(), 
index_length_bytes->size()));
+    return FinishEntry(entry_pos);
+}
+
+Result<std::unique_ptr<InputStream>> BlobFormatWriter::OpenBlobInputStream(
+    std::string_view blob_data, std::optional<int32_t> element_index) {
     // Whether blob_data is a serialized BlobDescriptor is detected by its 
magic header rather
-    // than taken from a blob_as_descriptor option, so each row may hold 
either form.
-    std::unique_ptr<InputStream> in;
+    // than taken from a blob_as_descriptor option, so each value may hold 
either form.
     PAIMON_ASSIGN_OR_RAISE(bool is_descriptor,
                            BlobDescriptor::IsBlobDescriptor(blob_data.data(), 
blob_data.size()));
     if (is_descriptor) {
-        PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<InputStream> descriptor_in,
-                               OpenDescriptorInputStream(blob_data));
-        // A null stream means a write-null option already converted the 
failure.
-        if (descriptor_in == nullptr) {
-            bin_lengths_.push_back(BlobDefs::kNullBinLength);
-            return Status::OK();
+        return OpenDescriptorInputStream(blob_data, element_index);
+    }
+    return std::unique_ptr<InputStream>(
+        std::make_unique<ByteArrayInputStream>(blob_data.data(), 
blob_data.size()));
+}
+
+Result<int64_t> BlobFormatWriter::CopyBlobData(InputStream* in) {
+    PAIMON_ASSIGN_OR_RAISE(int64_t length, in->Length());
+    int64_t copied = 0;
+    while (copied < length) {
+        const int64_t request =
+            std::min(length - copied, 
static_cast<int64_t>(tmp_buffer_->size()));
+        PAIMON_ASSIGN_OR_RAISE(int64_t read_len, in->Read(tmp_buffer_->data(), 
request));
+        if (read_len <= 0) {
+            return Status::IOError(
+                fmt::format("unexpected end of blob data after {} of {} bytes: 
read returned {}",
+                            copied, length, read_len));
+        }
+        if (read_len > request) {
+            return Status::Invalid(fmt::format("read returned {} bytes, more 
than the {} requested",
+                                               read_len, request));
         }

Review Comment:
   We can keep the previous logic here. Paimon Input assumes `read_len` must 
equal `request`, unless it reaches EOF. But request has already been aligned 
with EOF, so this issue shouldn’t happen.



##########
src/paimon/format/blob/blob_format_writer_test.cpp:
##########
@@ -933,6 +1077,46 @@ TEST_F(BlobFormatWriterWriteNullTest, 
TestWriteNullOnExistsCheckFailure) {
     }
 }
 
+TEST_F(BlobFormatWriterWriteNullTest, TestCopyWithShortReads) {
+    const std::string data = "0123456789";
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Blob> blob,
+                         Blob::FromPath(WriteSourceFile("source.bin", data)));
+    ASSERT_OK_AND_ASSIGN(auto array, PrepareDescriptorArray(blob));
+
+    auto short_read_fs = 
std::make_shared<ShortReadFileSystem>(/*max_read_size=*/3);
+    ASSERT_OK_AND_ASSIGN(
+        std::shared_ptr<BlobFormatWriter> writer,
+        BlobFormatWriter::Create(output_stream_, struct_type_,
+                                 /*write_null_on_missing_file=*/false,
+                                 /*write_null_on_fetch_failure=*/false,
+                                 /*write_placeholder=*/false, short_read_fs, 
pool_));
+    ASSERT_OK(AddBatchOnce(writer, array));
+    ASSERT_OK(writer->Finish());
+    ASSERT_EQ(short_read_fs->ReadCallCount(), 4);
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::StructArray> result_struct, 
ReadBackAsData());
+    ASSERT_EQ(result_struct->length(), 1);
+    auto binary_array = 
checked_pointer_cast<arrow::LargeBinaryArray>(result_struct->field(0));
+    ASSERT_EQ(binary_array->GetString(0), data);
+
+    const std::vector<std::pair<Status, std::string>> cases = {
+        {Status::IOError("mock read error"), "mock read error"},
+        {Status::OK(), "unexpected end of blob data after 7 of 10 bytes: read 
returned 0"}};
+    for (const auto& [end_status, expected_error] : cases) {
+        SCOPED_TRACE(expected_error);
+        ASSERT_OK_AND_ASSIGN(std::shared_ptr<BlobFormatWriter> failing_writer,
+                             BlobFormatWriter::Create(
+                                 output_stream_, struct_type_, 
/*write_null_on_missing_file=*/true,
+                                 /*write_null_on_fetch_failure=*/true, 
/*write_placeholder=*/false,
+                                 std::make_shared<ShortReadFileSystem>(
+                                     /*max_read_size=*/3, 
/*readable_length=*/7, end_status),
+                                 pool_));
+        Status status = AddBatchOnce(failing_writer, array);
+        ASSERT_NOK_WITH_MSG(status, "failed to copy BLOB field blob_col in row 
0 of blob file");
+        ASSERT_NOK_WITH_MSG(status, expected_error);
+    }

Review Comment:
   If we assume that input always reads back the exact number of bytes needed, 
then this test may not be necessary. In the future, we plan to remove the 
return value of `input->Read` to avoid ambiguity.



##########
src/paimon/common/data/blob_defs.h:
##########
@@ -86,13 +87,15 @@ class BlobDefs {
     /// Only the data-evolution blob fallback read path sets this.
     static constexpr char kEmitPlaceholderSentinelKey[] = 
"blob.internal.emit-placeholder-sentinel";
     /// Internal (non user-facing) format option, "false" by default: when 
"true", the blob
-    /// format writer persists a value exactly equal to kPlaceholderSentinel 
as a bin_length -2
-    /// entry. Only set for data-evolution partial updates, i.e. blob-only 
column writes of a
-    /// table with data evolution enabled; all other writes store bytes 
verbatim.
+    /// format writer persists a value exactly equal to kPlaceholderSentinel 
(or an ARRAY<BLOB>
+    /// whose only element is kPlaceholderSentinel) as a bin_length -2 entry. 
Only set for
+    /// data-evolution partial updates, i.e. blob-only column writes of a 
table with data evolution
+    /// enabled; no other write interprets a value as a placeholder.
     static constexpr char kWritePlaceholderKey[] = 
"blob.internal.write-placeholder";

Review Comment:
   I found a historical issue: if the write schema is `[id, blob]`, Java will 
unconditionally recognize the placeholder, meaning it supports fallback, but 
C++ requires the write schema to contain only blob in order to recognize the 
placeholder. This may cause compatibility issues. Please make C++ recognize the 
placeholder without depending on the write schema, similar to Java, and fix 
this in this PR as well. Also please add integration tests.



-- 
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]

Reply via email to