github-actions[bot] commented on code in PR #66206:
URL: https://github.com/apache/doris/pull/66206#discussion_r3671049071
##########
be/src/format_v2/parquet/parquet_statistics.cpp:
##########
@@ -452,6 +455,316 @@ std::optional<format::LocalColumnId>
file_column_id_by_block_position(
return std::nullopt;
}
+enum class VariantComparisonOp { EQ, NE, LT, LE, GT, GE };
+
+struct VariantShreddedPredicate {
+ int slot_index = -1;
+ std::vector<std::string> path;
+ DataTypePtr comparison_type;
+ DataTypePtr literal_type;
+ Field literal;
+ VariantComparisonOp op = VariantComparisonOp::EQ;
+};
+
+std::string callable_name(const VExprSPtr& expr) {
+ if (const auto function =
std::dynamic_pointer_cast<VectorizedFnCall>(expr);
+ function != nullptr) {
+ return function->function_name();
+ }
+ return expr == nullptr ? std::string {} : expr->expr_name();
+}
+
+std::optional<VariantComparisonOp> variant_comparison_op(std::string_view
name) {
+ if (name == "eq") {
+ return VariantComparisonOp::EQ;
+ }
+ if (name == "ne") {
+ return VariantComparisonOp::NE;
+ }
+ if (name == "lt") {
+ return VariantComparisonOp::LT;
+ }
+ if (name == "le") {
+ return VariantComparisonOp::LE;
+ }
+ if (name == "gt") {
+ return VariantComparisonOp::GT;
+ }
+ if (name == "ge") {
+ return VariantComparisonOp::GE;
+ }
+ return std::nullopt;
+}
+
+VariantComparisonOp reverse_variant_comparison(VariantComparisonOp op) {
+ switch (op) {
+ case VariantComparisonOp::EQ:
+ case VariantComparisonOp::NE:
+ return op;
+ case VariantComparisonOp::LT:
+ return VariantComparisonOp::GT;
+ case VariantComparisonOp::LE:
+ return VariantComparisonOp::GE;
+ case VariantComparisonOp::GT:
+ return VariantComparisonOp::LT;
+ case VariantComparisonOp::GE:
+ return VariantComparisonOp::LE;
+ }
+ __builtin_unreachable();
+}
+
+std::optional<std::pair<Field, DataTypePtr>> variant_literal(const VExprSPtr&
expr) {
+ const auto literal = std::dynamic_pointer_cast<VLiteral>(expr);
+ if (literal == nullptr || !literal->get_column_ptr() ||
literal->get_column_ptr()->empty()) {
+ return std::nullopt;
+ }
+ Field value;
+ literal->get_column_ptr()->get(0, value);
+ if (value.is_null()) {
+ return std::nullopt;
+ }
+ return std::make_pair(std::move(value), literal->get_data_type());
+}
+
+std::optional<VariantShreddedPredicate> extract_variant_shredded_predicate(
+ const VExprContextSPtr& conjunct) {
+ if (conjunct == nullptr || conjunct->root() == nullptr ||
+ conjunct->root()->get_num_children() != 2) {
+ return std::nullopt;
+ }
+ auto op = variant_comparison_op(callable_name(conjunct->root()));
+ if (!op.has_value()) {
+ return std::nullopt;
+ }
+
+ VExprSPtr value_expr;
+ std::optional<std::pair<Field, DataTypePtr>> literal;
+ if ((literal =
variant_literal(conjunct->root()->get_child(1))).has_value()) {
+ value_expr = conjunct->root()->get_child(0);
+ } else if ((literal =
variant_literal(conjunct->root()->get_child(0))).has_value()) {
+ value_expr = conjunct->root()->get_child(1);
+ op = reverse_variant_comparison(*op);
+ } else {
+ return std::nullopt;
+ }
+
+ const auto comparison_type = value_expr->data_type();
+ while (value_expr->node_type() == TExprNodeType::CAST_EXPR &&
+ value_expr->get_num_children() == 1) {
+ if (!expr_zonemap::data_types_compatible(value_expr->data_type(),
comparison_type)) {
+ // Every removed cast must preserve the comparison domain.
Otherwise bounds for the
+ // raw typed leaf could skip rows whose value changes in an
intermediate narrowing cast.
+ return std::nullopt;
+ }
+ value_expr = value_expr->get_child(0);
+ }
+
+ std::vector<std::string> reverse_path;
+ while (callable_name(value_expr) == "element_at" &&
value_expr->get_num_children() == 2) {
+ const auto key = variant_literal(value_expr->get_child(1));
+ if (!key.has_value() || key->first.get_type() != TYPE_STRING) {
+ // Repeated array shredding has no single scalar page range, so
only object keys are
+ // eligible for this file-level optimization.
+ return std::nullopt;
+ }
+ reverse_path.push_back(key->first.get<TYPE_STRING>());
+ value_expr = value_expr->get_child(0);
+ }
+ const auto slot = std::dynamic_pointer_cast<VSlotRef>(value_expr);
+ if (slot == nullptr || reverse_path.empty() || comparison_type == nullptr
||
+ remove_nullable(slot->data_type())->get_primitive_type() !=
TYPE_VARIANT ||
+ !expr_zonemap::data_types_compatible(comparison_type,
literal->second)) {
+ return std::nullopt;
+ }
+ std::ranges::reverse(reverse_path);
+ return VariantShreddedPredicate {.slot_index = slot->column_id(),
+ .path = std::move(reverse_path),
+ .comparison_type = comparison_type,
+ .literal_type = literal->second,
+ .literal = std::move(literal->first),
+ .op = *op};
+}
+
+bool has_variant_shredded_filter(const format::FileScanRequest& request) {
+ return std::ranges::any_of(request.conjuncts, [](const auto& conjunct) {
+ return extract_variant_shredded_predicate(conjunct).has_value();
+ });
+}
+
+const ParquetColumnSchema* child_named(const ParquetColumnSchema& parent,
std::string_view name) {
+ const auto it = std::ranges::find_if(parent.children, [&](const auto&
child) {
+ return child != nullptr && child->name == name;
+ });
+ return it == parent.children.end() ? nullptr : it->get();
+}
+
+struct ResolvedVariantShredding {
+ const ParquetColumnSchema* fallback_value = nullptr;
+ const ParquetColumnSchema* typed_value = nullptr;
+};
+
+bool metadata_cast_is_order_preserving(const DataTypePtr& source, const
DataTypePtr& target) {
+ if (expr_zonemap::data_types_compatible(source, target)) {
+ return true;
+ }
+ const auto source_type = remove_nullable(source);
+ const auto target_type = remove_nullable(target);
+ const auto source_primitive = source_type->get_primitive_type();
+ const auto target_primitive = target_type->get_primitive_type();
+ // Metadata bounds may cross only exact widening domains. This mirrors the
residual CAST while
+ // excluding rounding, overflow, and narrowing cases that could reverse a
pruning decision.
+ if (source_primitive == TYPE_FLOAT && target_primitive == TYPE_DOUBLE) {
+ return true;
+ }
+ if (is_int(source_primitive) && source_primitive != TYPE_LARGEINT &&
+ is_decimalv3(target_primitive)) {
+ const uint32_t required_integer_digits = source_primitive ==
TYPE_TINYINT ? 3
+ : source_primitive ==
TYPE_SMALLINT ? 5
+ : source_primitive ==
TYPE_INT ? 10
+
: 19;
+ return target_type->get_precision() >= target_type->get_scale() &&
+ target_type->get_precision() - target_type->get_scale() >=
required_integer_digits;
+ }
+ if (is_decimalv3(source_primitive) && is_decimalv3(target_primitive)) {
+ const uint32_t source_integer_digits =
+ source_type->get_precision() - source_type->get_scale();
+ const uint32_t target_integer_digits =
+ target_type->get_precision() - target_type->get_scale();
+ return target_integer_digits >= source_integer_digits &&
+ target_type->get_scale() >= source_type->get_scale();
+ }
+ return false;
+}
+
+std::optional<Field> cast_metadata_field(const Field& value, const
DataTypePtr& source,
+ const DataTypePtr& target) {
+ if (expr_zonemap::data_types_compatible(source, target)) {
+ return value;
+ }
+ const auto source_type = remove_nullable(source);
+ const auto target_type = remove_nullable(target);
+ if (source_type->get_primitive_type() == TYPE_FLOAT &&
+ target_type->get_primitive_type() == TYPE_DOUBLE) {
+ return
Field::create_field<TYPE_DOUBLE>(static_cast<double>(value.get<TYPE_FLOAT>()));
+ }
+ try {
+ auto source_column = source_type->create_column();
+ source_column->insert(value);
+ DataTypeSerDe::FormatOptions options =
DataTypeSerDe::get_default_format_options();
+ options.converted_from_string = true;
+ std::string text = source_type->to_string(*source_column, 0, options);
+ StringRef input(text.data(), text.size());
+ auto target_column = target_type->create_column();
+ if (!target_type->get_serde()
+ ->from_string_strict_mode(input, *target_column, options)
+ .ok() ||
+ target_column->size() != 1) {
+ return std::nullopt;
+ }
+ Field result;
+ target_column->get(0, result);
+ return result;
+ } catch (...) {
+ return std::nullopt;
+ }
+}
+
+std::optional<ParquetColumnStatistics> normalize_variant_statistics(
+ const VariantShreddedPredicate& predicate, const ParquetColumnSchema&
typed_value,
+ const ParquetColumnStatistics& statistics) {
+ if (!statistics.has_min_max ||
+ expr_zonemap::data_types_compatible(typed_value.type,
predicate.comparison_type)) {
+ return statistics;
+ }
+ auto min_value =
+ cast_metadata_field(statistics.min_value, typed_value.type,
predicate.comparison_type);
+ auto max_value =
+ cast_metadata_field(statistics.max_value, typed_value.type,
predicate.comparison_type);
+ if (!min_value.has_value() || !max_value.has_value()) {
+ return std::nullopt;
+ }
+ auto normalized = statistics;
+ normalized.min_value = std::move(*min_value);
+ normalized.max_value = std::move(*max_value);
+ return normalized;
+}
+
+std::optional<ResolvedVariantShredding> resolve_variant_shredding(
+ const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
+ const format::FileScanRequest& request, const
VariantShreddedPredicate& predicate) {
+ const auto local_id = file_column_id_by_block_position(request,
predicate.slot_index);
+ if (!local_id.has_value() || local_id->value() < 0 ||
+ local_id->value() >= static_cast<int>(file_schema.size())) {
+ return std::nullopt;
+ }
+ const ParquetColumnSchema* wrapper = file_schema[local_id->value()].get();
+ if (wrapper == nullptr || wrapper->kind !=
ParquetColumnSchemaKind::VARIANT) {
+ return std::nullopt;
+ }
+ for (const auto& component : predicate.path) {
+ const auto* typed_object = child_named(*wrapper, "typed_value");
+ if (typed_object == nullptr || typed_object->kind !=
ParquetColumnSchemaKind::STRUCT) {
+ return std::nullopt;
+ }
+ wrapper = child_named(*typed_object, component);
+ if (wrapper == nullptr || wrapper->kind !=
ParquetColumnSchemaKind::STRUCT) {
+ return std::nullopt;
+ }
+ }
+ const auto* fallback = child_named(*wrapper, "value");
+ const auto* typed = child_named(*wrapper, "typed_value");
+ if (fallback == nullptr || typed == nullptr ||
+ fallback->kind != ParquetColumnSchemaKind::PRIMITIVE ||
+ typed->kind != ParquetColumnSchemaKind::PRIMITIVE ||
typed->max_repetition_level != 0 ||
+ !metadata_cast_is_order_preserving(typed->type,
predicate.comparison_type) ||
+ !expr_zonemap::data_types_compatible(predicate.comparison_type,
predicate.literal_type)) {
+ return std::nullopt;
+ }
+ return ResolvedVariantShredding {.fallback_value = fallback, .typed_value
= typed};
+}
+
+bool fallback_is_all_null(const tparquet::RowGroup& row_group,
+ const ParquetColumnSchema& fallback) {
+ if (fallback.max_repetition_level != 0 || fallback.leaf_column_id < 0 ||
+ fallback.leaf_column_id >= static_cast<int>(row_group.columns.size()))
{
+ return false;
+ }
+ const auto& chunk = row_group.columns[fallback.leaf_column_id];
+ return row_group.num_rows >= 0 && chunk.__isset.meta_data &&
+ chunk.meta_data.num_values == row_group.num_rows &&
chunk.meta_data.__isset.statistics &&
+ chunk.meta_data.statistics.__isset.null_count &&
+ chunk.meta_data.statistics.null_count == chunk.meta_data.num_values;
+}
+
+bool variant_statistics_exclude(const VariantShreddedPredicate& predicate,
+ const ParquetColumnStatistics& statistics) {
+ if (!statistics.has_any_statistics()) {
+ return false;
+ }
+ if (!statistics.has_not_null) {
+ return true;
+ }
+ if (!statistics.has_min_max) {
+ return false;
+ }
+ const auto& literal = predicate.literal;
+ switch (predicate.op) {
+ case VariantComparisonOp::EQ:
+ return literal < statistics.min_value || statistics.max_value <
literal;
+ case VariantComparisonOp::NE:
+ return statistics.min_value == literal && statistics.max_value ==
literal;
+ case VariantComparisonOp::LT:
+ return statistics.min_value >= literal;
+ case VariantComparisonOp::LE:
+ return statistics.min_value > literal;
+ case VariantComparisonOp::GT:
Review Comment:
[P1] Do not use floating bounds without proving that NaN is absent.
Legacy Parquet FLOAT/DOUBLE min/max excludes NaNs, while Doris comparisons
order NaN above ordinary values. For a valid shredded page containing `[0.0,
NaN]`, `max=0.0` makes this `GT` branch discard `v['x'] > 1.0`, even though the
residual Doris predicate retains the NaN row. The new footer and page paths
accept TYPE_ORDER bounds but never establish `nan_count == 0`, and a missing
count must be treated as unknown. Please disable shredded floating min/max
pruning unless supported metadata proves no NaNs (or evaluate the advertised
floating order exactly), with differential footer/page tests for NaN
inequalities.
##########
be/src/format_v2/parquet/reader/variant_column_reader.cpp:
##########
@@ -0,0 +1,818 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include "format_v2/parquet/reader/variant_column_reader.h"
+
+#include <algorithm>
+#include <array>
+#include <cmath>
+#include <cstdint>
+#include <cstring>
+#include <limits>
+#include <mutex>
+#include <optional>
+#include <string_view>
+
+#include "common/exception.h"
+#include "core/assert_cast.h"
+#include "core/column/column_array.h"
+#include "core/column/column_decimal.h"
+#include "core/column/column_map.h"
+#include "core/column/column_nullable.h"
+#include "core/column/column_struct.h"
+#include "core/column/column_vector.h"
+#include "core/column/variant_v2/column_variant_v2.h"
+#include "core/column/variant_v2/column_variant_v2_typed_column.h"
+#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_variant_v2.h"
+#include "core/value/variant/variant_batch_builder.h"
+#include "core/value/variant/variant_metadata.h"
+#include "format_v2/parquet/parquet_column_schema.h"
+
+namespace doris::format::parquet {
+namespace {
+
+struct Cell {
+ const IColumn* column = nullptr;
+ bool is_null = false;
+};
+
+Cell cell_at(const IColumn& column, size_t row) {
+ if (row >= column.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant row {} exceeds
column size {}", row,
+ column.size());
+ }
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(column)) {
+ return {.column = &nullable->get_nested_column(),
+ .is_null = nullable->get_null_map_data()[row] != 0};
+ }
+ return {.column = &column, .is_null = false};
+}
+
+const ParquetColumnSchema* find_child(const ParquetColumnSchema& schema,
std::string_view name,
+ size_t* index) {
+ for (size_t i = 0; i < schema.children.size(); ++i) {
+ if (schema.children[i]->name == name) {
+ if (index != nullptr) {
+ *index = i;
+ }
+ return schema.children[i].get();
+ }
+ }
+ return nullptr;
+}
+
+Cell struct_child_at(const ParquetColumnSchema& schema, const IColumn&
physical, size_t row,
+ std::string_view name, const ParquetColumnSchema**
child_schema) {
+ const auto& structure = assert_cast<const ColumnStruct&>(physical);
+ size_t index = 0;
+ const auto* child = find_child(schema, name, &index);
+ if (child == nullptr || index >= structure.tuple_size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} has no
physical child {}",
+ schema.name, name);
+ }
+ if (child_schema != nullptr) {
+ *child_schema = child;
+ }
+ return cell_at(structure.get_column(index), row);
+}
+
+uint8_t decimal_width(int precision) {
+ if (precision <= 0 || precision > 38) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant decimal precision {} is outside [1,
38]", precision);
+ }
+ return precision <= 9 ? 4 : (precision <= 18 ? 8 : 16);
+}
+
+uint8_t integer_width(const ParquetColumnSchema& schema, PrimitiveType type) {
+ if (schema.type_descriptor.is_unsigned_integer) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Unsigned integers are not valid Parquet Variant typed
values");
+ }
+ if (schema.type_descriptor.integer_bit_width > 0) {
+ switch (schema.type_descriptor.integer_bit_width) {
+ case 8:
+ return 1;
+ case 16:
+ return 2;
+ case 32:
+ return 4;
+ case 64:
+ return 8;
+ default:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
integer width {}",
+ schema.type_descriptor.integer_bit_width);
+ }
+ }
+ switch (type) {
+ case TYPE_TINYINT:
+ return 1;
+ case TYPE_SMALLINT:
+ return 2;
+ case TYPE_INT:
+ return 4;
+ case TYPE_BIGINT:
+ return 8;
+ default:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
integer type {}", type);
+ }
+}
+
+void append_typed_scalar(const ParquetColumnSchema& schema, const IColumn&
column, size_t row,
+ VariantBatchBuilder::Row& builder) {
+ const PrimitiveType type =
remove_nullable(schema.type)->get_primitive_type();
+ switch (type) {
+ case TYPE_BOOLEAN:
+ builder.add_bool(assert_cast<const
ColumnUInt8&>(column).get_data()[row] != 0);
+ return;
+ case TYPE_TINYINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt8&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_SMALLINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt16&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_INT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt32&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_BIGINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt64&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_FLOAT:
+ builder.add_float(assert_cast<const
ColumnFloat32&>(column).get_data()[row]);
+ return;
+ case TYPE_DOUBLE:
+ builder.add_double(assert_cast<const
ColumnFloat64&>(column).get_data()[row]);
+ return;
+ case TYPE_DECIMAL128I: {
+ const auto value = assert_cast<const
ColumnDecimal128V3&>(column).get_data()[row].value;
+ builder.add_decimal(value,
static_cast<uint8_t>(schema.type_descriptor.decimal_scale),
+
decimal_width(schema.type_descriptor.decimal_precision));
+ return;
+ }
+ case TYPE_TIMEV2: {
+ const double seconds = assert_cast<const
ColumnTimeV2&>(column).get_data()[row];
+ if (!std::isfinite(seconds) ||
+ std::abs(seconds) >
static_cast<double>(std::numeric_limits<int64_t>::max()) / 1e6) {
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
TIME value");
+ }
+ builder.add_time_ntz_micros(static_cast<int64_t>(std::llround(seconds
* 1e6)));
+ return;
+ }
+ case TYPE_DATETIMEV2: {
+ if (schema.type_descriptor.time_unit == ParquetTimeUnit::NANOS) {
Review Comment:
[P1] Keep NANOS timestamps consistent across root and path reads.
`TIMESTAMP(..., NANOS)` is a valid Parquet Variant shredded type, but root
materialization reaches this `NOT_IMPLEMENTED` branch while the zero-copy path
bypasses it after native decoding has already floored the value to
microseconds. A query can therefore either fail or silently lose
sub-microsecond precision solely based on projection shape. Doris already has
`TIMESTAMP_NANOS`/`TIMESTAMP_NTZ_NANOS` Variant IDs and
`add_timestamp_nanos()`, so preserve the raw INT64 carrier and emit the
matching UTC/local nanos primitive (or reject before either access path can
return a value). Please add root and narrow-path tests for both UTC-adjusted
and local NANOS data.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java:
##########
@@ -713,7 +713,11 @@ public static Type
icebergTypeToDorisType(org.apache.iceberg.types.Type type, bo
.collect(Collectors.toCollection(ArrayList::new));
return new StructType(nestedTypes);
case VARIANT:
- return Type.UNSUPPORTED;
+ // Iceberg Variant uses the Parquet Variant encoding directly.
Mark it compute-only
+ // so BE scanners materialize ColumnVariantV2 without changing
persisted Doris
+ // table metadata semantics.
+ return new org.apache.doris.catalog.VariantType(
Review Comment:
[P1] Gate this new native reader capability during smooth upgrade.
This mapping unconditionally exposes root and nested Iceberg Variant to
planning, but all Parquet Variant decode support is new in the BE side of this
patch. `IcebergScanNode` can still assign such a split to a selected
`Backend.isSmoothUpgradeSrc()` process, so the same query succeeds or fails
according to old/new BE placement during a rolling upgrade. The scan node
already rejects smooth-upgrade source BEs for its newer native position-deletes
path; please add an equivalent Variant capability/version check before
scheduling semantic Variant projections. Cover mixed old/new backend selection
in the FE tests.
##########
be/src/exec/scan/file_scanner_v2.cpp:
##########
@@ -758,6 +809,22 @@ Status FileScannerV2::_build_projected_columns(const
format::TableReader& table_
const bool prefer_exact_name_match =
!_params->__isset.history_schema_info ||
supports_iceberg_scan_semantics_v1(_params);
+ const auto format_type = get_range_format_type(*_params, _current_range);
+ if (table_format_name(_current_range) == "iceberg" &&
+ format_type != TFileFormatType::FORMAT_PARQUET) {
+ const bool projects_variant =
Review Comment:
[P1] Exclude a retained COUNT(*) placeholder from this Variant-format guard.
An explicitly empty `push_down_count_slot_ids` list means COUNT(*) has no
semantic column, but Nereids can retain one minimum-width output slot.
`TableReader` later marks such non-predicate slots as
`count_star_placeholder_columns` and deliberately omits them from validation
and physical I/O; this earlier check instead treats a retained VARIANT slot as
a real projection. On an Iceberg table whose only retained slot is optional
VARIANT, `SELECT COUNT(*)` over ORC or a mixed split now fails with the
Variant/Parquet error even though no Variant value is read. Please derive this
check from semantic, non-placeholder projections while continuing to reject
COUNT(v), and add an ORC/mixed COUNT(*) case with a retained Variant slot.
##########
be/src/core/column/variant_v2/column_variant_v2.cpp:
##########
@@ -369,7 +379,26 @@ const DataTypePtr& ColumnVariantV2::typed_type() const {
return _typed_type;
}
+std::optional<VariantShreddedTypedValue>
ColumnVariantV2::find_shredded_typed_value(
+ std::span<const VariantShreddedPathSegment> path) const {
+ if (!_shredded) {
+ return std::nullopt;
+ }
+ return _shredded->find_typed_value(path);
+}
+
void ColumnVariantV2::ensure_encoded() {
+ if (_shredded) {
+ const ColumnVariantV2& materialized = _shredded->materialized_column();
+ DORIS_CHECK(!materialized.is_shredded())
+ << "shredded state materializer returned another shredded
column";
+ _metadatas = materialized._metadatas;
Review Comment:
[P1] Detach these materialized buffers before mutating a shredded copy.
This assigns the cached column's subcolumn pointers and resets only the
current object's `_shredded` reference. Several new paths can reach it while
another live owner still retains that state and cache: the copy constructor
used by generic `clone()`, same-size `clone_resized()`, and the full-source
fast path in `insert_range_from()`. The object being mutated therefore still
shares `_meta_ids`/`_values`, so `pop_back()`, growth, or a second range
insertion reaches `require_exclusive()` and fails instead of mutating. Please
prevent mutable results from sharing the state, or detach the adopted buffers
here when the cache is shared, and cover ordinary clone, same-size clone, and
two-source append cases.
##########
be/src/format_v2/parquet/reader/native_column_reader.cpp:
##########
@@ -114,6 +185,12 @@ void collect_projected_ids(const ParquetColumnSchema&
schema,
const format::LocalColumnIndex* projection,
const NativeFieldSchema& native_field,
std::set<uint64_t>* ids) {
DORIS_CHECK(ids != nullptr);
+ if (schema.kind == ParquetColumnSchemaKind::VARIANT) {
Review Comment:
[P2] Preserve the requested Variant path in the physical projection.
This branch expands every Variant access to the complete physical shredding
subtree. As a result, `SELECT v['hot_key']` or a predicate-only access still
fetches and decodes the root residual plus every unrelated sibling's fallback
and typed chunks before the zero-copy extractor runs; one large cold sibling
can dominate remote bytes and CPU. That removes the field-projection benefit of
Parquet Variant shredding and conflicts with the ScannerV2
predicate-first/lazy-I/O path. Please retain only the presence/metadata
carriers and requested path for narrow access, reserving full-subtree
reconstruction for a true root-Variant projection, and add a profile test
proving an unreferenced large sibling is not requested.
##########
be/src/format_v2/parquet/reader/variant_column_reader.cpp:
##########
@@ -0,0 +1,818 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include "format_v2/parquet/reader/variant_column_reader.h"
+
+#include <algorithm>
+#include <array>
+#include <cmath>
+#include <cstdint>
+#include <cstring>
+#include <limits>
+#include <mutex>
+#include <optional>
+#include <string_view>
+
+#include "common/exception.h"
+#include "core/assert_cast.h"
+#include "core/column/column_array.h"
+#include "core/column/column_decimal.h"
+#include "core/column/column_map.h"
+#include "core/column/column_nullable.h"
+#include "core/column/column_struct.h"
+#include "core/column/column_vector.h"
+#include "core/column/variant_v2/column_variant_v2.h"
+#include "core/column/variant_v2/column_variant_v2_typed_column.h"
+#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_variant_v2.h"
+#include "core/value/variant/variant_batch_builder.h"
+#include "core/value/variant/variant_metadata.h"
+#include "format_v2/parquet/parquet_column_schema.h"
+
+namespace doris::format::parquet {
+namespace {
+
+struct Cell {
+ const IColumn* column = nullptr;
+ bool is_null = false;
+};
+
+Cell cell_at(const IColumn& column, size_t row) {
+ if (row >= column.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant row {} exceeds
column size {}", row,
+ column.size());
+ }
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(column)) {
+ return {.column = &nullable->get_nested_column(),
+ .is_null = nullable->get_null_map_data()[row] != 0};
+ }
+ return {.column = &column, .is_null = false};
+}
+
+const ParquetColumnSchema* find_child(const ParquetColumnSchema& schema,
std::string_view name,
+ size_t* index) {
+ for (size_t i = 0; i < schema.children.size(); ++i) {
+ if (schema.children[i]->name == name) {
+ if (index != nullptr) {
+ *index = i;
+ }
+ return schema.children[i].get();
+ }
+ }
+ return nullptr;
+}
+
+Cell struct_child_at(const ParquetColumnSchema& schema, const IColumn&
physical, size_t row,
+ std::string_view name, const ParquetColumnSchema**
child_schema) {
+ const auto& structure = assert_cast<const ColumnStruct&>(physical);
+ size_t index = 0;
+ const auto* child = find_child(schema, name, &index);
+ if (child == nullptr || index >= structure.tuple_size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} has no
physical child {}",
+ schema.name, name);
+ }
+ if (child_schema != nullptr) {
+ *child_schema = child;
+ }
+ return cell_at(structure.get_column(index), row);
+}
+
+uint8_t decimal_width(int precision) {
+ if (precision <= 0 || precision > 38) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant decimal precision {} is outside [1,
38]", precision);
+ }
+ return precision <= 9 ? 4 : (precision <= 18 ? 8 : 16);
+}
+
+uint8_t integer_width(const ParquetColumnSchema& schema, PrimitiveType type) {
+ if (schema.type_descriptor.is_unsigned_integer) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Unsigned integers are not valid Parquet Variant typed
values");
+ }
+ if (schema.type_descriptor.integer_bit_width > 0) {
+ switch (schema.type_descriptor.integer_bit_width) {
+ case 8:
+ return 1;
+ case 16:
+ return 2;
+ case 32:
+ return 4;
+ case 64:
+ return 8;
+ default:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
integer width {}",
+ schema.type_descriptor.integer_bit_width);
+ }
+ }
+ switch (type) {
+ case TYPE_TINYINT:
+ return 1;
+ case TYPE_SMALLINT:
+ return 2;
+ case TYPE_INT:
+ return 4;
+ case TYPE_BIGINT:
+ return 8;
+ default:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
integer type {}", type);
+ }
+}
+
+void append_typed_scalar(const ParquetColumnSchema& schema, const IColumn&
column, size_t row,
+ VariantBatchBuilder::Row& builder) {
+ const PrimitiveType type =
remove_nullable(schema.type)->get_primitive_type();
+ switch (type) {
+ case TYPE_BOOLEAN:
+ builder.add_bool(assert_cast<const
ColumnUInt8&>(column).get_data()[row] != 0);
+ return;
+ case TYPE_TINYINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt8&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_SMALLINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt16&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_INT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt32&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_BIGINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt64&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_FLOAT:
+ builder.add_float(assert_cast<const
ColumnFloat32&>(column).get_data()[row]);
+ return;
+ case TYPE_DOUBLE:
+ builder.add_double(assert_cast<const
ColumnFloat64&>(column).get_data()[row]);
+ return;
+ case TYPE_DECIMAL128I: {
+ const auto value = assert_cast<const
ColumnDecimal128V3&>(column).get_data()[row].value;
+ builder.add_decimal(value,
static_cast<uint8_t>(schema.type_descriptor.decimal_scale),
+
decimal_width(schema.type_descriptor.decimal_precision));
+ return;
+ }
+ case TYPE_TIMEV2: {
+ const double seconds = assert_cast<const
ColumnTimeV2&>(column).get_data()[row];
+ if (!std::isfinite(seconds) ||
+ std::abs(seconds) >
static_cast<double>(std::numeric_limits<int64_t>::max()) / 1e6) {
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
TIME value");
+ }
+ builder.add_time_ntz_micros(static_cast<int64_t>(std::llround(seconds
* 1e6)));
+ return;
+ }
+ case TYPE_DATETIMEV2: {
+ if (schema.type_descriptor.time_unit == ParquetTimeUnit::NANOS) {
+ // Native DATETIMEV2 is microsecond based. Reject before returning
a silently truncated
+ // value; a raw INT64 nanos decoder can be added independently.
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant TIMESTAMP(NANOS) is not
supported");
+ }
+ const auto& value = assert_cast<const
ColumnDateTimeV2&>(column).get_data()[row];
+ builder.add_timestamp_micros(
+ variant_timestamp_micros(value, row, "Parquet Variant
TIMESTAMP"),
+ schema.type_descriptor.timestamp_is_adjusted_to_utc);
+ return;
+ }
+ case TYPE_TIMESTAMPTZ: {
+ if (schema.type_descriptor.time_unit == ParquetTimeUnit::NANOS) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant TIMESTAMP(NANOS) is not
supported");
+ }
+ const auto& value = assert_cast<const
ColumnTimeStampTz&>(column).get_data()[row];
+ builder.add_timestamp_micros(
+ variant_timestamp_micros(value, row, "Parquet Variant
TIMESTAMP"), true);
+ return;
+ }
+ case TYPE_VARBINARY: {
+ const StringRef value = column.get_data_at(row);
+ if (!schema.type_descriptor.is_uuid) {
+ builder.add_binary(value);
+ return;
+ }
+ if (value.size != 16) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant UUID has {} bytes instead of 16",
value.size);
+ }
+ std::array<uint8_t, 16> uuid {};
+ std::memcpy(uuid.data(), value.data, uuid.size());
+ builder.add_uuid(uuid);
+ return;
+ }
+ case TYPE_STRING: {
+ const StringRef value = column.get_data_at(row);
+ if (schema.type_descriptor.is_uuid) {
+ if (value.size != 16) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant UUID has {} bytes instead of
16", value.size);
+ }
+ std::array<uint8_t, 16> uuid {};
+ std::memcpy(uuid.data(), value.data, uuid.size());
+ builder.add_uuid(uuid);
+ } else if (schema.type_descriptor.is_string_annotation) {
+ builder.add_string(value);
+ } else {
+ builder.add_binary(value);
+ }
+ return;
+ }
+ default:
+ if (!is_supported_variant_typed_identity(type)) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant typed value {} is not supported",
+ remove_nullable(schema.type)->get_name());
+ }
+ dispatch_variant_typed_column(
+ column, type, [&]<PrimitiveType Type>(const auto&
typed_column) {
+ with_variant_typed_scalar<Type>(
+ typed_column, row,
+
static_cast<uint8_t>(remove_nullable(schema.type)->get_scale()),
+ [&](const VariantScalarRef& scalar) {
builder.add_scalar(scalar); });
+ });
+ }
+}
+
+enum class WrapperContext { ROOT, ARRAY_ELEMENT, OBJECT_FIELD };
+
+bool append_wrapper(const ParquetColumnSchema& schema, const IColumn& wrapper,
size_t row,
+ VariantMetadataRef metadata, VariantBatchBuilder::Row&
builder,
+ WrapperContext context);
+
+void append_typed_value(const ParquetColumnSchema& schema, const IColumn&
column, size_t row,
+ VariantMetadataRef metadata, const VariantRef*
residual,
+ VariantBatchBuilder::Row& builder) {
+ switch (schema.kind) {
+ case ParquetColumnSchemaKind::PRIMITIVE:
+ if (static_cast<bool>(residual)) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant scalar typed_value cannot have
residual value bytes");
+ }
+ append_typed_scalar(schema, column, row, builder);
+ return;
+ case ParquetColumnSchemaKind::STRUCT: {
+ if (static_cast<bool>(residual) && residual->basic_type() !=
VariantBasicType::OBJECT) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant object typed_value has non-object
residual value");
+ }
+ const auto& structure = assert_cast<const ColumnStruct&>(column);
+ if (structure.tuple_size() != schema.children.size()) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant object {} physical field count
mismatch", schema.name);
+ }
+ auto object = builder.start_object();
+ if (static_cast<bool>(residual)) {
+ for (uint32_t i = 0; i < residual->num_elements(); ++i) {
+ uint32_t field_id = 0;
+ const VariantRef child = residual->object_value_at(i,
&field_id);
+ object.add_key(residual->metadata.key_at(field_id));
+ builder.add_value(child);
+ }
+ }
+ for (size_t i = 0; i < schema.children.size(); ++i) {
+ const auto& child_schema = *schema.children[i];
+ const Cell child = cell_at(structure.get_column(i), row);
+ if (child.is_null) {
+ // Shredded object fields are optional wrapper groups. A
missing group means the
+ // key is absent, which differs from a present wrapper
encoding a Variant null.
+ continue;
+ }
+ // A null/null wrapper means this object field is absent. Delay
add_key until its
+ // presence is known so absent shredded fields do not turn into
Variant nulls.
+ size_t value_index = 0;
+ const auto* value_schema = find_child(child_schema, "value",
&value_index);
+ const auto& child_struct = assert_cast<const
ColumnStruct&>(*child.column);
+ const bool value_present = value_schema != nullptr &&
+
!cell_at(child_struct.get_column(value_index), row).is_null;
+ size_t typed_index = 0;
+ const auto* typed_schema = find_child(child_schema, "typed_value",
&typed_index);
+ const bool typed_present = typed_schema != nullptr &&
+
!cell_at(child_struct.get_column(typed_index), row).is_null;
+ if (!value_present && !typed_present) {
+ continue;
+ }
+ object.add_key(StringRef(child_schema.name));
+ (void)append_wrapper(child_schema, *child.column, row, metadata,
builder,
+ WrapperContext::OBJECT_FIELD);
+ }
+ object.finish();
+ return;
+ }
+ case ParquetColumnSchemaKind::LIST: {
+ if (static_cast<bool>(residual)) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant array typed_value cannot have
residual value bytes");
+ }
+ if (schema.children.size() != 1) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant array {} has invalid element
schema", schema.name);
+ }
+ const auto& array = assert_cast<const ColumnArray&>(column);
+ const size_t begin = array.offset_at(static_cast<ssize_t>(row));
+ const size_t end = array.get_offsets()[row];
+ auto scope = builder.start_array();
+ for (size_t element = begin; element < end; ++element) {
+ const Cell cell = cell_at(array.get_data(), element);
+ if (cell.is_null) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant shredded array element
wrapper is null");
+ }
+ (void)append_wrapper(*schema.children[0], *cell.column, element,
metadata, builder,
+ WrapperContext::ARRAY_ELEMENT);
+ }
+ scope.finish();
+ return;
+ }
+ case ParquetColumnSchemaKind::MAP:
+ case ParquetColumnSchemaKind::VARIANT:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
typed_value schema {}",
+ schema.name);
+ }
+}
+
+bool append_wrapper(const ParquetColumnSchema& schema, const IColumn& wrapper,
size_t row,
+ VariantMetadataRef metadata, VariantBatchBuilder::Row&
builder,
+ WrapperContext context) {
+ Cell value;
+ if (find_child(schema, "value", nullptr) != nullptr) {
+ value = struct_child_at(schema, wrapper, row, "value", nullptr);
+ } else {
+ value.is_null = true;
+ }
+ const ParquetColumnSchema* typed_schema = nullptr;
+ Cell typed;
+ if (find_child(schema, "typed_value", nullptr) != nullptr) {
+ typed = struct_child_at(schema, wrapper, row, "typed_value",
&typed_schema);
+ } else {
+ typed.is_null = true;
+ }
+
+ if (find_child(schema, "value", nullptr) == nullptr && typed_schema ==
nullptr) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant wrapper {} has neither value nor
typed_value",
+ schema.name);
+ }
+ if (value.is_null && typed.is_null) {
+ if (context == WrapperContext::OBJECT_FIELD) {
+ return false;
+ }
+ if (context == WrapperContext::ARRAY_ELEMENT) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant array
element is missing");
+ }
+ builder.add_null();
+ return true;
+ }
+
+ VariantRef residual {.metadata = metadata, .value = {}};
+ if (!value.is_null) {
+ residual.value = value.column->get_data_at(row);
+ }
+ if (typed.is_null) {
+ builder.add_value(residual);
+ return true;
+ }
+ append_typed_value(*typed_schema, *typed.column, row, metadata,
+ value.is_null ? nullptr : &residual, builder);
+ return true;
+}
+
+void encode_variant_range(const ParquetColumnSchema& schema, const IColumn&
wrapper,
+ const ColumnNullable* outer_nullable, size_t begin,
size_t end,
+ ColumnVariantV2& variants) {
+ try {
+ VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows =
end - begin});
+ for (size_t row = begin; row < end; ++row) {
+ auto output_row = builder.begin_row();
+ if (outer_nullable != nullptr &&
outer_nullable->get_null_map_data()[row] != 0) {
+ output_row.add_null();
+ output_row.finish();
+ continue;
+ }
+ const Cell metadata_cell = struct_child_at(schema, wrapper, row,
"metadata", nullptr);
+ if (metadata_cell.is_null) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant {} has null metadata at row
{}", schema.name, row);
+ }
+ const StringRef metadata_bytes =
metadata_cell.column->get_data_at(row);
+ const VariantMetadataRef metadata {metadata_bytes.data,
metadata_bytes.size};
+ metadata.validate();
+ (void)append_wrapper(schema, wrapper, row, metadata, output_row,
WrapperContext::ROOT);
+ output_row.finish();
+ }
+ VariantBatchBuilder batch = builder.finish_batch();
+ variants.insert_encoded_batch(batch);
+ } catch (...) {
+ if (end - begin <= 1) {
+ throw;
+ }
+ // A single builder has one metadata dictionary. If heterogeneous file
rows cannot fit in
+ // that dictionary, split without changing the destination column's
already-valid batches.
+ // Corrupt input still reaches a one-row range and propagates its
original exception.
+ const size_t middle = begin + (end - begin) / 2;
+ encode_variant_range(schema, wrapper, outer_nullable, begin, middle,
variants);
+ encode_variant_range(schema, wrapper, outer_nullable, middle, end,
variants);
+ }
+}
+
+ColumnVariantV2::MutablePtr encode_variant_column(const ParquetColumnSchema&
schema,
+ const IColumn& physical) {
+ if (schema.kind != ParquetColumnSchemaKind::VARIANT) {
+ throw Exception(ErrorCode::INVALID_ARGUMENT, "Parquet column {} is not
Variant",
+ schema.name);
+ }
+ const auto* outer_nullable =
check_and_get_column<ColumnNullable>(physical);
+ const IColumn& wrapper =
+ outer_nullable == nullptr ? physical :
outer_nullable->get_nested_column();
+ const auto& structure = assert_cast<const ColumnStruct&>(wrapper);
+ if (structure.tuple_size() != schema.children.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} physical
field count mismatch",
+ schema.name);
+ }
+
+ auto variants = ColumnVariantV2::create();
+ constexpr size_t MAX_RECONSTRUCTION_BATCH_ROWS = 4096;
+ for (size_t begin = 0; begin < physical.size(); begin +=
MAX_RECONSTRUCTION_BATCH_ROWS) {
+ encode_variant_range(schema, wrapper, outer_nullable, begin,
+ std::min(physical.size(), begin +
MAX_RECONSTRUCTION_BATCH_ROWS),
+ *variants);
+ }
+ return variants;
+}
+
+std::unique_ptr<ParquetColumnSchema> clone_schema(const ParquetColumnSchema&
source) {
+ auto result = std::make_unique<ParquetColumnSchema>();
+ result->local_id = source.local_id;
+ result->parquet_field_id = source.parquet_field_id;
+ result->name = source.name;
+ result->type = source.type;
+ result->variant_physical_type = source.variant_physical_type;
+ result->leaf_column_id = source.leaf_column_id;
+ result->type_descriptor = source.type_descriptor;
+ result->kind = source.kind;
+ result->max_definition_level = source.max_definition_level;
+ result->max_repetition_level = source.max_repetition_level;
+ result->nullable_definition_level = source.nullable_definition_level;
+ result->definition_level = source.definition_level;
+ result->repetition_level = source.repetition_level;
+ result->repeated_ancestor_definition_level =
source.repeated_ancestor_definition_level;
+ result->repeated_repetition_level = source.repeated_repetition_level;
+ result->children.reserve(source.children.size());
+ for (const auto& child : source.children) {
+ result->children.push_back(clone_schema(*child));
+ }
+ return result;
+}
+
+ColumnPtr unwrap_nullable(ColumnPtr column) {
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(*column)) {
+ return nullable->get_nested_column_ptr();
+ }
+ return column;
+}
+
+ColumnPtr struct_child(const ParquetColumnSchema& schema, ColumnPtr column,
std::string_view name,
+ const ParquetColumnSchema** child_schema) {
+ column = unwrap_nullable(std::move(column));
+ const auto* structure = check_and_get_column<ColumnStruct>(*column);
+ if (structure == nullptr) {
+ return nullptr;
+ }
+ size_t index = 0;
+ const auto* child = find_child(schema, name, &index);
+ if (child == nullptr || index >= structure->tuple_size()) {
+ return nullptr;
+ }
+ if (child_schema != nullptr) {
+ *child_schema = child;
+ }
+ return structure->get_column_ptr(index);
+}
+
+bool has_present_value(const ColumnPtr& column) {
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(*column)) {
+ return std::ranges::any_of(nullable->get_null_map_data(),
+ [](uint8_t is_null) { return is_null == 0;
});
+ }
+ return !column->empty();
+}
+
+class ParquetVariantShreddedState final : public VariantShreddedState {
+public:
+ ParquetVariantShreddedState(const ParquetColumnSchema& schema, ColumnPtr
physical)
+ : _schema(clone_schema(schema)), _physical(std::move(physical)) {
+ DORIS_CHECK(static_cast<bool>(_physical));
+ const ColumnPtr wrapper = unwrap_nullable(_physical);
+ const auto* structure = check_and_get_column<ColumnStruct>(*wrapper);
+ if (structure == nullptr || structure->tuple_size() !=
_schema->children.size()) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant {} physical field count
mismatch", _schema->name);
+ }
+ }
+
+ size_t size() const override { return _physical->size(); }
+ size_t byte_size() const override { return _physical->byte_size(); }
+ size_t allocated_bytes() const override { return
_physical->allocated_bytes(); }
Review Comment:
[P1] Include the cached canonical column in memory accounting.
`materialized_column()` keeps `_materialized` beside `_physical` for the
lifetime of this shared state, but these methods continue to report only the
physical tree. After a root cast, hash, grouping, or serialization forces
reconstruction, queues, exchange buffers, join/aggregation spill thresholds,
operator counters, and adaptive batch sizing can understate the live block by
roughly another full Variant representation. Please account for the cached
column under the same synchronization (without double-counting shared buffers),
or release/replace the physical representation once canonical bytes are
retained. A test should force materialization and verify both `byte_size()` and
`allocated_bytes()`.
##########
be/src/exec/operator/file_scan_operator.cpp:
##########
@@ -152,6 +181,16 @@ Status
FileScanLocalState::_init_scanners(std::list<ScannerSPtr>* scanners) {
const bool use_file_scanner_v2 =
_should_use_file_scanner_v2(state()->query_options(), is_load,
*scan_params);
_operator_profile->add_info_string("UseScannerV2", use_file_scanner_v2 ?
"true" : "false");
+ const auto* output_tuple_desc =
state()->desc_tbl().get_tuple_descriptor(_output_tuple_id);
+ DORIS_CHECK(output_tuple_desc != nullptr);
+ if (!is_load && !use_file_scanner_v2 &&
Review Comment:
[P1] Exclude a retained COUNT(*) placeholder from the legacy-scanner guard.
An explicitly empty `push_down_count_slot_ids` list identifies COUNT(*), but
Nereids may retain one minimum-width output slot only as a row-shaped
placeholder. This check runs before ScannerV1 constructs its reader and rejects
that slot if it is VARIANT, even though V1's `CountReader` later emits
table/file-metadata row counts and never consumes the retained value. Thus
`SELECT COUNT(*)` on a no-delete Iceberg table whose only column is VARIANT
fails whenever ScannerV2 is disabled. This is separate from the non-Parquet
ScannerV2 guard. Please derive the rejection from semantic, non-placeholder
projections while continuing to reject COUNT(v) and real Variant output, and
add a ScannerV1 COUNT(*) case for a Variant-only table.
##########
be/src/format_v2/parquet/reader/variant_column_reader.cpp:
##########
@@ -0,0 +1,818 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include "format_v2/parquet/reader/variant_column_reader.h"
+
+#include <algorithm>
+#include <array>
+#include <cmath>
+#include <cstdint>
+#include <cstring>
+#include <limits>
+#include <mutex>
+#include <optional>
+#include <string_view>
+
+#include "common/exception.h"
+#include "core/assert_cast.h"
+#include "core/column/column_array.h"
+#include "core/column/column_decimal.h"
+#include "core/column/column_map.h"
+#include "core/column/column_nullable.h"
+#include "core/column/column_struct.h"
+#include "core/column/column_vector.h"
+#include "core/column/variant_v2/column_variant_v2.h"
+#include "core/column/variant_v2/column_variant_v2_typed_column.h"
+#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_variant_v2.h"
+#include "core/value/variant/variant_batch_builder.h"
+#include "core/value/variant/variant_metadata.h"
+#include "format_v2/parquet/parquet_column_schema.h"
+
+namespace doris::format::parquet {
+namespace {
+
+struct Cell {
+ const IColumn* column = nullptr;
+ bool is_null = false;
+};
+
+Cell cell_at(const IColumn& column, size_t row) {
+ if (row >= column.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant row {} exceeds
column size {}", row,
+ column.size());
+ }
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(column)) {
+ return {.column = &nullable->get_nested_column(),
+ .is_null = nullable->get_null_map_data()[row] != 0};
+ }
+ return {.column = &column, .is_null = false};
+}
+
+const ParquetColumnSchema* find_child(const ParquetColumnSchema& schema,
std::string_view name,
+ size_t* index) {
+ for (size_t i = 0; i < schema.children.size(); ++i) {
+ if (schema.children[i]->name == name) {
+ if (index != nullptr) {
+ *index = i;
+ }
+ return schema.children[i].get();
+ }
+ }
+ return nullptr;
+}
+
+Cell struct_child_at(const ParquetColumnSchema& schema, const IColumn&
physical, size_t row,
+ std::string_view name, const ParquetColumnSchema**
child_schema) {
+ const auto& structure = assert_cast<const ColumnStruct&>(physical);
+ size_t index = 0;
+ const auto* child = find_child(schema, name, &index);
+ if (child == nullptr || index >= structure.tuple_size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} has no
physical child {}",
+ schema.name, name);
+ }
+ if (child_schema != nullptr) {
+ *child_schema = child;
+ }
+ return cell_at(structure.get_column(index), row);
+}
+
+uint8_t decimal_width(int precision) {
+ if (precision <= 0 || precision > 38) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant decimal precision {} is outside [1,
38]", precision);
+ }
+ return precision <= 9 ? 4 : (precision <= 18 ? 8 : 16);
+}
+
+uint8_t integer_width(const ParquetColumnSchema& schema, PrimitiveType type) {
+ if (schema.type_descriptor.is_unsigned_integer) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Unsigned integers are not valid Parquet Variant typed
values");
+ }
+ if (schema.type_descriptor.integer_bit_width > 0) {
+ switch (schema.type_descriptor.integer_bit_width) {
+ case 8:
+ return 1;
+ case 16:
+ return 2;
+ case 32:
+ return 4;
+ case 64:
+ return 8;
+ default:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
integer width {}",
+ schema.type_descriptor.integer_bit_width);
+ }
+ }
+ switch (type) {
+ case TYPE_TINYINT:
+ return 1;
+ case TYPE_SMALLINT:
+ return 2;
+ case TYPE_INT:
+ return 4;
+ case TYPE_BIGINT:
+ return 8;
+ default:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
integer type {}", type);
+ }
+}
+
+void append_typed_scalar(const ParquetColumnSchema& schema, const IColumn&
column, size_t row,
+ VariantBatchBuilder::Row& builder) {
+ const PrimitiveType type =
remove_nullable(schema.type)->get_primitive_type();
+ switch (type) {
+ case TYPE_BOOLEAN:
+ builder.add_bool(assert_cast<const
ColumnUInt8&>(column).get_data()[row] != 0);
+ return;
+ case TYPE_TINYINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt8&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_SMALLINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt16&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_INT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt32&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_BIGINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt64&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_FLOAT:
+ builder.add_float(assert_cast<const
ColumnFloat32&>(column).get_data()[row]);
+ return;
+ case TYPE_DOUBLE:
+ builder.add_double(assert_cast<const
ColumnFloat64&>(column).get_data()[row]);
+ return;
+ case TYPE_DECIMAL128I: {
+ const auto value = assert_cast<const
ColumnDecimal128V3&>(column).get_data()[row].value;
+ builder.add_decimal(value,
static_cast<uint8_t>(schema.type_descriptor.decimal_scale),
+
decimal_width(schema.type_descriptor.decimal_precision));
+ return;
+ }
+ case TYPE_TIMEV2: {
+ const double seconds = assert_cast<const
ColumnTimeV2&>(column).get_data()[row];
+ if (!std::isfinite(seconds) ||
+ std::abs(seconds) >
static_cast<double>(std::numeric_limits<int64_t>::max()) / 1e6) {
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
TIME value");
+ }
+ builder.add_time_ntz_micros(static_cast<int64_t>(std::llround(seconds
* 1e6)));
+ return;
+ }
+ case TYPE_DATETIMEV2: {
+ if (schema.type_descriptor.time_unit == ParquetTimeUnit::NANOS) {
+ // Native DATETIMEV2 is microsecond based. Reject before returning
a silently truncated
+ // value; a raw INT64 nanos decoder can be added independently.
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant TIMESTAMP(NANOS) is not
supported");
+ }
+ const auto& value = assert_cast<const
ColumnDateTimeV2&>(column).get_data()[row];
+ builder.add_timestamp_micros(
+ variant_timestamp_micros(value, row, "Parquet Variant
TIMESTAMP"),
+ schema.type_descriptor.timestamp_is_adjusted_to_utc);
+ return;
+ }
+ case TYPE_TIMESTAMPTZ: {
+ if (schema.type_descriptor.time_unit == ParquetTimeUnit::NANOS) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant TIMESTAMP(NANOS) is not
supported");
+ }
+ const auto& value = assert_cast<const
ColumnTimeStampTz&>(column).get_data()[row];
+ builder.add_timestamp_micros(
+ variant_timestamp_micros(value, row, "Parquet Variant
TIMESTAMP"), true);
+ return;
+ }
+ case TYPE_VARBINARY: {
+ const StringRef value = column.get_data_at(row);
+ if (!schema.type_descriptor.is_uuid) {
+ builder.add_binary(value);
+ return;
+ }
+ if (value.size != 16) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant UUID has {} bytes instead of 16",
value.size);
+ }
+ std::array<uint8_t, 16> uuid {};
+ std::memcpy(uuid.data(), value.data, uuid.size());
+ builder.add_uuid(uuid);
+ return;
+ }
+ case TYPE_STRING: {
+ const StringRef value = column.get_data_at(row);
+ if (schema.type_descriptor.is_uuid) {
+ if (value.size != 16) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant UUID has {} bytes instead of
16", value.size);
+ }
+ std::array<uint8_t, 16> uuid {};
+ std::memcpy(uuid.data(), value.data, uuid.size());
+ builder.add_uuid(uuid);
+ } else if (schema.type_descriptor.is_string_annotation) {
+ builder.add_string(value);
+ } else {
+ builder.add_binary(value);
+ }
+ return;
+ }
+ default:
+ if (!is_supported_variant_typed_identity(type)) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant typed value {} is not supported",
+ remove_nullable(schema.type)->get_name());
+ }
+ dispatch_variant_typed_column(
+ column, type, [&]<PrimitiveType Type>(const auto&
typed_column) {
+ with_variant_typed_scalar<Type>(
+ typed_column, row,
+
static_cast<uint8_t>(remove_nullable(schema.type)->get_scale()),
+ [&](const VariantScalarRef& scalar) {
builder.add_scalar(scalar); });
+ });
+ }
+}
+
+enum class WrapperContext { ROOT, ARRAY_ELEMENT, OBJECT_FIELD };
+
+bool append_wrapper(const ParquetColumnSchema& schema, const IColumn& wrapper,
size_t row,
+ VariantMetadataRef metadata, VariantBatchBuilder::Row&
builder,
+ WrapperContext context);
+
+void append_typed_value(const ParquetColumnSchema& schema, const IColumn&
column, size_t row,
+ VariantMetadataRef metadata, const VariantRef*
residual,
+ VariantBatchBuilder::Row& builder) {
+ switch (schema.kind) {
+ case ParquetColumnSchemaKind::PRIMITIVE:
+ if (static_cast<bool>(residual)) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant scalar typed_value cannot have
residual value bytes");
+ }
+ append_typed_scalar(schema, column, row, builder);
+ return;
+ case ParquetColumnSchemaKind::STRUCT: {
+ if (static_cast<bool>(residual) && residual->basic_type() !=
VariantBasicType::OBJECT) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant object typed_value has non-object
residual value");
+ }
+ const auto& structure = assert_cast<const ColumnStruct&>(column);
+ if (structure.tuple_size() != schema.children.size()) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant object {} physical field count
mismatch", schema.name);
+ }
+ auto object = builder.start_object();
+ if (static_cast<bool>(residual)) {
+ for (uint32_t i = 0; i < residual->num_elements(); ++i) {
+ uint32_t field_id = 0;
+ const VariantRef child = residual->object_value_at(i,
&field_id);
+ object.add_key(residual->metadata.key_at(field_id));
+ builder.add_value(child);
+ }
+ }
+ for (size_t i = 0; i < schema.children.size(); ++i) {
+ const auto& child_schema = *schema.children[i];
+ const Cell child = cell_at(structure.get_column(i), row);
+ if (child.is_null) {
+ // Shredded object fields are optional wrapper groups. A
missing group means the
+ // key is absent, which differs from a present wrapper
encoding a Variant null.
+ continue;
+ }
+ // A null/null wrapper means this object field is absent. Delay
add_key until its
+ // presence is known so absent shredded fields do not turn into
Variant nulls.
+ size_t value_index = 0;
+ const auto* value_schema = find_child(child_schema, "value",
&value_index);
+ const auto& child_struct = assert_cast<const
ColumnStruct&>(*child.column);
+ const bool value_present = value_schema != nullptr &&
+
!cell_at(child_struct.get_column(value_index), row).is_null;
+ size_t typed_index = 0;
+ const auto* typed_schema = find_child(child_schema, "typed_value",
&typed_index);
+ const bool typed_present = typed_schema != nullptr &&
+
!cell_at(child_struct.get_column(typed_index), row).is_null;
+ if (!value_present && !typed_present) {
+ continue;
+ }
+ object.add_key(StringRef(child_schema.name));
+ (void)append_wrapper(child_schema, *child.column, row, metadata,
builder,
+ WrapperContext::OBJECT_FIELD);
+ }
+ object.finish();
+ return;
+ }
+ case ParquetColumnSchemaKind::LIST: {
+ if (static_cast<bool>(residual)) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant array typed_value cannot have
residual value bytes");
+ }
+ if (schema.children.size() != 1) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant array {} has invalid element
schema", schema.name);
+ }
+ const auto& array = assert_cast<const ColumnArray&>(column);
+ const size_t begin = array.offset_at(static_cast<ssize_t>(row));
+ const size_t end = array.get_offsets()[row];
+ auto scope = builder.start_array();
+ for (size_t element = begin; element < end; ++element) {
+ const Cell cell = cell_at(array.get_data(), element);
+ if (cell.is_null) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant shredded array element
wrapper is null");
+ }
+ (void)append_wrapper(*schema.children[0], *cell.column, element,
metadata, builder,
+ WrapperContext::ARRAY_ELEMENT);
+ }
+ scope.finish();
+ return;
+ }
+ case ParquetColumnSchemaKind::MAP:
+ case ParquetColumnSchemaKind::VARIANT:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
typed_value schema {}",
+ schema.name);
+ }
+}
+
+bool append_wrapper(const ParquetColumnSchema& schema, const IColumn& wrapper,
size_t row,
+ VariantMetadataRef metadata, VariantBatchBuilder::Row&
builder,
+ WrapperContext context) {
+ Cell value;
+ if (find_child(schema, "value", nullptr) != nullptr) {
+ value = struct_child_at(schema, wrapper, row, "value", nullptr);
+ } else {
+ value.is_null = true;
+ }
+ const ParquetColumnSchema* typed_schema = nullptr;
+ Cell typed;
+ if (find_child(schema, "typed_value", nullptr) != nullptr) {
+ typed = struct_child_at(schema, wrapper, row, "typed_value",
&typed_schema);
+ } else {
+ typed.is_null = true;
+ }
+
+ if (find_child(schema, "value", nullptr) == nullptr && typed_schema ==
nullptr) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant wrapper {} has neither value nor
typed_value",
+ schema.name);
+ }
+ if (value.is_null && typed.is_null) {
+ if (context == WrapperContext::OBJECT_FIELD) {
+ return false;
+ }
+ if (context == WrapperContext::ARRAY_ELEMENT) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant array
element is missing");
+ }
+ builder.add_null();
+ return true;
+ }
+
+ VariantRef residual {.metadata = metadata, .value = {}};
+ if (!value.is_null) {
+ residual.value = value.column->get_data_at(row);
+ }
+ if (typed.is_null) {
+ builder.add_value(residual);
+ return true;
+ }
+ append_typed_value(*typed_schema, *typed.column, row, metadata,
+ value.is_null ? nullptr : &residual, builder);
+ return true;
+}
+
+void encode_variant_range(const ParquetColumnSchema& schema, const IColumn&
wrapper,
+ const ColumnNullable* outer_nullable, size_t begin,
size_t end,
+ ColumnVariantV2& variants) {
+ try {
+ VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows =
end - begin});
+ for (size_t row = begin; row < end; ++row) {
+ auto output_row = builder.begin_row();
+ if (outer_nullable != nullptr &&
outer_nullable->get_null_map_data()[row] != 0) {
+ output_row.add_null();
+ output_row.finish();
+ continue;
+ }
+ const Cell metadata_cell = struct_child_at(schema, wrapper, row,
"metadata", nullptr);
+ if (metadata_cell.is_null) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant {} has null metadata at row
{}", schema.name, row);
+ }
+ const StringRef metadata_bytes =
metadata_cell.column->get_data_at(row);
+ const VariantMetadataRef metadata {metadata_bytes.data,
metadata_bytes.size};
+ metadata.validate();
+ (void)append_wrapper(schema, wrapper, row, metadata, output_row,
WrapperContext::ROOT);
+ output_row.finish();
+ }
+ VariantBatchBuilder batch = builder.finish_batch();
+ variants.insert_encoded_batch(batch);
+ } catch (...) {
+ if (end - begin <= 1) {
+ throw;
+ }
+ // A single builder has one metadata dictionary. If heterogeneous file
rows cannot fit in
+ // that dictionary, split without changing the destination column's
already-valid batches.
+ // Corrupt input still reaches a one-row range and propagates its
original exception.
+ const size_t middle = begin + (end - begin) / 2;
+ encode_variant_range(schema, wrapper, outer_nullable, begin, middle,
variants);
+ encode_variant_range(schema, wrapper, outer_nullable, middle, end,
variants);
+ }
+}
+
+ColumnVariantV2::MutablePtr encode_variant_column(const ParquetColumnSchema&
schema,
+ const IColumn& physical) {
+ if (schema.kind != ParquetColumnSchemaKind::VARIANT) {
+ throw Exception(ErrorCode::INVALID_ARGUMENT, "Parquet column {} is not
Variant",
+ schema.name);
+ }
+ const auto* outer_nullable =
check_and_get_column<ColumnNullable>(physical);
+ const IColumn& wrapper =
+ outer_nullable == nullptr ? physical :
outer_nullable->get_nested_column();
+ const auto& structure = assert_cast<const ColumnStruct&>(wrapper);
+ if (structure.tuple_size() != schema.children.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} physical
field count mismatch",
+ schema.name);
+ }
+
+ auto variants = ColumnVariantV2::create();
+ constexpr size_t MAX_RECONSTRUCTION_BATCH_ROWS = 4096;
+ for (size_t begin = 0; begin < physical.size(); begin +=
MAX_RECONSTRUCTION_BATCH_ROWS) {
+ encode_variant_range(schema, wrapper, outer_nullable, begin,
+ std::min(physical.size(), begin +
MAX_RECONSTRUCTION_BATCH_ROWS),
+ *variants);
+ }
+ return variants;
+}
+
+std::unique_ptr<ParquetColumnSchema> clone_schema(const ParquetColumnSchema&
source) {
+ auto result = std::make_unique<ParquetColumnSchema>();
+ result->local_id = source.local_id;
+ result->parquet_field_id = source.parquet_field_id;
+ result->name = source.name;
+ result->type = source.type;
+ result->variant_physical_type = source.variant_physical_type;
+ result->leaf_column_id = source.leaf_column_id;
+ result->type_descriptor = source.type_descriptor;
+ result->kind = source.kind;
+ result->max_definition_level = source.max_definition_level;
+ result->max_repetition_level = source.max_repetition_level;
+ result->nullable_definition_level = source.nullable_definition_level;
+ result->definition_level = source.definition_level;
+ result->repetition_level = source.repetition_level;
+ result->repeated_ancestor_definition_level =
source.repeated_ancestor_definition_level;
+ result->repeated_repetition_level = source.repeated_repetition_level;
+ result->children.reserve(source.children.size());
+ for (const auto& child : source.children) {
+ result->children.push_back(clone_schema(*child));
+ }
+ return result;
+}
+
+ColumnPtr unwrap_nullable(ColumnPtr column) {
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(*column)) {
+ return nullable->get_nested_column_ptr();
+ }
+ return column;
+}
+
+ColumnPtr struct_child(const ParquetColumnSchema& schema, ColumnPtr column,
std::string_view name,
+ const ParquetColumnSchema** child_schema) {
+ column = unwrap_nullable(std::move(column));
+ const auto* structure = check_and_get_column<ColumnStruct>(*column);
+ if (structure == nullptr) {
+ return nullptr;
+ }
+ size_t index = 0;
+ const auto* child = find_child(schema, name, &index);
+ if (child == nullptr || index >= structure->tuple_size()) {
+ return nullptr;
+ }
+ if (child_schema != nullptr) {
+ *child_schema = child;
+ }
+ return structure->get_column_ptr(index);
+}
+
+bool has_present_value(const ColumnPtr& column) {
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(*column)) {
+ return std::ranges::any_of(nullable->get_null_map_data(),
+ [](uint8_t is_null) { return is_null == 0;
});
+ }
+ return !column->empty();
+}
+
+class ParquetVariantShreddedState final : public VariantShreddedState {
+public:
+ ParquetVariantShreddedState(const ParquetColumnSchema& schema, ColumnPtr
physical)
+ : _schema(clone_schema(schema)), _physical(std::move(physical)) {
+ DORIS_CHECK(static_cast<bool>(_physical));
+ const ColumnPtr wrapper = unwrap_nullable(_physical);
+ const auto* structure = check_and_get_column<ColumnStruct>(*wrapper);
+ if (structure == nullptr || structure->tuple_size() !=
_schema->children.size()) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant {} physical field count
mismatch", _schema->name);
+ }
+ }
+
+ size_t size() const override { return _physical->size(); }
+ size_t byte_size() const override { return _physical->byte_size(); }
+ size_t allocated_bytes() const override { return
_physical->allocated_bytes(); }
+ void sanity_check() const override { _physical->sanity_check(); }
+
+ void for_each_subcolumn(const IColumn::ColumnCallback& callback) const
override {
+ callback(*_physical);
+ }
+
+ std::optional<VariantShreddedTypedValue> find_typed_value(
+ std::span<const VariantShreddedPathSegment> path) const override {
+ if (path.empty()) {
+ return std::nullopt;
+ }
+
+ const ParquetColumnSchema* typed_schema = nullptr;
+ ColumnPtr typed = struct_child(*_schema, _physical, "typed_value",
&typed_schema);
+ if (!typed || typed_schema->kind != ParquetColumnSchemaKind::STRUCT) {
+ return std::nullopt;
+ }
+
+ for (size_t position = 0; position < path.size(); ++position) {
+ if (path[position].kind !=
VariantShreddedPathSegment::Kind::OBJECT_KEY) {
+ return std::nullopt;
+ }
+
+ const std::string_view key(path[position].key.data,
path[position].key.size);
+ const ParquetColumnSchema* wrapper_schema = nullptr;
+ ColumnPtr wrapper = struct_child(*typed_schema, typed, key,
&wrapper_schema);
+ if (!wrapper) {
+ return std::nullopt;
+ }
+
+ if (ColumnPtr residual = struct_child(*wrapper_schema, wrapper,
"value", nullptr);
+ static_cast<bool>(residual) && has_present_value(residual)) {
+ // A residual value can contribute data to the same logical
object. Reconstructing
+ // is required in that case; returning only the typed leaf
would drop information.
+ return std::nullopt;
+ }
+
+ typed = struct_child(*wrapper_schema, wrapper, "typed_value",
&typed_schema);
+ if (!typed) {
+ return std::nullopt;
+ }
+ if (position + 1 == path.size()) {
+ if (typed_schema->kind != ParquetColumnSchemaKind::PRIMITIVE ||
+ check_and_get_column<ColumnNullable>(*typed) == nullptr) {
+ return std::nullopt;
+ }
+ return VariantShreddedTypedValue {.column = std::move(typed),
+ .type =
remove_nullable(typed_schema->type)};
+ }
+ if (typed_schema->kind != ParquetColumnSchemaKind::STRUCT) {
+ return std::nullopt;
+ }
+ }
+ return std::nullopt;
+ }
+
+ const ColumnVariantV2& materialized_column() const override {
+ std::lock_guard lock(_materialization_lock);
+ if (!_materialized) {
+ _materialized = encode_variant_column(*_schema, *_physical);
+ }
+ return *_materialized;
+ }
+
+private:
+ std::unique_ptr<ParquetColumnSchema> _schema;
+ ColumnPtr _physical;
+ mutable std::mutex _materialization_lock;
+ mutable ColumnVariantV2::MutablePtr _materialized;
+};
+
+MutableColumnPtr build_variant_column(const ParquetColumnSchema& schema,
ColumnPtr physical) {
+ if (schema.kind != ParquetColumnSchemaKind::VARIANT) {
+ throw Exception(ErrorCode::INVALID_ARGUMENT, "Parquet column {} is not
Variant",
+ schema.name);
+ }
+
+ const auto* outer_nullable =
check_and_get_column<ColumnNullable>(*physical);
+ MutableColumnPtr variants = ColumnVariantV2::create_shredded(
+ std::make_shared<ParquetVariantShreddedState>(schema, physical));
+ if (outer_nullable == nullptr) {
+ return variants;
+ }
+ auto nulls =
outer_nullable->get_null_map_column().clone_resized(physical->size());
+ return ColumnNullable::create(std::move(variants), std::move(nulls));
+}
+
+MutableColumnPtr transform_node(const VariantMaterializationNode& plan,
ColumnPtr physical);
+
+MutableColumnPtr transform_non_nullable(const VariantMaterializationNode& plan,
Review Comment:
[P2] Avoid three full copies of unaffected nested buffers.
For a Variant inside STRUCT/LIST/MAP, this transform clones every primitive
sibling, enclosing null map, and collection offset.
`append_materialized_column()` then copies that entire converted tree into a
compatible scratch column and copies it again into the real output, performing
three full copy passes and leaving up to three batch representations of each
unaffected buffer live transiently in the scan path. A large non-Variant
sibling therefore pays repeated memory and CPU cost merely because another
child is Variant. Please rebuild only Variant-containing ancestors while moving
or sharing unchanged buffers, and preserve atomic validation without two
additional full-tree appends; add an allocation/copy benchmark for a nested
Variant beside a large sibling.
##########
be/src/format_v2/parquet/reader/variant_column_reader.cpp:
##########
@@ -0,0 +1,818 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include "format_v2/parquet/reader/variant_column_reader.h"
+
+#include <algorithm>
+#include <array>
+#include <cmath>
+#include <cstdint>
+#include <cstring>
+#include <limits>
+#include <mutex>
+#include <optional>
+#include <string_view>
+
+#include "common/exception.h"
+#include "core/assert_cast.h"
+#include "core/column/column_array.h"
+#include "core/column/column_decimal.h"
+#include "core/column/column_map.h"
+#include "core/column/column_nullable.h"
+#include "core/column/column_struct.h"
+#include "core/column/column_vector.h"
+#include "core/column/variant_v2/column_variant_v2.h"
+#include "core/column/variant_v2/column_variant_v2_typed_column.h"
+#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_variant_v2.h"
+#include "core/value/variant/variant_batch_builder.h"
+#include "core/value/variant/variant_metadata.h"
+#include "format_v2/parquet/parquet_column_schema.h"
+
+namespace doris::format::parquet {
+namespace {
+
+struct Cell {
+ const IColumn* column = nullptr;
+ bool is_null = false;
+};
+
+Cell cell_at(const IColumn& column, size_t row) {
+ if (row >= column.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant row {} exceeds
column size {}", row,
+ column.size());
+ }
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(column)) {
+ return {.column = &nullable->get_nested_column(),
+ .is_null = nullable->get_null_map_data()[row] != 0};
+ }
+ return {.column = &column, .is_null = false};
+}
+
+const ParquetColumnSchema* find_child(const ParquetColumnSchema& schema,
std::string_view name,
+ size_t* index) {
+ for (size_t i = 0; i < schema.children.size(); ++i) {
+ if (schema.children[i]->name == name) {
+ if (index != nullptr) {
+ *index = i;
+ }
+ return schema.children[i].get();
+ }
+ }
+ return nullptr;
+}
+
+Cell struct_child_at(const ParquetColumnSchema& schema, const IColumn&
physical, size_t row,
+ std::string_view name, const ParquetColumnSchema**
child_schema) {
+ const auto& structure = assert_cast<const ColumnStruct&>(physical);
+ size_t index = 0;
+ const auto* child = find_child(schema, name, &index);
+ if (child == nullptr || index >= structure.tuple_size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} has no
physical child {}",
+ schema.name, name);
+ }
+ if (child_schema != nullptr) {
+ *child_schema = child;
+ }
+ return cell_at(structure.get_column(index), row);
+}
+
+uint8_t decimal_width(int precision) {
+ if (precision <= 0 || precision > 38) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant decimal precision {} is outside [1,
38]", precision);
+ }
+ return precision <= 9 ? 4 : (precision <= 18 ? 8 : 16);
+}
+
+uint8_t integer_width(const ParquetColumnSchema& schema, PrimitiveType type) {
+ if (schema.type_descriptor.is_unsigned_integer) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Unsigned integers are not valid Parquet Variant typed
values");
+ }
+ if (schema.type_descriptor.integer_bit_width > 0) {
+ switch (schema.type_descriptor.integer_bit_width) {
+ case 8:
+ return 1;
+ case 16:
+ return 2;
+ case 32:
+ return 4;
+ case 64:
+ return 8;
+ default:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
integer width {}",
+ schema.type_descriptor.integer_bit_width);
+ }
+ }
+ switch (type) {
+ case TYPE_TINYINT:
+ return 1;
+ case TYPE_SMALLINT:
+ return 2;
+ case TYPE_INT:
+ return 4;
+ case TYPE_BIGINT:
+ return 8;
+ default:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
integer type {}", type);
+ }
+}
+
+void append_typed_scalar(const ParquetColumnSchema& schema, const IColumn&
column, size_t row,
+ VariantBatchBuilder::Row& builder) {
+ const PrimitiveType type =
remove_nullable(schema.type)->get_primitive_type();
+ switch (type) {
+ case TYPE_BOOLEAN:
+ builder.add_bool(assert_cast<const
ColumnUInt8&>(column).get_data()[row] != 0);
+ return;
+ case TYPE_TINYINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt8&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_SMALLINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt16&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_INT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt32&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_BIGINT:
+ builder.add_scalar(
+ VariantScalarRef::integer(assert_cast<const
ColumnInt64&>(column).get_data()[row],
+ integer_width(schema, type)));
+ return;
+ case TYPE_FLOAT:
+ builder.add_float(assert_cast<const
ColumnFloat32&>(column).get_data()[row]);
+ return;
+ case TYPE_DOUBLE:
+ builder.add_double(assert_cast<const
ColumnFloat64&>(column).get_data()[row]);
+ return;
+ case TYPE_DECIMAL128I: {
+ const auto value = assert_cast<const
ColumnDecimal128V3&>(column).get_data()[row].value;
+ builder.add_decimal(value,
static_cast<uint8_t>(schema.type_descriptor.decimal_scale),
+
decimal_width(schema.type_descriptor.decimal_precision));
+ return;
+ }
+ case TYPE_TIMEV2: {
+ const double seconds = assert_cast<const
ColumnTimeV2&>(column).get_data()[row];
+ if (!std::isfinite(seconds) ||
+ std::abs(seconds) >
static_cast<double>(std::numeric_limits<int64_t>::max()) / 1e6) {
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
TIME value");
+ }
+ builder.add_time_ntz_micros(static_cast<int64_t>(std::llround(seconds
* 1e6)));
+ return;
+ }
+ case TYPE_DATETIMEV2: {
+ if (schema.type_descriptor.time_unit == ParquetTimeUnit::NANOS) {
+ // Native DATETIMEV2 is microsecond based. Reject before returning
a silently truncated
+ // value; a raw INT64 nanos decoder can be added independently.
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant TIMESTAMP(NANOS) is not
supported");
+ }
+ const auto& value = assert_cast<const
ColumnDateTimeV2&>(column).get_data()[row];
+ builder.add_timestamp_micros(
+ variant_timestamp_micros(value, row, "Parquet Variant
TIMESTAMP"),
+ schema.type_descriptor.timestamp_is_adjusted_to_utc);
+ return;
+ }
+ case TYPE_TIMESTAMPTZ: {
+ if (schema.type_descriptor.time_unit == ParquetTimeUnit::NANOS) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant TIMESTAMP(NANOS) is not
supported");
+ }
+ const auto& value = assert_cast<const
ColumnTimeStampTz&>(column).get_data()[row];
+ builder.add_timestamp_micros(
+ variant_timestamp_micros(value, row, "Parquet Variant
TIMESTAMP"), true);
+ return;
+ }
+ case TYPE_VARBINARY: {
+ const StringRef value = column.get_data_at(row);
+ if (!schema.type_descriptor.is_uuid) {
+ builder.add_binary(value);
+ return;
+ }
+ if (value.size != 16) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant UUID has {} bytes instead of 16",
value.size);
+ }
+ std::array<uint8_t, 16> uuid {};
+ std::memcpy(uuid.data(), value.data, uuid.size());
+ builder.add_uuid(uuid);
+ return;
+ }
+ case TYPE_STRING: {
+ const StringRef value = column.get_data_at(row);
+ if (schema.type_descriptor.is_uuid) {
+ if (value.size != 16) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant UUID has {} bytes instead of
16", value.size);
+ }
+ std::array<uint8_t, 16> uuid {};
+ std::memcpy(uuid.data(), value.data, uuid.size());
+ builder.add_uuid(uuid);
+ } else if (schema.type_descriptor.is_string_annotation) {
+ builder.add_string(value);
+ } else {
+ builder.add_binary(value);
+ }
+ return;
+ }
+ default:
+ if (!is_supported_variant_typed_identity(type)) {
+ throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
+ "Parquet Variant typed value {} is not supported",
+ remove_nullable(schema.type)->get_name());
+ }
+ dispatch_variant_typed_column(
+ column, type, [&]<PrimitiveType Type>(const auto&
typed_column) {
+ with_variant_typed_scalar<Type>(
+ typed_column, row,
+
static_cast<uint8_t>(remove_nullable(schema.type)->get_scale()),
+ [&](const VariantScalarRef& scalar) {
builder.add_scalar(scalar); });
+ });
+ }
+}
+
+enum class WrapperContext { ROOT, ARRAY_ELEMENT, OBJECT_FIELD };
+
+bool append_wrapper(const ParquetColumnSchema& schema, const IColumn& wrapper,
size_t row,
+ VariantMetadataRef metadata, VariantBatchBuilder::Row&
builder,
+ WrapperContext context);
+
+void append_typed_value(const ParquetColumnSchema& schema, const IColumn&
column, size_t row,
+ VariantMetadataRef metadata, const VariantRef*
residual,
+ VariantBatchBuilder::Row& builder) {
+ switch (schema.kind) {
+ case ParquetColumnSchemaKind::PRIMITIVE:
+ if (static_cast<bool>(residual)) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant scalar typed_value cannot have
residual value bytes");
+ }
+ append_typed_scalar(schema, column, row, builder);
+ return;
+ case ParquetColumnSchemaKind::STRUCT: {
+ if (static_cast<bool>(residual) && residual->basic_type() !=
VariantBasicType::OBJECT) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant object typed_value has non-object
residual value");
+ }
+ const auto& structure = assert_cast<const ColumnStruct&>(column);
+ if (structure.tuple_size() != schema.children.size()) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant object {} physical field count
mismatch", schema.name);
+ }
+ auto object = builder.start_object();
+ if (static_cast<bool>(residual)) {
+ for (uint32_t i = 0; i < residual->num_elements(); ++i) {
+ uint32_t field_id = 0;
+ const VariantRef child = residual->object_value_at(i,
&field_id);
+ object.add_key(residual->metadata.key_at(field_id));
+ builder.add_value(child);
+ }
+ }
+ for (size_t i = 0; i < schema.children.size(); ++i) {
+ const auto& child_schema = *schema.children[i];
+ const Cell child = cell_at(structure.get_column(i), row);
+ if (child.is_null) {
+ // Shredded object fields are optional wrapper groups. A
missing group means the
+ // key is absent, which differs from a present wrapper
encoding a Variant null.
+ continue;
+ }
+ // A null/null wrapper means this object field is absent. Delay
add_key until its
+ // presence is known so absent shredded fields do not turn into
Variant nulls.
+ size_t value_index = 0;
+ const auto* value_schema = find_child(child_schema, "value",
&value_index);
+ const auto& child_struct = assert_cast<const
ColumnStruct&>(*child.column);
+ const bool value_present = value_schema != nullptr &&
+
!cell_at(child_struct.get_column(value_index), row).is_null;
+ size_t typed_index = 0;
+ const auto* typed_schema = find_child(child_schema, "typed_value",
&typed_index);
+ const bool typed_present = typed_schema != nullptr &&
+
!cell_at(child_struct.get_column(typed_index), row).is_null;
+ if (!value_present && !typed_present) {
+ continue;
+ }
+ object.add_key(StringRef(child_schema.name));
+ (void)append_wrapper(child_schema, *child.column, row, metadata,
builder,
+ WrapperContext::OBJECT_FIELD);
+ }
+ object.finish();
+ return;
+ }
+ case ParquetColumnSchemaKind::LIST: {
+ if (static_cast<bool>(residual)) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant array typed_value cannot have
residual value bytes");
+ }
+ if (schema.children.size() != 1) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant array {} has invalid element
schema", schema.name);
+ }
+ const auto& array = assert_cast<const ColumnArray&>(column);
+ const size_t begin = array.offset_at(static_cast<ssize_t>(row));
+ const size_t end = array.get_offsets()[row];
+ auto scope = builder.start_array();
+ for (size_t element = begin; element < end; ++element) {
+ const Cell cell = cell_at(array.get_data(), element);
+ if (cell.is_null) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant shredded array element
wrapper is null");
+ }
+ (void)append_wrapper(*schema.children[0], *cell.column, element,
metadata, builder,
+ WrapperContext::ARRAY_ELEMENT);
+ }
+ scope.finish();
+ return;
+ }
+ case ParquetColumnSchemaKind::MAP:
+ case ParquetColumnSchemaKind::VARIANT:
+ throw Exception(ErrorCode::CORRUPTION, "Invalid Parquet Variant
typed_value schema {}",
+ schema.name);
+ }
+}
+
+bool append_wrapper(const ParquetColumnSchema& schema, const IColumn& wrapper,
size_t row,
+ VariantMetadataRef metadata, VariantBatchBuilder::Row&
builder,
+ WrapperContext context) {
+ Cell value;
+ if (find_child(schema, "value", nullptr) != nullptr) {
+ value = struct_child_at(schema, wrapper, row, "value", nullptr);
+ } else {
+ value.is_null = true;
+ }
+ const ParquetColumnSchema* typed_schema = nullptr;
+ Cell typed;
+ if (find_child(schema, "typed_value", nullptr) != nullptr) {
+ typed = struct_child_at(schema, wrapper, row, "typed_value",
&typed_schema);
+ } else {
+ typed.is_null = true;
+ }
+
+ if (find_child(schema, "value", nullptr) == nullptr && typed_schema ==
nullptr) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant wrapper {} has neither value nor
typed_value",
+ schema.name);
+ }
+ if (value.is_null && typed.is_null) {
+ if (context == WrapperContext::OBJECT_FIELD) {
+ return false;
+ }
+ if (context == WrapperContext::ARRAY_ELEMENT) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant array
element is missing");
+ }
+ builder.add_null();
+ return true;
+ }
+
+ VariantRef residual {.metadata = metadata, .value = {}};
+ if (!value.is_null) {
+ residual.value = value.column->get_data_at(row);
+ }
+ if (typed.is_null) {
+ builder.add_value(residual);
+ return true;
+ }
+ append_typed_value(*typed_schema, *typed.column, row, metadata,
+ value.is_null ? nullptr : &residual, builder);
+ return true;
+}
+
+void encode_variant_range(const ParquetColumnSchema& schema, const IColumn&
wrapper,
+ const ColumnNullable* outer_nullable, size_t begin,
size_t end,
+ ColumnVariantV2& variants) {
+ try {
+ VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows =
end - begin});
+ for (size_t row = begin; row < end; ++row) {
+ auto output_row = builder.begin_row();
+ if (outer_nullable != nullptr &&
outer_nullable->get_null_map_data()[row] != 0) {
+ output_row.add_null();
+ output_row.finish();
+ continue;
+ }
+ const Cell metadata_cell = struct_child_at(schema, wrapper, row,
"metadata", nullptr);
+ if (metadata_cell.is_null) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant {} has null metadata at row
{}", schema.name, row);
+ }
+ const StringRef metadata_bytes =
metadata_cell.column->get_data_at(row);
+ const VariantMetadataRef metadata {metadata_bytes.data,
metadata_bytes.size};
+ metadata.validate();
+ (void)append_wrapper(schema, wrapper, row, metadata, output_row,
WrapperContext::ROOT);
+ output_row.finish();
+ }
+ VariantBatchBuilder batch = builder.finish_batch();
+ variants.insert_encoded_batch(batch);
+ } catch (...) {
+ if (end - begin <= 1) {
+ throw;
+ }
+ // A single builder has one metadata dictionary. If heterogeneous file
rows cannot fit in
+ // that dictionary, split without changing the destination column's
already-valid batches.
+ // Corrupt input still reaches a one-row range and propagates its
original exception.
+ const size_t middle = begin + (end - begin) / 2;
+ encode_variant_range(schema, wrapper, outer_nullable, begin, middle,
variants);
+ encode_variant_range(schema, wrapper, outer_nullable, middle, end,
variants);
+ }
+}
+
+ColumnVariantV2::MutablePtr encode_variant_column(const ParquetColumnSchema&
schema,
+ const IColumn& physical) {
+ if (schema.kind != ParquetColumnSchemaKind::VARIANT) {
+ throw Exception(ErrorCode::INVALID_ARGUMENT, "Parquet column {} is not
Variant",
+ schema.name);
+ }
+ const auto* outer_nullable =
check_and_get_column<ColumnNullable>(physical);
+ const IColumn& wrapper =
+ outer_nullable == nullptr ? physical :
outer_nullable->get_nested_column();
+ const auto& structure = assert_cast<const ColumnStruct&>(wrapper);
+ if (structure.tuple_size() != schema.children.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} physical
field count mismatch",
+ schema.name);
+ }
+
+ auto variants = ColumnVariantV2::create();
+ constexpr size_t MAX_RECONSTRUCTION_BATCH_ROWS = 4096;
+ for (size_t begin = 0; begin < physical.size(); begin +=
MAX_RECONSTRUCTION_BATCH_ROWS) {
+ encode_variant_range(schema, wrapper, outer_nullable, begin,
+ std::min(physical.size(), begin +
MAX_RECONSTRUCTION_BATCH_ROWS),
+ *variants);
+ }
+ return variants;
+}
+
+std::unique_ptr<ParquetColumnSchema> clone_schema(const ParquetColumnSchema&
source) {
+ auto result = std::make_unique<ParquetColumnSchema>();
+ result->local_id = source.local_id;
+ result->parquet_field_id = source.parquet_field_id;
+ result->name = source.name;
+ result->type = source.type;
+ result->variant_physical_type = source.variant_physical_type;
+ result->leaf_column_id = source.leaf_column_id;
+ result->type_descriptor = source.type_descriptor;
+ result->kind = source.kind;
+ result->max_definition_level = source.max_definition_level;
+ result->max_repetition_level = source.max_repetition_level;
+ result->nullable_definition_level = source.nullable_definition_level;
+ result->definition_level = source.definition_level;
+ result->repetition_level = source.repetition_level;
+ result->repeated_ancestor_definition_level =
source.repeated_ancestor_definition_level;
+ result->repeated_repetition_level = source.repeated_repetition_level;
+ result->children.reserve(source.children.size());
+ for (const auto& child : source.children) {
+ result->children.push_back(clone_schema(*child));
+ }
+ return result;
+}
+
+ColumnPtr unwrap_nullable(ColumnPtr column) {
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(*column)) {
+ return nullable->get_nested_column_ptr();
+ }
+ return column;
+}
+
+ColumnPtr struct_child(const ParquetColumnSchema& schema, ColumnPtr column,
std::string_view name,
+ const ParquetColumnSchema** child_schema) {
+ column = unwrap_nullable(std::move(column));
+ const auto* structure = check_and_get_column<ColumnStruct>(*column);
+ if (structure == nullptr) {
+ return nullptr;
+ }
+ size_t index = 0;
+ const auto* child = find_child(schema, name, &index);
+ if (child == nullptr || index >= structure->tuple_size()) {
+ return nullptr;
+ }
+ if (child_schema != nullptr) {
+ *child_schema = child;
+ }
+ return structure->get_column_ptr(index);
+}
+
+bool has_present_value(const ColumnPtr& column) {
+ if (const auto* nullable = check_and_get_column<ColumnNullable>(*column)) {
+ return std::ranges::any_of(nullable->get_null_map_data(),
+ [](uint8_t is_null) { return is_null == 0;
});
+ }
+ return !column->empty();
+}
+
+class ParquetVariantShreddedState final : public VariantShreddedState {
+public:
+ ParquetVariantShreddedState(const ParquetColumnSchema& schema, ColumnPtr
physical)
+ : _schema(clone_schema(schema)), _physical(std::move(physical)) {
+ DORIS_CHECK(static_cast<bool>(_physical));
+ const ColumnPtr wrapper = unwrap_nullable(_physical);
+ const auto* structure = check_and_get_column<ColumnStruct>(*wrapper);
+ if (structure == nullptr || structure->tuple_size() !=
_schema->children.size()) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant {} physical field count
mismatch", _schema->name);
+ }
+ }
+
+ size_t size() const override { return _physical->size(); }
+ size_t byte_size() const override { return _physical->byte_size(); }
+ size_t allocated_bytes() const override { return
_physical->allocated_bytes(); }
+ void sanity_check() const override { _physical->sanity_check(); }
+
+ void for_each_subcolumn(const IColumn::ColumnCallback& callback) const
override {
+ callback(*_physical);
+ }
+
+ std::optional<VariantShreddedTypedValue> find_typed_value(
+ std::span<const VariantShreddedPathSegment> path) const override {
+ if (path.empty()) {
+ return std::nullopt;
+ }
+
+ const ParquetColumnSchema* typed_schema = nullptr;
+ ColumnPtr typed = struct_child(*_schema, _physical, "typed_value",
&typed_schema);
+ if (!typed || typed_schema->kind != ParquetColumnSchemaKind::STRUCT) {
+ return std::nullopt;
+ }
+
+ for (size_t position = 0; position < path.size(); ++position) {
+ if (path[position].kind !=
VariantShreddedPathSegment::Kind::OBJECT_KEY) {
+ return std::nullopt;
+ }
+
+ const std::string_view key(path[position].key.data,
path[position].key.size);
+ const ParquetColumnSchema* wrapper_schema = nullptr;
+ ColumnPtr wrapper = struct_child(*typed_schema, typed, key,
&wrapper_schema);
+ if (!wrapper) {
+ return std::nullopt;
+ }
+
+ if (ColumnPtr residual = struct_child(*wrapper_schema, wrapper,
"value", nullptr);
+ static_cast<bool>(residual) && has_present_value(residual)) {
+ // A residual value can contribute data to the same logical
object. Reconstructing
+ // is required in that case; returning only the typed leaf
would drop information.
+ return std::nullopt;
+ }
+
+ typed = struct_child(*wrapper_schema, wrapper, "typed_value",
&typed_schema);
+ if (!typed) {
+ return std::nullopt;
+ }
+ if (position + 1 == path.size()) {
+ if (typed_schema->kind != ParquetColumnSchemaKind::PRIMITIVE ||
+ check_and_get_column<ColumnNullable>(*typed) == nullptr) {
+ return std::nullopt;
+ }
+ return VariantShreddedTypedValue {.column = std::move(typed),
Review Comment:
[P1] Preserve the Parquet Variant primitive identity on this fast path.
Returning only the decoded Doris type drops descriptor facts that are part
of the logical Variant value: integer/decimal width, raw BINARY versus
STRING/UUID, TIME, and timestamp unit. The caller immediately feeds this pair
to `ColumnVariantV2::create_typed()`, so valid TIME/VARBINARY leaves hit the
unsupported-identity check, mapping-disabled BINARY/UUID serialize as Variant
STRING, and integer or DECIMAL4/8 values can hash/serialize under a different
primitive ID than full reconstruction. The same file can therefore fail or
produce a different canonical Variant depending on whether the query extracts a
path or materializes the root. Please carry an exact Variant primitive
descriptor through the typed state, or decline this fast path for identities it
cannot represent exactly, and cover extraction followed by canonical
hash/serialization for every valid shredded scalar.
--
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]