Gabriel39 commented on code in PR #66302:
URL: https://github.com/apache/doris/pull/66302#discussion_r3689497523


##########
be/src/format_v2/column_mapper.cpp:
##########
@@ -2004,6 +2125,7 @@ Status 
TableColumnMapper::_create_mapping_for_column(const ColumnDefinition& tab
     mapping->global_index = global_index;
     mapping->table_column_name = table_column.name;
     mapping->table_type = table_column.type;
+    mapping->variant_access_paths = table_column.variant_access_paths;

Review Comment:
   Fixed. Nested Variant access paths are propagated through recursive 
mappings, nested projection finalization preserves physical typed leaves, and 
mapper/reader tests cover the nested path.



##########
be/src/format_v2/parquet/native_schema_desc.cpp:
##########
@@ -55,6 +56,209 @@ static bool is_map_node(const tparquet::SchemaElement& 
schema) {
            (schema.__isset.logicalType && schema.logicalType.__isset.MAP);
 }
 
+static bool is_variant_node(const tparquet::SchemaElement& schema) {
+    return schema.__isset.logicalType && schema.logicalType.__isset.VARIANT;
+}
+
+class ScopedBoolOverride {
+public:
+    ScopedBoolOverride(bool& target, bool value) : _target(target), 
_original(target) {
+        _target = value;
+    }
+    ~ScopedBoolOverride() { _target = _original; }
+
+private:
+    bool& _target;
+    bool _original;
+};
+
+static Status validate_variant_layout(const tparquet::SchemaElement& 
group_schema,
+                                      const NativeFieldSchema& group_field) {
+    const auto& annotation = group_schema.logicalType.VARIANT;
+    if (annotation.__isset.specification_version && 
annotation.specification_version != 1) {
+        return Status::NotSupported("Parquet Variant specification version {} 
is not supported",
+                                    annotation.specification_version);
+    }
+    if (group_field.children.size() < 2 || group_field.children.size() > 3) {
+        return Status::Corruption(
+                "Parquet Variant {} must contain metadata, value, and optional 
typed_value",
+                group_schema.name);
+    }
+
+    const NativeFieldSchema* metadata = nullptr;
+    const NativeFieldSchema* value = nullptr;
+    const NativeFieldSchema* typed_value = nullptr;
+    for (const auto& child : group_field.children) {
+        const NativeFieldSchema** target = nullptr;
+        if (child.name == "metadata") {
+            target = &metadata;
+        } else if (child.name == "value") {
+            target = &value;
+        } else if (child.name == "typed_value") {
+            target = &typed_value;
+        } else {
+            return Status::Corruption("Parquet Variant {} has unexpected child 
{}",
+                                      group_schema.name, child.name);
+        }
+        if (*target != nullptr) {
+            return Status::Corruption("Parquet Variant {} has duplicate child 
{}",
+                                      group_schema.name, child.name);
+        }
+        *target = &child;
+    }
+    if (metadata == nullptr || value == nullptr) {
+        return Status::Corruption("Parquet Variant {} requires metadata and 
value children",
+                                  group_schema.name);
+    }
+    if (!metadata->children.empty() || metadata->physical_type != 
tparquet::Type::BYTE_ARRAY ||
+        metadata->parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::REQUIRED) {
+        return Status::Corruption("Parquet Variant {} metadata must be a 
required BYTE_ARRAY",
+                                  group_schema.name);
+    }
+    const auto expected_value_repetition = typed_value == nullptr
+                                                   ? 
tparquet::FieldRepetitionType::REQUIRED
+                                                   : 
tparquet::FieldRepetitionType::OPTIONAL;
+    // SQL nullability belongs to the outer Variant group. Only shredding 
makes value optional,
+    // because typed_value may carry all or part of the logical value instead.
+    if (!value->children.empty() || value->physical_type != 
tparquet::Type::BYTE_ARRAY ||
+        value->parquet_schema.repetition_type != expected_value_repetition) {
+        return Status::Corruption("Parquet Variant {} value must be a {} 
BYTE_ARRAY",
+                                  group_schema.name,
+                                  typed_value == nullptr ? "required" : 
"optional");
+    }
+    if (typed_value != nullptr &&
+        typed_value->parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::OPTIONAL) {
+        return Status::Corruption("Parquet Variant {} typed_value must be 
optional",
+                                  group_schema.name);
+    }
+
+    enum class WrapperContext : uint8_t { OBJECT_FIELD, ARRAY_ELEMENT };
+    std::function<Status(const NativeFieldSchema&)> validate_typed_value;
+    std::function<Status(const NativeFieldSchema&, WrapperContext)> 
validate_wrapper;
+    validate_wrapper = [&](const NativeFieldSchema& wrapper, WrapperContext 
context) -> Status {
+        if (!wrapper.parquet_schema.__isset.repetition_type ||
+            wrapper.parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::REQUIRED) {
+            return Status::Corruption("Parquet Variant shredded wrapper {} 
must be required",
+                                      wrapper.name);
+        }
+        const NativeFieldSchema* fallback = nullptr;
+        const NativeFieldSchema* typed = nullptr;
+        for (const auto& child : wrapper.children) {
+            if (child.name == "value") {
+                if (fallback != nullptr) {
+                    return Status::Corruption(
+                            "Parquet Variant wrapper {} has duplicate value 
child", wrapper.name);
+                }
+                fallback = &child;
+            } else if (child.name == "typed_value") {
+                if (typed != nullptr) {
+                    return Status::Corruption(
+                            "Parquet Variant wrapper {} has duplicate 
typed_value child",
+                            wrapper.name);
+                }
+                typed = &child;
+            } else {
+                return Status::Corruption("Parquet Variant wrapper {} has 
unexpected child {}",
+                                          wrapper.name, child.name);
+            }
+        }
+        if (fallback == nullptr && typed == nullptr) {
+            return Status::Corruption(
+                    "Parquet Variant shredded wrapper {} requires at least one 
of value or "
+                    "typed_value",
+                    wrapper.name);
+        }
+        // Object fields always retain the fallback value carrier; only 
typed_value is optional.
+        // Array elements may omit either carrier when every element uses the 
remaining one.
+        if (context == WrapperContext::OBJECT_FIELD && fallback == nullptr) {
+            return Status::Corruption(
+                    "Parquet Variant object wrapper {} requires an optional 
value child",
+                    wrapper.name);
+        }
+        if (fallback != nullptr &&
+            (!fallback->children.empty() || fallback->physical_type != 
tparquet::Type::BYTE_ARRAY ||
+             !fallback->parquet_schema.__isset.repetition_type ||
+             fallback->parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::OPTIONAL)) {
+            return Status::Corruption(
+                    "Parquet Variant wrapper {} value must be an optional 
BYTE_ARRAY",
+                    wrapper.name);
+        }
+        if (typed != nullptr) {
+            if (!typed->parquet_schema.__isset.repetition_type ||
+                typed->parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::OPTIONAL) {
+                return Status::Corruption("Parquet Variant wrapper {} 
typed_value must be optional",
+                                          wrapper.name);
+            }
+            return validate_typed_value(*typed);
+        }
+        return Status::OK();
+    };
+    validate_typed_value = [&](const NativeFieldSchema& typed) -> Status {
+        if (!typed.unsupported_reason.empty()) {
+            return Status::NotSupported("Parquet Variant typed value {} is not 
supported: {}",
+                                        typed.name, typed.unsupported_reason);
+        }
+        if (typed.children.empty()) {
+            const auto& physical = typed.parquet_schema;
+            if (physical.__isset.logicalType && 
physical.logicalType.__isset.INTEGER &&
+                !physical.logicalType.INTEGER.isSigned) {
+                return Status::Corruption(
+                        "Parquet Variant unsigned integers are not valid typed 
values");
+            }
+            if (physical.__isset.converted_type &&
+                (physical.converted_type == tparquet::ConvertedType::UINT_8 ||
+                 physical.converted_type == tparquet::ConvertedType::UINT_16 ||
+                 physical.converted_type == tparquet::ConvertedType::UINT_32 ||
+                 physical.converted_type == tparquet::ConvertedType::UINT_64)) 
{
+                return Status::Corruption(
+                        "Parquet Variant unsigned integers are not valid typed 
values");
+            }
+            if (physical.__isset.logicalType && 
physical.logicalType.__isset.TIME) {
+                const auto& time = physical.logicalType.TIME;
+                // Variant v1 has one canonical TIME representation: local 
wall-clock MICROS.
+                // Accepting adjusted or lower-precision forms would make 
projection-dependent
+                // reconstruction disagree with the canonical Variant value.
+                if (time.isAdjustedToUTC) {
+                    return Status::Corruption(
+                            "Parquet Variant TIME must have 
isAdjustedToUTC=false");
+                }
+                if (!time.unit.__isset.MICROS) {
+                    return Status::Corruption(
+                            "Parquet Variant TIME(MILLIS) is not supported; 
use TIME(MICROS)");
+                }
+            }
+            if (physical.__isset.converted_type &&
+                physical.converted_type == 
tparquet::ConvertedType::TIME_MILLIS) {
+                return Status::Corruption(
+                        "Parquet Variant TIME(MILLIS) is not supported; use 
TIME(MICROS)");
+            }
+            if (physical.__isset.logicalType && 
physical.logicalType.__isset.TIMESTAMP &&
+                physical.logicalType.TIMESTAMP.unit.__isset.NANOS) {
+                // Reject at schema open so full reconstruction and direct 
typed-leaf access have
+                // the same precision contract instead of diverging after 
projection planning.
+                return Status::NotSupported("Parquet Variant TIMESTAMP(NANOS) 
is not supported");
+            }
+            return Status::OK();

Review Comment:
   Fixed. The reader now validates the complete Variant typed-value 
physical/logical type matrix, including width, annotation, signedness, decimal 
metadata, and time unit. Negative schema fixtures cover unsupported pairs.



##########
be/src/format_v2/parquet/native_schema_desc.cpp:
##########
@@ -55,6 +56,209 @@ static bool is_map_node(const tparquet::SchemaElement& 
schema) {
            (schema.__isset.logicalType && schema.logicalType.__isset.MAP);
 }
 
+static bool is_variant_node(const tparquet::SchemaElement& schema) {
+    return schema.__isset.logicalType && schema.logicalType.__isset.VARIANT;
+}
+
+class ScopedBoolOverride {
+public:
+    ScopedBoolOverride(bool& target, bool value) : _target(target), 
_original(target) {
+        _target = value;
+    }
+    ~ScopedBoolOverride() { _target = _original; }
+
+private:
+    bool& _target;
+    bool _original;
+};
+
+static Status validate_variant_layout(const tparquet::SchemaElement& 
group_schema,
+                                      const NativeFieldSchema& group_field) {
+    const auto& annotation = group_schema.logicalType.VARIANT;
+    if (annotation.__isset.specification_version && 
annotation.specification_version != 1) {
+        return Status::NotSupported("Parquet Variant specification version {} 
is not supported",
+                                    annotation.specification_version);
+    }
+    if (group_field.children.size() < 2 || group_field.children.size() > 3) {
+        return Status::Corruption(
+                "Parquet Variant {} must contain metadata, value, and optional 
typed_value",
+                group_schema.name);
+    }
+
+    const NativeFieldSchema* metadata = nullptr;
+    const NativeFieldSchema* value = nullptr;
+    const NativeFieldSchema* typed_value = nullptr;
+    for (const auto& child : group_field.children) {
+        const NativeFieldSchema** target = nullptr;
+        if (child.name == "metadata") {
+            target = &metadata;
+        } else if (child.name == "value") {
+            target = &value;
+        } else if (child.name == "typed_value") {
+            target = &typed_value;
+        } else {
+            return Status::Corruption("Parquet Variant {} has unexpected child 
{}",
+                                      group_schema.name, child.name);
+        }
+        if (*target != nullptr) {
+            return Status::Corruption("Parquet Variant {} has duplicate child 
{}",
+                                      group_schema.name, child.name);
+        }
+        *target = &child;
+    }
+    if (metadata == nullptr || value == nullptr) {
+        return Status::Corruption("Parquet Variant {} requires metadata and 
value children",
+                                  group_schema.name);
+    }
+    if (!metadata->children.empty() || metadata->physical_type != 
tparquet::Type::BYTE_ARRAY ||
+        metadata->parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::REQUIRED) {
+        return Status::Corruption("Parquet Variant {} metadata must be a 
required BYTE_ARRAY",
+                                  group_schema.name);
+    }
+    const auto expected_value_repetition = typed_value == nullptr
+                                                   ? 
tparquet::FieldRepetitionType::REQUIRED
+                                                   : 
tparquet::FieldRepetitionType::OPTIONAL;
+    // SQL nullability belongs to the outer Variant group. Only shredding 
makes value optional,
+    // because typed_value may carry all or part of the logical value instead.
+    if (!value->children.empty() || value->physical_type != 
tparquet::Type::BYTE_ARRAY ||
+        value->parquet_schema.repetition_type != expected_value_repetition) {
+        return Status::Corruption("Parquet Variant {} value must be a {} 
BYTE_ARRAY",
+                                  group_schema.name,
+                                  typed_value == nullptr ? "required" : 
"optional");
+    }
+    if (typed_value != nullptr &&
+        typed_value->parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::OPTIONAL) {
+        return Status::Corruption("Parquet Variant {} typed_value must be 
optional",
+                                  group_schema.name);
+    }
+
+    enum class WrapperContext : uint8_t { OBJECT_FIELD, ARRAY_ELEMENT };
+    std::function<Status(const NativeFieldSchema&)> validate_typed_value;
+    std::function<Status(const NativeFieldSchema&, WrapperContext)> 
validate_wrapper;
+    validate_wrapper = [&](const NativeFieldSchema& wrapper, WrapperContext 
context) -> Status {
+        if (!wrapper.parquet_schema.__isset.repetition_type ||
+            wrapper.parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::REQUIRED) {
+            return Status::Corruption("Parquet Variant shredded wrapper {} 
must be required",
+                                      wrapper.name);
+        }
+        const NativeFieldSchema* fallback = nullptr;
+        const NativeFieldSchema* typed = nullptr;
+        for (const auto& child : wrapper.children) {
+            if (child.name == "value") {
+                if (fallback != nullptr) {
+                    return Status::Corruption(
+                            "Parquet Variant wrapper {} has duplicate value 
child", wrapper.name);
+                }
+                fallback = &child;
+            } else if (child.name == "typed_value") {
+                if (typed != nullptr) {
+                    return Status::Corruption(
+                            "Parquet Variant wrapper {} has duplicate 
typed_value child",
+                            wrapper.name);
+                }
+                typed = &child;
+            } else {
+                return Status::Corruption("Parquet Variant wrapper {} has 
unexpected child {}",
+                                          wrapper.name, child.name);
+            }
+        }
+        if (fallback == nullptr && typed == nullptr) {
+            return Status::Corruption(
+                    "Parquet Variant shredded wrapper {} requires at least one 
of value or "
+                    "typed_value",
+                    wrapper.name);
+        }
+        // Object fields always retain the fallback value carrier; only 
typed_value is optional.
+        // Array elements may omit either carrier when every element uses the 
remaining one.
+        if (context == WrapperContext::OBJECT_FIELD && fallback == nullptr) {
+            return Status::Corruption(
+                    "Parquet Variant object wrapper {} requires an optional 
value child",
+                    wrapper.name);
+        }
+        if (fallback != nullptr &&
+            (!fallback->children.empty() || fallback->physical_type != 
tparquet::Type::BYTE_ARRAY ||
+             !fallback->parquet_schema.__isset.repetition_type ||
+             fallback->parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::OPTIONAL)) {
+            return Status::Corruption(
+                    "Parquet Variant wrapper {} value must be an optional 
BYTE_ARRAY",
+                    wrapper.name);
+        }
+        if (typed != nullptr) {
+            if (!typed->parquet_schema.__isset.repetition_type ||
+                typed->parquet_schema.repetition_type != 
tparquet::FieldRepetitionType::OPTIONAL) {
+                return Status::Corruption("Parquet Variant wrapper {} 
typed_value must be optional",
+                                          wrapper.name);
+            }
+            return validate_typed_value(*typed);
+        }
+        return Status::OK();
+    };
+    validate_typed_value = [&](const NativeFieldSchema& typed) -> Status {
+        if (!typed.unsupported_reason.empty()) {
+            return Status::NotSupported("Parquet Variant typed value {} is not 
supported: {}",
+                                        typed.name, typed.unsupported_reason);
+        }
+        if (typed.children.empty()) {
+            const auto& physical = typed.parquet_schema;
+            if (physical.__isset.logicalType && 
physical.logicalType.__isset.INTEGER &&
+                !physical.logicalType.INTEGER.isSigned) {
+                return Status::Corruption(
+                        "Parquet Variant unsigned integers are not valid typed 
values");
+            }
+            if (physical.__isset.converted_type &&
+                (physical.converted_type == tparquet::ConvertedType::UINT_8 ||
+                 physical.converted_type == tparquet::ConvertedType::UINT_16 ||
+                 physical.converted_type == tparquet::ConvertedType::UINT_32 ||
+                 physical.converted_type == tparquet::ConvertedType::UINT_64)) 
{
+                return Status::Corruption(
+                        "Parquet Variant unsigned integers are not valid typed 
values");
+            }
+            if (physical.__isset.logicalType && 
physical.logicalType.__isset.TIME) {
+                const auto& time = physical.logicalType.TIME;
+                // Variant v1 has one canonical TIME representation: local 
wall-clock MICROS.
+                // Accepting adjusted or lower-precision forms would make 
projection-dependent
+                // reconstruction disagree with the canonical Variant value.
+                if (time.isAdjustedToUTC) {
+                    return Status::Corruption(
+                            "Parquet Variant TIME must have 
isAdjustedToUTC=false");
+                }
+                if (!time.unit.__isset.MICROS) {
+                    return Status::Corruption(
+                            "Parquet Variant TIME(MILLIS) is not supported; 
use TIME(MICROS)");
+                }
+            }
+            if (physical.__isset.converted_type &&
+                physical.converted_type == 
tparquet::ConvertedType::TIME_MILLIS) {
+                return Status::Corruption(
+                        "Parquet Variant TIME(MILLIS) is not supported; use 
TIME(MICROS)");
+            }
+            if (physical.__isset.logicalType && 
physical.logicalType.__isset.TIMESTAMP &&
+                physical.logicalType.TIMESTAMP.unit.__isset.NANOS) {
+                // Reject at schema open so full reconstruction and direct 
typed-leaf access have
+                // the same precision contract instead of diverging after 
projection planning.
+                return Status::NotSupported("Parquet Variant TIMESTAMP(NANOS) 
is not supported");
+            }
+            return Status::OK();
+        }
+
+        const PrimitiveType primitive = 
remove_nullable(typed.data_type)->get_primitive_type();
+        if (primitive == TYPE_STRUCT) {
+            for (const auto& child : typed.children) {

Review Comment:
   Fixed. Shredded object wrapper names are validated for uniqueness before 
projection planning, with absent-first/present-second and both-present 
duplicate fixtures.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to