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]