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]