diff --git a/src/paimon/common/data/blob_utils.cpp b/src/paimon/common/data/blob_utils.cpp index ad2ad315f..d7d8f336c 100644 --- a/src/paimon/common/data/blob_utils.cpp +++ b/src/paimon/common/data/blob_utils.cpp @@ -108,6 +108,26 @@ bool BlobUtils::IsBlobField(const std::shared_ptr& field) { return IsBlobMetadata(field->metadata()); } +bool BlobUtils::IsMapBlobField(const std::shared_ptr& field) { + if (field == nullptr || field->type()->id() != arrow::Type::MAP) { + return false; + } + const auto& map_type = checked_cast(*field->type()); + // Arrow's C schema bridge does not retain nested field metadata for MapType. Paimon's + // ordinary binary type is BINARY, so LARGE_BINARY uniquely identifies a BLOB value here. + return map_type.item_type()->id() == arrow::Type::LARGE_BINARY; +} + +Status BlobUtils::ValidateMapBlobWriteSchema(const std::shared_ptr& 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."); + } + } + return Status::OK(); +} + bool BlobUtils::IsBlobMetadata(const std::shared_ptr& metadata) { if (!metadata) { return false; diff --git a/src/paimon/common/data/blob_utils.h b/src/paimon/common/data/blob_utils.h index c13e24f91..86e0ef375 100644 --- a/src/paimon/common/data/blob_utils.h +++ b/src/paimon/common/data/blob_utils.h @@ -73,6 +73,10 @@ class PAIMON_EXPORT BlobUtils { const std::set& inline_fields); static bool IsBlobField(const std::shared_ptr& field); + /// Returns whether the field is a top-level MAP whose values are BLOBs. + static bool IsMapBlobField(const std::shared_ptr& field); + /// Rejects schemas that the C++ writer cannot safely mutate. + static Status ValidateMapBlobWriteSchema(const std::shared_ptr& schema); static bool IsBlobMetadata(const std::shared_ptr& metadata); static bool IsBlobFile(const std::string& file_name); diff --git a/src/paimon/common/reader/blob_fallback_batch_reader.cpp b/src/paimon/common/reader/blob_fallback_batch_reader.cpp index f3bcdd611..1a23ec2d1 100644 --- a/src/paimon/common/reader/blob_fallback_batch_reader.cpp +++ b/src/paimon/common/reader/blob_fallback_batch_reader.cpp @@ -56,7 +56,8 @@ Result> 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))) { + if (BlobUtils::IsBlobField(read_schema->field(i)) || + BlobUtils::IsMapBlobField(read_schema->field(i))) { if (blob_field_idx != -1) { return Status::Invalid( "Blob fallback read schema should contain exactly one blob field."); @@ -193,10 +194,29 @@ Result> BlobFallbackBatchReader::ComputePlaceholderFlags( std::fill(flags.begin() + pos, flags.begin() + pos + chunk.length, true); } else { std::shared_ptr blob_col = chunk.array->field(blob_field_idx_); - if (!blob_col || blob_col->type_id() != arrow::Type::LARGE_BINARY) { - return Status::Invalid(fmt::format( - "Blob fallback expects the blob column to be large binary, but got {}", - blob_col ? blob_col->type()->ToString() : "null")); + if (!blob_col) { + return Status::Invalid("Blob fallback got a null blob column."); + } + if (blob_col->type_id() == arrow::Type::MAP) { + auto map_col = checked_pointer_cast(blob_col); + const std::shared_ptr& keys = map_col->keys(); + const std::shared_ptr& 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); + } + } + 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())); } auto binary_col = checked_pointer_cast(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 011b3ff1b..3664b0c18 100644 --- a/src/paimon/common/reader/blob_fallback_batch_reader.h +++ b/src/paimon/common/reader/blob_fallback_batch_reader.h @@ -54,9 +54,8 @@ namespace paimon { /// vector has to reach every group the same way, through the file segments' readers and /// 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. Placeholder rows are identified by exact equality with the -/// BlobDefs::kPlaceholderSentinel bytes, emitted by the blob format reader when -/// BlobDefs::kEmitPlaceholderSentinelKey is set. +/// 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>. /// 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 b6afc15c4..8ff580f3c 100644 --- a/src/paimon/core/append/append_compact_coordinator.cpp +++ b/src/paimon/core/append/append_compact_coordinator.cpp @@ -25,6 +25,7 @@ #include #include "paimon/common/data/binary_row.h" +#include "paimon/common/data/blob_utils.h" #include "paimon/common/types/data_field.h" #include "paimon/common/utils/linked_hash_map.h" #include "paimon/core/append/append_compact_task.h" @@ -198,7 +199,9 @@ Result, CoreOptions>> LoadSchemaAndOption /// Validate that the table is an append-only unaware-bucket table without DV. Status ValidateTable(const std::shared_ptr& table_schema, + const std::shared_ptr& arrow_schema, const CoreOptions& core_options) { + PAIMON_RETURN_NOT_OK(BlobUtils::ValidateMapBlobWriteSchema(arrow_schema)); if (!table_schema->PrimaryKeys().empty() || core_options.GetBucket() != -1) { return Status::Invalid( "AppendCompactCoordinator only supports append-only tables " @@ -320,12 +323,12 @@ Result>> AppendCompactCoordinator::Ru PAIMON_ASSIGN_OR_RAISE(schema_and_options, LoadSchemaAndOptions(table_path, options, file_system)); const auto& [table_schema, core_options] = schema_and_options; + auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(table_schema->Fields()); // Validate table type - PAIMON_RETURN_NOT_OK(ValidateTable(table_schema, core_options)); + PAIMON_RETURN_NOT_OK(ValidateTable(table_schema, arrow_schema, core_options)); // Build shared objects - auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(table_schema->Fields()); PAIMON_ASSIGN_OR_RAISE( std::shared_ptr partition_schema, FieldMapping::GetPartitionSchema(arrow_schema, table_schema->PartitionKeys())); diff --git a/src/paimon/core/append/append_compact_coordinator_test.cpp b/src/paimon/core/append/append_compact_coordinator_test.cpp index 1c8c7affa..bb166fe56 100644 --- a/src/paimon/core/append/append_compact_coordinator_test.cpp +++ b/src/paimon/core/append/append_compact_coordinator_test.cpp @@ -511,6 +511,31 @@ TEST_F(AppendCompactCoordinatorTest, TestValidateFailsOnDvTable) { "not support for dv in UNAWARE_BUCKET mode"); } +TEST_F(AppendCompactCoordinatorTest, TestValidateFailsOnLoadedMapBlobTable) { + auto fs = dir_->GetFileSystem(); + SchemaManager schema_manager(fs, TablePath()); + std::string schema_json = R"json({ + "version" : 3, + "id" : 0, + "fields" : [ { + "id" : 0, + "name" : "blob_map", + "type" : {"type":"MAP", "key":"STRING", "value":"BLOB"} + } ], + "highestFieldId" : 0, + "partitionKeys" : [], + "primaryKeys" : [], + "options" : {"bucket":"-1"}, + "timeMillis" : 1721614341162 + })json"; + ASSERT_OK(fs->AtomicStore(PathUtil::JoinPath(schema_manager.SchemaDirectory(), "schema-0"), + schema_json)); + + ASSERT_NOK_WITH_MSG( + AppendCompactCoordinator::Run(TablePath(), /*options=*/{}, /*partitions=*/{}, fs, pool_), + "Writing a table with MAP<..., BLOB> is not supported by the C++ writer"); +} + /// Test that compact output files are written to external path when configured. TEST_F(AppendCompactCoordinatorTest, TestCompactWithExternalPath) { auto external_dir = UniqueTestDirectory::Create("local"); diff --git a/src/paimon/core/operation/file_store_write.cpp b/src/paimon/core/operation/file_store_write.cpp index fa1294d1b..74f56dc90 100644 --- a/src/paimon/core/operation/file_store_write.cpp +++ b/src/paimon/core/operation/file_store_write.cpp @@ -23,6 +23,7 @@ #include #include "fmt/format.h" +#include "paimon/common/data/blob_utils.h" #include "paimon/common/types/data_field.h" #include "paimon/common/utils/fields_comparator.h" #include "paimon/core/core_options.h" @@ -109,6 +110,8 @@ Result> FileStoreWrite::Create(std::unique_ptrFields()); + PAIMON_RETURN_NOT_OK(BlobUtils::ValidateMapBlobWriteSchema(arrow_schema)); auto opts = schema->Options(); for (const auto& [key, value] : ctx->GetOptions()) { opts[key] = value; @@ -116,7 +119,6 @@ Result> FileStoreWrite::Create(std::unique_ptrGetSpecificFileSystem(), ctx->GetFileSystemSchemeToIdentifierMap())); - auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(schema->Fields()); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr partition_schema, FieldMapping::GetPartitionSchema(arrow_schema, schema->PartitionKeys())); diff --git a/src/paimon/core/operation/file_store_write_test.cpp b/src/paimon/core/operation/file_store_write_test.cpp index 3e1b3ea4d..23633488e 100644 --- a/src/paimon/core/operation/file_store_write_test.cpp +++ b/src/paimon/core/operation/file_store_write_test.cpp @@ -75,6 +75,34 @@ TEST(FileStoreWriteTest, TestCreateAppendTable) { FileStoreWrite::Create(std::move(write_context))); } +TEST(FileStoreWriteTest, TestCreateWriterForLoadedMapBlobTable) { + auto dir = UniqueTestDirectory::Create(); + std::string table_path = PathUtil::JoinPath(dir->Str(), "foo.db/bar"); + auto fs = std::make_shared(); + SchemaManager schema_manager(fs, table_path); + std::string schema_json = R"json({ + "version" : 3, + "id" : 0, + "fields" : [ { + "id" : 0, + "name" : "blob_map", + "type" : {"type":"MAP", "key":"STRING", "value":"BLOB"} + } ], + "highestFieldId" : 0, + "partitionKeys" : [], + "primaryKeys" : [], + "options" : {}, + "timeMillis" : 1721614341162 + })json"; + ASSERT_OK(fs->AtomicStore(PathUtil::JoinPath(schema_manager.SchemaDirectory(), "schema-0"), + schema_json)); + + WriteContextBuilder context_builder(table_path, "commit_user_1"); + ASSERT_OK_AND_ASSIGN(std::unique_ptr write_context, context_builder.Finish()); + ASSERT_NOK_WITH_MSG(FileStoreWrite::Create(std::move(write_context)), + "Writing a table with MAP<..., BLOB> is not supported by the C++ writer"); +} + TEST(FileStoreWriteTest, TestCreateAppendTableWithInvalidBucket) { auto dir = UniqueTestDirectory::Create(); arrow::FieldVector fields = { diff --git a/src/paimon/core/schema/arrow_schema_validator.cpp b/src/paimon/core/schema/arrow_schema_validator.cpp index 0db1145f7..009979365 100644 --- a/src/paimon/core/schema/arrow_schema_validator.cpp +++ b/src/paimon/core/schema/arrow_schema_validator.cpp @@ -164,8 +164,10 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId( const auto& item_field = checked_cast(type.get())->item_field(); PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId( key_field->type(), key_field->metadata(), /*allow_blob=*/false, field_id_set)); + bool allow_direct_blob_item = + allow_blob && item_field->type()->id() == arrow::Type::LARGE_BINARY; PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId( - item_field->type(), item_field->metadata(), /*allow_blob=*/false, field_id_set)); + item_field->type(), item_field->metadata(), allow_direct_blob_item, field_id_set)); break; } case arrow::Type::type::LARGE_BINARY: { @@ -248,7 +250,9 @@ Status ArrowSchemaValidator::ValidateField(const std::shared_ptr& const auto& item_field = checked_cast(*field->type()).item_field(); PAIMON_RETURN_NOT_OK(ValidateField(key_field, /*allow_blob=*/false)); - PAIMON_RETURN_NOT_OK(ValidateField(item_field, /*allow_blob=*/false)); + bool allow_direct_blob_item = + allow_blob && item_field->type()->id() == arrow::Type::LARGE_BINARY; + PAIMON_RETURN_NOT_OK(ValidateField(item_field, allow_direct_blob_item)); break; } case arrow::Type::type::LARGE_BINARY: { diff --git a/src/paimon/core/schema/arrow_schema_validator_test.cpp b/src/paimon/core/schema/arrow_schema_validator_test.cpp index ed56b6370..d644e5635 100644 --- a/src/paimon/core/schema/arrow_schema_validator_test.cpp +++ b/src/paimon/core/schema/arrow_schema_validator_test.cpp @@ -195,6 +195,16 @@ TEST(ArrowSchemaValidatorTest, TestBlobFieldMustBeTopLevel) { ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema), "Blob field must be a top-level field."); } + { + auto map_blob_field = arrow::field( + "map_blob", arrow::map(arrow::utf8(), BlobUtils::ToArrowField("value", true))); + auto arrow_schema = arrow::schema(arrow::FieldVector({map_blob_field})); + ASSERT_OK(ArrowSchemaValidator::ValidateSchema(*arrow_schema)); + + std::vector fields = {DataField(0, map_blob_field)}; + arrow_schema = DataField::ConvertDataFieldsToArrowSchema(fields); + ASSERT_OK(ArrowSchemaValidator::ValidateSchemaWithFieldId(*arrow_schema)); + } { auto map_blob_field = arrow::field( "map_blob", @@ -203,6 +213,15 @@ TEST(ArrowSchemaValidatorTest, TestBlobFieldMustBeTopLevel) { ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema), "Blob field must be a top-level field."); } + { + auto nested_map_blob_field = arrow::field( + "nested_map_blob", + arrow::map(arrow::utf8(), + 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."); + } { std::vector nested_fields = { DataField(1, BlobUtils::ToArrowField("blob", true))}; diff --git a/src/paimon/core/schema/schema_validation_test.cpp b/src/paimon/core/schema/schema_validation_test.cpp index 054f76da3..3e708f87a 100644 --- a/src/paimon/core/schema/schema_validation_test.cpp +++ b/src/paimon/core/schema/schema_validation_test.cpp @@ -25,6 +25,7 @@ #include "gtest/gtest.h" #include "paimon/common/data/blob_utils.h" #include "paimon/common/data/variant/variant_type_utils.h" +#include "paimon/common/utils/checked_cast.h" #include "paimon/core/schema/table_schema.h" #include "paimon/defs.h" #include "paimon/testing/utils/testharness.h" @@ -1159,8 +1160,6 @@ TEST(SchemaValidationTest, TestMapSharedShreddingRequiresNonNullableKey) { } TEST(SchemaValidationTest, TestMapSharedShreddingRejectsBlobValue) { - auto direct_blob_map = - arrow::map(arrow::utf8(), BlobUtils::ToArrowField("value", /*nullable=*/true)); auto nested_blob_map = arrow::map( arrow::utf8(), arrow::field("value", arrow::struct_({BlobUtils::ToArrowField("blob")}))); std::map options = { @@ -1169,15 +1168,38 @@ TEST(SchemaValidationTest, TestMapSharedShreddingRejectsBlobValue) { {"fields.f1.map.storage-layout", "shared-shredding"}, }; - for (const auto& map_type : {direct_blob_map, nested_blob_map}) { - auto schema = arrow::schema({ - arrow::field("f0", arrow::utf8()), - arrow::field("f1", map_type), - }); - ASSERT_NOK_WITH_MSG(TableSchema::Create(/*schema_id=*/0, schema, /*partition_keys=*/{}, - /*primary_keys=*/{}, options), - "Blob field must be a top-level field."); - } + const std::string loaded_schema = R"json({ + "version": 3, + "id": 0, + "fields": [ + {"id": 0, "name": "f0", "type": "STRING"}, + {"id": 1, "name": "f1", + "type": {"type": "MAP", "key": "STRING", "value": "BLOB"}} + ], + "highestFieldId": 1, + "partitionKeys": [], + "primaryKeys": [], + "options": { + "bucket": "1", + "bucket-key": "f0", + "fields.f1.map.storage-layout": "shared-shredding" + }, + "timeMillis": 0 + })json"; + ASSERT_OK_AND_ASSIGN(std::shared_ptr table_schema, + TableSchema::CreateFromJson(loaded_schema)); + auto loaded_map = checked_pointer_cast(table_schema->Fields()[1].Type()); + ASSERT_TRUE(BlobUtils::IsBlobField(loaded_map->item_field())); + ASSERT_NOK_WITH_MSG(SchemaValidation::ValidateTableSchema(*table_schema), + "MAP shared-shredding currently cannot contain BLOB fields."); + + auto nested_schema = arrow::schema({ + arrow::field("f0", arrow::utf8()), + arrow::field("f1", nested_blob_map), + }); + ASSERT_NOK_WITH_MSG(TableSchema::Create(/*schema_id=*/0, nested_schema, + /*partition_keys=*/{}, /*primary_keys=*/{}, options), + "Blob field must be a top-level field."); } TEST(SchemaValidationTest, TestMapSharedShreddingCompression) { diff --git a/src/paimon/core/schema/table_schema.cpp b/src/paimon/core/schema/table_schema.cpp index 6e2d8747e..9d2405601 100644 --- a/src/paimon/core/schema/table_schema.cpp +++ b/src/paimon/core/schema/table_schema.cpp @@ -27,6 +27,7 @@ #include "arrow/api.h" #include "arrow/c/bridge.h" #include "fmt/format.h" +#include "paimon/common/data/blob_utils.h" #include "paimon/common/data/variant/variant_type_utils.h" #include "paimon/common/utils/arrow/status_utils.h" #include "paimon/common/utils/checked_cast.h" @@ -56,6 +57,7 @@ Result> TableSchema::Create( for (const auto& primary_key : primary_keys) { primary_key_set.insert(primary_key); } + PAIMON_RETURN_NOT_OK(BlobUtils::ValidateMapBlobWriteSchema(schema)); for (const auto& field : schema->fields()) { PAIMON_ASSIGN_OR_RAISE(std::shared_ptr 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 0acbee0b2..e195802c0 100644 --- a/src/paimon/core/schema/table_schema_test.cpp +++ b/src/paimon/core/schema/table_schema_test.cpp @@ -23,6 +23,7 @@ #include "arrow/api.h" #include "gtest/gtest.h" +#include "paimon/common/data/blob_utils.h" #include "paimon/common/data/variant/variant_type_utils.h" #include "paimon/common/utils/checked_cast.h" #include "paimon/common/utils/date_time_utils.h" @@ -1331,6 +1332,40 @@ TEST_F(TableSchemaTest, NullableMapKeySchemaIsSupported) { ASSERT_TRUE(direct_map_type->key_field()->nullable()); } +TEST_F(TableSchemaTest, MapBlobSchemaLoadsFromJson) { + std::string table_schema_str = R"json({ + "version" : 3, + "id" : 0, + "fields" : [ { + "id" : 0, + "name" : "string_blob_map", + "type" : {"type":"MAP", "key":"STRING", "value":"BLOB"} + } ], + "highestFieldId" : 0, + "partitionKeys" : [], + "primaryKeys" : [], + "options" : {}, + "timeMillis" : 1721614341162 + })json"; + ASSERT_OK_AND_ASSIGN(std::unique_ptr table_schema, + TableSchema::CreateFromJson(table_schema_str)); + ASSERT_TRUE(BlobUtils::IsMapBlobField( + DataField::ConvertDataFieldToArrowField(table_schema->Fields()[0]))); + ASSERT_OK_AND_ASSIGN(std::string serialized, table_schema->ToJsonString()); + ASSERT_OK_AND_ASSIGN(std::unique_ptr restored, + TableSchema::CreateFromJson(serialized)); + ASSERT_OK_AND_ASSIGN(std::string restored_json, restored->ToJsonString()); + ASSERT_EQ(restored_json, serialized); +} + +TEST_F(TableSchemaTest, CreatingMapBlobSchemaIsRejected) { + auto map_type = arrow::map(arrow::utf8(), BlobUtils::ToArrowField("value", /*nullable=*/true)); + ASSERT_NOK_WITH_MSG( + TableSchema::Create(/*schema_id=*/0, arrow::schema({arrow::field("blob_map", map_type)}), + /*partition_keys=*/{}, /*primary_keys=*/{}, /*options=*/{}), + "not supported by the C++ writer"); +} + TEST_F(TableSchemaTest, MapKeysSortedIsNormalized) { auto sorted_map = std::make_shared(arrow::field("key", arrow::utf8(), /*nullable=*/false), diff --git a/src/paimon/core/utils/nested_projection_utils.cpp b/src/paimon/core/utils/nested_projection_utils.cpp index f956ca8ff..03712e93a 100644 --- a/src/paimon/core/utils/nested_projection_utils.cpp +++ b/src/paimon/core/utils/nested_projection_utils.cpp @@ -33,6 +33,7 @@ #include "arrow/compute/cast.h" #include "arrow/type.h" #include "fmt/format.h" +#include "paimon/common/data/blob_utils.h" #include "paimon/common/data/variant/variant_access_utils.h" #include "paimon/common/data/variant/variant_type_utils.h" #include "paimon/common/utils/checked_cast.h" @@ -434,6 +435,10 @@ Result> NestedProjectionUtils::GetMapSelectedKeys( if (!get_result.ok()) { return result; } + if (BlobUtils::IsMapBlobField(field)) { + return Status::NotImplemented( + "paimon.map.selected-keys is not supported for MAP<..., BLOB>"); + } auto tokens = StringUtils::Split(get_result.ValueUnsafe(), ",", /*ignore_empty=*/false); std::unordered_set deduplicated; deduplicated.reserve(tokens.size()); diff --git a/src/paimon/core/utils/nested_projection_utils_test.cpp b/src/paimon/core/utils/nested_projection_utils_test.cpp index 153d90e17..8adf7b242 100644 --- a/src/paimon/core/utils/nested_projection_utils_test.cpp +++ b/src/paimon/core/utils/nested_projection_utils_test.cpp @@ -29,6 +29,7 @@ #include "arrow/memory_pool.h" #include "arrow/type.h" #include "gtest/gtest.h" +#include "paimon/common/data/blob_utils.h" #include "paimon/common/data/variant/variant_access_utils.h" #include "paimon/common/data/variant/variant_type_utils.h" #include "paimon/common/types/data_field.h" @@ -548,6 +549,16 @@ TEST(NestedProjectionUtilsTest, GetMapSelectedKeysDuplicateKey) { "Duplicate selected key 'a'"); } +TEST(NestedProjectionUtilsTest, GetMapSelectedKeysRejectsMapBlob) { + auto metadata = arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a"}); + auto map_type = arrow::map(arrow::utf8(), BlobUtils::ToArrowField("value")); + auto field = arrow::field("m", map_type, /*nullable=*/true, metadata); + ASSERT_NOK_WITH_MSG(NestedProjectionUtils::GetMapSelectedKeys(field), + "paimon.map.selected-keys is not supported for MAP<..., BLOB>"); + ASSERT_NOK_WITH_MSG(NestedProjectionUtils::HasMapSelectedKeysRecursively(field), + "paimon.map.selected-keys is not supported for MAP<..., BLOB>"); +} + // ============== MapSharedShreddingAccessField ============== TEST(NestedProjectionUtilsTest, IsMapSharedShreddingAccessField) { diff --git a/src/paimon/format/blob/blob_file_batch_reader.cpp b/src/paimon/format/blob/blob_file_batch_reader.cpp index f5e3ef9ce..acbf152ee 100644 --- a/src/paimon/format/blob/blob_file_batch_reader.cpp +++ b/src/paimon/format/blob/blob_file_batch_reader.cpp @@ -19,14 +19,20 @@ #include "paimon/format/blob/blob_file_batch_reader.h" #include +#include #include #include +#include +#include #include "arrow/api.h" #include "arrow/array/builder_dict.h" #include "arrow/array/builder_nested.h" #include "arrow/c/bridge.h" #include "arrow/util/bit_util.h" +#include "arrow/util/decimal.h" +#include "arrow/util/ubsan.h" +#include "arrow/util/utf8.h" #include "fmt/format.h" #include "paimon/common/data/blob_utils.h" #include "paimon/common/io/offset_input_stream.h" @@ -35,10 +41,128 @@ #include "paimon/common/utils/arrow/status_utils.h" #include "paimon/common/utils/checked_cast.h" #include "paimon/common/utils/delta_varint_compressor.h" +#include "paimon/common/utils/math.h" #include "paimon/common/utils/stream_utils.h" #include "paimon/data/blob.h" namespace paimon::blob { +namespace { + +constexpr int32_t kMapBlobMagicNumber = 0x4D424342; +constexpr int8_t kMapBlobVersion = 1; +constexpr int32_t kMapBlobHeaderLength = 9; +constexpr int32_t kMapBlobIndexLengthsSize = 8; +constexpr int32_t kMapBlobMinPayloadLength = kMapBlobHeaderLength + kMapBlobIndexLengthsSize; + +template +T ReadLittleEndian(const uint8_t* data) { + return FromLittleEndian(arrow::util::SafeLoadAs(data)); +} + +Result GetMapBlobFixedKeyLength(const std::shared_ptr& key_type) { + switch (key_type->id()) { + case arrow::Type::BOOL: + case arrow::Type::INT8: + return 1; + case arrow::Type::INT16: + return 2; + case arrow::Type::INT32: + case arrow::Type::DATE32: + return 4; + case arrow::Type::INT64: + return 8; + case arrow::Type::DECIMAL128: { + const auto& decimal_type = static_cast(*key_type); + return decimal_type.precision() <= 18 ? 8 : -1; + } + case arrow::Type::STRING: + case arrow::Type::BINARY: + return -1; + default: + return Status::Invalid( + fmt::format("unsupported MAP<..., BLOB> key type: {}", key_type->ToString())); + } +} + +Status AppendMapBlobKey(const std::shared_ptr& key_type, const uint8_t* data, + int32_t length, arrow::ArrayBuilder* builder) { + switch (key_type->id()) { + case arrow::Type::BOOL: { + if (data[0] != 0 && data[0] != 1) { + return Status::Invalid("invalid MAP<..., BLOB> boolean key"); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW( + checked_cast(builder)->Append(data[0] == 1)); + return Status::OK(); + } + case arrow::Type::INT8: { + PAIMON_RETURN_NOT_OK_FROM_ARROW( + checked_cast(builder)->Append(static_cast(data[0]))); + return Status::OK(); + } + case arrow::Type::INT16: { + PAIMON_RETURN_NOT_OK_FROM_ARROW(checked_cast(builder)->Append( + ReadLittleEndian(data))); + return Status::OK(); + } + case arrow::Type::INT32: { + PAIMON_RETURN_NOT_OK_FROM_ARROW(checked_cast(builder)->Append( + ReadLittleEndian(data))); + return Status::OK(); + } + case arrow::Type::INT64: { + PAIMON_RETURN_NOT_OK_FROM_ARROW(checked_cast(builder)->Append( + ReadLittleEndian(data))); + return Status::OK(); + } + case arrow::Type::DATE32: { + PAIMON_RETURN_NOT_OK_FROM_ARROW(checked_cast(builder)->Append( + ReadLittleEndian(data))); + return Status::OK(); + } + case arrow::Type::STRING: { + if (!arrow::util::ValidateUTF8(data, length)) { + return Status::Invalid("invalid UTF-8 in MAP key"); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW( + checked_cast(builder)->Append(data, length)); + return Status::OK(); + } + case arrow::Type::BINARY: { + PAIMON_RETURN_NOT_OK_FROM_ARROW( + checked_cast(builder)->Append(data, length)); + return Status::OK(); + } + case arrow::Type::DECIMAL128: { + const auto& decimal_type = static_cast(*key_type); + arrow::Decimal128 value; + if (decimal_type.precision() <= 18) { + value = arrow::Decimal128(ReadLittleEndian(data)); + } else { + // Java BigInteger.toByteArray() uses the shortest big-endian two's-complement + // representation. This makes byte-wise duplicate detection equivalent to + // decoded decimal equality. + if (length <= 0 || (length > 1 && ((data[0] == 0x00 && (data[1] & 0x80) == 0) || + (data[0] == 0xFF && (data[1] & 0x80) != 0)))) { + return Status::Invalid("invalid MAP<..., BLOB> non-canonical decimal key"); + } + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(value, + arrow::Decimal128::FromBigEndian(data, length)); + } + if (!value.FitsInPrecision(decimal_type.precision())) { + return Status::Invalid("MAP<..., BLOB> decimal key exceeds declared precision"); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW( + checked_cast(builder)->Append(value)); + return Status::OK(); + } + default: + return Status::Invalid( + fmt::format("unsupported MAP<..., BLOB> key type: {}", key_type->ToString())); + } +} + +} // namespace Result> BlobFileBatchReader::Create( const std::shared_ptr& input_stream, int32_t batch_size, bool blob_as_descriptor, @@ -133,9 +257,15 @@ Status BlobFileBatchReader::SetReadSchema(::ArrowSchema* read_schema, return Status::Invalid( fmt::format("read schema field number {} is not 1", arrow_schema->num_fields())); } - if (!BlobUtils::IsBlobField(arrow_schema->field(0))) { + std::shared_ptr read_field = arrow_schema->field(0); + if (!BlobUtils::IsBlobField(read_field) && !BlobUtils::IsMapBlobField(read_field)) { return Status::Invalid( - fmt::format("field {} is not BLOB", arrow_schema->field(0)->ToString())); + fmt::format("field {} must be BLOB or MAP<..., BLOB>", read_field->ToString())); + } + if (BlobUtils::IsMapBlobField(read_field)) { + const auto& map_type = static_cast(*read_field->type()); + PAIMON_ASSIGN_OR_RAISE([[maybe_unused]] int32_t fixed_key_length, + GetMapBlobFixedKeyLength(map_type.key_type())); } if (selection_bitmap != std::nullopt) { int32_t cardinality = selection_bitmap->Cardinality(); @@ -160,7 +290,7 @@ Status BlobFileBatchReader::SetReadSchema(::ArrowSchema* read_schema, target_blob_offsets_ = new_offsets; target_blob_row_indexes_ = new_row_indexes; } - target_type_ = arrow::struct_(arrow_schema->fields()); + target_type_ = arrow::struct_({read_field}); current_pos_ = 0; previous_batch_start_pos_ = std::numeric_limits::max(); previous_batch_row_count_ = 0; @@ -252,9 +382,247 @@ Result> BlobFileBatchReader::BuildContentArray( return std::make_shared(struct_array_data); } +Result BlobFileBatchReader::ReadMapBlobPayload( + size_t row_index, int32_t fixed_key_length) const { + if (target_blob_lengths_[row_index] < 0) { + return Status::Invalid(fmt::format("unsupported MAP<..., 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 < kMapBlobMinPayloadLength) { + return Status::Invalid( + fmt::format("invalid MAP<..., BLOB> payload length: {}", payload_length)); + } + + std::array header; + PAIMON_RETURN_NOT_OK(ReadBlobContentAt(payload_offset, header.size(), header.data())); + const auto magic_number = ReadLittleEndian(header.data()); + if (magic_number != kMapBlobMagicNumber) { + return Status::Invalid( + fmt::format("invalid MAP<..., BLOB> payload magic number: {}", magic_number)); + } + const auto version = static_cast(header[4]); + if (version != kMapBlobVersion) { + return Status::NotImplemented( + fmt::format("unsupported MAP<..., BLOB> payload version: {}", version)); + } + const auto entry_count = ReadLittleEndian(header.data() + 5); + if (entry_count < 0) { + return Status::Invalid(fmt::format("invalid MAP<..., BLOB> entry count: {}", entry_count)); + } + + const int64_t index_lengths_offset = payload_offset + payload_length - kMapBlobIndexLengthsSize; + std::array index_lengths; + PAIMON_RETURN_NOT_OK( + ReadBlobContentAt(index_lengths_offset, index_lengths.size(), index_lengths.data())); + const auto key_index_length = ReadLittleEndian(index_lengths.data()); + const auto value_index_length = + ReadLittleEndian(index_lengths.data() + sizeof(int32_t)); + const int64_t maximum_indexes_length = payload_length - kMapBlobMinPayloadLength; + if (key_index_length < 0 || key_index_length > maximum_indexes_length) { + return Status::Invalid( + fmt::format("invalid MAP<..., BLOB> key index length: {}", key_index_length)); + } + if (value_index_length < 0 || value_index_length > maximum_indexes_length) { + return Status::Invalid( + fmt::format("invalid MAP<..., BLOB> value index length: {}", value_index_length)); + } + if (static_cast(key_index_length) + value_index_length > maximum_indexes_length) { + return Status::Invalid("MAP<..., BLOB> indexes exceed the payload length"); + } + if (entry_count > key_index_length || entry_count > value_index_length) { + return Status::Invalid("MAP<..., BLOB> entry count exceeds index length"); + } + + const int64_t value_index_offset = index_lengths_offset - value_index_length; + const int64_t key_index_offset = value_index_offset - key_index_length; + std::vector key_index_bytes(key_index_length); + std::vector value_index_bytes(value_index_length); + PAIMON_RETURN_NOT_OK(ReadBlobContentAt(key_index_offset, key_index_length, + reinterpret_cast(key_index_bytes.data()))); + PAIMON_RETURN_NOT_OK(ReadBlobContentAt(value_index_offset, value_index_length, + reinterpret_cast(value_index_bytes.data()))); + PAIMON_ASSIGN_OR_RAISE(std::vector key_lengths, + DeltaVarintCompressor::Decompress(key_index_bytes)); + PAIMON_ASSIGN_OR_RAISE(std::vector value_lengths, + DeltaVarintCompressor::Decompress(value_index_bytes)); + if (key_lengths.size() != static_cast(entry_count)) { + return Status::Invalid("MAP<..., BLOB> entry count does not match key index length"); + } + if (value_lengths.size() != static_cast(entry_count)) { + return Status::Invalid("MAP<..., BLOB> entry count does not match value index length"); + } + + const int64_t data_offset = payload_offset + kMapBlobHeaderLength; + const int64_t data_length = key_index_offset - data_offset; + int64_t key_data_length = 0; + for (int64_t key_length : key_lengths) { + if (key_length < 0) { + return Status::Invalid("MAP<..., BLOB> keys cannot be null"); + } + if (key_length > std::numeric_limits::max()) { + return Status::Invalid(fmt::format("MAP<..., BLOB> key is too large: {}", key_length)); + } + if (fixed_key_length >= 0 && key_length != fixed_key_length) { + return Status::Invalid( + fmt::format("invalid MAP<..., BLOB> fixed-width key length: {}", key_length)); + } + if (key_length > data_length - key_data_length) { + return Status::Invalid("MAP<..., BLOB> key lengths exceed the payload data length"); + } + key_data_length += key_length; + } + + const int64_t maximum_value_data_length = data_length - key_data_length; + int64_t value_data_length = 0; + for (int64_t value_length : value_lengths) { + if (value_length == BlobDefs::kNullBinLength) { + continue; + } + if (value_length < 0) { + return Status::Invalid( + fmt::format("invalid MAP<..., BLOB> value length: {}", value_length)); + } + if (!blob_as_descriptor_ && value_length > std::numeric_limits::max()) { + return Status::Invalid( + fmt::format("MAP<..., BLOB> inline value is too large: {}", value_length)); + } + if (value_length > maximum_value_data_length - value_data_length) { + return Status::Invalid("MAP<..., BLOB> value lengths exceed the payload data length"); + } + value_data_length += value_length; + } + if (value_data_length != maximum_value_data_length) { + return Status::Invalid( + "MAP<..., BLOB> key/value lengths do not match the payload data length"); + } + return MapBlobPayload{std::move(key_lengths), std::move(value_lengths), data_offset, + key_data_length}; +} + +Status BlobFileBatchReader::AppendMapBlobKeys(const MapBlobPayload& payload, + const std::shared_ptr& key_type, + arrow::ArrayBuilder* key_builder) const { + int64_t key_offset = payload.data_offset; + std::set serialized_keys; + for (int64_t key_length_64 : payload.key_lengths) { + const auto key_length = static_cast(key_length_64); + PAIMON_UNIQUE_PTR key_bytes = + Bytes::AllocateBytes(static_cast(key_length), pool_.get()); + if (key_length > 0) { + PAIMON_RETURN_NOT_OK(ReadBlobContentAt(key_offset, key_length, + reinterpret_cast(key_bytes->data()))); + } + std::string serialized_key; + if (key_length > 0) { + serialized_key.assign(key_bytes->data(), key_length); + } + if (!serialized_keys.emplace(std::move(serialized_key)).second) { + return Status::Invalid("invalid MAP<..., BLOB> payload: duplicate key"); + } + const uint8_t empty_key = 0; + const uint8_t* key_data = + key_length == 0 ? &empty_key : reinterpret_cast(key_bytes->data()); + PAIMON_RETURN_NOT_OK(AppendMapBlobKey(key_type, key_data, key_length, key_builder)); + key_offset += key_length; + } + return Status::OK(); +} + +Status BlobFileBatchReader::AppendMapBlobValues(const MapBlobPayload& payload, + arrow::LargeBinaryBuilder* blob_builder) const { + int64_t value_offset = payload.data_offset + payload.key_data_length; + for (int64_t value_length : payload.value_lengths) { + if (value_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::FromPath(file_path_, value_offset, value_length)); + PAIMON_UNIQUE_PTR descriptor = blob->ToDescriptor(pool_); + PAIMON_RETURN_NOT_OK_FROM_ARROW( + blob_builder->Append(descriptor->data(), descriptor->size())); + } else { + PAIMON_UNIQUE_PTR value_bytes = + Bytes::AllocateBytes(static_cast(value_length), pool_.get()); + if (value_length > 0) { + PAIMON_RETURN_NOT_OK(ReadBlobContentAt( + value_offset, value_length, reinterpret_cast(value_bytes->data()))); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW( + blob_builder->Append(value_bytes->data(), value_length)); + } + value_offset += value_length; + } + return Status::OK(); +} + +Result> BlobFileBatchReader::BuildMapBlobArray( + int32_t rows_to_read) const { + const auto& struct_type = static_cast(*target_type_); + const std::shared_ptr& map_field = struct_type.field(0); + auto map_type = checked_pointer_cast(map_field->type()); + const std::shared_ptr& key_type = map_type->key_type(); + if (key_type->id() == arrow::Type::STRING) { + arrow::util::InitializeUTF8(); + } + PAIMON_ASSIGN_OR_RAISE(int32_t fixed_key_length, GetMapBlobFixedKeyLength(key_type)); + + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::unique_ptr key_builder_unique, + arrow::MakeBuilder(key_type, arrow_pool_.get())); + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::unique_ptr item_builder_unique, + arrow::MakeBuilder(map_type->item_type(), arrow_pool_.get())); + std::shared_ptr key_builder(std::move(key_builder_unique)); + std::shared_ptr item_builder(std::move(item_builder_unique)); + if (!item_builder || !item_builder->type() || + item_builder->type()->id() != arrow::Type::LARGE_BINARY) { + return Status::Invalid("cast MAP<..., BLOB> item builder to large binary builder failed"); + } + auto* blob_builder = checked_cast(item_builder.get()); + arrow::MapBuilder map_builder(arrow_pool_.get(), key_builder, item_builder, map_type); + + for (int32_t k = 0; k < rows_to_read; ++k) { + const size_t row_index = current_pos_ + k; + if (IsTargetNull(row_index)) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(map_builder.AppendNull()); + continue; + } + if (IsTargetPlaceholder(row_index)) { + // Duplicate map keys cannot occur in a valid Paimon map, so two empty/default keys + // with null values form an unambiguous, Arrow-valid internal sentinel. + PAIMON_RETURN_NOT_OK_FROM_ARROW(map_builder.Append()); + PAIMON_RETURN_NOT_OK_FROM_ARROW(key_builder->AppendEmptyValues(2)); + PAIMON_RETURN_NOT_OK_FROM_ARROW(blob_builder->AppendNulls(2)); + continue; + } + + PAIMON_ASSIGN_OR_RAISE(MapBlobPayload payload, + ReadMapBlobPayload(row_index, fixed_key_length)); + PAIMON_RETURN_NOT_OK_FROM_ARROW(map_builder.Append()); + PAIMON_RETURN_NOT_OK(AppendMapBlobKeys(payload, key_type, key_builder.get())); + PAIMON_RETURN_NOT_OK(AppendMapBlobValues(payload, blob_builder)); + } + + std::shared_ptr built_map_array; + PAIMON_RETURN_NOT_OK_FROM_ARROW(map_builder.Finish(&built_map_array)); + auto map_array = std::make_shared( + map_type, built_map_array->length(), built_map_array->value_offsets(), + built_map_array->keys(), built_map_array->items(), built_map_array->null_bitmap(), + built_map_array->null_count(), built_map_array->offset()); + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr struct_array, + arrow::StructArray::Make({map_array}, {map_field})); + return struct_array; +} + Result> BlobFileBatchReader::BuildTargetArray( int32_t rows_to_read) const { - std::shared_ptr blob_array; + const auto& struct_type = static_cast(*target_type_); + if (struct_type.field(0)->type()->id() == arrow::Type::MAP) { + return BuildMapBlobArray(rows_to_read); + } if (!blob_as_descriptor_) { return BuildContentArray(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 9a05c15b1..b1512f790 100644 --- a/src/paimon/format/blob/blob_file_batch_reader.h +++ b/src/paimon/format/blob/blob_file_batch_reader.h @@ -36,6 +36,11 @@ #include "paimon/result.h" #include "paimon/utils/roaring_bitmap32.h" +namespace arrow { +class ArrayBuilder; +class LargeBinaryBuilder; +} // namespace arrow + namespace paimon::blob { /// Binary Blob File Layout Specification @@ -99,9 +104,8 @@ 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 that path to merge away. Stored values - /// are returned verbatim; see BlobDefs::kPlaceholderSentinel for the accepted collision - /// with a user value exactly equal to the sentinel. + /// 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. static Result> Create( const std::shared_ptr& input_stream, int32_t batch_size, bool blob_as_descriptor, bool emit_placeholder_sentinel, @@ -150,6 +154,13 @@ class BlobFileBatchReader : public FileBatchReader { } private: + struct MapBlobPayload { + std::vector key_lengths; + std::vector value_lengths; + int64_t data_offset; + int64_t key_data_length; + }; + static constexpr uint64_t kDefaultReadChunkSize = 1024 * 1024; static int32_t GetIndexLength(const int8_t* bytes, int32_t offset); @@ -168,6 +179,13 @@ class BlobFileBatchReader : public FileBatchReader { /// Builds a null bitmap buffer for the given rows. Returns nullptr if no nulls. Result> BuildNullBitmap(int32_t rows_to_read) const; Result> BuildContentArray(int32_t rows_to_read) const; + Result ReadMapBlobPayload(size_t row_index, int32_t fixed_key_length) const; + Status AppendMapBlobKeys(const MapBlobPayload& payload, + const std::shared_ptr& key_type, + arrow::ArrayBuilder* key_builder) const; + Status AppendMapBlobValues(const MapBlobPayload& payload, + arrow::LargeBinaryBuilder* blob_builder) const; + Result> BuildMapBlobArray(int32_t rows_to_read) const; Result> BuildTargetArray(int32_t rows_to_read) const; /// Returns true if the blob at the given index is null (bin_length == kNullBinLength). 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 17482af77..beefdd8f2 100644 --- a/src/paimon/format/blob/blob_file_batch_reader_test.cpp +++ b/src/paimon/format/blob/blob_file_batch_reader_test.cpp @@ -18,10 +18,17 @@ #include "paimon/format/blob/blob_file_batch_reader.h" +#include +#include + #include "arrow/api.h" +#include "arrow/c/bridge.h" #include "arrow/c/helpers.h" +#include "arrow/ipc/json_simple.h" #include "gtest/gtest.h" +#include "paimon/common/data/blob_defs.h" #include "paimon/common/data/blob_utils.h" +#include "paimon/common/reader/blob_fallback_batch_reader.h" #include "paimon/common/utils/arrow/mem_utils.h" #include "paimon/data/blob.h" #include "paimon/format/blob/blob_format_writer.h" @@ -32,6 +39,32 @@ #include "paimon/testing/utils/testharness.h" namespace paimon::blob::test { +namespace { + +std::string HexToBytes(std::string_view hex) { + auto hex_value = [](char c) -> uint8_t { + return c <= '9' ? static_cast(c - '0') : static_cast(c - 'a' + 10); + }; + std::string bytes; + bytes.reserve(hex.size() / 2); + for (size_t i = 0; i < hex.size(); i += 2) { + bytes.push_back(static_cast((hex_value(hex[i]) << 4) | hex_value(hex[i + 1]))); + } + return bytes; +} + +std::string MapBlobGoldenBytes() { + // Java-compatible MAP golden file. Rows are: + // {alpha: "hello", empty: "", missing: null}, null, {}, {omega: "world"}. + return HexToBytes( + "cf114e584243424d0103000000616c706861656d7074796d697373696e676865" + "6c6c6f0a00040a090103000000030000003d000000000000002a64aaabcf114e" + "584243424d0100000000000000000000000021000000000000008360591ecf11" + "4e584243424d01010000006f6d656761776f726c640a0a01000000010000002d" + "00000000000000248fe4237a7b44180400000001"); +} + +} // namespace TEST(BlobReaderBuilderTest, RejectsNullMemoryPool) { BlobReaderBuilder builder(/*batch_size=*/10, /*options=*/{}); @@ -108,6 +141,44 @@ class BlobFileBatchReaderTest : public testing::Test, public ::testing::WithPara } } + Result ReadMapBlobValue(const std::shared_ptr& blob_array, + int64_t index, bool blob_as_descriptor, + const std::shared_ptr& file_system) { + std::string stored_value = blob_array->GetString(index); + if (!blob_as_descriptor) { + return stored_value; + } + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr blob, + Blob::FromDescriptor(stored_value.data(), stored_value.size())); + PAIMON_ASSIGN_OR_RAISE(PAIMON_UNIQUE_PTR value, blob->ToData(file_system, pool_)); + if (value->size() == 0) { + return std::string(); + } + return std::string(value->data(), value->size()); + } + + void CheckMapBlobReadFails(const std::string& file_bytes, + const std::shared_ptr& key_type, + const std::string& expected_message) { + auto dir = paimon::test::UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + const std::string file_path = dir->Str() + "/corrupt-map.blob"; + std::shared_ptr file_system = std::make_shared(); + ASSERT_OK(file_system->WriteFile(file_path, file_bytes, /*overwrite=*/true)); + + auto map_type = arrow::map(key_type, BlobUtils::ToArrowField("value", /*nullable=*/true)); + auto schema = arrow::schema({arrow::field("blob_map", map_type)}); + ::ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportSchema(*schema, &c_schema).ok()); + ASSERT_OK_AND_ASSIGN(std::shared_ptr input, file_system->Open(file_path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + BlobFileBatchReader::Create( + input, /*batch_size=*/1, /*blob_as_descriptor=*/false, + /*emit_placeholder_sentinel=*/false, 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 pool_; @@ -132,6 +203,284 @@ TEST_P(BlobFileBatchReaderTest, TestSimple) { {"blob_9_f54d253c.bin"}, blob_as_descriptor); } +TEST_P(BlobFileBatchReaderTest, TestMapBlob) { + auto dir = paimon::test::UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + const std::string file_path = dir->Str() + "/map-blob.blob"; + std::shared_ptr file_system = std::make_shared(); + + const std::string file_bytes = MapBlobGoldenBytes(); + ASSERT_OK(file_system->WriteFile(file_path, file_bytes, /*overwrite=*/true)); + + std::shared_ptr blob_item = BlobUtils::ToArrowField("value", true); + auto key_field = arrow::field("key", arrow::utf8(), false); + auto map_type = std::make_shared(key_field, blob_item); + ASSERT_TRUE(BlobUtils::IsBlobField(map_type->item_field())); + auto map_field = arrow::field("blob_map", map_type, true); + auto schema = arrow::schema({map_field}); + ::ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportSchema(*schema, &c_schema).ok()); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr input_stream, file_system->Open(file_path)); + const bool blob_as_descriptor = GetParam(); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + BlobFileBatchReader::Create( + input_stream, /*batch_size=*/2, blob_as_descriptor, + /*emit_placeholder_sentinel=*/false, pool_, GetArrowPool(pool_))); + ASSERT_OK(reader->SetReadSchema(&c_schema, nullptr, std::nullopt)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr chunked_array, + paimon::test::ReadResultCollector::CollectResult(std::move(reader))); + std::shared_ptr combined_array = + arrow::Concatenate(chunked_array->chunks()).ValueOrDie(); + + auto struct_array = std::dynamic_pointer_cast(combined_array); + ASSERT_TRUE(struct_array); + auto map_array = std::dynamic_pointer_cast(struct_array->field(0)); + ASSERT_TRUE(map_array); + auto values = std::dynamic_pointer_cast(map_array->items()); + 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; + } + ASSERT_OK_AND_ASSIGN(std::string value, + ReadMapBlobValue(values, i, blob_as_descriptor, file_system)); + ASSERT_TRUE(normalized_values_builder.Append(value).ok()); + } + std::shared_ptr normalized_values; + ASSERT_TRUE(normalized_values_builder.Finish(&normalized_values).ok()); + auto normalized_map = std::make_shared( + map_type, map_array->length(), map_array->value_offsets(), map_array->keys(), + normalized_values, map_array->null_bitmap(), map_array->null_count(), map_array->offset()); + std::shared_ptr expected = arrow::ipc::internal::json::ArrayFromJSON(map_type, + R"json([ + [["alpha", "hello"], ["empty", ""], ["missing", null]], + null, + [], + [["omega", "world"]] + ])json") + .ValueOrDie(); + ASSERT_TRUE(expected->Equals(normalized_map)) + << "expected: " << expected->ToString() << "\nactual: " << normalized_map->ToString(); +} + +TEST_P(BlobFileBatchReaderTest, MapBlobFallbackAcrossSequenceLayers) { + auto dir = paimon::test::UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + std::shared_ptr file_system = std::make_shared(); + const std::string old_file_path = dir->Str() + "/old-map.blob"; + const std::string new_file_path = dir->Str() + "/new-placeholder.blob"; + + const std::string old_bytes = MapBlobGoldenBytes(); + ASSERT_OK(file_system->WriteFile(old_file_path, old_bytes, /*overwrite=*/true)); + + // Generate four genuine -2 outer-file entries through the scalar writer. The outer blob + // index is type-independent; the map reader turns them into its map placeholder sentinel. + std::shared_ptr scalar_blob_field = BlobUtils::ToArrowField("blob_map", true); + auto scalar_struct_type = arrow::struct_({scalar_blob_field}); + ASSERT_OK_AND_ASSIGN(std::shared_ptr new_output, + file_system->Create(new_file_path, /*overwrite=*/true)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr writer, + BlobFormatWriter::Create(new_output, scalar_struct_type, + /*write_null_on_missing_file=*/false, + /*write_null_on_fetch_failure=*/false, + /*write_placeholder=*/true, file_system, pool_)); + arrow::LargeBinaryBuilder scalar_builder; + const std::string sentinel(BlobDefs::PlaceholderSentinelView()); + for (int32_t i = 0; i < 4; i++) { + ASSERT_TRUE(scalar_builder.Append(sentinel).ok()); + } + std::shared_ptr scalar_values; + ASSERT_TRUE(scalar_builder.Finish(&scalar_values).ok()); + std::shared_ptr scalar_rows = + arrow::StructArray::Make({scalar_values}, {scalar_blob_field}).ValueOrDie(); + for (int32_t i = 0; i < 4; i++) { + ::ArrowArray c_array; + ASSERT_TRUE(arrow::ExportArray(*scalar_rows->Slice(i, 1), &c_array).ok()); + ASSERT_OK(writer->AddBatch(&c_array)); + } + ASSERT_OK(writer->Finish()); + ASSERT_OK(new_output->Close()); + + auto map_type = std::make_shared(arrow::field("key", arrow::utf8(), false), + BlobUtils::ToArrowField("value", true)); + std::shared_ptr map_field = arrow::field("blob_map", map_type, true); + auto map_schema = arrow::schema({map_field}); + const bool blob_as_descriptor = GetParam(); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr new_input, file_system->Open(new_file_path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr new_reader, + BlobFileBatchReader::Create( + new_input, /*batch_size=*/2, blob_as_descriptor, + /*emit_placeholder_sentinel=*/true, pool_, GetArrowPool(pool_))); + ::ArrowSchema new_schema; + ASSERT_TRUE(arrow::ExportSchema(*map_schema, &new_schema).ok()); + ASSERT_OK(new_reader->SetReadSchema(&new_schema, nullptr, std::nullopt)); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr old_input, file_system->Open(old_file_path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr old_reader, + BlobFileBatchReader::Create( + old_input, /*batch_size=*/2, blob_as_descriptor, + /*emit_placeholder_sentinel=*/true, pool_, GetArrowPool(pool_))); + ::ArrowSchema old_schema; + ASSERT_TRUE(arrow::ExportSchema(*map_schema, &old_schema).ok()); + ASSERT_OK(old_reader->SetReadSchema(&old_schema, nullptr, std::nullopt)); + + std::vector> groups(2); + groups[0].push_back({std::move(new_reader), {}}); + groups[1].push_back({std::move(old_reader), {}}); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr fallback, + BlobFallbackBatchReader::Create(std::move(groups), map_schema, /*read_batch_size=*/2, + GetArrowPool(pool_))); + ASSERT_OK_AND_ASSIGN(std::shared_ptr result, + paimon::test::ReadResultCollector::CollectResult(std::move(fallback))); + std::shared_ptr combined = arrow::Concatenate(result->chunks()).ValueOrDie(); + auto struct_array = std::dynamic_pointer_cast(combined); + ASSERT_TRUE(struct_array); + auto map_array = std::dynamic_pointer_cast(struct_array->field(0)); + ASSERT_TRUE(map_array); + ASSERT_EQ(4, map_array->length()); + ASSERT_EQ(3, map_array->value_length(0)); + ASSERT_TRUE(map_array->IsNull(1)); + ASSERT_EQ(0, map_array->value_length(2)); + ASSERT_EQ(1, map_array->value_length(3)); + auto keys = std::dynamic_pointer_cast(map_array->keys()); + auto values = std::dynamic_pointer_cast(map_array->items()); + 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)); + ASSERT_OK_AND_ASSIGN(std::string last_value, + ReadMapBlobValue(values, 3, blob_as_descriptor, file_system)); + ASSERT_EQ("hello", first_value); + ASSERT_EQ("world", last_value); +} + +TEST_F(BlobFileBatchReaderTest, RejectsNonCanonicalDecimalMapKey) { + auto dir = paimon::test::UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + const std::string file_path = dir->Str() + "/bad-decimal-map.blob"; + std::shared_ptr file_system = std::make_shared(); + // The two raw keys 00 and 0000 both decode to decimal zero. The second is not the shortest + // Java BigInteger two's-complement representation and must not bypass duplicate detection. + const std::string file_bytes = HexToBytes( + "cf114e584243424d010200000000000002020000020000000200000028000000" + "0000000000000000d0000200000001"); + ASSERT_OK(file_system->WriteFile(file_path, file_bytes, /*overwrite=*/true)); + + auto map_type = + std::make_shared(arrow::field("key", arrow::decimal128(20, 0), false), + BlobUtils::ToArrowField("value", true)); + auto schema = arrow::schema({arrow::field("blob_map", map_type, true)}); + ::ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportSchema(*schema, &c_schema).ok()); + ASSERT_OK_AND_ASSIGN(std::shared_ptr input, file_system->Open(file_path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + BlobFileBatchReader::Create( + input, /*batch_size=*/1, /*blob_as_descriptor=*/false, + /*emit_placeholder_sentinel=*/false, pool_, GetArrowPool(pool_))); + ASSERT_OK(reader->SetReadSchema(&c_schema, nullptr, std::nullopt)); + ASSERT_NOK_WITH_MSG(reader->NextBatch(), "non-canonical decimal key"); +} + +TEST_F(BlobFileBatchReaderTest, RejectsInvalidUtf8StringMapKey) { + auto dir = paimon::test::UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + const std::string file_path = dir->Str() + "/bad-string-map.blob"; + std::shared_ptr file_system = std::make_shared(); + std::string file_bytes = MapBlobGoldenBytes(); + ASSERT_GT(file_bytes.size(), 13); + // The first key starts after the outer magic and the map header. + file_bytes[13] = static_cast(0xFF); + ASSERT_OK(file_system->WriteFile(file_path, file_bytes, /*overwrite=*/true)); + + auto map_type = arrow::map(arrow::utf8(), BlobUtils::ToArrowField("value", /*nullable=*/true)); + auto schema = arrow::schema({arrow::field("blob_map", map_type)}); + ::ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportSchema(*schema, &c_schema).ok()); + ASSERT_OK_AND_ASSIGN(std::shared_ptr input, file_system->Open(file_path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + BlobFileBatchReader::Create( + input, /*batch_size=*/1, /*blob_as_descriptor=*/false, + /*emit_placeholder_sentinel=*/false, pool_, GetArrowPool(pool_))); + ASSERT_OK(reader->SetReadSchema(&c_schema, nullptr, std::nullopt)); + ASSERT_NOK_WITH_MSG(reader->NextBatch(), "invalid UTF-8"); +} + +TEST_F(BlobFileBatchReaderTest, RejectsNullMapKey) { + auto dir = paimon::test::UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + const std::string file_path = dir->Str() + "/null-key-map.blob"; + std::shared_ptr file_system = std::make_shared(); + std::string file_bytes = MapBlobGoldenBytes(); + const std::string key_index = HexToBytes("0a0004"); + const size_t key_index_offset = file_bytes.find(key_index); + ASSERT_NE(std::string::npos, key_index_offset); + // The first delta-varint changes from key length 5 to Java's null marker -1. + file_bytes[key_index_offset] = static_cast(0x01); + ASSERT_OK(file_system->WriteFile(file_path, file_bytes, /*overwrite=*/true)); + + auto map_type = arrow::map(arrow::utf8(), BlobUtils::ToArrowField("value", /*nullable=*/true)); + auto schema = arrow::schema({arrow::field("blob_map", map_type)}); + ::ArrowSchema c_schema; + ASSERT_TRUE(arrow::ExportSchema(*schema, &c_schema).ok()); + ASSERT_OK_AND_ASSIGN(std::shared_ptr input, file_system->Open(file_path)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + BlobFileBatchReader::Create( + input, /*batch_size=*/1, /*blob_as_descriptor=*/false, + /*emit_placeholder_sentinel=*/false, pool_, GetArrowPool(pool_))); + ASSERT_OK(reader->SetReadSchema(&c_schema, nullptr, std::nullopt)); + ASSERT_NOK_WITH_MSG(reader->NextBatch(), "MAP<..., BLOB> keys cannot be null"); +} + +TEST_F(BlobFileBatchReaderTest, RejectsCorruptMapPayloadMetadata) { + // The Java golden's first bin starts with the outer BLOB magic (4 bytes), followed by a + // 45-byte MAP payload: header [4, 13), key/value data [13, 35), indexes [35, 41), and the + // two index lengths [41, 49). Mutating the payload does not require updating the outer CRC, + // which BlobFileBatchReader intentionally does not validate. + const std::string golden = MapBlobGoldenBytes(); + + std::string corrupted = golden; + corrupted[4] = 0; + CheckMapBlobReadFails(corrupted, arrow::utf8(), "invalid MAP<..., BLOB> payload magic number"); + + corrupted = golden; + corrupted[8] = 2; + CheckMapBlobReadFails(corrupted, arrow::utf8(), "unsupported MAP<..., BLOB> payload version"); + + corrupted = golden; + std::fill(corrupted.begin() + 9, corrupted.begin() + 13, static_cast(0xFF)); + CheckMapBlobReadFails(corrupted, arrow::utf8(), "invalid MAP<..., BLOB> entry count"); + + corrupted = golden; + corrupted[41] = static_cast(0xFF); + CheckMapBlobReadFails(corrupted, arrow::utf8(), "invalid MAP<..., BLOB> key index length"); + + corrupted = golden; + corrupted[9] = 2; + CheckMapBlobReadFails(corrupted, arrow::utf8(), "entry count does not match key index length"); + + CheckMapBlobReadFails(golden, arrow::int32(), "invalid MAP<..., BLOB> fixed-width key length"); + + corrupted = golden; + corrupted[38] = 0x0C; + corrupted[39] = 0x0B; + CheckMapBlobReadFails(corrupted, arrow::utf8(), "value lengths exceed the payload data length"); + + corrupted = golden; + corrupted[38] = 0x08; + corrupted[39] = 0x07; + CheckMapBlobReadFails(corrupted, arrow::utf8(), + "key/value lengths do not match the payload data length"); + + corrupted = golden; + std::copy_n(corrupted.begin() + 13, 5, corrupted.begin() + 18); + CheckMapBlobReadFails(corrupted, arrow::utf8(), "payload: duplicate key"); +} + 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(); @@ -391,7 +740,7 @@ 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 is not BLOB"); + "field my_blob_field: large_binary must be BLOB or MAP<..., BLOB>"); } { auto schema = arrow::schema({BlobUtils::ToArrowField("my_blob_field", false)}); diff --git a/test/inte/blob_table_inte_test.cpp b/test/inte/blob_table_inte_test.cpp index 4b76cef16..4419b3567 100644 --- a/test/inte/blob_table_inte_test.cpp +++ b/test/inte/blob_table_inte_test.cpp @@ -535,6 +535,50 @@ class BlobTableInteTest : public testing::Test, public ::testing::WithParamInter }); } + Result> NormalizeMapBlobValues( + const std::shared_ptr& map_array, bool blob_as_descriptor) const { + const auto& values = checked_cast(*map_array->items()); + auto fs = std::make_shared(); + arrow::LargeBinaryBuilder builder; + for (int64_t i = 0; i < values.length(); ++i) { + if (values.IsNull(i)) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.AppendNull()); + 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::FromDescriptor(stored.data(), static_cast(stored.size()))); + PAIMON_ASSIGN_OR_RAISE(PAIMON_UNIQUE_PTR data, blob->ToData(fs, pool_)); + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.Append(data->data(), data->size())); + } + std::shared_ptr normalized_values; + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.Finish(&normalized_values)); + return std::make_shared(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()); + } + + void CheckMapBlobColumn(const std::shared_ptr& rows, + const std::string& field_name, const std::string& expected_json, + bool blob_as_descriptor) const { + auto map_array = + std::dynamic_pointer_cast(rows->GetFieldByName(field_name)); + ASSERT_TRUE(map_array) << field_name; + ASSERT_OK_AND_ASSIGN(auto normalized, + NormalizeMapBlobValues(map_array, blob_as_descriptor)); + auto expected = arrow::ipc::internal::json::ArrayFromJSON(map_array->type(), expected_json) + .ValueOrDie(); + ASSERT_TRUE(expected->Equals(normalized)) + << field_name << " expected: " << expected->ToString() + << " actual: " << normalized->ToString(); + } + /// Verify DataFileMeta properties from a scan plan. /// Each vector element corresponds to one expected DataFileMeta (ordered by file index). static void VerifyDataFileMetas( @@ -4334,4 +4378,111 @@ TEST_P(BlobTableInteTest, TestReadBlobDescriptorFieldFromJava) { ASSERT_TRUE(resolved->Equals(expected_with_rk)); } +TEST_P(BlobTableInteTest, TestReadMapBlobTableFromJava) { + if (GetParam() != "parquet") { + GTEST_SKIP() << "the Java fixture uses Parquet"; + } + const std::string table_path = GetDataDir() + "/parquet/map_blob_java.db/map_blob_java"; + const std::vector read_fields = {"id", + "string_payloads", + "boolean_payloads", + "tinyint_payloads", + "smallint_payloads", + "int_payloads", + "bigint_payloads", + "date_payloads", + "binary_payloads", + "compact_decimal_payloads", + "large_decimal_payloads"}; + + for (int64_t snapshot_id : {1, 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 string_layer_count = 0; + for (const auto& split : plan->Splits()) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + for (const auto& file : data_split->DataFiles()) { + if (file->write_cols == + std::optional>({"string_payloads"})) { + ++string_layer_count; + } + } + } + ASSERT_EQ(snapshot_id == 1 ? 1 : 3, string_layer_count); + + for (bool blob_as_descriptor : {false, true}) { + std::map 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(combined); + ASSERT_TRUE(rows); + ASSERT_EQ(4, rows->length()); + const auto& ids = checked_cast(*rows->GetFieldByName("id")); + for (int64_t i = 0; i < ids.length(); ++i) { + ASSERT_EQ(i + 1, ids.Value(i)); + } + + CheckMapBlobColumn( + rows, "string_payloads", + R"json([[["", "string-empty"], ["alpha", "string-alpha"]], [], null, [["omega", "string-omega"]]])json", + blob_as_descriptor); + CheckMapBlobColumn( + rows, "boolean_payloads", + R"json([[[false, "bool-false"], [true, "bool-true"]], null, null, null])json", + blob_as_descriptor); + CheckMapBlobColumn( + rows, "tinyint_payloads", + R"json([[[-128, "tiny-min"], [-1, "tiny-negative"], [127, "tiny-max"]], null, null, null])json", + blob_as_descriptor); + CheckMapBlobColumn( + rows, "smallint_payloads", + R"json([[[-32768, "small-min"], [-1, "small-negative"], [32767, "small-max"]], null, null, null])json", + blob_as_descriptor); + CheckMapBlobColumn( + rows, "int_payloads", + R"json([[[-2147483648, "int-min"], [-1, "int-negative"], [2147483647, "int-max"]], null, null, null])json", + blob_as_descriptor); + CheckMapBlobColumn( + rows, "bigint_payloads", + R"json([[[-9223372036854775808, "big-min"], [-1, "big-negative"], [9223372036854775807, "big-max"]], null, null, null])json", + blob_as_descriptor); + CheckMapBlobColumn( + rows, "date_payloads", + R"json([[[-1, "date-negative"], [0, "date-epoch"]], null, null, null])json", + blob_as_descriptor); + CheckMapBlobColumn( + rows, "compact_decimal_payloads", + R"json([[["-99999999.99", "compact-negative"], ["99999999.99", "compact-positive"]], null, null, null])json", + blob_as_descriptor); + CheckMapBlobColumn( + rows, "large_decimal_payloads", + R"json([[["-999999999999999999.99", "large-negative"], ["999999999999999999.99", "large-positive"]], null, null, null])json", + blob_as_descriptor); + + auto binary_map = + std::dynamic_pointer_cast(rows->GetFieldByName("binary_payloads")); + ASSERT_TRUE(binary_map); + ASSERT_OK_AND_ASSIGN(auto normalized_binary, + NormalizeMapBlobValues(binary_map, blob_as_descriptor)); + const auto& binary_keys = + checked_cast(*normalized_binary->keys()); + const auto& binary_values = + checked_cast(*normalized_binary->items()); + ASSERT_EQ(2, normalized_binary->value_length(0)); + ASSERT_EQ("", binary_keys.GetString(0)); + ASSERT_EQ(std::string("\0\xff\1\2", 4), binary_keys.GetString(1)); + ASSERT_EQ("binary-empty", binary_values.GetString(0)); + ASSERT_EQ("binary-bytes", binary_values.GetString(1)); + } + } +} + } // namespace paimon::test diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/README.md b/test/test_data/parquet/map_blob_java.db/map_blob_java/README.md new file mode 100644 index 000000000..402a48bbc --- /dev/null +++ b/test/test_data/parquet/map_blob_java.db/map_blob_java/README.md @@ -0,0 +1,48 @@ +Generated by Apache Paimon Java at commit a176eba1c6f9b0402eceea641bf435b05976470b. + +Schema: +id INT +string_payloads MAP +boolean_payloads MAP +tinyint_payloads MAP +smallint_payloads MAP +int_payloads MAP +bigint_payloads MAP +date_payloads MAP +binary_payloads MAP +compact_decimal_payloads MAP +large_decimal_payloads MAP + +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: +id 1: + string_payloads: {"": "string-empty", "alpha": "string-alpha"} + boolean_payloads: {false: "bool-false", true: "bool-true"} + tinyint_payloads: {-128: "tiny-min", -1: "tiny-negative", 127: "tiny-max"} + smallint_payloads: {-32768: "small-min", -1: "small-negative", 32767: "small-max"} + int_payloads: {-2147483648: "int-min", -1: "int-negative", 2147483647: "int-max"} + bigint_payloads: {-9223372036854775808: "big-min", -1: "big-negative", 9223372036854775807: "big-max"} + date_payloads: {-1: "date-negative", 0: "date-epoch"} + binary_payloads: {empty bytes: "binary-empty", 00 ff 01 02: "binary-bytes"} + compact_decimal_payloads: {-99999999.99: "compact-negative", 99999999.99: "compact-positive"} + large_decimal_payloads: {-999999999999999999.99: "large-negative", 999999999999999999.99: "large-positive"} +id 2: string_payloads is empty; all other map columns are null +id 3: all map columns are null +id 4: string_payloads is {"omega": "string-omega"}; all other map columns are null + +snapshot-2 +Create BatchTableWrite with the write type projected to string_payloads and write +BlobMapPlaceholder.INSTANCE for the four row positions. Set the first row id to 0 before +committing. This adds a second sequence layer without changing the logical values. + +snapshot-3 +Repeat the snapshot-2 write to add a third sequence layer. Reading this snapshot exercises +fallback through both newer placeholder layers to the values from snapshot-1. diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1a677b74-2f94-4d34-9166-82aaedb2b84b-0.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1a677b74-2f94-4d34-9166-82aaedb2b84b-0.blob new file mode 100644 index 000000000..4dc6a6a31 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1a677b74-2f94-4d34-9166-82aaedb2b84b-0.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-0.parquet b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-0.parquet new file mode 100644 index 000000000..830f7d858 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-0.parquet differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-1.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-1.blob new file mode 100644 index 000000000..11ce90ed6 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-1.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-10.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-10.blob new file mode 100644 index 000000000..e936ae192 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-10.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-2.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-2.blob new file mode 100644 index 000000000..71e7312c1 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-2.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-3.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-3.blob new file mode 100644 index 000000000..3a4a5bb91 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-3.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-4.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-4.blob new file mode 100644 index 000000000..d4a9c4969 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-4.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-5.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-5.blob new file mode 100644 index 000000000..263f8179b Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-5.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-6.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-6.blob new file mode 100644 index 000000000..5cee202c9 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-6.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-7.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-7.blob new file mode 100644 index 000000000..d0c44d4f5 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-7.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-8.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-8.blob new file mode 100644 index 000000000..d101f8d0c Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-8.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-9.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-9.blob new file mode 100644 index 000000000..9f84915e7 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-1be78d56-71bd-453f-ab1f-4312f3bf480e-9.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-d5380318-987b-4af0-9073-9c8e8162b568-0.blob b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-d5380318-987b-4af0-9073-9c8e8162b568-0.blob new file mode 100644 index 000000000..4dc6a6a31 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/bucket-0/data-d5380318-987b-4af0-9073-9c8e8162b568-0.blob differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-0e14efd4-f0a2-46cf-99ba-993506575ea4-0 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-0e14efd4-f0a2-46cf-99ba-993506575ea4-0 new file mode 100644 index 000000000..a51f68820 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-0e14efd4-f0a2-46cf-99ba-993506575ea4-0 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-1856e00c-bc22-481b-a8ac-c6304d6307fb-0 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-1856e00c-bc22-481b-a8ac-c6304d6307fb-0 new file mode 100644 index 000000000..fc43e02d5 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-1856e00c-bc22-481b-a8ac-c6304d6307fb-0 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-cbd224f1-d718-466b-89c4-f6ee7b673e2f-0 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-cbd224f1-d718-466b-89c4-f6ee7b673e2f-0 new file mode 100644 index 000000000..8a188e7e3 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-cbd224f1-d718-466b-89c4-f6ee7b673e2f-0 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-437bd3e4-24fc-40e8-9695-532ce23df015-0 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-437bd3e4-24fc-40e8-9695-532ce23df015-0 new file mode 100644 index 000000000..8c989918b Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-437bd3e4-24fc-40e8-9695-532ce23df015-0 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-437bd3e4-24fc-40e8-9695-532ce23df015-1 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-437bd3e4-24fc-40e8-9695-532ce23df015-1 new file mode 100644 index 000000000..bc483799b Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-437bd3e4-24fc-40e8-9695-532ce23df015-1 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-7475337a-832d-4e57-969d-c0cedd1b28be-0 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-7475337a-832d-4e57-969d-c0cedd1b28be-0 new file mode 100644 index 000000000..27f752348 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-7475337a-832d-4e57-969d-c0cedd1b28be-0 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-7475337a-832d-4e57-969d-c0cedd1b28be-1 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-7475337a-832d-4e57-969d-c0cedd1b28be-1 new file mode 100644 index 000000000..2eb6a32c9 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-7475337a-832d-4e57-969d-c0cedd1b28be-1 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-da204778-ac4b-4d76-af5e-bd55d0287b14-0 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-da204778-ac4b-4d76-af5e-bd55d0287b14-0 new file mode 100644 index 000000000..5ce489a3b Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-da204778-ac4b-4d76-af5e-bd55d0287b14-0 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-da204778-ac4b-4d76-af5e-bd55d0287b14-1 b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-da204778-ac4b-4d76-af5e-bd55d0287b14-1 new file mode 100644 index 000000000..f6479b2b8 Binary files /dev/null and b/test/test_data/parquet/map_blob_java.db/map_blob_java/manifest/manifest-list-da204778-ac4b-4d76-af5e-bd55d0287b14-1 differ diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/schema/schema-0 b/test/test_data/parquet/map_blob_java.db/map_blob_java/schema/schema-0 new file mode 100644 index 000000000..53d44145b --- /dev/null +++ b/test/test_data/parquet/map_blob_java.db/map_blob_java/schema/schema-0 @@ -0,0 +1,99 @@ +{ + "version" : 3, + "id" : 0, + "fields" : [ { + "id" : 0, + "name" : "id", + "type" : "INT" + }, { + "id" : 1, + "name" : "string_payloads", + "type" : { + "type" : "MAP", + "key" : "STRING", + "value" : "BLOB" + } + }, { + "id" : 2, + "name" : "boolean_payloads", + "type" : { + "type" : "MAP", + "key" : "BOOLEAN", + "value" : "BLOB" + } + }, { + "id" : 3, + "name" : "tinyint_payloads", + "type" : { + "type" : "MAP", + "key" : "TINYINT", + "value" : "BLOB" + } + }, { + "id" : 4, + "name" : "smallint_payloads", + "type" : { + "type" : "MAP", + "key" : "SMALLINT", + "value" : "BLOB" + } + }, { + "id" : 5, + "name" : "int_payloads", + "type" : { + "type" : "MAP", + "key" : "INT", + "value" : "BLOB" + } + }, { + "id" : 6, + "name" : "bigint_payloads", + "type" : { + "type" : "MAP", + "key" : "BIGINT", + "value" : "BLOB" + } + }, { + "id" : 7, + "name" : "date_payloads", + "type" : { + "type" : "MAP", + "key" : "DATE", + "value" : "BLOB" + } + }, { + "id" : 8, + "name" : "binary_payloads", + "type" : { + "type" : "MAP", + "key" : "BINARY(4)", + "value" : "BLOB" + } + }, { + "id" : 9, + "name" : "compact_decimal_payloads", + "type" : { + "type" : "MAP", + "key" : "DECIMAL(10, 2)", + "value" : "BLOB" + } + }, { + "id" : 10, + "name" : "large_decimal_payloads", + "type" : { + "type" : "MAP", + "key" : "DECIMAL(20, 2)", + "value" : "BLOB" + } + } ], + "highestFieldId" : 10, + "partitionKeys" : [ ], + "primaryKeys" : [ ], + "options" : { + "bucket" : "-1", + "data-evolution.enabled" : "true", + "file.format" : "parquet", + "row-tracking.enabled" : "true" + }, + "timeMillis" : 1788529095182 +} \ No newline at end of file diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/EARLIEST b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/EARLIEST new file mode 100644 index 000000000..56a6051ca --- /dev/null +++ b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/EARLIEST @@ -0,0 +1 @@ +1 \ No newline at end of file diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/LATEST b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/LATEST new file mode 100644 index 000000000..e440e5c84 --- /dev/null +++ b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/LATEST @@ -0,0 +1 @@ +3 \ No newline at end of file diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-1 b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-1 new file mode 100644 index 000000000..4a8ddfbe1 --- /dev/null +++ b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-1 @@ -0,0 +1,18 @@ +{ + "version" : 3, + "uuid" : "deee4971-d5b5-4930-92cb-b7bbcba49b86", + "id" : 1, + "schemaId" : 0, + "baseManifestList" : "manifest-list-437bd3e4-24fc-40e8-9695-532ce23df015-0", + "baseManifestListSize" : 1006, + "deltaManifestList" : "manifest-list-437bd3e4-24fc-40e8-9695-532ce23df015-1", + "deltaManifestListSize" : 1109, + "commitUser" : "daa569f0-8a00-4b77-82ba-8d58749de860", + "writerVersion" : "java-2.1-SNAPSHOT-a176eba1c6f9b0402eceea641bf435b05976470b", + "commitIdentifier" : 9223372036854775807, + "commitKind" : "APPEND", + "timeMillis" : 1788529096198, + "totalRecordCount" : 44, + "deltaRecordCount" : 44, + "nextRowId" : 4 +} \ No newline at end of file diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-2 b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-2 new file mode 100644 index 000000000..9f7ea0f4e --- /dev/null +++ b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-2 @@ -0,0 +1,18 @@ +{ + "version" : 3, + "uuid" : "42ae0281-d813-4af7-9fba-a461f970b9e6", + "id" : 2, + "schemaId" : 0, + "baseManifestList" : "manifest-list-7475337a-832d-4e57-969d-c0cedd1b28be-0", + "baseManifestListSize" : 1109, + "deltaManifestList" : "manifest-list-7475337a-832d-4e57-969d-c0cedd1b28be-1", + "deltaManifestListSize" : 1111, + "commitUser" : "08992445-7b6a-4421-bbcb-9f0dfe92b4e7", + "writerVersion" : "java-2.1-SNAPSHOT-a176eba1c6f9b0402eceea641bf435b05976470b", + "commitIdentifier" : 9223372036854775807, + "commitKind" : "APPEND", + "timeMillis" : 1788529096366, + "totalRecordCount" : 48, + "deltaRecordCount" : 4, + "nextRowId" : 4 +} \ No newline at end of file diff --git a/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-3 b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-3 new file mode 100644 index 000000000..38a44d56c --- /dev/null +++ b/test/test_data/parquet/map_blob_java.db/map_blob_java/snapshot/snapshot-3 @@ -0,0 +1,18 @@ +{ + "version" : 3, + "uuid" : "7c4965a2-a674-4d64-a873-819f8365909e", + "id" : 3, + "schemaId" : 0, + "baseManifestList" : "manifest-list-da204778-ac4b-4d76-af5e-bd55d0287b14-0", + "baseManifestListSize" : 1140, + "deltaManifestList" : "manifest-list-da204778-ac4b-4d76-af5e-bd55d0287b14-1", + "deltaManifestListSize" : 1111, + "commitUser" : "24f1b152-e7a9-4cb2-b2e1-8bea095a2b62", + "writerVersion" : "java-2.1-SNAPSHOT-a176eba1c6f9b0402eceea641bf435b05976470b", + "commitIdentifier" : 9223372036854775807, + "commitKind" : "APPEND", + "timeMillis" : 1788529096385, + "totalRecordCount" : 52, + "deltaRecordCount" : 4, + "nextRowId" : 4 +} \ No newline at end of file