SteNicholas commented on code in PR #392:
URL: https://github.com/apache/paimon-cpp/pull/392#discussion_r4131835000


##########
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:
   Agreed, reverted to the previous logic: `CopyBlobData` requests at most the 
remaining length and fails when a read returns a different number of bytes.



##########
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:
   Instead of a TODO, this PR now implements the reuse following 
`ReusingBlobRefStreamProvider`: the stream of a file referenced by a descriptor 
with a known length stays open for the next such descriptor of the same file. 
It is released when a descriptor references another file or a range it cannot 
serve, after a failed copy and on `Finish()`, and a failure to close it fails 
the write regardless of the write-null options. As in Java, a descriptor with a 
dynamic length (`-1`) opens a stream of its own, which is closed after the 
copy. This is covered by `TestReuseSourceStreamAcrossDescriptors`, 
`TestDynamicLengthElementReadsGrownFile` and `TestSourceCloseFailureFailsWrite`.



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