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);
     }
 }
 

Reply via email to