This is an automated email from the ASF dual-hosted git repository.
lxy-9602 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 5a2fb5cb fix(variant): write parquet field ids in shredded physical
schema (#379)
5a2fb5cb is described below
commit 5a2fb5cb61e0577a9a52193c513d1d9cf89f88a3
Author: Nicholas Jiang <[email protected]>
AuthorDate: Mon Sep 21 18:47:40 2026 +0800
fix(variant): write parquet field ids in shredded physical schema (#379)
---
.../common/data/variant/variant_shredding_test.cpp | 106 ++++++++++++++
.../data/variant/variant_shredding_utils.cpp | 68 ++++++---
src/paimon/common/table/special_fields.h | 11 ++
src/paimon/common/table/special_fields_test.cpp | 7 +
src/paimon/format/parquet/variant_parquet_test.cpp | 161 ++++++++++++++++-----
5 files changed, 301 insertions(+), 52 deletions(-)
diff --git a/src/paimon/common/data/variant/variant_shredding_test.cpp
b/src/paimon/common/data/variant/variant_shredding_test.cpp
index cbf71f78..c0b54eb0 100644
--- a/src/paimon/common/data/variant/variant_shredding_test.cpp
+++ b/src/paimon/common/data/variant/variant_shredding_test.cpp
@@ -31,12 +31,22 @@
#include "paimon/common/data/variant/variant_schema.h"
#include "paimon/common/data/variant/variant_shredding_utils.h"
#include "paimon/common/data/variant/variant_shredding_writer.h"
+#include "paimon/common/types/data_field.h"
#include "paimon/common/utils/checked_cast.h"
#include "paimon/memory/memory_pool.h"
#include "paimon/testing/utils/testharness.h"
namespace paimon::test {
+namespace {
+
+int32_t ShreddedFieldId(const std::shared_ptr<arrow::Field>& field) {
+ Result<DataField> data_field =
DataField::ConvertArrowFieldToDataField(field);
+ return data_field.ok() ? data_field.value().Id() : -1;
+}
+
+} // namespace
+
class VariantShreddingTest : public ::testing::Test {
public:
// Shreds the given JSON documents (nullptr = null variant) with the given
shredding type,
@@ -217,6 +227,102 @@ TEST_F(VariantShreddingTest, ShreddingSchemaShape) {
VariantShreddingUtils::VariantShreddingSchema(arrow::map(arrow::utf8(),
arrow::int32())));
}
+TEST_F(VariantShreddingTest, ShreddingSchemaFieldIds) {
+ auto shredding_type = arrow::struct_(
+ {arrow::field("a", arrow::int32()), arrow::field("b",
arrow::list(arrow::utf8())),
+ arrow::field("c", arrow::struct_({arrow::field("d",
arrow::null())}))});
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::DataType> physical,
+
VariantShreddingUtils::VariantShreddingSchema(shredding_type));
+ ASSERT_EQ(physical->num_fields(), 3);
+ ASSERT_EQ(ShreddedFieldId(physical->field(0)), 0);
+ ASSERT_EQ(ShreddedFieldId(physical->field(1)), 1);
+ ASSERT_EQ(ShreddedFieldId(physical->field(2)), 2);
+
+ const std::shared_ptr<arrow::DataType>& object =
physical->field(2)->type();
+ ASSERT_EQ(object->num_fields(), 3);
+ ASSERT_EQ(ShreddedFieldId(object->field(0)), 0);
+ ASSERT_EQ(ShreddedFieldId(object->field(1)), 1);
+ ASSERT_EQ(ShreddedFieldId(object->field(2)), 2);
+
+ const std::shared_ptr<arrow::DataType>& scalar = object->field(0)->type();
+ ASSERT_EQ(scalar->num_fields(), 2);
+ ASSERT_EQ(ShreddedFieldId(scalar->field(0)), 0);
+ ASSERT_EQ(ShreddedFieldId(scalar->field(1)), 1);
+
+ const std::shared_ptr<arrow::DataType>& array = object->field(1)->type();
+ ASSERT_EQ(array->num_fields(), 2);
+ ASSERT_EQ(ShreddedFieldId(array->field(0)), 0);
+ ASSERT_EQ(ShreddedFieldId(array->field(1)), 1);
+ ASSERT_EQ(array->field(1)->type()->id(), arrow::Type::LIST);
+ const std::shared_ptr<arrow::Field>& element =
array->field(1)->type()->field(0);
+ ASSERT_EQ(ShreddedFieldId(element), 536871936);
+ ASSERT_EQ(element->type()->num_fields(), 2);
+ ASSERT_EQ(ShreddedFieldId(element->type()->field(0)), 0);
+ ASSERT_EQ(ShreddedFieldId(element->type()->field(1)), 1);
+
+ const std::shared_ptr<arrow::DataType>& nested = object->field(2)->type();
+ ASSERT_EQ(nested->num_fields(), 2);
+ ASSERT_EQ(ShreddedFieldId(nested->field(0)), 0);
+ ASSERT_EQ(ShreddedFieldId(nested->field(1)), 1);
+ const std::shared_ptr<arrow::DataType>& nested_object =
nested->field(1)->type();
+ ASSERT_EQ(nested_object->num_fields(), 1);
+ ASSERT_EQ(ShreddedFieldId(nested_object->field(0)), 0);
+ const std::shared_ptr<arrow::DataType>& untyped =
nested_object->field(0)->type();
+ ASSERT_EQ(untyped->num_fields(), 1);
+ ASSERT_EQ(ShreddedFieldId(untyped->field(0)), 0);
+
+ auto configured = arrow::struct_(
+ {arrow::field("a", arrow::int32(), true,
+ arrow::KeyValueMetadata::Make({DataField::FIELD_ID},
{"7"})),
+ arrow::field("b", arrow::utf8(), true,
+ arrow::KeyValueMetadata::Make({DataField::DESCRIPTION},
{"description"}))});
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::DataType> configured_physical,
+
VariantShreddingUtils::VariantShreddingSchema(configured));
+ ASSERT_EQ(configured_physical->num_fields(), 3);
+ ASSERT_EQ(configured_physical->field(2)->type()->num_fields(), 2);
+
ASSERT_EQ(ShreddedFieldId(configured_physical->field(2)->type()->field(0)), 7);
+
ASSERT_EQ(ShreddedFieldId(configured_physical->field(2)->type()->field(1)), 1);
+}
+
+TEST_F(VariantShreddingTest, ShreddingSchemaRejectsInvalidFieldIds) {
+ for (const std::string field_id : {"invalid", "", "2147483648"}) {
+ SCOPED_TRACE(field_id);
+ auto shredding_type = arrow::struct_(
+ {arrow::field("a", arrow::int32(), true,
+ arrow::KeyValueMetadata::Make({DataField::FIELD_ID},
{field_id}))});
+
ASSERT_NOK_WITH_MSG(VariantShreddingUtils::VariantShreddingSchema(shredding_type),
+ "cannot cast field id");
+ }
+}
+
+TEST_F(VariantShreddingTest, ShreddingArraySchemaFieldIds) {
+ for (const auto& leaf_type : {arrow::utf8(), arrow::null()}) {
+ SCOPED_TRACE(leaf_type->ToString());
+ ASSERT_OK_AND_ASSIGN(
+ std::shared_ptr<arrow::DataType> physical,
+
VariantShreddingUtils::VariantShreddingSchema(arrow::list(arrow::list(leaf_type))));
+ ASSERT_EQ(physical->num_fields(), 3);
+ ASSERT_EQ(ShreddedFieldId(physical->field(0)), 0);
+ ASSERT_EQ(ShreddedFieldId(physical->field(1)), 1);
+ ASSERT_EQ(ShreddedFieldId(physical->field(2)), 2);
+ ASSERT_EQ(physical->field(2)->type()->id(), arrow::Type::LIST);
+ const auto& outer_element = physical->field(2)->type()->field(0);
+ ASSERT_EQ(ShreddedFieldId(outer_element), 536872960);
+ ASSERT_EQ(outer_element->type()->num_fields(), 2);
+ ASSERT_EQ(ShreddedFieldId(outer_element->type()->field(0)), 0);
+ ASSERT_EQ(ShreddedFieldId(outer_element->type()->field(1)), 1);
+ const auto& inner_array = outer_element->type()->field(1)->type();
+ ASSERT_EQ(inner_array->id(), arrow::Type::LIST);
+ const auto& inner_element = inner_array->field(0);
+ ASSERT_EQ(ShreddedFieldId(inner_element), 536871936);
+ ASSERT_EQ(inner_element->type()->num_fields(), leaf_type->id() ==
arrow::Type::NA ? 1 : 2);
+ ASSERT_EQ(ShreddedFieldId(inner_element->type()->field(0)), 0);
+ if (leaf_type->id() != arrow::Type::NA) {
+ ASSERT_EQ(ShreddedFieldId(inner_element->type()->field(1)), 1);
+ }
+ }
+}
+
TEST_F(VariantShreddingTest, ShredObject) {
// Mirrors the Java GenericVariantTest#testShredding scenarios.
auto variant_json = R"({"a": 1, "b": "hello"})";
diff --git a/src/paimon/common/data/variant/variant_shredding_utils.cpp
b/src/paimon/common/data/variant/variant_shredding_utils.cpp
index 4d5c7dda..9796f905 100644
--- a/src/paimon/common/data/variant/variant_shredding_utils.cpp
+++ b/src/paimon/common/data/variant/variant_shredding_utils.cpp
@@ -31,6 +31,8 @@
#include "fmt/format.h"
#include "paimon/common/data/variant/variant_defs.h"
#include "paimon/common/data/variant/variant_type_utils.h"
+#include "paimon/common/table/special_fields.h"
+#include "paimon/common/types/data_field.h"
#include "paimon/common/utils/checked_cast.h"
namespace paimon {
@@ -42,14 +44,34 @@ Status InvalidVariantShreddingSchema(const
std::shared_ptr<arrow::DataType>& typ
fmt::format("Invalid variant shredding schema: {}", type ?
type->ToString() : "null"));
}
+std::shared_ptr<arrow::KeyValueMetadata> FieldIdMetadata(int32_t field_id) {
+ return arrow::KeyValueMetadata::Make({DataField::FIELD_ID},
{std::to_string(field_id)});
+}
+
+// Preserve configured IDs; inferred object fields use their positions, as in
Java RowType.
+Result<int32_t> GetObjectFieldId(const std::shared_ptr<arrow::Field>& field,
int32_t index) {
+ if (!field->metadata() ||
!field->metadata()->Contains(DataField::FIELD_ID)) {
+ return index;
+ }
+ PAIMON_ASSIGN_OR_RAISE(DataField data_field,
DataField::ConvertArrowFieldToDataField(field));
+ return data_field.Id();
+}
+
// Mirrors the Java `PaimonShreddingUtils.variantShreddingSchema(dataType,
isTopLevel,
// isObjectField)`.
Result<std::shared_ptr<arrow::DataType>> VariantShreddingSchemaImpl(
const std::shared_ptr<arrow::DataType>& data_type, bool is_top_level, bool
is_object_field) {
arrow::FieldVector fields;
+ // Java's Parquet reader requires field IDs; RowType.builder() numbers
each row from zero.
+ int32_t next_field_id = 0;
+ auto shredded_field = [&next_field_id](const std::string& name,
+ const
std::shared_ptr<arrow::DataType>& type,
+ bool nullable) {
+ return arrow::field(name, type, nullable,
FieldIdMetadata(next_field_id++));
+ };
if (is_top_level) {
- fields.push_back(arrow::field(VariantDefs::kMetadataFieldName,
arrow::binary(),
- /*nullable=*/false));
+ fields.push_back(shredded_field(VariantDefs::kMetadataFieldName,
arrow::binary(),
+ /*nullable=*/false));
}
switch (data_type->id()) {
case arrow::Type::LIST: {
@@ -58,10 +80,15 @@ Result<std::shared_ptr<arrow::DataType>>
VariantShreddingSchemaImpl(
VariantShreddingSchemaImpl(list_type->value_type(),
/*is_top_level=*/false,
/*is_object_field=*/false));
- fields.push_back(
- arrow::field(VariantDefs::kValueFieldName, arrow::binary(),
/*nullable=*/true));
- fields.push_back(arrow::field(VariantDefs::kTypedValueFieldName,
- arrow::list(element_type),
/*nullable=*/true));
+ fields.push_back(shredded_field(VariantDefs::kValueFieldName,
arrow::binary(),
+ /*nullable=*/true));
+ // Java resets collection depth at each ROW; every shredded array
element is a ROW.
+ const int32_t element_id =
+ SpecialFields::GetArrayElementFieldId(next_field_id,
/*depth=*/1);
+ std::shared_ptr<arrow::Field> element_field =
+
arrow::list(element_type)->field(0)->WithMetadata(FieldIdMetadata(element_id));
+ fields.push_back(shredded_field(VariantDefs::kTypedValueFieldName,
+ arrow::list(element_field),
/*nullable=*/true));
break;
}
case arrow::Type::STRUCT: {
@@ -70,18 +97,21 @@ Result<std::shared_ptr<arrow::DataType>>
VariantShreddingSchemaImpl(
// "value" and "typed_value" to null.
const auto& struct_type =
checked_pointer_cast<arrow::StructType>(data_type);
arrow::FieldVector shredded_fields;
- for (const auto& field : struct_type->fields()) {
+ for (int32_t index = 0; index < struct_type->num_fields();
++index) {
+ const std::shared_ptr<arrow::Field>& field =
struct_type->field(index);
+ PAIMON_ASSIGN_OR_RAISE(int32_t field_id,
GetObjectFieldId(field, index));
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::DataType>
field_type,
VariantShreddingSchemaImpl(field->type(),
/*is_top_level=*/false,
/*is_object_field=*/true));
- shredded_fields.push_back(
- arrow::field(field->name(), field_type,
/*nullable=*/false));
+ shredded_fields.push_back(arrow::field(field->name(),
field_type,
+ /*nullable=*/false,
+
FieldIdMetadata(field_id)));
}
- fields.push_back(
- arrow::field(VariantDefs::kValueFieldName, arrow::binary(),
/*nullable=*/true));
- fields.push_back(arrow::field(VariantDefs::kTypedValueFieldName,
- arrow::struct_(shredded_fields),
/*nullable=*/true));
+ fields.push_back(shredded_field(VariantDefs::kValueFieldName,
arrow::binary(),
+ /*nullable=*/true));
+ fields.push_back(shredded_field(VariantDefs::kTypedValueFieldName,
+ arrow::struct_(shredded_fields),
/*nullable=*/true));
break;
}
case arrow::Type::NA: {
@@ -89,8 +119,8 @@ Result<std::shared_ptr<arrow::DataType>>
VariantShreddingSchemaImpl(
// need a typed column. If there is no typed column, value is
required for array
// elements or top-level fields, but optional for objects (where a
null represents a
// missing field).
- fields.push_back(arrow::field(VariantDefs::kValueFieldName,
arrow::binary(),
- /*nullable=*/is_object_field));
+ fields.push_back(shredded_field(VariantDefs::kValueFieldName,
arrow::binary(),
+ /*nullable=*/is_object_field));
break;
}
case arrow::Type::STRING:
@@ -103,10 +133,10 @@ Result<std::shared_ptr<arrow::DataType>>
VariantShreddingSchemaImpl(
case arrow::Type::INT64:
case arrow::Type::FLOAT:
case arrow::Type::DOUBLE: {
- fields.push_back(
- arrow::field(VariantDefs::kValueFieldName, arrow::binary(),
/*nullable=*/true));
- fields.push_back(arrow::field(VariantDefs::kTypedValueFieldName,
data_type,
- /*nullable=*/true));
+ fields.push_back(shredded_field(VariantDefs::kValueFieldName,
arrow::binary(),
+ /*nullable=*/true));
+ fields.push_back(shredded_field(VariantDefs::kTypedValueFieldName,
data_type,
+ /*nullable=*/true));
break;
}
default:
diff --git a/src/paimon/common/table/special_fields.h
b/src/paimon/common/table/special_fields.h
index f23dafd8..e810fc1a 100644
--- a/src/paimon/common/table/special_fields.h
+++ b/src/paimon/common/table/special_fields.h
@@ -74,6 +74,17 @@ struct SpecialFields {
return data_field;
}
+ /// Field ID range and stride used by Java SpecialFields for structured
types.
+ static constexpr int32_t STRUCTURED_TYPE_FIELD_ID_BASE =
+ std::numeric_limits<int32_t>::max() / 4;
+ static constexpr int32_t STRUCTURED_TYPE_FIELD_DEPTH_LIMIT = 1 << 10;
+
+ /// Match Java's Parquet array-element IDs. Depth starts at 1 and resets
at each ROW.
+ static constexpr int32_t GetArrayElementFieldId(int32_t array_field_id,
int32_t depth) {
+ return STRUCTURED_TYPE_FIELD_ID_BASE + array_field_id *
STRUCTURED_TYPE_FIELD_DEPTH_LIMIT +
+ depth;
+ }
+
static bool IsSystemField(const std::string& field_name) {
if (StringUtils::StartsWith(field_name, KEY_FIELD_PREFIX)) {
return true;
diff --git a/src/paimon/common/table/special_fields_test.cpp
b/src/paimon/common/table/special_fields_test.cpp
index 58a025ba..a985e2b8 100644
--- a/src/paimon/common/table/special_fields_test.cpp
+++ b/src/paimon/common/table/special_fields_test.cpp
@@ -66,6 +66,13 @@ TEST(SpecialFieldsTest, TestKeyValueSpecialFieldCount) {
ASSERT_EQ(SpecialFields::KEY_VALUE_SPECIAL_FIELD_COUNT, 2);
}
+TEST(SpecialFieldsTest, TestGetArrayElementFieldId) {
+ ASSERT_EQ(SpecialFields::STRUCTURED_TYPE_FIELD_ID_BASE, 536870911);
+ ASSERT_EQ(SpecialFields::STRUCTURED_TYPE_FIELD_DEPTH_LIMIT, 1024);
+ ASSERT_EQ(SpecialFields::GetArrayElementFieldId(1, 1), 536871936);
+ ASSERT_EQ(SpecialFields::GetArrayElementFieldId(2, 2), 536872961);
+}
+
TEST(SpecialFieldsTest, TestIsSystemField) {
ASSERT_TRUE(SpecialFields::IsSystemField("_SEQUENCE_NUMBER"));
ASSERT_TRUE(SpecialFields::IsSystemField("_VALUE_KIND"));
diff --git a/src/paimon/format/parquet/variant_parquet_test.cpp
b/src/paimon/format/parquet/variant_parquet_test.cpp
index 46a63e38..2c1bab1d 100644
--- a/src/paimon/format/parquet/variant_parquet_test.cpp
+++ b/src/paimon/format/parquet/variant_parquet_test.cpp
@@ -484,6 +484,38 @@ constexpr const char* kAgeCityShreddingSchema = R"({
} ]
})";
+constexpr const char* kNestedShreddingSchema = R"({
+ "type": "ROW",
+ "fields": [ {
+ "id": 0,
+ "name": "v",
+ "type": {
+ "type": "ROW",
+ "fields": [
+ {"id": 3, "name": "age", "type": "INT"},
+ {"id": 4, "name": "addr", "type": {
+ "type": "ROW",
+ "fields": [ {"id": 5, "name": "city", "type": "STRING"} ]
+ }},
+ {"id": 6, "name": "tags", "type": {"type": "ARRAY", "element":
"STRING"}}
+ ]
+ }
+ } ]
+})";
+
+void CollectFieldIds(const ::parquet::schema::Node& node, const std::string&
prefix,
+ std::map<std::string, int32_t>* field_ids) {
+ std::string path = prefix.empty() ? node.name() : prefix + "." +
node.name();
+ (*field_ids)[path] = node.field_id();
+ if (!node.is_group()) {
+ return;
+ }
+ const auto& group = checked_cast<const
::parquet::schema::GroupNode&>(node);
+ for (int32_t i = 0; i < group.field_count(); ++i) {
+ CollectFieldIds(*group.field(i), path, field_ids);
+ }
+}
+
} // namespace
TEST_F(VariantParquetTest, PhysicalLayoutMatchesJava) {
@@ -611,46 +643,109 @@ TEST_F(VariantParquetTest, WriteAndReadRoundTrip) {
TEST_F(VariantParquetTest, ShreddedWriteAndReadRoundTrip) {
std::vector<const char*> jsons = {
- R"({"age": 35, "city": "Hangzhou"})",
+ R"({"age": 35, "addr": {"city": "Hangzhou"}, "tags": ["x", "y"]})",
nullptr,
R"({"age": "not a number", "extra": [1, 2]})",
"[\"top level array\"]",
};
- WriteShreddedFile(jsons, kAgeCityShreddingSchema);
+ for (const std::string mode : {"configured", "per-file", "adaptive"}) {
+ SCOPED_TRACE(mode);
+ const bool configured = mode == "configured";
+ std::map<std::string, std::string> option_map;
+ if (configured) {
+ option_map[Options::VARIANT_SHREDDING_SCHEMA] =
kNestedShreddingSchema;
+ } else {
+ option_map[Options::VARIANT_INFER_SHREDDING_SCHEMA] = "true";
+ option_map[Options::VARIANT_SHREDDING_INFERENCE_MODE] = mode;
+ }
+ ASSERT_OK_AND_ASSIGN(CoreOptions options,
CoreOptions::FromMap(option_map));
+ auto factory = VariantShreddingWritePlanFactory::Create(options,
paimon_schema_, pool_);
+ std::vector<std::shared_ptr<arrow::Array>> samples =
{BuildArray({jsons[0]})};
+ const int32_t file_count = mode == "adaptive" ? 2 : 1;
+ for (int32_t file_index = 0; file_index < file_count; ++file_index) {
+ SCOPED_TRACE(file_index);
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<ShreddingBatchConverter>
converter,
+ factory->CreateConverter("parquet", samples));
+ ASSERT_NE(converter, nullptr);
+ auto logical = BuildArray(jsons);
+ auto c_logical = std::make_unique<ArrowArray>();
+ ASSERT_TRUE(arrow::ExportArray(*logical, c_logical.get()).ok());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<ArrowArray> c_physical,
+ converter->Convert(c_logical.get()));
+ ASSERT_NO_FATAL_FAILURE(WriteFile(converter->GetPhysicalSchema(),
c_physical.get()));
+ ASSERT_OK(factory->OnFileCompleted(converter));
+ samples.clear();
+
+ {
+ std::unique_ptr<FileBatchReader> file_reader;
+ std::shared_ptr<arrow::Schema> file_schema;
+ ASSERT_NO_FATAL_FAILURE(OpenFile(&file_reader, &file_schema));
+ auto file_variant_field = file_schema->GetFieldByName("v");
+ ASSERT_NE(file_variant_field, nullptr);
+
ASSERT_TRUE(VariantShreddingUtils::IsShreddedFileType(file_variant_field->type()))
+ << file_variant_field->type()->ToString();
+ file_reader->Close();
+ }
- {
- std::unique_ptr<FileBatchReader> file_reader;
- std::shared_ptr<arrow::Schema> file_schema;
- OpenFile(&file_reader, &file_schema);
- auto file_variant_field = file_schema->GetFieldByName("v");
- ASSERT_NE(file_variant_field, nullptr);
-
ASSERT_TRUE(VariantShreddingUtils::IsShreddedFileType(file_variant_field->type()))
- << file_variant_field->type()->ToString();
- file_reader->Close();
- }
+ {
+ auto file = arrow::io::ReadableFile::Open(file_path_,
arrow_pool_.get());
+ ASSERT_TRUE(file.ok());
+ std::unique_ptr<::parquet::arrow::FileReader> reader;
+ auto status =
+ ::parquet::arrow::OpenFile(file.ValueOrDie(),
arrow_pool_.get(), &reader);
+ ASSERT_TRUE(status.ok()) << status.ToString();
+ const auto* root =
reader->parquet_reader()->metadata()->schema()->group_node();
+ ASSERT_EQ(root->field_count(), 2);
+ std::map<std::string, int32_t> field_ids;
+ CollectFieldIds(*root->field(1), "", &field_ids);
+ const std::map<std::string, int32_t> expected = {
+ {"v", 2},
+ {"v.metadata", 0},
+ {"v.value", 1},
+ {"v.typed_value", 2},
+ {"v.typed_value.age", configured ? 3 : 1},
+ {"v.typed_value.age.value", 0},
+ {"v.typed_value.age.typed_value", 1},
+ {"v.typed_value.addr", configured ? 4 : 0},
+ {"v.typed_value.addr.value", 0},
+ {"v.typed_value.addr.typed_value", 1},
+ {"v.typed_value.addr.typed_value.city", configured ? 5 :
0},
+ {"v.typed_value.addr.typed_value.city.value", 0},
+ {"v.typed_value.addr.typed_value.city.typed_value", 1},
+ {"v.typed_value.tags", configured ? 6 : 2},
+ {"v.typed_value.tags.value", 0},
+ {"v.typed_value.tags.typed_value", 1},
+ {"v.typed_value.tags.typed_value.list", -1},
+ {"v.typed_value.tags.typed_value.list.element", 536871936},
+ {"v.typed_value.tags.typed_value.list.element.value", 0},
+
{"v.typed_value.tags.typed_value.list.element.typed_value", 1},
+ };
+ ASSERT_EQ(field_ids, expected);
+ }
- // Reading the column as a plain VARIANT reassembles every physical shape
back to the
- // original logical value.
- std::shared_ptr<arrow::StructArray> variant_column;
- ReadVariantColumn(paimon_schema_, &variant_column);
- ASSERT_EQ(variant_column->length(), static_cast<int64_t>(jsons.size()));
- auto value_column =
checked_pointer_cast<arrow::BinaryArray>(variant_column->field(0));
- auto metadata_column =
checked_pointer_cast<arrow::BinaryArray>(variant_column->field(1));
- for (size_t i = 0; i < jsons.size(); ++i) {
- SCOPED_TRACE("row " + std::to_string(i));
- if (jsons[i] == nullptr) {
- ASSERT_TRUE(variant_column->IsNull(i));
- continue;
+ std::shared_ptr<arrow::StructArray> variant_column;
+ ASSERT_NO_FATAL_FAILURE(ReadVariantColumn(paimon_schema_,
&variant_column));
+ ASSERT_EQ(variant_column->length(),
static_cast<int64_t>(jsons.size()));
+ auto value_column =
checked_pointer_cast<arrow::BinaryArray>(variant_column->field(0));
+ auto metadata_column =
+
checked_pointer_cast<arrow::BinaryArray>(variant_column->field(1));
+ for (size_t i = 0; i < jsons.size(); ++i) {
+ SCOPED_TRACE(i);
+ if (jsons[i] == nullptr) {
+ ASSERT_TRUE(variant_column->IsNull(i));
+ continue;
+ }
+ ASSERT_FALSE(variant_column->IsNull(i));
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<GenericVariant> variant,
+
GenericVariant::Create(value_column->GetView(i),
+
metadata_column->GetView(i), pool_));
+ ASSERT_OK_AND_ASSIGN(std::string actual_json,
variant->ToJson());
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<GenericVariant> expected,
+ GenericVariant::FromJson(jsons[i],
pool_));
+ ASSERT_OK_AND_ASSIGN(std::string expected_json,
expected->ToJson());
+ ASSERT_EQ(actual_json, expected_json);
+ }
}
- ASSERT_FALSE(variant_column->IsNull(i));
- ASSERT_OK_AND_ASSIGN(
- std::shared_ptr<GenericVariant> variant,
- GenericVariant::Create(value_column->GetView(i),
metadata_column->GetView(i), pool_));
- ASSERT_OK_AND_ASSIGN(std::string actual_json, variant->ToJson());
- ASSERT_OK_AND_ASSIGN(std::shared_ptr<GenericVariant> expected,
- GenericVariant::FromJson(jsons[i], pool_));
- ASSERT_OK_AND_ASSIGN(std::string expected_json, expected->ToJson());
- ASSERT_EQ(actual_json, expected_json);
}
}