This is an automated email from the ASF dual-hosted git repository.
SteNicholas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-cpp.git
The following commit(s) were added to refs/heads/main by this push:
new e00d7996 feat(blob): support reading ARRAY<BLOB> fields (#388)
e00d7996 is described below
commit e00d799602067a5149e0b6b39525ca9c66be9d3f
Author: lxy <[email protected]>
AuthorDate: Thu Sep 24 20:23:31 2026 +0800
feat(blob): support reading ARRAY<BLOB> fields (#388)
---
src/paimon/common/data/blob_defs.h | 5 +-
src/paimon/common/data/blob_utils.cpp | 54 ++++++-
src/paimon/common/data/blob_utils.h | 15 +-
src/paimon/common/data/blob_utils_test.cpp | 51 +++++++
.../common/reader/blob_fallback_batch_reader.cpp | 31 ++--
.../common/reader/blob_fallback_batch_reader.h | 3 +-
.../core/append/append_compact_coordinator.cpp | 2 +-
src/paimon/core/operation/file_store_write.cpp | 2 +-
src/paimon/core/schema/arrow_schema_validator.cpp | 19 ++-
.../core/schema/arrow_schema_validator_test.cpp | 24 +--
src/paimon/core/schema/schema_validation.cpp | 29 ++--
src/paimon/core/schema/schema_validation_test.cpp | 129 +++++++++++++++-
src/paimon/core/schema/table_schema.cpp | 2 +-
src/paimon/core/schema/table_schema_test.cpp | 9 ++
src/paimon/format/blob/blob_file_batch_reader.cpp | 162 ++++++++++++++++++++-
src/paimon/format/blob/blob_file_batch_reader.h | 14 +-
.../format/blob/blob_file_batch_reader_test.cpp | 159 ++++++++++++++++++--
src/paimon/format/blob/blob_stats_extractor.cpp | 4 +-
.../format/blob/blob_stats_extractor_test.cpp | 2 +-
test/inte/blob_table_inte_test.cpp | 99 ++++++++++++-
test/inte/paimon_read_compat_inte_test.cpp | 34 ++++-
.../rust_array_blob_types/README.md | 3 +-
.../array_blob_java.db/array_blob_java/README.md | 29 ++++
...ata-0e375e23-85cb-4455-a0bd-511a331c4dc8-0.blob | Bin 0 -> 9 bytes
...ata-53163ebc-5467-47c9-ad31-228f7915c4db-0.blob | Bin 0 -> 52 bytes
...-91236ae6-4891-42dd-b2f9-a35751786cf2-0.parquet | Bin 0 -> 309 bytes
...ata-91236ae6-4891-42dd-b2f9-a35751786cf2-1.blob | Bin 0 -> 135 bytes
...manifest-1f8acd08-a84a-40d4-8069-86e92070df21-0 | Bin 0 -> 2188 bytes
...manifest-814974cc-e019-45eb-b1cc-5aa4b6cf9987-0 | Bin 0 -> 2189 bytes
...manifest-fd42fd8f-82ac-4a8f-acd0-771e4db82f25-0 | Bin 0 -> 2237 bytes
...est-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-0 | Bin 0 -> 1158 bytes
...est-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-1 | Bin 0 -> 1261 bytes
...est-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-0 | Bin 0 -> 1261 bytes
...est-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-1 | Bin 0 -> 1261 bytes
...est-list-f082e141-fa95-4685-9b7c-5555137b5310-0 | Bin 0 -> 1295 bytes
...est-list-f082e141-fa95-4685-9b7c-5555137b5310-1 | Bin 0 -> 1258 bytes
.../array_blob_java/schema/schema-0 | 26 ++++
.../array_blob_java/snapshot/EARLIEST | 1 +
.../array_blob_java/snapshot/LATEST | 1 +
.../array_blob_java/snapshot/snapshot-1 | 18 +++
.../array_blob_java/snapshot/snapshot-2 | 18 +++
.../array_blob_java/snapshot/snapshot-3 | 18 +++
42 files changed, 875 insertions(+), 88 deletions(-)
diff --git a/src/paimon/common/data/blob_defs.h
b/src/paimon/common/data/blob_defs.h
index deacd1d0..07828e82 100644
--- a/src/paimon/common/data/blob_defs.h
+++ b/src/paimon/common/data/blob_defs.h
@@ -65,8 +65,9 @@ class BlobDefs {
/// such write arrays. Outside that mode the writer never interprets
values, so arbitrary
/// user bytes can never be turned into a placeholder entry.
/// - Read channel: a placeholder-aware reader (see
kEmitPlaceholderSentinelKey) emits these
- /// bytes for -2 entries so the fallback merge can identify placeholders
after the batch
- /// has passed through schema-mapping readers.
+ /// bytes directly for a scalar BLOB, or as the only element of an
ARRAY<BLOB>, so the
+ /// fallback merge can identify -2 entries after the batch has passed
through schema-mapping
+ /// readers.
///
/// Both channels identify a placeholder by exact byte equality with this
internal reserved
/// value (IsPlaceholderSentinel), and the fallback merge byte-compares
every layer of a
diff --git a/src/paimon/common/data/blob_utils.cpp
b/src/paimon/common/data/blob_utils.cpp
index 6c52a9c0..11da794c 100644
--- a/src/paimon/common/data/blob_utils.cpp
+++ b/src/paimon/common/data/blob_utils.cpp
@@ -21,6 +21,7 @@
#include <cstddef>
#include <set>
+#include <string_view>
#include <vector>
#include "arrow/api.h"
@@ -98,6 +99,9 @@ Result<BlobUtils::SeparatedStructArrays>
BlobUtils::SeparateBlobArray(
}
bool BlobUtils::IsBlobField(const std::shared_ptr<arrow::Field>& field) {
+ if (field == nullptr) {
+ return false;
+ }
const auto& type = field->type();
if (type->id() != arrow::Type::LARGE_BINARY) {
return false;
@@ -108,6 +112,16 @@ bool BlobUtils::IsBlobField(const
std::shared_ptr<arrow::Field>& field) {
return IsBlobMetadata(field->metadata());
}
+bool BlobUtils::IsArrayBlobField(const std::shared_ptr<arrow::Field>& field) {
+ if (field == nullptr || field->type()->id() != arrow::Type::LIST) {
+ return false;
+ }
+ const auto& list_type = checked_cast<const
arrow::ListType&>(*field->type());
+ // Arrow's C schema importer passes MakeChildField(0) directly to
ListType, retaining the
+ // element field's metadata.
+ return IsBlobField(list_type.value_field());
+}
+
bool BlobUtils::IsMapBlobField(const std::shared_ptr<arrow::Field>& field) {
if (field == nullptr || field->type()->id() != arrow::Type::MAP) {
return false;
@@ -118,12 +132,50 @@ bool BlobUtils::IsMapBlobField(const
std::shared_ptr<arrow::Field>& field) {
return map_type.item_type()->id() == arrow::Type::LARGE_BINARY;
}
-Status BlobUtils::ValidateMapBlobWriteSchema(const
std::shared_ptr<arrow::Schema>& schema) {
+bool BlobUtils::IsAnyBlobField(const std::shared_ptr<arrow::Field>& field) {
+ return IsBlobField(field) || IsArrayBlobField(field) ||
IsMapBlobField(field);
+}
+
+bool BlobUtils::IsArrayBlobPlaceholder(const arrow::ListArray& array, int64_t
row) {
+ if (array.IsNull(row) || array.value_length(row) != 1) {
+ return false;
+ }
+ const std::shared_ptr<arrow::Array>& values = array.values();
+ if (values->type_id() != arrow::Type::LARGE_BINARY) {
+ return false;
+ }
+ const int64_t value_index = array.value_offset(row);
+ if (values->IsNull(value_index)) {
+ return false;
+ }
+ const auto& binary_values = checked_cast<const
arrow::LargeBinaryArray&>(*values);
+ const std::string_view value = binary_values.GetView(value_index);
+ return BlobDefs::IsPlaceholderSentinel(value.data(), value.size());
+}
+
+bool BlobUtils::IsMapBlobPlaceholder(const arrow::MapArray& array, int64_t
row) {
+ if (array.IsNull(row) || array.value_length(row) != 2) {
+ return false;
+ }
+ const int64_t entry_index = array.value_offset(row);
+ const std::shared_ptr<arrow::Array>& items = array.items();
+ if (!items->IsNull(entry_index) || !items->IsNull(entry_index + 1)) {
+ return false;
+ }
+ const std::shared_ptr<arrow::Array>& keys = array.keys();
+ return keys->RangeEquals(entry_index, entry_index + 1, entry_index + 1,
*keys);
+}
+
+Status BlobUtils::ValidateContainerBlobWriteSchema(const
std::shared_ptr<arrow::Schema>& schema) {
for (const auto& field : schema->fields()) {
if (IsMapBlobField(field)) {
return Status::NotImplemented(
"Writing a table with MAP<..., BLOB> is not supported by the
C++ writer.");
}
+ if (IsArrayBlobField(field)) {
+ return Status::NotImplemented(
+ "Writing a table with ARRAY<BLOB> is not supported by the C++
writer.");
+ }
}
return Status::OK();
}
diff --git a/src/paimon/common/data/blob_utils.h
b/src/paimon/common/data/blob_utils.h
index 86e0ef37..18f55862 100644
--- a/src/paimon/common/data/blob_utils.h
+++ b/src/paimon/common/data/blob_utils.h
@@ -19,6 +19,7 @@
#pragma once
+#include <cstdint>
#include <memory>
#include <set>
#include <string>
@@ -31,6 +32,8 @@
namespace arrow {
class Field;
class KeyValueMetadata;
+class ListArray;
+class MapArray;
class Schema;
class StructArray;
} // namespace arrow
@@ -73,10 +76,18 @@ class PAIMON_EXPORT BlobUtils {
const std::set<std::string>& inline_fields);
static bool IsBlobField(const std::shared_ptr<arrow::Field>& field);
+ /// Returns whether the field is a top-level ARRAY whose elements are
BLOBs.
+ static bool IsArrayBlobField(const std::shared_ptr<arrow::Field>& field);
/// Returns whether the field is a top-level MAP whose values are BLOBs.
static bool IsMapBlobField(const std::shared_ptr<arrow::Field>& field);
- /// Rejects schemas that the C++ writer cannot safely mutate.
- static Status ValidateMapBlobWriteSchema(const
std::shared_ptr<arrow::Schema>& schema);
+ /// Returns whether the field is stored in a standalone blob file.
+ static bool IsAnyBlobField(const std::shared_ptr<arrow::Field>& field);
+ /// Returns whether an ARRAY<BLOB> row is the internal fallback sentinel.
+ static bool IsArrayBlobPlaceholder(const arrow::ListArray& array, int64_t
row);
+ /// Returns whether a MAP<..., BLOB> row is the internal fallback sentinel.
+ static bool IsMapBlobPlaceholder(const arrow::MapArray& array, int64_t
row);
+ /// Rejects container BLOB types, which are currently supported by the C++
reader only.
+ static Status ValidateContainerBlobWriteSchema(const
std::shared_ptr<arrow::Schema>& schema);
static bool IsBlobMetadata(const std::shared_ptr<const
arrow::KeyValueMetadata>& metadata);
static bool IsBlobFile(const std::string& file_name);
diff --git a/src/paimon/common/data/blob_utils_test.cpp
b/src/paimon/common/data/blob_utils_test.cpp
index e523e5f7..67c357b4 100644
--- a/src/paimon/common/data/blob_utils_test.cpp
+++ b/src/paimon/common/data/blob_utils_test.cpp
@@ -21,6 +21,7 @@
#include "arrow/api.h"
#include "arrow/c/bridge.h"
+#include "arrow/ipc/json_simple.h"
#include "gtest/gtest.h"
#include "paimon/catalog/identifier.h"
#include "paimon/common/data/blob_defs.h"
@@ -78,6 +79,56 @@ TEST_F(BlobUtilsTest, IsBlobField) {
ASSERT_FALSE(BlobUtils::IsBlobField(binary_field_wrong_meta));
}
+TEST_F(BlobUtilsTest, IsArrayBlobField) {
+ auto array_blob_field =
+ arrow::field("array_blob", arrow::list(BlobUtils::ToArrowField("item",
true)));
+ ASSERT_TRUE(BlobUtils::IsArrayBlobField(array_blob_field));
+ ASSERT_TRUE(BlobUtils::IsAnyBlobField(array_blob_field));
+
+ ASSERT_FALSE(BlobUtils::IsArrayBlobField(nullptr));
+ ASSERT_FALSE(BlobUtils::IsArrayBlobField(
+ arrow::field("plain_binary", arrow::list(arrow::large_binary()))));
+ ASSERT_FALSE(BlobUtils::IsArrayBlobField(
+ arrow::field("nested",
arrow::list(arrow::list(BlobUtils::ToArrowField("item", true))))));
+}
+
+TEST_F(BlobUtilsTest, IsArrayBlobPlaceholder) {
+ const std::string json = R"([
+ ["_PAIMON_BLOB_PLACEHOLDER"],
+ null,
+ [],
+ [null],
+ ["_PAIMON_BLOB_PLACEHOLDER", "_PAIMON_BLOB_PLACEHOLDER"],
+ ["user-value"]
+ ])";
+ std::shared_ptr<arrow::Array> array =
+
arrow::ipc::internal::json::ArrayFromJSON(arrow::list(arrow::large_binary()),
json)
+ .ValueOrDie();
+ const auto& list_array = checked_cast<const arrow::ListArray&>(*array);
+ ASSERT_TRUE(BlobUtils::IsArrayBlobPlaceholder(list_array, 0));
+ for (int64_t i = 1; i < list_array.length(); ++i) {
+ ASSERT_FALSE(BlobUtils::IsArrayBlobPlaceholder(list_array, i));
+ }
+}
+
+TEST_F(BlobUtilsTest, IsMapBlobPlaceholder) {
+ auto map_type = arrow::map(arrow::utf8(), arrow::large_binary());
+ std::shared_ptr<arrow::Array> array =
arrow::ipc::internal::json::ArrayFromJSON(map_type, R"([
+ [["same", null], ["same", null]],
+ null,
+ [],
+ [["same", null]],
+ [["left", null], ["right", null]],
+ [["same", "value"], ["same", null]]
+ ])")
+ .ValueOrDie();
+ const auto& map_array = checked_cast<const arrow::MapArray&>(*array);
+ ASSERT_TRUE(BlobUtils::IsMapBlobPlaceholder(map_array, 0));
+ for (int64_t i = 1; i < map_array.length(); ++i) {
+ ASSERT_FALSE(BlobUtils::IsMapBlobPlaceholder(map_array, i));
+ }
+}
+
TEST_F(BlobUtilsTest, SeparateBlobSchema) {
auto int_field = arrow::field("f1_int", arrow::int32());
auto string_field = arrow::field("f2_string", arrow::utf8());
diff --git a/src/paimon/common/reader/blob_fallback_batch_reader.cpp
b/src/paimon/common/reader/blob_fallback_batch_reader.cpp
index 1a23ec2d..830b0df4 100644
--- a/src/paimon/common/reader/blob_fallback_batch_reader.cpp
+++ b/src/paimon/common/reader/blob_fallback_batch_reader.cpp
@@ -56,8 +56,7 @@ Result<std::unique_ptr<BlobFallbackBatchReader>>
BlobFallbackBatchReader::Create
}
int32_t blob_field_idx = -1;
for (int32_t i = 0; i < read_schema->num_fields(); i++) {
- if (BlobUtils::IsBlobField(read_schema->field(i)) ||
- BlobUtils::IsMapBlobField(read_schema->field(i))) {
+ if (BlobUtils::IsAnyBlobField(read_schema->field(i))) {
if (blob_field_idx != -1) {
return Status::Invalid(
"Blob fallback read schema should contain exactly one blob
field.");
@@ -113,7 +112,8 @@ Result<int64_t> BlobFallbackBatchReader::FillWindow(size_t
group_idx, int64_t wa
const std::shared_ptr<arrow::StructArray>& front =
cursor.pending.front();
int64_t available = front->length() - cursor.pending_pos;
int64_t take = std::min(available, want - collected);
- chunks->push_back(Chunk{front, cursor.pending_pos, take, {}});
+ chunks->push_back(
+ Chunk{front, cursor.pending_pos, take, {}}); //
NOLINT(modernize-use-emplace)
cursor.pending_pos += take;
collected += take;
if (cursor.pending_pos == front->length()) {
@@ -199,24 +199,25 @@ Result<std::vector<bool>>
BlobFallbackBatchReader::ComputePlaceholderFlags(
}
if (blob_col->type_id() == arrow::Type::MAP) {
auto map_col = checked_pointer_cast<arrow::MapArray>(blob_col);
- const std::shared_ptr<arrow::Array>& keys = map_col->keys();
- const std::shared_ptr<arrow::Array>& items = map_col->items();
for (int64_t k = 0; k < chunk.length; k++) {
- int64_t idx = chunk.offset + k;
- if (!map_col->IsNull(idx) && map_col->value_length(idx) ==
2) {
- int64_t entry_idx = map_col->value_offset(idx);
- flags[pos + k] =
- items->IsNull(entry_idx) &&
items->IsNull(entry_idx + 1) &&
- keys->RangeEquals(entry_idx, entry_idx + 1,
entry_idx + 1, *keys);
- }
+ flags[pos + k] = BlobUtils::IsMapBlobPlaceholder(*map_col,
chunk.offset + k);
+ }
+ pos += chunk.length;
+ continue;
+ }
+ if (blob_col->type_id() == arrow::Type::LIST) {
+ auto list_col =
checked_pointer_cast<arrow::ListArray>(blob_col);
+ for (int64_t k = 0; k < chunk.length; k++) {
+ flags[pos + k] =
BlobUtils::IsArrayBlobPlaceholder(*list_col, chunk.offset + k);
}
pos += chunk.length;
continue;
}
if (blob_col->type_id() != arrow::Type::LARGE_BINARY) {
- return Status::Invalid(
- fmt::format("Blob fallback expects a BLOB or MAP<...,
BLOB> column, but got {}",
- blob_col->type()->ToString()));
+ return Status::Invalid(fmt::format(
+ "Blob fallback expects a BLOB, ARRAY<BLOB> or MAP<...,
BLOB> column, but "
+ "got {}",
+ blob_col->type()->ToString()));
}
auto binary_col =
checked_pointer_cast<arrow::LargeBinaryArray>(blob_col);
for (int64_t k = 0; k < chunk.length; k++) {
diff --git a/src/paimon/common/reader/blob_fallback_batch_reader.h
b/src/paimon/common/reader/blob_fallback_batch_reader.h
index 3664b0c1..e15667b4 100644
--- a/src/paimon/common/reader/blob_fallback_batch_reader.h
+++ b/src/paimon/common/reader/blob_fallback_batch_reader.h
@@ -55,7 +55,8 @@ namespace paimon {
/// through the row ids the caller leaves in a gap segment's
`gap_selected_ranges`.
/// 3. Each output row takes the first group, in max-sequence order, whose row
is not a
/// placeholder. The blob format reader emits placeholders as
BlobDefs::kPlaceholderSentinel
-/// bytes for scalar BLOBs, or as a two-entry map with duplicate keys for
MAP<..., BLOB>.
+/// bytes for scalar BLOBs, as a one-element sentinel array for
ARRAY<BLOB>, or as a two-entry
+/// map with duplicate keys for MAP<..., BLOB>.
/// 4. A row that is a placeholder in every group degrades to a null blob: it
keeps its
/// _ROW_ID, reports -1 as its _SEQUENCE_NUMBER, and returns null for every
other field.
class BlobFallbackBatchReader : public BatchReader {
diff --git a/src/paimon/core/append/append_compact_coordinator.cpp
b/src/paimon/core/append/append_compact_coordinator.cpp
index eb031c05..d25bd7b6 100644
--- a/src/paimon/core/append/append_compact_coordinator.cpp
+++ b/src/paimon/core/append/append_compact_coordinator.cpp
@@ -202,7 +202,7 @@ Result<std::pair<std::shared_ptr<TableSchema>,
CoreOptions>> LoadSchemaAndOption
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::ValidateMapBlobWriteSchema(arrow_schema));
+
PAIMON_RETURN_NOT_OK(BlobUtils::ValidateContainerBlobWriteSchema(arrow_schema));
if (!table_schema->PrimaryKeys().empty() || core_options.GetBucket() !=
-1) {
return Status::Invalid(
"AppendCompactCoordinator only supports append-only tables "
diff --git a/src/paimon/core/operation/file_store_write.cpp
b/src/paimon/core/operation/file_store_write.cpp
index 8e25f69f..6056ba1e 100644
--- a/src/paimon/core/operation/file_store_write.cpp
+++ b/src/paimon/core/operation/file_store_write.cpp
@@ -239,7 +239,7 @@ Result<std::unique_ptr<FileStoreWrite>>
FileStoreWrite::Create(std::unique_ptr<W
}
const std::shared_ptr<TableSchema>& schema = latest_schema;
auto arrow_schema =
DataField::ConvertDataFieldsToArrowSchema(schema->Fields());
- PAIMON_RETURN_NOT_OK(BlobUtils::ValidateMapBlobWriteSchema(arrow_schema));
+
PAIMON_RETURN_NOT_OK(BlobUtils::ValidateContainerBlobWriteSchema(arrow_schema));
auto opts = schema->Options();
for (const auto& [key, value] : ctx->GetOptions()) {
opts[key] = value;
diff --git a/src/paimon/core/schema/arrow_schema_validator.cpp
b/src/paimon/core/schema/arrow_schema_validator.cpp
index daf364c6..c0447c31 100644
--- a/src/paimon/core/schema/arrow_schema_validator.cpp
+++ b/src/paimon/core/schema/arrow_schema_validator.cpp
@@ -128,8 +128,11 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
return Status::OK();
case arrow::Type::type::LIST: {
const auto& value_field =
checked_cast<arrow::BaseListType*>(type.get())->value_field();
- PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
- value_field->type(), value_field->metadata(),
/*allow_blob=*/false, field_id_set));
+ bool allow_direct_blob_element =
+ allow_blob && value_field->type()->id() ==
arrow::Type::LARGE_BINARY;
+ PAIMON_RETURN_NOT_OK(
+ ValidateDataTypeWithFieldId(value_field->type(),
value_field->metadata(),
+ allow_direct_blob_element,
field_id_set));
break;
}
case arrow::Type::type::FIXED_SIZE_LIST: {
@@ -176,8 +179,8 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
if (BlobUtils::IsBlobMetadata(key_value_metadata)) {
if (!allow_blob) {
return Status::Invalid(
- "BLOB field must be a top-level field or the direct
value of a "
- "top-level MAP field.");
+ "BLOB field must be a top-level field or the direct
element/value of a "
+ "top-level ARRAY/MAP field.");
}
break;
}
@@ -219,7 +222,9 @@ Status ArrowSchemaValidator::ValidateField(const
std::shared_ptr<arrow::Field>&
case arrow::Type::type::LIST: {
const auto& value_field =
checked_cast<const
arrow::BaseListType&>(*field->type()).value_field();
- PAIMON_RETURN_NOT_OK(ValidateField(value_field,
/*allow_blob=*/false));
+ bool allow_direct_blob_element =
+ allow_blob && value_field->type()->id() ==
arrow::Type::LARGE_BINARY;
+ PAIMON_RETURN_NOT_OK(ValidateField(value_field,
allow_direct_blob_element));
break;
}
case arrow::Type::type::FIXED_SIZE_LIST: {
@@ -266,8 +271,8 @@ Status ArrowSchemaValidator::ValidateField(const
std::shared_ptr<arrow::Field>&
if (BlobUtils::IsBlobField(field)) {
if (!allow_blob) {
return Status::Invalid(
- "BLOB field must be a top-level field or the direct
value of a "
- "top-level MAP field.");
+ "BLOB field must be a top-level field or the direct
element/value of a "
+ "top-level ARRAY/MAP field.");
}
break;
}
diff --git a/src/paimon/core/schema/arrow_schema_validator_test.cpp
b/src/paimon/core/schema/arrow_schema_validator_test.cpp
index 6a9a55f3..46f30d3c 100644
--- a/src/paimon/core/schema/arrow_schema_validator_test.cpp
+++ b/src/paimon/core/schema/arrow_schema_validator_test.cpp
@@ -220,16 +220,18 @@ TEST(ArrowSchemaValidatorTest, TestBlobFieldPlacement) {
arrow::field("nested",
arrow::struct_({BlobUtils::ToArrowField("blob", true)}));
auto arrow_schema =
arrow::schema(arrow::FieldVector({nested_blob_field}));
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
- "BLOB field must be a top-level field or the
direct value of a "
- "top-level MAP field.");
+ "BLOB field must be a top-level field or the
direct element/value of a "
+ "top-level ARRAY/MAP field.");
}
{
auto array_blob_field =
arrow::field("array_blob",
arrow::list(BlobUtils::ToArrowField("item", true)));
auto arrow_schema =
arrow::schema(arrow::FieldVector({array_blob_field}));
-
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
- "BLOB field must be a top-level field or the
direct value of a "
- "top-level MAP field.");
+ ASSERT_OK(ArrowSchemaValidator::ValidateSchema(*arrow_schema));
+
+ std::vector<DataField> fields = {DataField(0, array_blob_field)};
+ arrow_schema = DataField::ConvertDataFieldsToArrowSchema(fields);
+
ASSERT_OK(ArrowSchemaValidator::ValidateSchemaWithFieldId(*arrow_schema));
}
{
auto map_blob_field = arrow::field(
@@ -247,8 +249,8 @@ TEST(ArrowSchemaValidatorTest, TestBlobFieldPlacement) {
arrow::map(arrow::utf8(),
arrow::struct_({BlobUtils::ToArrowField("blob", true)})));
auto arrow_schema =
arrow::schema(arrow::FieldVector({map_blob_field}));
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
- "BLOB field must be a top-level field or the
direct value of a "
- "top-level MAP field.");
+ "BLOB field must be a top-level field or the
direct element/value of a "
+ "top-level ARRAY/MAP field.");
}
{
auto nested_map_blob_field = arrow::field(
@@ -257,8 +259,8 @@ TEST(ArrowSchemaValidatorTest, TestBlobFieldPlacement) {
arrow::map(arrow::utf8(),
BlobUtils::ToArrowField("value", true))));
auto arrow_schema =
arrow::schema(arrow::FieldVector({nested_map_blob_field}));
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
- "BLOB field must be a top-level field or the
direct value of a "
- "top-level MAP field.");
+ "BLOB field must be a top-level field or the
direct element/value of a "
+ "top-level ARRAY/MAP field.");
}
{
std::vector<DataField> nested_fields = {
@@ -268,8 +270,8 @@ TEST(ArrowSchemaValidatorTest, TestBlobFieldPlacement) {
arrow::field("nested",
DataField::ConvertDataFieldsToArrowStructType(nested_fields)))};
auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(fields);
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchemaWithFieldId(*arrow_schema),
- "BLOB field must be a top-level field or the
direct value of a "
- "top-level MAP field.");
+ "BLOB field must be a top-level field or the
direct element/value of a "
+ "top-level ARRAY/MAP field.");
}
}
diff --git a/src/paimon/core/schema/schema_validation.cpp
b/src/paimon/core/schema/schema_validation.cpp
index f6d96294..9a5b2983 100644
--- a/src/paimon/core/schema/schema_validation.cpp
+++ b/src/paimon/core/schema/schema_validation.cpp
@@ -775,7 +775,7 @@ Status SchemaValidation::ValidateRowTracking(const
TableSchema& table_schema,
std::vector<std::string> blob_names;
for (const auto& field : table_schema.Fields()) {
- if (BlobUtils::IsBlobField(field.ArrowField())) {
+ if (BlobUtils::IsAnyBlobField(field.ArrowField())) {
blob_names.push_back(field.Name());
}
}
@@ -809,26 +809,31 @@ Status SchemaValidation::ValidateBlobFields(const
TableSchema& schema, const Cor
}
auto validate_blob_fields = [&](const std::vector<std::string>&
field_names,
- const std::string& option_key) -> Status {
+ const std::string& option_key,
+ bool allow_container_blob) -> Status {
if (field_names.empty()) {
return Status::OK();
}
PAIMON_RETURN_NOT_OK(ValidateNoDuplicateField(field_names,
option_key));
PAIMON_ASSIGN_OR_RAISE(std::vector<DataField> blob_fields,
schema.GetFields(field_names));
for (const auto& blob_field : blob_fields) {
- if (!BlobUtils::IsBlobField(blob_field.ArrowField())) {
+ bool is_blob = allow_container_blob ?
BlobUtils::IsAnyBlobField(blob_field.ArrowField())
+ :
BlobUtils::IsBlobField(blob_field.ArrowField());
+ if (!is_blob) {
+ const std::string expected_type =
+ allow_container_blob ? "BLOB, ARRAY<BLOB> or MAP<...,
BLOB>" : "BLOB";
return Status::Invalid(
- fmt::format("Field '{}' in '{}' must be a BLOB field in
table schema.",
- blob_field.Name(), option_key));
+ fmt::format("Field '{}' in '{}' must be a {} field in
table schema.",
+ blob_field.Name(), option_key, expected_type));
}
}
return Status::OK();
};
- PAIMON_RETURN_NOT_OK(validate_blob_fields(configured_blob_names,
Options::BLOB_FIELD));
+ PAIMON_RETURN_NOT_OK(validate_blob_fields(configured_blob_names,
Options::BLOB_FIELD, true));
PAIMON_RETURN_NOT_OK(
- validate_blob_fields(blob_descriptor_names,
Options::BLOB_DESCRIPTOR_FIELD));
- PAIMON_RETURN_NOT_OK(validate_blob_fields(blob_view_names,
Options::BLOB_VIEW_FIELD));
+ validate_blob_fields(blob_descriptor_names,
Options::BLOB_DESCRIPTOR_FIELD, false));
+ PAIMON_RETURN_NOT_OK(validate_blob_fields(blob_view_names,
Options::BLOB_VIEW_FIELD, false));
std::set<std::string>
blob_descriptor_name_set(blob_descriptor_names.begin(),
blob_descriptor_names.end());
@@ -900,10 +905,10 @@ Status SchemaValidation::ValidateMosaicDataFields(const
TableSchema& schema,
const std::set<std::string>
inline_blob_field_set(inline_blob_fields.begin(),
inline_blob_fields.end());
// Match Java SchemaValidation by validating only fields stored in the
normal data file.
- // Top-level BLOB fields stored in separate files are skipped; descriptor
and view fields are
- // inline, so Mosaic must reject them here.
+ // Top-level blob-file fields stored in separate files are skipped;
descriptor and view fields
+ // are inline, so Mosaic must reject them here.
for (const DataField& field : schema.Fields()) {
- if (BlobUtils::IsBlobField(field.ArrowField()) &&
+ if (BlobUtils::IsAnyBlobField(field.ArrowField()) &&
inline_blob_field_set.count(field.Name()) == 0) {
continue;
}
@@ -969,7 +974,7 @@ Status SchemaValidation::ValidateLanceDataFields(const
TableSchema& schema,
return Status::OK();
}
for (const DataField& field : schema.Fields()) {
- if (BlobUtils::IsBlobField(field.ArrowField()) &&
+ if (BlobUtils::IsAnyBlobField(field.ArrowField()) &&
inline_blob_field_set.count(field.Name()) == 0) {
continue;
}
diff --git a/src/paimon/core/schema/schema_validation_test.cpp
b/src/paimon/core/schema/schema_validation_test.cpp
index d3ee43f7..bccbad8e 100644
--- a/src/paimon/core/schema/schema_validation_test.cpp
+++ b/src/paimon/core/schema/schema_validation_test.cpp
@@ -20,6 +20,8 @@
#include "paimon/core/schema/schema_validation.h"
#include <map>
+#include <string>
+#include <vector>
#include "arrow/api.h"
#include "gtest/gtest.h"
@@ -44,6 +46,40 @@ Result<std::unique_ptr<TableSchema>>
MakePrimaryKeyBTreeSchema(
options);
}
+const std::vector<std::string>& ContainerBlobTypeJsons() {
+ static const std::vector<std::string> container_blob_types = {
+ R"({"type":"ARRAY","element":"BLOB"})",
+ R"({"type":"MAP","key":"STRING","value":"BLOB"})",
+ };
+ return container_blob_types;
+}
+
+Result<std::unique_ptr<TableSchema>> MakeContainerBlobSchema(const
std::string& blob_type_json,
+ const
std::string& partition_keys_json,
+ const
std::string& options_json) {
+ const std::string schema_json =
+ R"({
+ "version": 3,
+ "id": 0,
+ "fields": [
+ {"id": 0, "name": "id", "type": "INT"},
+ {"id": 1, "name": "blob", "type": )" +
+ blob_type_json +
+ R"(}
+ ],
+ "highestFieldId": 1,
+ "partitionKeys": )" +
+ partition_keys_json +
+ R"(,
+ "primaryKeys": [],
+ "options": )" +
+ options_json +
+ R"(,
+ "timeMillis": 0
+ })";
+ return TableSchema::CreateFromJson(schema_json);
+}
+
} // namespace
TEST(SchemaValidationTest, TestSimple) {
@@ -255,6 +291,18 @@ TEST(SchemaValidationTest, TestMosaicDataTypes) {
/*partition_keys=*/{}, /*primary_keys=*/{},
blob_options));
ASSERT_OK(SchemaValidation::ValidateTableSchema(*table_schema));
+ for (const std::string& blob_type_json : ContainerBlobTypeJsons()) {
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<TableSchema>
container_blob_schema,
+ MakeContainerBlobSchema(blob_type_json, "[]",
+ R"({
+ "bucket": "-1",
+ "file.format": "mosaic",
+ "row-tracking.enabled": "true",
+ "data-evolution.enabled": "true"
+ })"));
+
ASSERT_OK(SchemaValidation::ValidateTableSchema(*container_blob_schema));
+ }
+
blob_options[Options::BLOB_DESCRIPTOR_FIELD] = "blob";
ASSERT_OK_AND_ASSIGN(
table_schema,
@@ -299,6 +347,18 @@ TEST(SchemaValidationTest, TestLanceDataTypes) {
/*partition_keys=*/{},
/*primary_keys=*/{}, options));
ASSERT_OK(SchemaValidation::ValidateTableSchema(*table_schema));
+ for (const std::string& blob_type_json : ContainerBlobTypeJsons()) {
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<TableSchema>
container_blob_schema,
+ MakeContainerBlobSchema(blob_type_json, "[]",
+ R"({
+ "bucket": "-1",
+ "file.format": "lance",
+ "row-tracking.enabled": "true",
+ "data-evolution.enabled": "true"
+ })"));
+
ASSERT_OK(SchemaValidation::ValidateTableSchema(*container_blob_schema));
+ }
+
arrow::FieldVector unsupported_fields = {
arrow::field("map", arrow::map(arrow::int32(), arrow::utf8())),
arrow::field("ltz", arrow::timestamp(arrow::TimeUnit::MICRO, "UTC")),
@@ -373,6 +433,32 @@ TEST(SchemaValidationTest, TestRowTracking) {
ASSERT_OK(SchemaValidation::ValidateTableSchema(*deletion_vector_table_schema));
}
+TEST(SchemaValidationTest, TestContainerBlobRowTracking) {
+ for (const std::string& blob_type_json : ContainerBlobTypeJsons()) {
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<TableSchema> table_schema,
+ MakeContainerBlobSchema(blob_type_json, "[]",
+ R"({
+ "bucket": "-1",
+ "row-tracking.enabled": "true",
+ "data-evolution.enabled": "false"
+ })"));
+ ASSERT_OK_AND_ASSIGN(CoreOptions options,
CoreOptions::FromMap(table_schema->Options()));
+ ASSERT_NOK_WITH_MSG(
+ SchemaValidation::ValidateRowTracking(*table_schema, options),
+ "Data evolution config must be enabled for table with BLOB type
column.");
+
+ ASSERT_OK_AND_ASSIGN(table_schema,
MakeContainerBlobSchema(blob_type_json, R"(["blob"])",
+ R"({
+ "bucket": "-1",
+ "row-tracking.enabled": "true",
+ "data-evolution.enabled": "true"
+ })"));
+ ASSERT_OK_AND_ASSIGN(options,
CoreOptions::FromMap(table_schema->Options()));
+
ASSERT_NOK_WITH_MSG(SchemaValidation::ValidateRowTracking(*table_schema,
options),
+ "Blob field blob cannot be a partition key.");
+ }
+}
+
TEST(SchemaValidationTest, TestWithBlobField) {
auto f0 = arrow::field("f0", arrow::utf8());
auto f1 = arrow::field("f1", arrow::int32());
@@ -393,6 +479,32 @@ TEST(SchemaValidationTest, TestWithBlobField) {
TableSchema::Create(/*schema_id=*/0, schema, partition_keys,
primary_keys, options));
ASSERT_OK(SchemaValidation::ValidateTableSchema(*table_schema));
}
+ for (const std::string& blob_type_json : ContainerBlobTypeJsons()) {
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<TableSchema> table_schema,
+ MakeContainerBlobSchema(blob_type_json, "[]",
+ R"({
+ "bucket": "-1",
+ "row-tracking.enabled": "true",
+ "data-evolution.enabled": "true",
+ "blob-field": "blob"
+ })"));
+ ASSERT_OK(SchemaValidation::ValidateTableSchema(*table_schema));
+ }
+ const std::vector<std::pair<std::string, std::string>> inline_blob_options
= {
+ {std::string(Options::BLOB_DESCRIPTOR_FIELD),
+
R"({"bucket":"-1","row-tracking.enabled":"true","data-evolution.enabled":"true","blob-descriptor-field":"blob"})"},
+ {std::string(Options::BLOB_VIEW_FIELD),
+
R"({"bucket":"-1","row-tracking.enabled":"true","data-evolution.enabled":"true","blob-view-field":"blob"})"},
+ };
+ for (const std::string& blob_type_json : ContainerBlobTypeJsons()) {
+ for (const auto& [option_key, options_json] : inline_blob_options) {
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<TableSchema> table_schema,
+ MakeContainerBlobSchema(blob_type_json, "[]",
options_json));
+ ASSERT_NOK_WITH_MSG(
+ SchemaValidation::ValidateTableSchema(*table_schema),
+ "Field 'blob' in '" + option_key + "' must be a BLOB field in
table schema.");
+ }
+ }
{
// a blob table with data evolution may also enable deletion vectors
arrow::FieldVector fields = {f0, f1, f2, f3};
@@ -530,7 +642,8 @@ TEST(SchemaValidationTest, TestWithBlobField) {
std::shared_ptr<TableSchema> table_schema,
TableSchema::Create(/*schema_id=*/0, schema, partition_keys,
primary_keys, options));
ASSERT_NOK_WITH_MSG(SchemaValidation::ValidateTableSchema(*table_schema),
- "Field 'f0' in 'blob-field' must be a BLOB field
in table schema.");
+ "Field 'f0' in 'blob-field' must be a BLOB,
ARRAY<BLOB> or "
+ "MAP<..., BLOB> field in table schema.");
}
{
arrow::FieldVector fields = {f0, f1, f2, f3};
@@ -1428,8 +1541,9 @@ TEST(SchemaValidationTest,
TestMapSharedShreddingRejectsBlobValue) {
auto nested_blob_map = arrow::map(
arrow::utf8(), arrow::field("value",
arrow::struct_({BlobUtils::ToArrowField("blob")})));
std::map<std::string, std::string> options = {
- {Options::BUCKET, "1"},
- {Options::BUCKET_KEY, "f0"},
+ {Options::BUCKET, "-1"},
+ {Options::ROW_TRACKING_ENABLED, "true"},
+ {Options::DATA_EVOLUTION_ENABLED, "true"},
{"fields.f1.map.storage-layout", "shared-shredding"},
};
@@ -1445,8 +1559,9 @@ TEST(SchemaValidationTest,
TestMapSharedShreddingRejectsBlobValue) {
"partitionKeys": [],
"primaryKeys": [],
"options": {
- "bucket": "1",
- "bucket-key": "f0",
+ "bucket": "-1",
+ "row-tracking.enabled": "true",
+ "data-evolution.enabled": "true",
"fields.f1.map.storage-layout": "shared-shredding"
},
"timeMillis": 0
@@ -1464,8 +1579,8 @@ TEST(SchemaValidationTest,
TestMapSharedShreddingRejectsBlobValue) {
});
ASSERT_NOK_WITH_MSG(TableSchema::Create(/*schema_id=*/0, nested_schema,
/*partition_keys=*/{},
/*primary_keys=*/{}, options),
- "BLOB field must be a top-level field or the direct
value of a "
- "top-level MAP field.");
+ "BLOB field must be a top-level field or the direct
element/value of a "
+ "top-level ARRAY/MAP field.");
}
TEST(SchemaValidationTest, TestMapSharedShreddingCompression) {
diff --git a/src/paimon/core/schema/table_schema.cpp
b/src/paimon/core/schema/table_schema.cpp
index 9d240560..5e70ff5a 100644
--- a/src/paimon/core/schema/table_schema.cpp
+++ b/src/paimon/core/schema/table_schema.cpp
@@ -57,7 +57,7 @@ Result<std::unique_ptr<TableSchema>> TableSchema::Create(
for (const auto& primary_key : primary_keys) {
primary_key_set.insert(primary_key);
}
- PAIMON_RETURN_NOT_OK(BlobUtils::ValidateMapBlobWriteSchema(schema));
+ PAIMON_RETURN_NOT_OK(BlobUtils::ValidateContainerBlobWriteSchema(schema));
for (const auto& field : schema->fields()) {
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Field> field_with_id,
AssignFieldIdsRecursively(field,
/*set_field_id=*/true, &field_id));
diff --git a/src/paimon/core/schema/table_schema_test.cpp
b/src/paimon/core/schema/table_schema_test.cpp
index e195802c..966c6e27 100644
--- a/src/paimon/core/schema/table_schema_test.cpp
+++ b/src/paimon/core/schema/table_schema_test.cpp
@@ -1366,6 +1366,15 @@ TEST_F(TableSchemaTest, CreatingMapBlobSchemaIsRejected)
{
"not supported by the C++ writer");
}
+TEST_F(TableSchemaTest, CreatingArrayBlobSchemaIsRejected) {
+ auto array_type = arrow::list(BlobUtils::ToArrowField("item",
/*nullable=*/true));
+ ASSERT_NOK_WITH_MSG(
+ TableSchema::Create(/*schema_id=*/0,
+ arrow::schema({arrow::field("blob_array",
array_type)}),
+ /*partition_keys=*/{}, /*primary_keys=*/{},
/*options=*/{}),
+ "Writing a table with ARRAY<BLOB> is not supported by the C++ writer");
+}
+
TEST_F(TableSchemaTest, MapKeysSortedIsNormalized) {
auto sorted_map =
std::make_shared<arrow::MapType>(arrow::field("key", arrow::utf8(),
/*nullable=*/false),
diff --git a/src/paimon/format/blob/blob_file_batch_reader.cpp
b/src/paimon/format/blob/blob_file_batch_reader.cpp
index acbf152e..531b32a0 100644
--- a/src/paimon/format/blob/blob_file_batch_reader.cpp
+++ b/src/paimon/format/blob/blob_file_batch_reader.cpp
@@ -48,6 +48,12 @@
namespace paimon::blob {
namespace {
+constexpr int32_t kArrayBlobMagicNumber = 1094861634;
+constexpr int8_t kArrayBlobVersion = 1;
+constexpr int32_t kArrayBlobHeaderLength = 9;
+constexpr int32_t kArrayBlobIndexLengthSize = 4;
+constexpr int32_t kArrayBlobMinPayloadLength = kArrayBlobHeaderLength +
kArrayBlobIndexLengthSize;
+
constexpr int32_t kMapBlobMagicNumber = 0x4D424342;
constexpr int8_t kMapBlobVersion = 1;
constexpr int32_t kMapBlobHeaderLength = 9;
@@ -258,9 +264,9 @@ Status BlobFileBatchReader::SetReadSchema(::ArrowSchema*
read_schema,
fmt::format("read schema field number {} is not 1",
arrow_schema->num_fields()));
}
std::shared_ptr<arrow::Field> read_field = arrow_schema->field(0);
- if (!BlobUtils::IsBlobField(read_field) &&
!BlobUtils::IsMapBlobField(read_field)) {
- return Status::Invalid(
- fmt::format("field {} must be BLOB or MAP<..., BLOB>",
read_field->ToString()));
+ if (!BlobUtils::IsAnyBlobField(read_field)) {
+ return Status::Invalid(fmt::format("field {} must be BLOB, ARRAY<BLOB>
or MAP<..., BLOB>",
+ read_field->ToString()));
}
if (BlobUtils::IsMapBlobField(read_field)) {
const auto& map_type = static_cast<const
arrow::MapType&>(*read_field->type());
@@ -382,6 +388,153 @@ Result<std::shared_ptr<arrow::Array>>
BlobFileBatchReader::BuildContentArray(
return std::make_shared<arrow::StructArray>(struct_array_data);
}
+Result<BlobFileBatchReader::ArrayBlobPayload>
BlobFileBatchReader::ReadArrayBlobPayload(
+ size_t row_index) const {
+ if (target_blob_lengths_[row_index] < 0) {
+ return Status::Invalid(fmt::format("unsupported ARRAY<BLOB> record
length: {}",
+ target_blob_lengths_[row_index]));
+ }
+
+ const int64_t payload_offset = GetTargetContentOffset(row_index);
+ const int64_t payload_length = GetTargetContentLength(row_index);
+ if (payload_length < kArrayBlobMinPayloadLength) {
+ return Status::Invalid(
+ fmt::format("invalid ARRAY<BLOB> payload length: {}",
payload_length));
+ }
+
+ std::array<uint8_t, kArrayBlobHeaderLength> header;
+ PAIMON_RETURN_NOT_OK(ReadBlobContentAt(payload_offset, header.size(),
header.data()));
+ const auto magic_number = ReadLittleEndian<int32_t>(header.data());
+ if (magic_number != kArrayBlobMagicNumber) {
+ return Status::Invalid(
+ fmt::format("invalid ARRAY<BLOB> payload magic number: {}",
magic_number));
+ }
+ const auto version = static_cast<int8_t>(header[4]);
+ if (version != kArrayBlobVersion) {
+ return Status::NotImplemented(
+ fmt::format("unsupported ARRAY<BLOB> payload version: {}",
version));
+ }
+ const auto element_count = ReadLittleEndian<int32_t>(header.data() + 5);
+ if (element_count < 0) {
+ return Status::Invalid(fmt::format("invalid ARRAY<BLOB> element count:
{}", element_count));
+ }
+
+ const int64_t index_length_offset = payload_offset + payload_length -
kArrayBlobIndexLengthSize;
+ std::array<uint8_t, kArrayBlobIndexLengthSize> index_length_bytes;
+ PAIMON_RETURN_NOT_OK(ReadBlobContentAt(index_length_offset,
index_length_bytes.size(),
+ index_length_bytes.data()));
+ const auto index_length =
ReadLittleEndian<int32_t>(index_length_bytes.data());
+ const int64_t maximum_index_length = payload_length -
kArrayBlobMinPayloadLength;
+ if (index_length < 0 || index_length > maximum_index_length) {
+ return Status::Invalid(
+ fmt::format("invalid ARRAY<BLOB> element index length: {}",
index_length));
+ }
+ if (element_count > index_length) {
+ return Status::Invalid("ARRAY<BLOB> element count exceeds element
index length");
+ }
+
+ const int64_t index_offset = index_length_offset - index_length;
+ std::vector<char> index_bytes(index_length);
+ PAIMON_RETURN_NOT_OK(ReadBlobContentAt(index_offset, index_length,
+
reinterpret_cast<uint8_t*>(index_bytes.data())));
+ PAIMON_ASSIGN_OR_RAISE(std::vector<int64_t> element_lengths,
+ DeltaVarintCompressor::Decompress(index_bytes));
+ if (element_lengths.size() != static_cast<size_t>(element_count)) {
+ return Status::Invalid("ARRAY<BLOB> element count does not match
element index length");
+ }
+
+ const int64_t data_offset = payload_offset + kArrayBlobHeaderLength;
+ const int64_t data_length = index_offset - data_offset;
+ int64_t total_element_length = 0;
+ for (int64_t element_length : element_lengths) {
+ if (element_length == BlobDefs::kNullBinLength) {
+ continue;
+ }
+ if (element_length < 0) {
+ return Status::Invalid(
+ fmt::format("invalid ARRAY<BLOB> element length: {}",
element_length));
+ }
+ if (!blob_as_descriptor_ && element_length >
std::numeric_limits<int32_t>::max()) {
+ return Status::Invalid(
+ fmt::format("ARRAY<BLOB> inline element is too large: {}",
element_length));
+ }
+ if (element_length > data_length - total_element_length) {
+ return Status::Invalid("ARRAY<BLOB> element lengths exceed the
payload data length");
+ }
+ total_element_length += element_length;
+ }
+ if (total_element_length != data_length) {
+ return Status::Invalid("ARRAY<BLOB> element lengths do not match the
payload data length");
+ }
+ return ArrayBlobPayload{std::move(element_lengths), data_offset};
+}
+
+Status BlobFileBatchReader::AppendArrayBlobValues(const ArrayBlobPayload&
payload,
+ arrow::LargeBinaryBuilder*
blob_builder) const {
+ int64_t element_offset = payload.data_offset;
+ for (int64_t element_length : payload.element_lengths) {
+ if (element_length == BlobDefs::kNullBinLength) {
+ PAIMON_RETURN_NOT_OK_FROM_ARROW(blob_builder->AppendNull());
+ continue;
+ }
+ if (blob_as_descriptor_) {
+ PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<Blob> blob,
+ Blob::FromPath(file_path_, element_offset,
element_length));
+ PAIMON_UNIQUE_PTR<Bytes> descriptor = blob->ToDescriptor(pool_);
+ PAIMON_RETURN_NOT_OK_FROM_ARROW(
+ blob_builder->Append(descriptor->data(), descriptor->size()));
+ } else {
+ PAIMON_UNIQUE_PTR<Bytes> element_bytes =
+ Bytes::AllocateBytes(static_cast<size_t>(element_length),
pool_.get());
+ if (element_length > 0) {
+ PAIMON_RETURN_NOT_OK(
+ ReadBlobContentAt(element_offset, element_length,
+
reinterpret_cast<uint8_t*>(element_bytes->data())));
+ }
+ PAIMON_RETURN_NOT_OK_FROM_ARROW(
+ blob_builder->Append(element_bytes->data(), element_length));
+ }
+ element_offset += element_length;
+ }
+ return Status::OK();
+}
+
+Result<std::shared_ptr<arrow::Array>> BlobFileBatchReader::BuildArrayBlobArray(
+ int32_t rows_to_read) const {
+ const auto& struct_type = checked_cast<const
arrow::StructType&>(*target_type_);
+ const std::shared_ptr<arrow::Field>& list_field = struct_type.field(0);
+ const std::shared_ptr<arrow::ListType> list_type =
+ checked_pointer_cast<arrow::ListType>(list_field->type());
+ if (list_type->value_type()->id() != arrow::Type::LARGE_BINARY) {
+ return Status::Invalid("ARRAY<BLOB> element type must be large
binary");
+ }
+
+ std::unique_ptr<arrow::ArrayBuilder> array_builder;
+ PAIMON_RETURN_NOT_OK_FROM_ARROW(
+ arrow::MakeBuilder(arrow_pool_.get(), list_type, &array_builder));
+ auto* list_builder =
checked_cast<arrow::ListBuilder*>(array_builder.get());
+ auto* blob_builder =
checked_cast<arrow::LargeBinaryBuilder*>(list_builder->value_builder());
+ PAIMON_RETURN_NOT_OK_FROM_ARROW(list_builder->Reserve(rows_to_read));
+ for (int32_t k = 0; k < rows_to_read; ++k) {
+ const size_t row_index = current_pos_ + k;
+ const bool is_null = IsTargetNull(row_index);
+ PAIMON_RETURN_NOT_OK_FROM_ARROW(list_builder->Append(!is_null));
+ if (IsTargetPlaceholder(row_index)) {
+ PAIMON_RETURN_NOT_OK_FROM_ARROW(blob_builder->Append(
+ BlobDefs::kPlaceholderSentinel,
BlobDefs::kPlaceholderSentinelLength));
+ } else if (!is_null) {
+ PAIMON_ASSIGN_OR_RAISE(ArrayBlobPayload payload,
ReadArrayBlobPayload(row_index));
+ PAIMON_RETURN_NOT_OK(AppendArrayBlobValues(payload, blob_builder));
+ }
+ }
+
+ std::shared_ptr<arrow::Array> list_array;
+ PAIMON_RETURN_NOT_OK_FROM_ARROW(list_builder->Finish(&list_array));
+ PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::StructArray>
struct_array,
+ arrow::StructArray::Make({list_array},
{list_field}));
+ return struct_array;
+}
+
Result<BlobFileBatchReader::MapBlobPayload>
BlobFileBatchReader::ReadMapBlobPayload(
size_t row_index, int32_t fixed_key_length) const {
if (target_blob_lengths_[row_index] < 0) {
@@ -620,6 +773,9 @@ Result<std::shared_ptr<arrow::Array>>
BlobFileBatchReader::BuildMapBlobArray(
Result<std::shared_ptr<arrow::Array>> BlobFileBatchReader::BuildTargetArray(
int32_t rows_to_read) const {
const auto& struct_type = static_cast<const
arrow::StructType&>(*target_type_);
+ if (struct_type.field(0)->type()->id() == arrow::Type::LIST) {
+ return BuildArrayBlobArray(rows_to_read);
+ }
if (struct_type.field(0)->type()->id() == arrow::Type::MAP) {
return BuildMapBlobArray(rows_to_read);
}
diff --git a/src/paimon/format/blob/blob_file_batch_reader.h
b/src/paimon/format/blob/blob_file_batch_reader.h
index b1512f79..e49a2180 100644
--- a/src/paimon/format/blob/blob_file_batch_reader.h
+++ b/src/paimon/format/blob/blob_file_batch_reader.h
@@ -104,8 +104,9 @@ class BlobFileBatchReader : public FileBatchReader {
/// `emit_placeholder_sentinel` controls how placeholder entries
(bin_length ==
/// BlobDefs::kPlaceholderBinLength) are read: when false they fail the
read, as resolving
/// them requires the data-evolution blob fallback path; when true they
are returned as the
- /// non-null BlobDefs::kPlaceholderSentinel bytes for scalar BLOB, or a
two-entry map with
- /// duplicate keys for MAP<..., BLOB>. The fallback path removes these
internal values.
+ /// non-null BlobDefs::kPlaceholderSentinel bytes for scalar BLOB, a
one-element sentinel
+ /// array for ARRAY<BLOB>, or a two-entry map with duplicate keys for
MAP<..., BLOB>. The
+ /// fallback path removes these internal values.
static Result<std::unique_ptr<BlobFileBatchReader>> Create(
const std::shared_ptr<InputStream>& input_stream, int32_t batch_size,
bool blob_as_descriptor, bool emit_placeholder_sentinel,
@@ -154,6 +155,11 @@ class BlobFileBatchReader : public FileBatchReader {
}
private:
+ struct ArrayBlobPayload {
+ std::vector<int64_t> element_lengths;
+ int64_t data_offset;
+ };
+
struct MapBlobPayload {
std::vector<int64_t> key_lengths;
std::vector<int64_t> value_lengths;
@@ -179,6 +185,10 @@ class BlobFileBatchReader : public FileBatchReader {
/// Builds a null bitmap buffer for the given rows. Returns nullptr if no
nulls.
Result<std::shared_ptr<arrow::Buffer>> BuildNullBitmap(int32_t
rows_to_read) const;
Result<std::shared_ptr<arrow::Array>> BuildContentArray(int32_t
rows_to_read) const;
+ Result<ArrayBlobPayload> ReadArrayBlobPayload(size_t row_index) const;
+ Status AppendArrayBlobValues(const ArrayBlobPayload& payload,
+ arrow::LargeBinaryBuilder* blob_builder)
const;
+ Result<std::shared_ptr<arrow::Array>> BuildArrayBlobArray(int32_t
rows_to_read) const;
Result<MapBlobPayload> ReadMapBlobPayload(size_t row_index, int32_t
fixed_key_length) const;
Status AppendMapBlobKeys(const MapBlobPayload& payload,
const std::shared_ptr<arrow::DataType>& key_type,
diff --git a/src/paimon/format/blob/blob_file_batch_reader_test.cpp
b/src/paimon/format/blob/blob_file_batch_reader_test.cpp
index beefdd8f..80dde03f 100644
--- a/src/paimon/format/blob/blob_file_batch_reader_test.cpp
+++ b/src/paimon/format/blob/blob_file_batch_reader_test.cpp
@@ -64,6 +64,15 @@ std::string MapBlobGoldenBytes() {
"00000000000000248fe4237a7b44180400000001");
}
+std::string ArrayBlobGoldenBytes() {
+ // Java-compatible ARRAY<BLOB> golden file. Rows are:
+ // [], ["inline", null, "", "descriptor"], null, placeholder.
+ return HexToBytes(
+ "cf114e58424342410100000000000000001d000000000000009bd49157cf114e"
+ "58424342410104000000696e6c696e6564657363726970746f720c0d02140400"
+ "00003100000000000000d08307713a2863010400000001");
+}
+
} // namespace
TEST(BlobReaderBuilderTest, RejectsNullMemoryPool) {
@@ -141,9 +150,9 @@ class BlobFileBatchReaderTest : public testing::Test,
public ::testing::WithPara
}
}
- Result<std::string> ReadMapBlobValue(const
std::shared_ptr<arrow::LargeBinaryArray>& blob_array,
- int64_t index, bool
blob_as_descriptor,
- const std::shared_ptr<FileSystem>&
file_system) {
+ Result<std::string> ReadBlobValue(const
std::shared_ptr<arrow::LargeBinaryArray>& blob_array,
+ int64_t index, bool blob_as_descriptor,
+ const std::shared_ptr<FileSystem>&
file_system) {
std::string stored_value = blob_array->GetString(index);
if (!blob_as_descriptor) {
return stored_value;
@@ -179,6 +188,27 @@ class BlobFileBatchReaderTest : public testing::Test,
public ::testing::WithPara
ASSERT_NOK_WITH_MSG(reader->NextBatch(), expected_message);
}
+ void CheckArrayBlobReadFails(const std::string& file_bytes,
+ const std::string& expected_message) {
+ auto dir = paimon::test::UniqueTestDirectory::Create();
+ ASSERT_TRUE(dir);
+ const std::string file_path = dir->Str() + "/corrupt-array.blob";
+ std::shared_ptr<FileSystem> file_system =
std::make_shared<LocalFileSystem>();
+ ASSERT_OK(file_system->WriteFile(file_path, file_bytes,
/*overwrite=*/true));
+
+ auto array_type = arrow::list(BlobUtils::ToArrowField("item",
/*nullable=*/true));
+ auto schema = arrow::schema({arrow::field("blob_array", array_type)});
+ ::ArrowSchema c_schema;
+ ASSERT_TRUE(arrow::ExportSchema(*schema, &c_schema).ok());
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> input,
file_system->Open(file_path));
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<BlobFileBatchReader> reader,
+ BlobFileBatchReader::Create(
+ input, /*batch_size=*/16,
/*blob_as_descriptor=*/false,
+ /*emit_placeholder_sentinel=*/true, pool_,
GetArrowPool(pool_)));
+ ASSERT_OK(reader->SetReadSchema(&c_schema, nullptr, std::nullopt));
+ ASSERT_NOK_WITH_MSG(reader->NextBatch(), expected_message);
+ }
+
private:
std::string blob_field_name_;
std::shared_ptr<MemoryPool> pool_;
@@ -246,7 +276,7 @@ TEST_P(BlobFileBatchReaderTest, TestMapBlob) {
continue;
}
ASSERT_OK_AND_ASSIGN(std::string value,
- ReadMapBlobValue(values, i, blob_as_descriptor,
file_system));
+ ReadBlobValue(values, i, blob_as_descriptor,
file_system));
ASSERT_TRUE(normalized_values_builder.Append(value).ok());
}
std::shared_ptr<arrow::Array> normalized_values;
@@ -266,6 +296,73 @@ TEST_P(BlobFileBatchReaderTest, TestMapBlob) {
<< "expected: " << expected->ToString() << "\nactual: " <<
normalized_map->ToString();
}
+TEST_P(BlobFileBatchReaderTest, TestArrayBlob) {
+ auto dir = paimon::test::UniqueTestDirectory::Create();
+ ASSERT_TRUE(dir);
+ const std::string file_path = dir->Str() + "/array-blob.blob";
+ std::shared_ptr<FileSystem> file_system =
std::make_shared<LocalFileSystem>();
+ ASSERT_OK(file_system->WriteFile(file_path, ArrayBlobGoldenBytes(),
/*overwrite=*/true));
+
+ auto array_type = arrow::list(BlobUtils::ToArrowField("item",
/*nullable=*/true));
+ auto schema = arrow::schema({arrow::field("blob_array", array_type)});
+ ::ArrowSchema c_schema;
+ ASSERT_TRUE(arrow::ExportSchema(*schema, &c_schema).ok());
+
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> input,
file_system->Open(file_path));
+ const bool blob_as_descriptor = GetParam();
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<BlobFileBatchReader> reader,
+ BlobFileBatchReader::Create(input, /*batch_size=*/2,
blob_as_descriptor,
+
/*emit_placeholder_sentinel=*/true, pool_,
+ GetArrowPool(pool_)));
+ ASSERT_OK(reader->SetReadSchema(&c_schema, nullptr, std::nullopt));
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::ChunkedArray> chunked_array,
+
paimon::test::ReadResultCollector::CollectResult(std::move(reader)));
+ std::shared_ptr<arrow::Array> combined_array =
+ arrow::Concatenate(chunked_array->chunks()).ValueOrDie();
+
+ auto struct_array =
std::dynamic_pointer_cast<arrow::StructArray>(combined_array);
+ ASSERT_TRUE(struct_array);
+ auto list_array =
std::dynamic_pointer_cast<arrow::ListArray>(struct_array->field(0));
+ ASSERT_TRUE(list_array);
+ std::shared_ptr<arrow::ListArray> normalized_list = list_array;
+ if (blob_as_descriptor) {
+ auto values =
std::dynamic_pointer_cast<arrow::LargeBinaryArray>(list_array->values());
+ ASSERT_TRUE(values);
+ arrow::LargeBinaryBuilder normalized_values_builder;
+ for (int64_t i = 0; i < values->length(); ++i) {
+ if (values->IsNull(i)) {
+ ASSERT_TRUE(normalized_values_builder.AppendNull().ok());
+ continue;
+ }
+ const std::string_view stored_value = values->GetView(i);
+ if (BlobDefs::IsPlaceholderSentinel(stored_value.data(),
stored_value.size())) {
+
ASSERT_TRUE(normalized_values_builder.Append(stored_value).ok());
+ continue;
+ }
+ ASSERT_OK_AND_ASSIGN(
+ std::string value,
+ ReadBlobValue(values, i, /*blob_as_descriptor=*/true,
file_system));
+ ASSERT_TRUE(normalized_values_builder.Append(value).ok());
+ }
+ std::shared_ptr<arrow::Array> normalized_values;
+ ASSERT_TRUE(normalized_values_builder.Finish(&normalized_values).ok());
+ normalized_list = std::make_shared<arrow::ListArray>(
+ list_array->type(), list_array->length(),
list_array->value_offsets(),
+ normalized_values, list_array->null_bitmap(),
list_array->null_count(),
+ list_array->offset());
+ }
+ const std::string expected_json = R"([
+ [],
+ ["inline", null, "", "descriptor"],
+ null,
+ ["_PAIMON_BLOB_PLACEHOLDER"]
+ ])";
+ std::shared_ptr<arrow::Array> expected =
+ arrow::ipc::internal::json::ArrayFromJSON(array_type,
expected_json).ValueOrDie();
+ ASSERT_TRUE(expected->Equals(normalized_list))
+ << "expected: " << expected->ToString() << "\nactual: " <<
normalized_list->ToString();
+}
+
TEST_P(BlobFileBatchReaderTest, MapBlobFallbackAcrossSequenceLayers) {
auto dir = paimon::test::UniqueTestDirectory::Create();
ASSERT_TRUE(dir);
@@ -329,8 +426,8 @@ TEST_P(BlobFileBatchReaderTest,
MapBlobFallbackAcrossSequenceLayers) {
ASSERT_OK(old_reader->SetReadSchema(&old_schema, nullptr, std::nullopt));
std::vector<std::vector<BlobFallbackBatchReader::Segment>> groups(2);
- groups[0].push_back({std::move(new_reader), {}});
- groups[1].push_back({std::move(old_reader), {}});
+ groups[0].push_back({std::move(new_reader), {}}); //
NOLINT(modernize-use-emplace)
+ groups[1].push_back({std::move(old_reader), {}}); //
NOLINT(modernize-use-emplace)
ASSERT_OK_AND_ASSIGN(
std::unique_ptr<BlobFallbackBatchReader> fallback,
BlobFallbackBatchReader::Create(std::move(groups), map_schema,
/*read_batch_size=*/2,
@@ -352,9 +449,9 @@ TEST_P(BlobFileBatchReaderTest,
MapBlobFallbackAcrossSequenceLayers) {
ASSERT_EQ("alpha", keys->GetString(0));
ASSERT_EQ("omega", keys->GetString(3));
ASSERT_OK_AND_ASSIGN(std::string first_value,
- ReadMapBlobValue(values, 0, blob_as_descriptor,
file_system));
+ ReadBlobValue(values, 0, blob_as_descriptor,
file_system));
ASSERT_OK_AND_ASSIGN(std::string last_value,
- ReadMapBlobValue(values, 3, blob_as_descriptor,
file_system));
+ ReadBlobValue(values, 3, blob_as_descriptor,
file_system));
ASSERT_EQ("hello", first_value);
ASSERT_EQ("world", last_value);
}
@@ -481,6 +578,49 @@ TEST_F(BlobFileBatchReaderTest,
RejectsCorruptMapPayloadMetadata) {
CheckMapBlobReadFails(corrupted, arrow::utf8(), "payload: duplicate key");
}
+TEST_F(BlobFileBatchReaderTest, RejectsCorruptArrayPayloadMetadata) {
+ // The first empty-array record starts at byte 0. Its ARRAY payload
occupies [4, 17): header
+ // [4, 13) and the index length [13, 17). The second record starts at byte
29; its payload has
+ // header [33, 42), data [42, 58), index [58, 62), and index length [62,
66).
+ const std::string golden = ArrayBlobGoldenBytes();
+
+ std::string corrupted = golden;
+ corrupted[4] = 0;
+ CheckArrayBlobReadFails(corrupted, "invalid ARRAY<BLOB> payload magic
number");
+
+ corrupted = golden;
+ corrupted[8] = 2;
+ CheckArrayBlobReadFails(corrupted, "unsupported ARRAY<BLOB> payload
version");
+
+ corrupted = golden;
+ std::fill(corrupted.begin() + 9, corrupted.begin() + 13,
static_cast<char>(0xFF));
+ CheckArrayBlobReadFails(corrupted, "invalid ARRAY<BLOB> element count");
+
+ corrupted = golden;
+ std::fill(corrupted.begin() + 13, corrupted.begin() + 17,
static_cast<char>(0xFF));
+ CheckArrayBlobReadFails(corrupted, "invalid ARRAY<BLOB> element index
length");
+
+ corrupted = golden;
+ corrupted[9] = 1;
+ CheckArrayBlobReadFails(corrupted, "element count exceeds element index
length");
+
+ corrupted = golden;
+ corrupted[38] = 3;
+ CheckArrayBlobReadFails(corrupted, "element count does not match element
index length");
+
+ corrupted = golden;
+ corrupted[59] = 0x0F;
+ CheckArrayBlobReadFails(corrupted, "invalid ARRAY<BLOB> element length");
+
+ corrupted = golden;
+ corrupted[61] = 0x16;
+ CheckArrayBlobReadFails(corrupted, "element lengths exceed the payload
data length");
+
+ corrupted = golden;
+ corrupted[61] = 0x12;
+ CheckArrayBlobReadFails(corrupted, "element lengths do not match the
payload data length");
+}
+
TEST_P(BlobFileBatchReaderTest, TestPushdownBitmap) {
std::string test_data_path = paimon::test::GetDataDir() +
"/db_with_blob.db/table_with_blob/";
auto dir = paimon::test::UniqueTestDirectory::Create();
@@ -740,7 +880,8 @@ TEST_F(BlobFileBatchReaderTest,
SetReadSchemaWithInvalidInputs) {
GetArrowPool(pool_)));
ASSERT_NOK_WITH_MSG(reader->SetReadSchema(&c_schema,
/*predicate=*/nullptr,
/*selection_bitmap=*/std::nullopt),
- "field my_blob_field: large_binary must be BLOB or
MAP<..., BLOB>");
+ "field my_blob_field: large_binary must be BLOB,
ARRAY<BLOB> or "
+ "MAP<..., BLOB>");
}
{
auto schema = arrow::schema({BlobUtils::ToArrowField("my_blob_field",
false)});
diff --git a/src/paimon/format/blob/blob_stats_extractor.cpp
b/src/paimon/format/blob/blob_stats_extractor.cpp
index f4bc757f..5e71121f 100644
--- a/src/paimon/format/blob/blob_stats_extractor.cpp
+++ b/src/paimon/format/blob/blob_stats_extractor.cpp
@@ -46,9 +46,9 @@ BlobStatsExtractor::ExtractWithFileInfo(const
std::shared_ptr<FileSystem>& file_
return Status::Invalid(
fmt::format("schema field number {} is not 1",
write_schema_->num_fields()));
}
- if (!BlobUtils::IsBlobField(write_schema_->field(0))) {
+ if (!BlobUtils::IsAnyBlobField(write_schema_->field(0))) {
return Status::Invalid(
- fmt::format("field {} is not BLOB",
write_schema_->field(0)->ToString()));
+ fmt::format("field {} is not a blob-file field",
write_schema_->field(0)->ToString()));
}
// The reader only serves footer metadata (GetNumberOfRows); NextBatch is
never called, so
diff --git a/src/paimon/format/blob/blob_stats_extractor_test.cpp
b/src/paimon/format/blob/blob_stats_extractor_test.cpp
index befbfa52..6b5f4954 100644
--- a/src/paimon/format/blob/blob_stats_extractor_test.cpp
+++ b/src/paimon/format/blob/blob_stats_extractor_test.cpp
@@ -97,7 +97,7 @@ TEST_F(BlobStatsExtractorTest, TestInvalidCase) {
auto non_blob_schema = arrow::schema({string_field});
BlobStatsExtractor extractor(non_blob_schema);
ASSERT_NOK_WITH_MSG(extractor.ExtractWithFileInfo(fs_, blob_file_path,
pool_),
- "field string_field: string is not BLOB");
+ "field string_field: string is not a blob-file
field");
}
// Should fail because file doesn't exist
{
diff --git a/test/inte/blob_table_inte_test.cpp
b/test/inte/blob_table_inte_test.cpp
index 40f0625f..5e36c3aa 100644
--- a/test/inte/blob_table_inte_test.cpp
+++ b/test/inte/blob_table_inte_test.cpp
@@ -535,9 +535,9 @@ class BlobTableInteTest : public testing::Test, public
::testing::WithParamInter
});
}
- Result<std::shared_ptr<arrow::MapArray>> NormalizeMapBlobValues(
- const std::shared_ptr<arrow::MapArray>& map_array, bool
blob_as_descriptor) const {
- const auto& values = checked_cast<const
arrow::LargeBinaryArray&>(*map_array->items());
+ Result<std::shared_ptr<arrow::Array>> ResolveBlobDescriptors(
+ const std::shared_ptr<arrow::Array>& values_array) const {
+ const auto& values = checked_cast<const
arrow::LargeBinaryArray&>(*values_array);
auto fs = std::make_shared<LocalFileSystem>();
arrow::LargeBinaryBuilder builder;
for (int64_t i = 0; i < values.length(); ++i) {
@@ -546,10 +546,6 @@ class BlobTableInteTest : public testing::Test, public
::testing::WithParamInter
continue;
}
std::string_view stored = values.GetView(i);
- if (!blob_as_descriptor) {
- PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.Append(stored));
- continue;
- }
PAIMON_ASSIGN_OR_RAISE(
std::unique_ptr<Blob> blob,
Blob::FromDescriptor(stored.data(),
static_cast<int64_t>(stored.size())));
@@ -558,12 +554,50 @@ class BlobTableInteTest : public testing::Test, public
::testing::WithParamInter
}
std::shared_ptr<arrow::Array> normalized_values;
PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.Finish(&normalized_values));
+ return normalized_values;
+ }
+
+ Result<std::shared_ptr<arrow::MapArray>> NormalizeMapBlobValues(
+ const std::shared_ptr<arrow::MapArray>& map_array, bool
blob_as_descriptor) const {
+ if (!blob_as_descriptor) {
+ return map_array;
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Array> normalized_values,
+ ResolveBlobDescriptors(map_array->items()));
return std::make_shared<arrow::MapArray>(map_array->type(),
map_array->length(),
map_array->value_offsets(),
map_array->keys(),
normalized_values,
map_array->null_bitmap(),
map_array->null_count(),
map_array->offset());
}
+ Result<std::shared_ptr<arrow::ListArray>> NormalizeArrayBlobValues(
+ const std::shared_ptr<arrow::ListArray>& list_array, bool
blob_as_descriptor) const {
+ if (!blob_as_descriptor) {
+ return list_array;
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Array> normalized_values,
+ ResolveBlobDescriptors(list_array->values()));
+ return std::make_shared<arrow::ListArray>(list_array->type(),
list_array->length(),
+ list_array->value_offsets(),
normalized_values,
+ list_array->null_bitmap(),
+ list_array->null_count(),
list_array->offset());
+ }
+
+ void CheckArrayBlobColumn(const std::shared_ptr<arrow::StructArray>& rows,
+ const std::string& field_name, const
std::string& expected_json,
+ bool blob_as_descriptor) const {
+ auto list_array =
+
std::dynamic_pointer_cast<arrow::ListArray>(rows->GetFieldByName(field_name));
+ ASSERT_TRUE(list_array) << field_name;
+ ASSERT_OK_AND_ASSIGN(auto normalized,
+ NormalizeArrayBlobValues(list_array,
blob_as_descriptor));
+ auto expected =
arrow::ipc::internal::json::ArrayFromJSON(list_array->type(), expected_json)
+ .ValueOrDie();
+ ASSERT_TRUE(expected->Equals(normalized))
+ << field_name << " expected: " << expected->ToString()
+ << " actual: " << normalized->ToString();
+ }
+
void CheckMapBlobColumn(const std::shared_ptr<arrow::StructArray>& rows,
const std::string& field_name, const std::string&
expected_json,
bool blob_as_descriptor) const {
@@ -4378,6 +4412,57 @@ TEST_P(BlobTableInteTest,
TestReadBlobDescriptorFieldFromJava) {
ASSERT_TRUE(resolved->Equals(expected_with_rk));
}
+TEST_P(BlobTableInteTest, TestReadArrayBlobTableFromJava) {
+ if (GetParam() != "parquet") {
+ GTEST_SKIP() << "the Java fixture uses Parquet";
+ }
+ const std::string table_path = GetDataDir() +
"/parquet/array_blob_java.db/array_blob_java";
+ const std::vector<std::string> read_fields = {"id", "array_payloads"};
+
+ for (int64_t snapshot_id : {1, 2, 3}) {
+ ScanContextBuilder scan_builder(table_path);
+ scan_builder.AddOption(Options::SCAN_SNAPSHOT_ID,
std::to_string(snapshot_id));
+ ASSERT_OK_AND_ASSIGN(auto scan_context, scan_builder.Finish());
+ ASSERT_OK_AND_ASSIGN(auto table_scan,
TableScan::Create(std::move(scan_context)));
+ ASSERT_OK_AND_ASSIGN(auto plan, table_scan->CreatePlan());
+
+ size_t array_layer_count = 0;
+ for (const auto& split : plan->Splits()) {
+ auto data_split = std::dynamic_pointer_cast<DataSplitImpl>(split);
+ ASSERT_TRUE(data_split);
+ for (const auto& file : data_split->DataFiles()) {
+ if (file->write_cols ==
+
std::optional<std::vector<std::string>>({"array_payloads"})) {
+ ++array_layer_count;
+ }
+ }
+ }
+ ASSERT_EQ(snapshot_id, static_cast<int64_t>(array_layer_count));
+
+ for (bool blob_as_descriptor : {false, true}) {
+ std::map<std::string, std::string> read_options = {
+ {Options::BLOB_AS_DESCRIPTOR, blob_as_descriptor ? "true" :
"false"}};
+ ASSERT_OK_AND_ASSIGN(auto result, ReadTable(table_path,
read_fields, plan,
+ /*predicate=*/nullptr,
read_options));
+ ASSERT_TRUE(result);
+ auto combined = arrow::Concatenate(result->chunks()).ValueOrDie();
+ auto rows =
std::dynamic_pointer_cast<arrow::StructArray>(combined);
+ ASSERT_TRUE(rows);
+ ASSERT_EQ(4, rows->length());
+ const auto& ids = checked_cast<const
arrow::Int32Array&>(*rows->GetFieldByName("id"));
+ for (int64_t i = 0; i < ids.length(); ++i) {
+ ASSERT_EQ(i + 1, ids.Value(i));
+ }
+
+ const std::string expected_json =
+ snapshot_id == 1
+ ? R"json([["array-alpha", null, "", "array-omega"], null,
[], ["array-single"]])json"
+ : R"json([["array-alpha", null, "", "array-omega"],
["array-updated"], [], ["array-single"]])json";
+ CheckArrayBlobColumn(rows, "array_payloads", expected_json,
blob_as_descriptor);
+ }
+ }
+}
+
TEST_P(BlobTableInteTest, TestReadMapBlobTableFromJava) {
if (GetParam() != "parquet") {
GTEST_SKIP() << "the Java fixture uses Parquet";
diff --git a/test/inte/paimon_read_compat_inte_test.cpp
b/test/inte/paimon_read_compat_inte_test.cpp
index f6468b6f..5761720d 100644
--- a/test/inte/paimon_read_compat_inte_test.cpp
+++ b/test/inte/paimon_read_compat_inte_test.cpp
@@ -488,6 +488,32 @@ TEST_P(PaimonReadCompatInteTest, ReadsBlobValues) {
ASSERT_EQ(inline_descriptor->Offset(), 7);
ASSERT_EQ(inline_descriptor->Length(), 11);
+ // Python and Java store ARRAY<BLOB> values in standard separate BLOB
files. Validate resolved
+ // payloads, a null element, a null array, and an empty array. Rust stores
raw values inline in
+ // Parquet, which is not a compatible Paimon BLOB representation and is
asserted separately.
+ if (param.writer_prefix != "rust") {
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::ChunkedArray>
array_blob_result,
+ ReadTable(param, param.writer_prefix +
"_array_blob_types",
+ {"id", "f_array_blob"},
blob_value_options));
+ std::shared_ptr<arrow::StructArray> array_blob_rows =
GetOnlyStructChunk(array_blob_result);
+ ASSERT_TRUE(array_blob_rows);
+ ASSERT_EQ(array_blob_rows->length(), 3);
+ AssertFieldEqualsJson(array_blob_rows, "id", arrow::int32(), "[1, 2,
3]");
+
+ std::string expected_array_blob_json;
+ if (param.file_format == "parquet") {
+ expected_array_blob_json = R"([["blob-array-0", null,
"blob-array-2"], null, []])";
+ } else if (param.file_format == "orc") {
+ expected_array_blob_json =
+ R"([["blob-array-left", null, "blob-array-right"], null, []])";
+ } else {
+ ASSERT_EQ(param.file_format, "avro");
+ expected_array_blob_json = R"([["array-blob-value", null, ""],
null, []])";
+ }
+ AssertFieldEqualsJson(array_blob_rows, "f_array_blob",
arrow::list(arrow::large_binary()),
+ expected_array_blob_json);
+ }
+
// Python and Java store MAP<STRING, BLOB> values in standard separate
BLOB files. Validate
// resolved payloads, null values, and an empty map. Rust stores raw
values inline in Parquet,
// which is not a compatible Paimon BLOB representation and is asserted
separately below.
@@ -512,6 +538,12 @@ TEST(PaimonReadCompatInteStandaloneTest,
RejectsNonStandardRustMapBlob) {
"Parquet does not support partial projection inside list/map: src
map<string, binary");
}
+TEST(PaimonReadCompatInteStandaloneTest, RejectsNonStandardRustArrayBlob) {
+ ASSERT_NOK_WITH_MSG(
+ ReadTable("parquet", "rust_array_blob_types", {"f_array_blob"}),
+ "Parquet does not support partial projection inside list/map: src
list<element: binary");
+}
+
TEST_P(PaimonReadCompatInteTest, ReadsVectorValues) {
const CompatibilityParam& param = GetParam();
if (!param.supports_vector) {
@@ -581,8 +613,6 @@ TEST_P(PaimonUnsupportedTypeInteTest, ReportsExpectedError)
{
std::vector<UnsupportedReadParam> UnsupportedReadParams() {
const std::vector<UnsupportedReadCase> read_cases = {
- {"ArrayBlob", "array_blob_types", "f_array_blob",
- "BLOB field must be a top-level field or the direct value of a
top-level MAP field"},
{"TimePrecision0", "time_types", "f_time_0", ""},
{"TimePrecision3", "time_types", "f_time_3", ""},
{"TimePrecision6", "time_types", "f_time_6", ""},
diff --git
a/test/test_data/parquet/append_types_compatibility.db/rust_array_blob_types/README.md
b/test/test_data/parquet/append_types_compatibility.db/rust_array_blob_types/README.md
index c4f2883b..bb256391 100644
---
a/test/test_data/parquet/append_types_compatibility.db/rust_array_blob_types/README.md
+++
b/test/test_data/parquet/append_types_compatibility.db/rust_array_blob_types/README.md
@@ -14,4 +14,5 @@ Add: (1, ["blob-array-0", null, "blob-array-2"])
Add: (2, null)
Add: (3, [])
-Paimon C++ is expected to reject this table while ARRAY<BLOB> is unsupported.
+Paimon Rust stores these values inline as Parquet binary; Java and C++
intentionally reject that
+representation as a Paimon BLOB array.
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/README.md
b/test/test_data/parquet/array_blob_java.db/array_blob_java/README.md
new file mode 100644
index 00000000..d5b42734
--- /dev/null
+++ b/test/test_data/parquet/array_blob_java.db/array_blob_java/README.md
@@ -0,0 +1,29 @@
+Schema:
+id INT
+array_payloads ARRAY<BLOB>
+
+Options:
+bucket = -1
+data-evolution.enabled = true
+file.format = parquet
+row-tracking.enabled = true
+
+Msgs:
+snapshot-1
+Commit four rows with BatchTableWrite over the full row type. BLOB values
below are UTF-8 byte
+payloads shown as text; "" is a zero-length BLOB:
+id 1: array_payloads is ["array-alpha", null, "", "array-omega"]
+id 2: array_payloads is null
+id 3: array_payloads is empty
+id 4: array_payloads is ["array-single"]
+
+snapshot-2
+Create BatchTableWrite with the write type projected to array_payloads. Write
+BlobArrayPlaceholder.INSTANCE for row positions 0, 2, and 3, and
["array-updated"] for row
+position 1. Set the first row id to 0 before committing. This adds a second
sequence layer while
+updating id 2 and falling back to snapshot-1 values for the other rows.
+
+snapshot-3
+Repeat a projected write with BlobArrayPlaceholder.INSTANCE for all four row
positions and set the
+first row id to 0. Reading this snapshot falls back to snapshot 2 for id 2 and
through both newer
+placeholder layers to snapshot 1 for ids 1, 3, and 4.
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-0e375e23-85cb-4455-a0bd-511a331c4dc8-0.blob
b/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-0e375e23-85cb-4455-a0bd-511a331c4dc8-0.blob
new file mode 100644
index 00000000..4dc6a6a3
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-0e375e23-85cb-4455-a0bd-511a331c4dc8-0.blob
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-53163ebc-5467-47c9-ad31-228f7915c4db-0.blob
b/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-53163ebc-5467-47c9-ad31-228f7915c4db-0.blob
new file mode 100644
index 00000000..2ed55e52
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-53163ebc-5467-47c9-ad31-228f7915c4db-0.blob
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-91236ae6-4891-42dd-b2f9-a35751786cf2-0.parquet
b/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-91236ae6-4891-42dd-b2f9-a35751786cf2-0.parquet
new file mode 100644
index 00000000..95a7c0ee
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-91236ae6-4891-42dd-b2f9-a35751786cf2-0.parquet
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-91236ae6-4891-42dd-b2f9-a35751786cf2-1.blob
b/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-91236ae6-4891-42dd-b2f9-a35751786cf2-1.blob
new file mode 100644
index 00000000..6b318752
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/bucket-0/data-91236ae6-4891-42dd-b2f9-a35751786cf2-1.blob
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-1f8acd08-a84a-40d4-8069-86e92070df21-0
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-1f8acd08-a84a-40d4-8069-86e92070df21-0
new file mode 100644
index 00000000..a95e0863
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-1f8acd08-a84a-40d4-8069-86e92070df21-0
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-814974cc-e019-45eb-b1cc-5aa4b6cf9987-0
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-814974cc-e019-45eb-b1cc-5aa4b6cf9987-0
new file mode 100644
index 00000000..85b1ca63
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-814974cc-e019-45eb-b1cc-5aa4b6cf9987-0
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-fd42fd8f-82ac-4a8f-acd0-771e4db82f25-0
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-fd42fd8f-82ac-4a8f-acd0-771e4db82f25-0
new file mode 100644
index 00000000..2c30fdd0
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-fd42fd8f-82ac-4a8f-acd0-771e4db82f25-0
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-0
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-0
new file mode 100644
index 00000000..3dfff774
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-0
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-1
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-1
new file mode 100644
index 00000000..aa347b91
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-1
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-0
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-0
new file mode 100644
index 00000000..81e1f6af
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-0
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-1
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-1
new file mode 100644
index 00000000..bc1052bb
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-1
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-f082e141-fa95-4685-9b7c-5555137b5310-0
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-f082e141-fa95-4685-9b7c-5555137b5310-0
new file mode 100644
index 00000000..83d0c9fc
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-f082e141-fa95-4685-9b7c-5555137b5310-0
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-f082e141-fa95-4685-9b7c-5555137b5310-1
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-f082e141-fa95-4685-9b7c-5555137b5310-1
new file mode 100644
index 00000000..2024b739
Binary files /dev/null and
b/test/test_data/parquet/array_blob_java.db/array_blob_java/manifest/manifest-list-f082e141-fa95-4685-9b7c-5555137b5310-1
differ
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/schema/schema-0
b/test/test_data/parquet/array_blob_java.db/array_blob_java/schema/schema-0
new file mode 100644
index 00000000..b99a52a4
--- /dev/null
+++ b/test/test_data/parquet/array_blob_java.db/array_blob_java/schema/schema-0
@@ -0,0 +1,26 @@
+{
+ "version" : 3,
+ "id" : 0,
+ "fields" : [ {
+ "id" : 0,
+ "name" : "id",
+ "type" : "INT"
+ }, {
+ "id" : 1,
+ "name" : "array_payloads",
+ "type" : {
+ "type" : "ARRAY",
+ "element" : "BLOB"
+ }
+ } ],
+ "highestFieldId" : 1,
+ "partitionKeys" : [ ],
+ "primaryKeys" : [ ],
+ "options" : {
+ "bucket" : "-1",
+ "data-evolution.enabled" : "true",
+ "file.format" : "parquet",
+ "row-tracking.enabled" : "true"
+ },
+ "timeMillis" : 1790220836832
+}
\ No newline at end of file
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/EARLIEST
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/EARLIEST
new file mode 100644
index 00000000..56a6051c
--- /dev/null
+++
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/EARLIEST
@@ -0,0 +1 @@
+1
\ No newline at end of file
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/LATEST
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/LATEST
new file mode 100644
index 00000000..e440e5c8
--- /dev/null
+++ b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/LATEST
@@ -0,0 +1 @@
+3
\ No newline at end of file
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-1
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-1
new file mode 100644
index 00000000..0242211f
--- /dev/null
+++
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-1
@@ -0,0 +1,18 @@
+{
+ "version" : 3,
+ "uuid" : "cd63c7b3-e3ec-450c-a2ca-10d755d1471b",
+ "id" : 1,
+ "schemaId" : 0,
+ "baseManifestList" : "manifest-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-0",
+ "baseManifestListSize" : 1158,
+ "deltaManifestList" : "manifest-list-ad3c543d-86c3-4e03-bf02-05fe25f6ec26-1",
+ "deltaManifestListSize" : 1261,
+ "commitUser" : "3b5ecbe8-697e-438e-97bb-6d51ec94d6e4",
+ "writerVersion" :
"java-2.2-SNAPSHOT-91ee13f838e066c399810ca8160f64f2c204fd01",
+ "commitIdentifier" : 9223372036854775807,
+ "commitKind" : "APPEND",
+ "timeMillis" : 1790220837732,
+ "totalRecordCount" : 8,
+ "deltaRecordCount" : 8,
+ "nextRowId" : 4
+}
\ No newline at end of file
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-2
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-2
new file mode 100644
index 00000000..2166950e
--- /dev/null
+++
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-2
@@ -0,0 +1,18 @@
+{
+ "version" : 3,
+ "uuid" : "26cf1411-8f75-44c5-a152-a2ca4cf9e6f3",
+ "id" : 2,
+ "schemaId" : 0,
+ "baseManifestList" : "manifest-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-0",
+ "baseManifestListSize" : 1261,
+ "deltaManifestList" : "manifest-list-c4a8e4a1-7c33-4c93-b1cf-5b1b47dd1b35-1",
+ "deltaManifestListSize" : 1261,
+ "commitUser" : "26618ef8-7e24-4ec7-a0e1-b49e30510a56",
+ "writerVersion" :
"java-2.2-SNAPSHOT-91ee13f838e066c399810ca8160f64f2c204fd01",
+ "commitIdentifier" : 9223372036854775807,
+ "commitKind" : "APPEND",
+ "timeMillis" : 1790220837869,
+ "totalRecordCount" : 12,
+ "deltaRecordCount" : 4,
+ "nextRowId" : 4
+}
\ No newline at end of file
diff --git
a/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-3
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-3
new file mode 100644
index 00000000..ca4d9ce7
--- /dev/null
+++
b/test/test_data/parquet/array_blob_java.db/array_blob_java/snapshot/snapshot-3
@@ -0,0 +1,18 @@
+{
+ "version" : 3,
+ "uuid" : "dbef5223-e0bf-458b-bc81-b9c9015ccf57",
+ "id" : 3,
+ "schemaId" : 0,
+ "baseManifestList" : "manifest-list-f082e141-fa95-4685-9b7c-5555137b5310-0",
+ "baseManifestListSize" : 1295,
+ "deltaManifestList" : "manifest-list-f082e141-fa95-4685-9b7c-5555137b5310-1",
+ "deltaManifestListSize" : 1258,
+ "commitUser" : "ae0cfda3-b97e-4136-a999-6d1488d2f591",
+ "writerVersion" :
"java-2.2-SNAPSHOT-91ee13f838e066c399810ca8160f64f2c204fd01",
+ "commitIdentifier" : 9223372036854775807,
+ "commitKind" : "APPEND",
+ "timeMillis" : 1790220837886,
+ "totalRecordCount" : 16,
+ "deltaRecordCount" : 4,
+ "nextRowId" : 4
+}
\ No newline at end of file