diff --git a/be/src/core/column/variant_v2/column_variant_v2.cpp b/be/src/core/column/variant_v2/column_variant_v2.cpp index 265e9f0e56dde8..adb7e947b91cc9 100644 --- a/be/src/core/column/variant_v2/column_variant_v2.cpp +++ b/be/src/core/column/variant_v2/column_variant_v2.cpp @@ -22,9 +22,11 @@ #include #include #include +#include #include #include #include +#include #include "common/check.h" #include "common/exception.h" @@ -313,6 +315,307 @@ ValidatedTypedInput validate_typed_input(ColumnPtr column, DataTypePtr scalar_ty "ColumnVariantV2::{} is intentionally unsupported for Variant values", method); } +class CompositeVariantShreddedState final : public VariantShreddedState { +public: + explicit CompositeVariantShreddedState( + std::vector> segments) + : _segments(std::move(segments)) { + DORIS_CHECK(std::ranges::all_of(_segments, [](const auto& segment) { + return segment != nullptr; + })) << "composite Variant shredded segments must not be null"; + for (const auto& segment : _segments) { + const size_t segment_rows = segment->size(); + DORIS_CHECK_LE(segment_rows, std::numeric_limits::max() - _rows) + << "composite Variant shredded row count overflows size_t"; + _rows += segment_rows; + } + } + + size_t size() const override { return _rows; } + + size_t recompute_size() const { + size_t rows = 0; + for (const auto& segment : _segments) { + const size_t segment_rows = segment->size(); + DORIS_CHECK_LE(segment_rows, std::numeric_limits::max() - rows) + << "composite Variant shredded row count overflows size_t"; + rows += segment_rows; + } + return rows; + } + + size_t byte_size() const override { + size_t bytes = 0; + for (const auto& segment : _segments) { + bytes += segment->byte_size(); + } + std::lock_guard lock(_materialization_lock); + return bytes + (_materialized ? _materialized->byte_size() : 0) + + (_serialized ? _serialized->byte_size() : 0); + } + + size_t allocated_bytes() const override { + size_t bytes = 0; + for (const auto& segment : _segments) { + bytes += segment->allocated_bytes(); + } + std::lock_guard lock(_materialization_lock); + return bytes + (_materialized ? _materialized->allocated_bytes() : 0) + + (_serialized ? _serialized->allocated_bytes() : 0); + } + + void sanity_check() const override { + // Row count is read in per-row expression loops, so cache it and reserve the full segment + // walk for invariant checks instead of making extraction quadratic in segment count. + DORIS_CHECK_EQ(recompute_size(), _rows) + << "cached composite Variant shredded row count is stale"; + for (const auto& segment : _segments) { + segment->sanity_check(); + } + } + + void for_each_subcolumn(const IColumn::ImutableColumnCallback& callback) const override { + for (const auto& segment : _segments) { + segment->for_each_subcolumn(callback); + } + } + + std::shared_ptr filter(const IColumn::Filter& filter, + ssize_t /*result_size_hint*/) const override { + DORIS_CHECK_EQ(filter.size(), size()) + << "composite Variant shredded filter size does not match row count"; + std::vector> selected; + selected.reserve(_segments.size()); + size_t offset = 0; + for (const auto& segment : _segments) { + IColumn::Filter segment_filter; + segment_filter.insert(filter.begin() + offset, + filter.begin() + offset + segment->size()); + auto filtered = segment->filter(segment_filter, -1); + if (filtered->size() != 0) { + selected.push_back(std::move(filtered)); + } + offset += segment->size(); + } + return pack(std::move(selected)); + } + + std::shared_ptr select_range(size_t start, size_t length) const override { + DORIS_CHECK_LE(start, size()) << "composite Variant range starts past source size"; + DORIS_CHECK_LE(length, size() - start) << "composite Variant range exceeds source size"; + std::vector> selected; + if (length == 0) { + return pack(std::move(selected)); + } + const size_t end = start + length; + size_t offset = 0; + for (const auto& segment : _segments) { + const size_t segment_end = offset + segment->size(); + const size_t overlap_begin = std::max(start, offset); + const size_t overlap_end = std::min(end, segment_end); + if (overlap_begin < overlap_end) { + selected.push_back( + segment->select_range(overlap_begin - offset, overlap_end - overlap_begin)); + } + offset = segment_end; + if (offset >= end) { + break; + } + } + return pack(std::move(selected)); + } + + std::shared_ptr select_indices( + const uint32_t* indices_begin, const uint32_t* indices_end) const override { + if (indices_begin == indices_end) { + return pack({}); + } + DORIS_CHECK(indices_begin != nullptr && indices_end != nullptr && + indices_begin < indices_end) + << "composite Variant indices are invalid"; + + std::vector segment_ends; + segment_ends.reserve(_segments.size()); + size_t rows = 0; + for (const auto& segment : _segments) { + rows += segment->size(); + segment_ends.push_back(rows); + } + + std::vector> selected; + const uint32_t* cursor = indices_begin; + while (cursor != indices_end) { + DORIS_CHECK_LT(*cursor, rows) << "composite Variant source index is out of range"; + const size_t segment_index = + std::upper_bound(segment_ends.begin(), segment_ends.end(), *cursor) - + segment_ends.begin(); + const size_t segment_begin = segment_index == 0 ? 0 : segment_ends[segment_index - 1]; + DorisVector local_indices; + while (cursor != indices_end && *cursor >= segment_begin && + *cursor < segment_ends[segment_index]) { + local_indices.push_back(static_cast(*cursor - segment_begin)); + ++cursor; + } + selected.push_back(_segments[segment_index]->select_indices( + local_indices.data(), local_indices.data() + local_indices.size())); + } + return pack(std::move(selected)); + } + + bool can_materialize() const override { + return std::ranges::all_of(_segments, + [](const auto& segment) { return segment->can_materialize(); }); + } + + bool try_append(const VariantShreddedState& source) override { + const size_t source_rows = source.size(); + DORIS_CHECK_LE(source_rows, std::numeric_limits::max() - _rows) + << "composite Variant shredded row count overflows size_t"; + if (const auto* composite = dynamic_cast(&source)) { + for (const auto& segment : composite->_segments) { + append(segment); + } + } else { + append(source.select_range(0, source.size())); + } + std::lock_guard lock(_materialization_lock); + _materialized.reset(); + _serialized.reset(); + _rows += source_rows; + return true; + } + + std::optional find_typed_value( + std::span path) const override { + if (_segments.empty()) { + return std::nullopt; + } + std::vector matches; + matches.reserve(_segments.size()); + bool all_direct = true; + for (const auto& segment : _segments) { + auto match = segment->find_typed_value(path); + if (!match.has_value()) { + all_direct = false; + break; + } + matches.push_back(std::move(*match)); + } + + const bool homogeneous = all_direct && matches.front().column && matches.front().type && + std::ranges::all_of(matches, [&](const auto& match) { + return match.column && match.type && !match.normalized && + exact_typed_identity(matches.front().type, match.type); + }); + if (homogeneous) { + MutableColumnPtr combined = matches.front().column->clone_empty(); + for (const auto& match : matches) { + combined->insert_range_from(*match.column, 0, match.column->size()); + } + return VariantShreddedTypedValue {.column = std::move(combined), + .type = matches.front().type, + .normalized = nullptr}; + } + + auto normalized = find_normalized_value(path); + if (!normalized.has_value()) { + return std::nullopt; + } + return VariantShreddedTypedValue { + .column = nullptr, .type = nullptr, .normalized = std::move(*normalized)}; + } + + std::optional find_normalized_value( + std::span path) const override { + auto values = ColumnVariantV2::create(); + auto nulls = ColumnUInt8::create(); + nulls->reserve(size()); + for (const auto& segment : _segments) { + auto normalized = segment->find_normalized_value(path); + if (!normalized.has_value()) { + return std::nullopt; + } + const auto& nullable = assert_cast(**normalized); + const auto& variants = + assert_cast(nullable.get_nested_column()); + values->insert_range_from(variants, 0, variants.size()); + nulls->insert_range_from(nullable.get_null_map_column(), 0, nullable.size()); + } + return ColumnNullable::create(std::move(values), std::move(nulls)); + } + + const ColumnVariantV2& materialized_column() const override { + std::lock_guard lock(_materialization_lock); + if (!_materialized) { + auto materialized = ColumnVariantV2::create(); + for (const auto& segment : _segments) { + const ColumnVariantV2& source = segment->materialized_column(); + materialized->insert_range_from(source, 0, source.size()); + } + _materialized = std::move(materialized); + } + return *_materialized; + } + + const ColumnVariantV2& serialized_column() const override { + std::lock_guard lock(_materialization_lock); + if (!_serialized) { + auto serialized = ColumnVariantV2::create(); + for (const auto& segment : _segments) { + const ColumnVariantV2& source = segment->serialized_column(); + serialized->insert_range_from(source, 0, source.size()); + } + _serialized = std::move(serialized); + } + return *_serialized; + } + +private: + static std::shared_ptr pack( + std::vector> segments) { + if (segments.size() == 1) { + return std::move(segments.front()); + } + return std::make_shared(std::move(segments)); + } + + void append(std::shared_ptr source) { + if (source->size() == 0) { + return; + } + if (const auto* composite = + dynamic_cast(source.get())) { + _segments.insert(_segments.end(), composite->_segments.begin(), + composite->_segments.end()); + return; + } + if (!_segments.empty()) { + auto& tail = _segments.back(); + if (tail.use_count() != 1) { + tail = tail->select_range(0, tail->size()); + } + if (tail->try_append(*source)) { + return; + } + } + _segments.push_back(std::move(source)); + } + + std::vector> _segments; + size_t _rows = 0; + mutable std::mutex _materialization_lock; + mutable ColumnVariantV2::MutablePtr _materialized; + mutable ColumnVariantV2::MutablePtr _serialized; +}; + +std::shared_ptr combine_shredded_states( + std::shared_ptr left, std::shared_ptr right) { + auto combined = std::make_shared( + std::vector> {std::move(left)}); + combined->try_append(*right); + return combined; +} + } // namespace #ifdef BE_TEST @@ -387,21 +690,21 @@ std::optional ColumnVariantV2::find_shredded_typed_va return _shredded->find_typed_value(path); } +const ColumnVariantV2& ColumnVariantV2::serialization_column() const { + if (!_shredded) { + return *this; + } + const ColumnVariantV2& serialized = _shredded->serialized_column(); + DORIS_CHECK(!serialized.is_shredded()) + << "shredded Variant wire materializer returned another shredded column"; + DORIS_CHECK_EQ(serialized.size(), size()) + << "shredded Variant wire materializer changed the row count"; + return serialized; +} + void ColumnVariantV2::ensure_encoded() { if (_shredded) { - const ColumnVariantV2& materialized = _shredded->materialized_column(); - DORIS_CHECK(!materialized.is_shredded()) - << "shredded state materializer returned another shredded column"; - // The shredded state may cache and share its canonical materialization across readers. - // Detach every mutable buffer before dropping that owner so later COW mutations stay legal. - _metadatas = materialized._metadatas->clone_resized(materialized._metadatas->size()); - _meta_ids = materialized._meta_ids->clone_resized(materialized._meta_ids->size()); - _values = materialized._values->clone_resized(materialized._values->size()); - _typed = materialized._typed == nullptr - ? nullptr - : materialized._typed->clone_resized(materialized._typed->size()); - _typed_type = materialized._typed_type; - _shredded.reset(); + _replace_shredded_state_with(_shredded->materialized_column()); } if (!_typed) { DCHECK(_typed_type == nullptr); @@ -426,6 +729,12 @@ void ColumnVariantV2::ensure_encoded() { _check_invariants(); } +void ColumnVariantV2::_ensure_serialized() { + if (_shredded) { + _replace_shredded_state_with(_shredded->serialized_column()); + } +} + std::string ColumnVariantV2::get_name() const { if (_shredded) { return "variant_v2(shredded)"; @@ -804,12 +1113,22 @@ void ColumnVariantV2::insert_range_from( // NOLINT(readability-function-size) _check_invariants(); return; } + if (!_shredded->can_materialize() || !selected_source->can_materialize()) { + // A projected segment cannot reconstruct omitted root fields. Preserve it beside any + // complete or projected neighbor so later path extraction can choose per-segment + // direct or canonical evaluation without forcing the incomplete state to materialize. + _shredded = combine_shredded_states(std::move(_shredded), std::move(selected_source)); + _check_invariants(); + return; + } } if (_shredded) { - ensure_encoded(); + // A merging exchange can combine a local projected state with its remotely serialized + // peer. Use the wire representation so omitted roots are never requested from either side. + _ensure_serialized(); } if (source._shredded) { - insert_range_from(source._shredded->materialized_column(), start, length); + insert_range_from(source.serialization_column(), start, length); return; } @@ -897,9 +1216,6 @@ void ColumnVariantV2::insert_indices_from( // NOLINT(readability-function-size) return; } - if (_shredded) { - ensure_encoded(); - } if (!_typed && empty() && _metadatas->empty() && source._shredded) { // Gather into the native shredded representation for the same reason as range selection: // row selection does not require, and may not have, a complete logical Variant value. @@ -907,8 +1223,29 @@ void ColumnVariantV2::insert_indices_from( // NOLINT(readability-function-size) _check_invariants(); return; } + if (_shredded && source._shredded) { + auto selected_source = source._shredded->select_indices(indices_begin, indices_end); + if (_shredded.use_count() != 1) { + _shredded = _shredded->select_range(0, size()); + } + if (_shredded->try_append(*selected_source)) { + _check_invariants(); + return; + } + if (!_shredded->can_materialize() || !selected_source->can_materialize()) { + // Exchange gathers may mix complete files with projected files. Keep their boundaries + // because only complete segments are allowed to reconstruct the full logical root. + _shredded = combine_shredded_states(std::move(_shredded), std::move(selected_source)); + _check_invariants(); + return; + } + } + if (_shredded) { + // Indexed gathers have the same local/remote representation boundary as range gathers. + _ensure_serialized(); + } if (source._shredded) { - insert_indices_from(source._shredded->materialized_column(), indices_begin, indices_end); + insert_indices_from(source.serialization_column(), indices_begin, indices_end); return; } @@ -1353,7 +1690,16 @@ MutableColumnPtr ColumnVariantV2::permute(const Permutation& permutation, size_t } if (_shredded) { - return _shredded->materialized_column().permute(permutation, limit); + DorisVector selected_indices(result_size); + for (size_t row = 0; row < result_size; ++row) { + DORIS_CHECK_LE(permutation[row], std::numeric_limits::max()) + << "shredded Variant permutation index exceeds uint32 domain"; + selected_indices[row] = static_cast(permutation[row]); + } + // Local TopN selection may run before exchange and projected states cannot reconstruct + // omitted root fields, so preserve the native shredded representation while gathering. + return ColumnVariantV2::create_shredded(_shredded->select_indices( + selected_indices.data(), selected_indices.data() + selected_indices.size())); } if (_typed) { @@ -1393,6 +1739,11 @@ MutableColumnPtr ColumnVariantV2::clone_resized(size_t new_size) const { result->_check_invariants(); return result; } + if (new_size < size()) { + // LIMIT truncation is a row selection and must not require complete Variant roots from + // a projected scanner state. + return ColumnVariantV2::create_shredded(_shredded->select_range(0, new_size)); + } return _shredded->materialized_column().clone_resized(new_size); } if (_typed) { @@ -1433,7 +1784,18 @@ MutableColumnPtr ColumnVariantV2::clone_resized(size_t new_size) const { void ColumnVariantV2::resize(size_t new_size) { const size_t old_size = size(); - if (_shredded && new_size != old_size) { + if (_shredded && new_size < old_size) { + // LIMIT truncation only selects existing rows, so keep it physical: projected states may + // omit roots that cannot be reconstructed merely to reduce the row count. + if (new_size == 0) { + _shredded.reset(); + } else { + _shredded = _shredded->select_range(0, new_size); + } + _check_invariants(); + return; + } + if (_shredded && new_size > old_size) { ensure_encoded(); } if (_typed) { @@ -1491,6 +1853,25 @@ uint32_t ColumnVariantV2::_find_or_insert_metadata(StringRef metadata) { return id; } +void ColumnVariantV2::_replace_shredded_state_with(const ColumnVariantV2& replacement) { + DORIS_CHECK(_shredded != nullptr) << "replacing shredded state requires a shredded destination"; + DORIS_CHECK(!replacement.is_shredded()) + << "shredded state replacement must be a non-shredded column"; + DORIS_CHECK_EQ(replacement.size(), size()) + << "shredded state replacement changed the row count"; + // Format states cache and share materialized columns. Detach every mutable buffer before + // dropping the state owner so later COW mutations cannot modify a cached representation. + _metadatas = replacement._metadatas->clone_resized(replacement._metadatas->size()); + _meta_ids = replacement._meta_ids->clone_resized(replacement._meta_ids->size()); + _values = replacement._values->clone_resized(replacement._values->size()); + _typed = replacement._typed == nullptr + ? nullptr + : replacement._typed->clone_resized(replacement._typed->size()); + _typed_type = replacement._typed_type; + _shredded.reset(); + _check_invariants(); +} + void ColumnVariantV2::_adopt_state_from(ColumnVariantV2& replacement) { DORIS_CHECK(this != &replacement) << "cannot adopt ColumnVariantV2 state from itself"; _metadatas = std::move(replacement._metadatas); diff --git a/be/src/core/column/variant_v2/column_variant_v2.h b/be/src/core/column/variant_v2/column_variant_v2.h index da0091d79fc165..7cd167663be0e4 100644 --- a/be/src/core/column/variant_v2/column_variant_v2.h +++ b/be/src/core/column/variant_v2/column_variant_v2.h @@ -53,6 +53,9 @@ struct VariantShreddedTypedValue { // retain the decoded leaf without copying it or depending on scanner lifetime. ColumnPtr column; DataTypePtr type; + // Physical identities such as binary annotations cannot use the typed scalar state. In that + // case the format reader may return an exact Nullable leaf instead. + ColumnPtr normalized; }; // Format readers keep their native shredded representation behind this interface. Core Variant @@ -76,15 +79,26 @@ class VariantShreddedState { size_t length) const = 0; virtual std::shared_ptr select_indices( const uint32_t* indices_begin, const uint32_t* indices_end) const = 0; + // False means the state contains only projected leaves and cannot reconstruct root values. + virtual bool can_materialize() const = 0; // Appends another state only when both format-owned physical layouts have identical semantics. // An incompatible source must leave this state unchanged and return false. virtual bool try_append(const VariantShreddedState& source) = 0; virtual std::optional find_typed_value( std::span path) const = 0; + // Produces an exact Variant representation of one requested path. Complete states may fall + // back to their canonical roots; projected states must preserve the format's physical scalar + // identity rather than inferring it from the decoded value. + virtual std::optional find_normalized_value( + std::span path) const = 0; // The returned column is cached and owned by this state, so borrowed VariantRef values remain // valid for the state lifetime. Implementations must not materialize before this is called. virtual const ColumnVariantV2& materialized_column() const = 0; + // Whole-column transport may encode only the retained projection because access-path planning + // guarantees that omitted fields have no downstream consumer. The returned column must be a + // self-contained, non-shredded wire representation with the same row count. + virtual const ColumnVariantV2& serialized_column() const = 0; }; // ColumnVariantV2 stores a whole column in exactly one state: encoded Variant bytes, one nullable @@ -146,6 +160,7 @@ class ColumnVariantV2 final : public COWHelper { const DataTypePtr& typed_type() const; std::optional find_shredded_typed_value( std::span path) const; + const ColumnVariantV2& serialization_column() const; void ensure_encoded(); ReadView read_view() const; @@ -235,6 +250,8 @@ class ColumnVariantV2 final : public COWHelper { ColumnVariantV2(const ColumnVariantV2& other); uint32_t _find_or_insert_metadata(StringRef metadata); + void _replace_shredded_state_with(const ColumnVariantV2& replacement); + void _ensure_serialized(); void _adopt_state_from(ColumnVariantV2& replacement); void _detach_metadata_for_write(); void _check_invariants() const; diff --git a/be/src/core/data_type_serde/data_type_variant_v2_serde.cpp b/be/src/core/data_type_serde/data_type_variant_v2_serde.cpp index a9a137e4fea85b..57521e7f34aa0d 100644 --- a/be/src/core/data_type_serde/data_type_variant_v2_serde.cpp +++ b/be/src/core/data_type_serde/data_type_variant_v2_serde.cpp @@ -181,7 +181,7 @@ DataTypeVariantV2SerDe::DataTypeVariantV2SerDe(int nesting_level) : DataTypeSerD int64_t DataTypeVariantV2SerDe::get_uncompressed_serialized_bytes(const IColumn& column, int be_exec_version) { - const auto& variant = get_variant_v2_column(column); + const auto& variant = get_variant_v2_column(column).serialization_column(); int64_t size = sizeof(bool) + sizeof(size_t) * 2 + sizeof(bool); if (variant.is_typed()) { const DataTypePtr nullable_type = make_nullable(variant._typed_type); @@ -199,7 +199,7 @@ char* DataTypeVariantV2SerDe::serialize(const IColumn& column, char* buf, int be const IColumn* physical = &column; size_t saved_rows = 0; buf = serialize_const_flag_and_row_num(&physical, buf, &saved_rows); - const auto& variant = assert_cast(*physical); + const auto& variant = assert_cast(*physical).serialization_column(); DCHECK_EQ(variant.size(), saved_rows); unaligned_store(buf, variant.is_typed()); buf += sizeof(bool); diff --git a/be/src/core/value/variant/variant_batch_builder.cpp b/be/src/core/value/variant/variant_batch_builder.cpp index ca03c1e670b10e..fedec1a5bd63f5 100644 --- a/be/src/core/value/variant/variant_batch_builder.cpp +++ b/be/src/core/value/variant/variant_batch_builder.cpp @@ -602,10 +602,18 @@ class VariantCollectionCore { add_bool(value.get_bool()); return; case VariantPrimitiveId::INT8: + add_scalar(VariantScalarRef::integer(value.get_int(), 1)); + return; case VariantPrimitiveId::INT16: + add_scalar(VariantScalarRef::integer(value.get_int(), 2)); + return; case VariantPrimitiveId::INT32: + add_scalar(VariantScalarRef::integer(value.get_int(), 4)); + return; case VariantPrimitiveId::INT64: - add_int(value.get_int()); + // Importing an existing VariantRef must retain its physical identity; external + // shredded schemas use the width to distinguish otherwise equal scalar values. + add_scalar(VariantScalarRef::integer(value.get_int(), 8)); return; case VariantPrimitiveId::FLOAT: add_scalar(VariantScalarRef::float32(value.get_float())); @@ -621,7 +629,7 @@ class VariantCollectionCore { throw Exception(ErrorCode::CORRUPTION, "Variant imported decimal exceeds precision 38"); } - add_scalar(VariantScalarRef::decimal(decimal.unscaled, decimal.scale)); + add_scalar(VariantScalarRef::decimal(decimal.unscaled, decimal.scale, decimal.width)); return; } case VariantPrimitiveId::DATE: diff --git a/be/src/exprs/function/function_variant_element_v2.cpp b/be/src/exprs/function/function_variant_element_v2.cpp index 90863fe86c5c72..2c506c05377e65 100644 --- a/be/src/exprs/function/function_variant_element_v2.cpp +++ b/be/src/exprs/function/function_variant_element_v2.cpp @@ -152,7 +152,8 @@ std::optional extract_shredded_typed_variant_element( if (!match.has_value()) { return std::nullopt; } - const auto& leaf = assert_cast(*match->column); + const ColumnPtr& matched_column = match->normalized ? match->normalized : match->column; + const auto& leaf = assert_cast(*matched_column); auto nulls = leaf.get_null_map_column().clone_resized(source.size()); auto& null_data = assert_cast(*nulls).get_data(); for (size_t row = 0; row < source.size(); ++row) { @@ -160,6 +161,10 @@ std::optional extract_shredded_typed_variant_element( static_cast(null_data[row] != 0 || is_outer_null(outer_nulls, row)); } + if (match->normalized) { + return ColumnNullable::create(leaf.get_nested_column_ptr(), std::move(nulls)); + } + // The typed ColumnVariantV2 retains the exact decoded Parquet leaf. Only the SQL result null // map is produced here, so predicates and casts can consume the leaf without reconstructing // canonical Variant rows. diff --git a/be/src/format_v2/file_reader.cpp b/be/src/format_v2/file_reader.cpp index a2ca4894044404..1b1f2f284405f9 100644 --- a/be/src/format_v2/file_reader.cpp +++ b/be/src/format_v2/file_reader.cpp @@ -74,7 +74,10 @@ std::string FileScanRequest::debug_string() const { out << column_id << ":" << block_position; } out << "}, conjunct_count=" << conjuncts.size() - << ", delete_conjunct_count=" << delete_conjuncts.size() + << ", delete_conjunct_count=" << delete_conjuncts.size() << ", variant_schema_overrides=" + << join_debug_strings( + variant_schema_overrides, + [](const LocalColumnIndex& projection) { return projection.debug_string(); }) << ", count_star_placeholder_columns={"; const char* delimiter = ""; for (const auto column_id : count_star_placeholder_columns) { diff --git a/be/src/format_v2/file_reader.h b/be/src/format_v2/file_reader.h index 0256b8c1eebcb2..c7060ce4e98852 100644 --- a/be/src/format_v2/file_reader.h +++ b/be/src/format_v2/file_reader.h @@ -94,6 +94,11 @@ struct FileScanRequest { // predicate_columns, the value is semantically required and must still be validated and read. std::vector count_star_placeholder_columns; + // Table formats may assign semantics that legacy physical files do not encode. Each path here + // identifies an unannotated Parquet group that the physical reader must validate and decode as + // Variant. Keeping this explicit prevents generic Parquet scans from guessing based on names. + std::vector variant_schema_overrides; + bool is_count_star_placeholder(LocalColumnId column_id) const { return std::ranges::find(count_star_placeholder_columns, column_id) != count_star_placeholder_columns.end(); diff --git a/be/src/format_v2/parquet/native_schema_desc.cpp b/be/src/format_v2/parquet/native_schema_desc.cpp index b56afb6f7170fb..16f4623bb1bafb 100644 --- a/be/src/format_v2/parquet/native_schema_desc.cpp +++ b/be/src/format_v2/parquet/native_schema_desc.cpp @@ -65,6 +65,8 @@ enum class VariantPrimitiveAnnotation : uint8_t { NONE, INT8, INT16, + INT32, + INT64, DECIMAL, DATE, TIME_MICROS, @@ -91,6 +93,15 @@ static VariantPrimitiveAnnotation variant_logical_annotation( if (logical.INTEGER.bitWidth == 16) { return VariantPrimitiveAnnotation::INT16; } + // Iceberg 1.11 writes full-width signed INTEGER annotations even though the Variant + // specification represents these widths without an annotation. Keep the widths distinct + // so validation only accepts them with the matching physical type. + if (logical.INTEGER.bitWidth == 32) { + return VariantPrimitiveAnnotation::INT32; + } + if (logical.INTEGER.bitWidth == 64) { + return VariantPrimitiveAnnotation::INT64; + } return VariantPrimitiveAnnotation::UNSUPPORTED; } if (logical.__isset.DECIMAL) { @@ -136,6 +147,11 @@ static VariantPrimitiveAnnotation variant_converted_annotation( return VariantPrimitiveAnnotation::INT8; case tparquet::ConvertedType::INT_16: return VariantPrimitiveAnnotation::INT16; + // Parquet Java mirrors full-width logical annotations into these legacy converted types. + case tparquet::ConvertedType::INT_32: + return VariantPrimitiveAnnotation::INT32; + case tparquet::ConvertedType::INT_64: + return VariantPrimitiveAnnotation::INT64; case tparquet::ConvertedType::DECIMAL: return VariantPrimitiveAnnotation::DECIMAL; case tparquet::ConvertedType::DATE: @@ -215,11 +231,13 @@ static Status validate_variant_primitive_type(const NativeFieldSchema& typed) { valid = annotation == VariantPrimitiveAnnotation::NONE || annotation == VariantPrimitiveAnnotation::INT8 || annotation == VariantPrimitiveAnnotation::INT16 || + annotation == VariantPrimitiveAnnotation::INT32 || annotation == VariantPrimitiveAnnotation::DECIMAL || annotation == VariantPrimitiveAnnotation::DATE; break; case tparquet::Type::INT64: valid = annotation == VariantPrimitiveAnnotation::NONE || + annotation == VariantPrimitiveAnnotation::INT64 || annotation == VariantPrimitiveAnnotation::DECIMAL || annotation == VariantPrimitiveAnnotation::TIME_MICROS || annotation == VariantPrimitiveAnnotation::TIMESTAMP_MICROS || @@ -266,17 +284,22 @@ class ScopedBoolOverride { bool _original; }; -static Status validate_variant_layout(const tparquet::SchemaElement& group_schema, - const NativeFieldSchema& group_field) { - const auto& annotation = group_schema.logicalType.VARIANT; - if (annotation.__isset.specification_version && annotation.specification_version != 1) { +Status validate_variant_layout(const NativeFieldSchema& group_field, + std::optional specification_version, + bool allow_optional_shredded_metadata) { + if (specification_version.has_value() && *specification_version != 1) { return Status::NotSupported("Parquet Variant specification version {} is not supported", - annotation.specification_version); + *specification_version); + } + if (group_field.parquet_schema.__isset.repetition_type && + group_field.parquet_schema.repetition_type == tparquet::FieldRepetitionType::REPEATED) { + return Status::NotSupported("repeated Parquet Variant group {} is not supported", + group_field.name); } if (group_field.children.size() < 2 || group_field.children.size() > 3) { return Status::Corruption( "Parquet Variant {} must contain metadata, value, and optional typed_value", - group_schema.name); + group_field.name); } const NativeFieldSchema* metadata = nullptr; @@ -292,22 +315,29 @@ static Status validate_variant_layout(const tparquet::SchemaElement& group_schem target = &typed_value; } else { return Status::Corruption("Parquet Variant {} has unexpected child {}", - group_schema.name, child.name); + group_field.name, child.name); } if (*target != nullptr) { - return Status::Corruption("Parquet Variant {} has duplicate child {}", - group_schema.name, child.name); + return Status::Corruption("Parquet Variant {} has duplicate child {}", group_field.name, + child.name); } *target = &child; } if (metadata == nullptr || value == nullptr) { return Status::Corruption("Parquet Variant {} requires metadata and value children", - group_schema.name); - } + group_field.name); + } + const auto metadata_repetition = metadata->parquet_schema.repetition_type; + // Paimon makes every field in its shredded carrier optional. Restrict that compatibility to + // unannotated overrides; row materialization still rejects null metadata for a non-null value. + const bool valid_metadata_repetition = + metadata_repetition == tparquet::FieldRepetitionType::REQUIRED || + (allow_optional_shredded_metadata && typed_value != nullptr && + metadata_repetition == tparquet::FieldRepetitionType::OPTIONAL); if (!metadata->children.empty() || metadata->physical_type != tparquet::Type::BYTE_ARRAY || - metadata->parquet_schema.repetition_type != tparquet::FieldRepetitionType::REQUIRED) { + !valid_metadata_repetition) { return Status::Corruption("Parquet Variant {} metadata must be a required BYTE_ARRAY", - group_schema.name); + group_field.name); } const auto expected_value_repetition = typed_value == nullptr ? tparquet::FieldRepetitionType::REQUIRED @@ -317,13 +347,13 @@ static Status validate_variant_layout(const tparquet::SchemaElement& group_schem if (!value->children.empty() || value->physical_type != tparquet::Type::BYTE_ARRAY || value->parquet_schema.repetition_type != expected_value_repetition) { return Status::Corruption("Parquet Variant {} value must be a {} BYTE_ARRAY", - group_schema.name, + group_field.name, typed_value == nullptr ? "required" : "optional"); } if (typed_value != nullptr && typed_value->parquet_schema.repetition_type != tparquet::FieldRepetitionType::OPTIONAL) { return Status::Corruption("Parquet Variant {} typed_value must be optional", - group_schema.name); + group_field.name); } enum class WrapperContext : uint8_t { OBJECT_FIELD, ARRAY_ELEMENT }; @@ -945,7 +975,12 @@ Status NativeFieldDescriptor::parse_group_field( ScopedBoolOverride timestamp_tz_mapping(_enable_mapping_timestamp_tz, true); RETURN_IF_ERROR(parse_struct_field(t_schemas, curr_pos, group_field)); } - RETURN_IF_ERROR(validate_variant_layout(group_schema, *group_field)); + const auto& annotation = group_schema.logicalType.VARIANT; + const auto specification_version = + annotation.__isset.specification_version + ? std::optional(annotation.specification_version) + : std::nullopt; + RETURN_IF_ERROR(validate_variant_layout(*group_field, specification_version)); group_field->variant_physical_type = group_field->data_type; // Native page readers dispatch groups from data_type, so preserve the physical STRUCT // here. The public Parquet schema maps it to logical Variant without losing this shape. diff --git a/be/src/format_v2/parquet/native_schema_desc.h b/be/src/format_v2/parquet/native_schema_desc.h index 918be6d2c65c58..be0e739f80b6ff 100644 --- a/be/src/format_v2/parquet/native_schema_desc.h +++ b/be/src/format_v2/parquet/native_schema_desc.h @@ -22,6 +22,7 @@ #include #include +#include #include #include #include @@ -86,6 +87,10 @@ struct NativeFieldSchema { uint64_t get_max_column_id() const; }; +Status validate_variant_layout(const NativeFieldSchema& group_field, + std::optional specification_version = std::nullopt, + bool allow_optional_shredded_metadata = false); + // V2 owns this schema tree and parser so footer/schema planning never invokes the V1 reader path. class NativeFieldDescriptor { private: diff --git a/be/src/format_v2/parquet/parquet_column_schema.cpp b/be/src/format_v2/parquet/parquet_column_schema.cpp index 71416e17dc9209..984e0259165074 100644 --- a/be/src/format_v2/parquet/parquet_column_schema.cpp +++ b/be/src/format_v2/parquet/parquet_column_schema.cpp @@ -178,6 +178,41 @@ void propagate_native_max_levels(ParquetColumnSchema* schema) { } } +void rebuild_logical_complex_type(ParquetColumnSchema* schema) { + DORIS_CHECK(schema != nullptr); + schema->contains_variant = schema->kind == ParquetColumnSchemaKind::VARIANT; + for (const auto& child : schema->children) { + schema->contains_variant |= child->contains_variant; + } + if (schema->kind == ParquetColumnSchemaKind::VARIANT || !schema->contains_variant) { + return; + } + + DataTypePtr logical_type; + if (schema->kind == ParquetColumnSchemaKind::LIST) { + DORIS_CHECK(schema->children.size() == 1); + logical_type = std::make_shared(schema->children[0]->type); + } else if (schema->kind == ParquetColumnSchemaKind::MAP) { + DORIS_CHECK(schema->children.size() == 2); + logical_type = std::make_shared(make_nullable(schema->children[0]->type), + make_nullable(schema->children[1]->type)); + } else { + DORIS_CHECK(schema->kind == ParquetColumnSchemaKind::STRUCT); + DataTypes child_types; + Strings child_names; + child_types.reserve(schema->children.size()); + child_names.reserve(schema->children.size()); + for (const auto& child : schema->children) { + child_types.push_back(child->type); + child_names.push_back(child->name); + } + logical_type = + std::make_shared(std::move(child_types), std::move(child_names)); + } + schema->type = schema->type->is_nullable() ? make_nullable(std::move(logical_type)) + : std::move(logical_type); +} + std::unique_ptr build_native_node_schema(const NativeFieldSchema& field, int32_t local_id) { auto result = std::make_unique(); @@ -228,34 +263,51 @@ std::unique_ptr build_native_node_schema(const NativeFieldS } // A nested Variant changes its public child type from the physical STRUCT carrier. Rebuild // every enclosing complex type so file-block columns keep the same logical shape as readers. - if (result->kind != ParquetColumnSchemaKind::VARIANT && result->contains_variant) { - DataTypePtr logical_type; - if (result->kind == ParquetColumnSchemaKind::LIST) { - DORIS_CHECK(result->children.size() == 1); - logical_type = std::make_shared(result->children[0]->type); - } else if (result->kind == ParquetColumnSchemaKind::MAP) { - DORIS_CHECK(result->children.size() == 2); - logical_type = std::make_shared(make_nullable(result->children[0]->type), - make_nullable(result->children[1]->type)); - } else { - DataTypes child_types; - Strings child_names; - child_types.reserve(result->children.size()); - child_names.reserve(result->children.size()); - for (const auto& child : result->children) { - child_types.push_back(child->type); - child_names.push_back(child->name); - } - logical_type = std::make_shared(std::move(child_types), - std::move(child_names)); - } - result->type = result->type->is_nullable() ? make_nullable(std::move(logical_type)) - : std::move(logical_type); - } + rebuild_logical_complex_type(result.get()); propagate_native_max_levels(result.get()); return result; } +Status apply_variant_schema_override(const NativeFieldSchema& native_schema, + const format::LocalColumnIndex& override, + ParquetColumnSchema* field) { + DORIS_CHECK(field != nullptr); + if (field->local_id != override.local_id()) { + return Status::InvalidArgument("Variant schema override local id {} does not match {}", + override.local_id(), field->local_id); + } + if (override.project_all_children) { + if (field->kind == ParquetColumnSchemaKind::VARIANT) { + return Status::OK(); + } + if (field->kind != ParquetColumnSchemaKind::STRUCT) { + return Status::Corruption("Parquet Variant {} must use a group carrier", + native_schema.name); + } + RETURN_IF_ERROR(validate_variant_layout(native_schema, std::nullopt, true)); + field->variant_physical_type = field->type; + DataTypePtr variant_type = std::make_shared(); + field->type = field->type->is_nullable() ? make_nullable(std::move(variant_type)) + : std::move(variant_type); + field->kind = ParquetColumnSchemaKind::VARIANT; + field->contains_variant = true; + return Status::OK(); + } + for (const auto& child_override : override.children) { + const auto child_idx = child_override.local_id(); + if (child_idx < 0 || child_idx >= static_cast(field->children.size()) || + child_idx >= static_cast(native_schema.children.size())) { + return Status::InvalidArgument("Invalid nested Variant schema override {} under {}", + child_idx, field->name); + } + RETURN_IF_ERROR(apply_variant_schema_override(native_schema.children[child_idx], + child_override, + field->children[child_idx].get())); + } + rebuild_logical_complex_type(field); + return Status::OK(); +} + } // namespace Status build_parquet_column_schema(const NativeFieldDescriptor& schema, @@ -277,4 +329,24 @@ Status build_parquet_column_schema(const NativeFieldDescriptor& schema, return Status::OK(); } +Status apply_variant_schema_overrides( + const NativeFieldDescriptor& native_schema, + const std::vector& variant_schema_overrides, + std::vector>* fields) { + if (fields == nullptr) { + return Status::InvalidArgument("fields is null"); + } + const auto& native_fields = native_schema.get_fields_schema(); + for (const auto& override : variant_schema_overrides) { + const auto local_id = override.local_id(); + if (local_id < 0 || local_id >= static_cast(fields->size()) || + local_id >= static_cast(native_fields.size())) { + return Status::InvalidArgument("Invalid Variant schema override root {}", local_id); + } + RETURN_IF_ERROR(apply_variant_schema_override(native_fields[local_id], override, + (*fields)[local_id].get())); + } + return Status::OK(); +} + } // namespace doris::format::parquet diff --git a/be/src/format_v2/parquet/parquet_column_schema.h b/be/src/format_v2/parquet/parquet_column_schema.h index 3103b97cfc253c..55de6f33e5fd35 100644 --- a/be/src/format_v2/parquet/parquet_column_schema.h +++ b/be/src/format_v2/parquet/parquet_column_schema.h @@ -21,6 +21,7 @@ #include "common/status.h" #include "core/data_type/data_type.h" +#include "format_v2/column_data.h" #include "format_v2/parquet/parquet_type.h" namespace doris::format::parquet { @@ -81,4 +82,9 @@ struct ParquetColumnSchema { Status build_parquet_column_schema(const NativeFieldDescriptor& schema, std::vector>* fields); +Status apply_variant_schema_overrides( + const NativeFieldDescriptor& native_schema, + const std::vector& variant_schema_overrides, + std::vector>* fields); + } // namespace doris::format::parquet diff --git a/be/src/format_v2/parquet/parquet_profile.cpp b/be/src/format_v2/parquet/parquet_profile.cpp index 5fd3100a2e5bcf..ec70db46840771 100644 --- a/be/src/format_v2/parquet/parquet_profile.cpp +++ b/be/src/format_v2/parquet/parquet_profile.cpp @@ -22,6 +22,19 @@ namespace doris::format::parquet { +namespace { + +std::shared_ptr add_persistent_counter(RuntimeProfile* profile, + const std::string& name, + TUnit::type type, + const std::string& parent) { + // A shredded Variant may be materialized after its scanner profile is destroyed. Keep the + // counter storage alive with the column state instead of retaining a dangling profile pointer. + return profile->add_shared_counter(name, type, parent, 1); +} + +} // namespace + void ParquetProfile::init(RuntimeProfile* profile) { if (profile == nullptr) { return; @@ -90,18 +103,18 @@ void ParquetProfile::init(RuntimeProfile* profile) { ADD_CHILD_TIMER_WITH_LEVEL(profile, "LevelOnlySkipTime", parquet_profile, 1); materialization_time = ADD_CHILD_TIMER_WITH_LEVEL(profile, "MaterializationTime", parquet_profile, 1); - variant_reconstruction_time = - ADD_CHILD_TIMER_WITH_LEVEL(profile, "VariantReconstructionTime", parquet_profile, 1); - variant_reconstructed_rows = ADD_CHILD_COUNTER_WITH_LEVEL(profile, "VariantReconstructedRows", - TUnit::UNIT, parquet_profile, 1); - variant_direct_leaf_rows = ADD_CHILD_COUNTER_WITH_LEVEL(profile, "VariantDirectLeafRows", - TUnit::UNIT, parquet_profile, 1); - variant_direct_leaf_path_misses = ADD_CHILD_COUNTER_WITH_LEVEL( - profile, "VariantDirectLeafPathMisses", TUnit::UNIT, parquet_profile, 1); - variant_direct_leaf_residual_fallbacks = ADD_CHILD_COUNTER_WITH_LEVEL( - profile, "VariantDirectLeafResidualFallbacks", TUnit::UNIT, parquet_profile, 1); - variant_direct_leaf_unsupported_fallbacks = ADD_CHILD_COUNTER_WITH_LEVEL( - profile, "VariantDirectLeafUnsupportedFallbacks", TUnit::UNIT, parquet_profile, 1); + variant_reconstruction_time = add_persistent_counter(profile, "VariantReconstructionTime", + TUnit::TIME_NS, parquet_profile); + variant_reconstructed_rows = add_persistent_counter(profile, "VariantReconstructedRows", + TUnit::UNIT, parquet_profile); + variant_direct_leaf_rows = + add_persistent_counter(profile, "VariantDirectLeafRows", TUnit::UNIT, parquet_profile); + variant_direct_leaf_path_misses = add_persistent_counter(profile, "VariantDirectLeafPathMisses", + TUnit::UNIT, parquet_profile); + variant_direct_leaf_residual_fallbacks = add_persistent_counter( + profile, "VariantDirectLeafResidualFallbacks", TUnit::UNIT, parquet_profile); + variant_direct_leaf_unsupported_fallbacks = add_persistent_counter( + profile, "VariantDirectLeafUnsupportedFallbacks", TUnit::UNIT, parquet_profile); hybrid_selection_batches = ADD_CHILD_COUNTER_WITH_LEVEL(profile, "HybridSelectionBatches", TUnit::UNIT, parquet_profile, 1); hybrid_selection_ranges = ADD_CHILD_COUNTER_WITH_LEVEL(profile, "HybridSelectionRanges", diff --git a/be/src/format_v2/parquet/parquet_profile.h b/be/src/format_v2/parquet/parquet_profile.h index ed1faa8f935134..764fef1d80c190 100644 --- a/be/src/format_v2/parquet/parquet_profile.h +++ b/be/src/format_v2/parquet/parquet_profile.h @@ -15,6 +15,8 @@ #pragma once +#include + #include "runtime/runtime_profile.h" namespace doris::format::parquet { @@ -38,12 +40,12 @@ struct ParquetColumnReaderProfile { RuntimeProfile::Counter* level_only_read_time = nullptr; RuntimeProfile::Counter* level_only_skip_time = nullptr; RuntimeProfile::Counter* materialization_time = nullptr; // value materialization time (ns) - RuntimeProfile::Counter* variant_reconstruction_time = nullptr; - RuntimeProfile::Counter* variant_reconstructed_rows = nullptr; - RuntimeProfile::Counter* variant_direct_leaf_rows = nullptr; - RuntimeProfile::Counter* variant_direct_leaf_path_misses = nullptr; - RuntimeProfile::Counter* variant_direct_leaf_residual_fallbacks = nullptr; - RuntimeProfile::Counter* variant_direct_leaf_unsupported_fallbacks = nullptr; + std::shared_ptr variant_reconstruction_time; + std::shared_ptr variant_reconstructed_rows; + std::shared_ptr variant_direct_leaf_rows; + std::shared_ptr variant_direct_leaf_path_misses; + std::shared_ptr variant_direct_leaf_residual_fallbacks; + std::shared_ptr variant_direct_leaf_unsupported_fallbacks; RuntimeProfile::Counter* hybrid_selection_batches = nullptr; RuntimeProfile::Counter* hybrid_selection_ranges = nullptr; RuntimeProfile::Counter* hybrid_selection_null_fallback_batches = nullptr; @@ -174,12 +176,12 @@ struct ParquetProfile { RuntimeProfile::Counter* level_only_read_time = nullptr; RuntimeProfile::Counter* level_only_skip_time = nullptr; RuntimeProfile::Counter* materialization_time = nullptr; - RuntimeProfile::Counter* variant_reconstruction_time = nullptr; - RuntimeProfile::Counter* variant_reconstructed_rows = nullptr; - RuntimeProfile::Counter* variant_direct_leaf_rows = nullptr; - RuntimeProfile::Counter* variant_direct_leaf_path_misses = nullptr; - RuntimeProfile::Counter* variant_direct_leaf_residual_fallbacks = nullptr; - RuntimeProfile::Counter* variant_direct_leaf_unsupported_fallbacks = nullptr; + std::shared_ptr variant_reconstruction_time; + std::shared_ptr variant_reconstructed_rows; + std::shared_ptr variant_direct_leaf_rows; + std::shared_ptr variant_direct_leaf_path_misses; + std::shared_ptr variant_direct_leaf_residual_fallbacks; + std::shared_ptr variant_direct_leaf_unsupported_fallbacks; RuntimeProfile::Counter* hybrid_selection_batches = nullptr; RuntimeProfile::Counter* hybrid_selection_ranges = nullptr; RuntimeProfile::Counter* hybrid_selection_null_fallback_batches = nullptr; diff --git a/be/src/format_v2/parquet/parquet_reader.cpp b/be/src/format_v2/parquet/parquet_reader.cpp index 8c6efe27558bbc..4316f2b21d701e 100644 --- a/be/src/format_v2/parquet/parquet_reader.cpp +++ b/be/src/format_v2/parquet/parquet_reader.cpp @@ -332,10 +332,32 @@ DataTypePtr apply_timestamp_tz_mapping(ParquetColumnSchema* column_schema) { } column_schema->type = nullable_like_original( column_schema->type, std::make_shared(child_types, child_names)); + } else if (column_schema->kind == ParquetColumnSchemaKind::VARIANT) { + Strings child_names; + child_names.reserve(column_schema->children.size()); + for (const auto& child : column_schema->children) { + child_names.push_back(child->name); + } + column_schema->variant_physical_type = + nullable_like_original(column_schema->variant_physical_type, + std::make_shared(child_types, child_names)); } return column_schema->type; } +void apply_timestamp_tz_mapping_in_variants(ParquetColumnSchema* column_schema) { + DORIS_CHECK(column_schema != nullptr); + if (column_schema->kind == ParquetColumnSchemaKind::VARIANT) { + // Shredded Variant timestamps always represent instants, even when the surrounding table + // did not request catalog-level TIMESTAMPTZ mapping for ordinary Parquet columns. + apply_timestamp_tz_mapping(column_schema); + return; + } + for (auto& child : column_schema->children) { + apply_timestamp_tz_mapping_in_variants(child.get()); + } +} + static Status find_projected_minmax_leaf(const ParquetColumnSchema& column_schema, const format::LocalColumnIndex& projection, const ParquetColumnSchema** leaf_schema) { @@ -526,6 +548,21 @@ Status ParquetReader::open(std::shared_ptr request) { } auto request_snapshot = request; DORIS_CHECK(request_snapshot != nullptr); + if (!request_snapshot->variant_schema_overrides.empty()) { + // Apply table-format semantics before Variant projection planning. The override is the + // explicit proof that an otherwise ordinary Parquet group is a Variant carrier. + RETURN_IF_ERROR(apply_variant_schema_overrides( + _state->file_context.native_metadata->schema(), + request_snapshot->variant_schema_overrides, &_state->file_schema)); + for (auto& column_schema : _state->file_schema) { + apply_timestamp_tz_mapping_in_variants(column_schema.get()); + } + _state->file_context.contains_variant = + std::ranges::any_of(_state->file_schema, [](const auto& column) { + DORIS_CHECK(column != nullptr); + return column->contains_variant; + }); + } size_t retained_variant_leaf_projections = 0; if (_state->file_context.contains_variant) { retained_variant_leaf_projections = diff --git a/be/src/format_v2/parquet/reader/variant_column_reader.cpp b/be/src/format_v2/parquet/reader/variant_column_reader.cpp index ba9f63b9356920..8e9d28720c4c27 100644 --- a/be/src/format_v2/parquet/reader/variant_column_reader.cpp +++ b/be/src/format_v2/parquet/reader/variant_column_reader.cpp @@ -398,7 +398,7 @@ bool append_wrapper(const ParquetColumnSchema& schema, const IColumn& wrapper, s void encode_variant_range(const ParquetColumnSchema& schema, const IColumn& wrapper, const ColumnNullable* outer_nullable, size_t begin, size_t end, - ColumnVariantV2& variants) { + bool require_metadata, ColumnVariantV2& variants) { try { VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows = end - begin}); for (size_t row = begin; row < end; ++row) { @@ -408,14 +408,22 @@ void encode_variant_range(const ParquetColumnSchema& schema, const IColumn& wrap 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); + VariantMetadataRef metadata; + if (find_child(schema, "metadata", nullptr) != nullptr) { + 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); + metadata = {metadata_bytes.data, metadata_bytes.size}; + metadata.validate(); + } else if (require_metadata) { + throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} has no root metadata", + schema.name); } - 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(); } @@ -429,13 +437,16 @@ void encode_variant_range(const ParquetColumnSchema& schema, const IColumn& wrap // 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); + encode_variant_range(schema, wrapper, outer_nullable, begin, middle, require_metadata, + variants); + encode_variant_range(schema, wrapper, outer_nullable, middle, end, require_metadata, + variants); } } ColumnVariantV2::MutablePtr encode_variant_column(const ParquetColumnSchema& schema, - const IColumn& physical) { + const IColumn& physical, + bool require_metadata = true) { if (schema.kind != ParquetColumnSchemaKind::VARIANT) { throw Exception(ErrorCode::INVALID_ARGUMENT, "Parquet column {} is not Variant", schema.name); @@ -454,7 +465,7 @@ ColumnVariantV2::MutablePtr encode_variant_column(const ParquetColumnSchema& sch 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); + require_metadata, *variants); } return variants; } @@ -553,6 +564,72 @@ bool supports_direct_typed_variant_state(const ParquetColumnSchema& schema) { } } +ColumnPtr normalize_projected_primitive_leaf(const ParquetColumnSchema& schema, + const ColumnPtr& typed) { + const auto& nullable = assert_cast(*typed); + VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows = nullable.size()}); + for (size_t row = 0; row < nullable.size(); ++row) { + auto output_row = builder.begin_row(); + if (nullable.get_null_map_data()[row] != 0) { + output_row.add_null(); + } else { + append_typed_scalar(schema, nullable.get_nested_column(), row, output_row); + } + output_row.finish(); + } + auto values = ColumnVariantV2::create(); + values->insert_encoded_batch(builder.finish_batch()); + auto nulls = nullable.get_null_map_column().clone_resized(nullable.size()); + return ColumnNullable::create(std::move(values), std::move(nulls)); +} + +bool find_materialized_path(VariantRef current, std::span path, + VariantRef* output) { + DORIS_CHECK(output != nullptr); + for (const auto& segment : path) { + if (segment.kind == VariantShreddedPathSegment::Kind::OBJECT_KEY) { + if (current.basic_type() != VariantBasicType::OBJECT || + !current.object_find(segment.key, ¤t)) { + return false; + } + continue; + } + if (current.basic_type() != VariantBasicType::ARRAY) { + return false; + } + const int64_t count = current.num_elements(); + const int64_t index = segment.index < 0 ? count + segment.index : segment.index; + if (index < 0 || index >= count) { + return false; + } + current = current.array_at(static_cast(index)); + } + *output = current; + return true; +} + +ColumnPtr normalize_materialized_path(const ColumnVariantV2& materialized, + std::span path) { + VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows = materialized.size()}); + auto nulls = ColumnUInt8::create(); + nulls->reserve(materialized.size()); + for (size_t row = 0; row < materialized.size(); ++row) { + auto output_row = builder.begin_row(); + VariantRef value; + if (find_materialized_path(materialized.get_value_ref(row), path, &value)) { + output_row.add_value(value); + nulls->insert_value(0); + } else { + output_row.add_null(); + nulls->insert_value(1); + } + output_row.finish(); + } + auto values = ColumnVariantV2::create(); + values->insert_encoded_batch(builder.finish_batch()); + return ColumnNullable::create(std::move(values), std::move(nulls)); +} + bool same_data_type(const DataTypePtr& left, const DataTypePtr& right) { return (!left && !right) || (left && right && left->equals(*right)); } @@ -612,12 +689,14 @@ class ParquetVariantShreddedState final : public VariantShreddedState { size_t size() const override { return _physical->size(); } size_t byte_size() const override { std::lock_guard lock(_materialization_lock); - return _physical->byte_size() + (_materialized ? _materialized->byte_size() : 0); + return _physical->byte_size() + (_materialized ? _materialized->byte_size() : 0) + + (_serialized ? _serialized->byte_size() : 0); } size_t allocated_bytes() const override { std::lock_guard lock(_materialization_lock); return _physical->allocated_bytes() + - (_materialized ? _materialized->allocated_bytes() : 0); + (_materialized ? _materialized->allocated_bytes() : 0) + + (_serialized ? _serialized->allocated_bytes() : 0); } void sanity_check() const override { _physical->sanity_check(); } @@ -648,6 +727,8 @@ class ParquetVariantShreddedState final : public VariantShreddedState { _complete, _profile); } + bool can_materialize() const override { return _complete; } + bool try_append(const VariantShreddedState& source) override { const auto* parquet_source = dynamic_cast(&source); if (parquet_source == nullptr || _complete != parquet_source->_complete || @@ -660,6 +741,7 @@ class ParquetVariantShreddedState final : public VariantShreddedState { _physical = std::move(mutable_physical); std::lock_guard lock(_materialization_lock); _materialized.reset(); + _serialized.reset(); return true; } @@ -705,15 +787,29 @@ class ParquetVariantShreddedState final : public VariantShreddedState { } if (position + 1 == path.size()) { if (typed_schema->kind != ParquetColumnSchemaKind::PRIMITIVE || - check_and_get_column(*typed) == nullptr || - !supports_direct_typed_variant_state(*typed_schema)) { + check_and_get_column(*typed) == nullptr) { update_counter(_profile.variant_direct_leaf_unsupported_fallbacks, 1); return std::nullopt; } + if (!supports_direct_typed_variant_state(*typed_schema)) { + if (_complete) { + update_counter(_profile.variant_direct_leaf_unsupported_fallbacks, 1); + return std::nullopt; + } + // A partial projection cannot reconstruct its root. Normalize only the exact + // requested leaf so Parquet annotations survive heterogeneous file schemas. + update_counter(_profile.variant_direct_leaf_rows, + static_cast(typed->size())); + return VariantShreddedTypedValue { + .column = nullptr, + .type = nullptr, + .normalized = normalize_projected_primitive_leaf(*typed_schema, typed)}; + } update_counter(_profile.variant_direct_leaf_rows, static_cast(typed->size())); return VariantShreddedTypedValue {.column = std::move(typed), - .type = remove_nullable(typed_schema->type)}; + .type = remove_nullable(typed_schema->type), + .normalized = nullptr}; } if (typed_schema->kind != ParquetColumnSchemaKind::STRUCT) { return path_miss(); @@ -722,6 +818,56 @@ class ParquetVariantShreddedState final : public VariantShreddedState { return std::nullopt; } + std::optional find_normalized_value( + std::span 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) { + bool direct = true; + for (size_t position = 0; position < path.size(); ++position) { + if (path[position].kind != VariantShreddedPathSegment::Kind::OBJECT_KEY) { + direct = false; + break; + } + 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) { + direct = false; + break; + } + ColumnPtr residual = struct_child(*wrapper_schema, wrapper, "value", nullptr); + if (residual && has_present_value(residual)) { + direct = false; + break; + } + typed = struct_child(*wrapper_schema, wrapper, "typed_value", &typed_schema); + if (!typed) { + direct = false; + break; + } + if (position + 1 == path.size()) { + direct = typed_schema->kind == ParquetColumnSchemaKind::PRIMITIVE && + check_and_get_column(*typed) != nullptr; + } else if (typed_schema->kind != ParquetColumnSchemaKind::STRUCT) { + direct = false; + break; + } + } + if (direct) { + return normalize_projected_primitive_leaf(*typed_schema, typed); + } + } + if (!_complete) { + return std::nullopt; + } + return normalize_materialized_path(materialized_column(), path); + } + const ColumnVariantV2& materialized_column() const override { std::lock_guard lock(_materialization_lock); if (!_complete) { @@ -730,7 +876,7 @@ class ParquetVariantShreddedState final : public VariantShreddedState { "A projected Parquet Variant can only serve its validated shredded leaves"); } if (!_materialized) { - SCOPED_TIMER(_profile.variant_reconstruction_time); + SCOPED_TIMER(_profile.variant_reconstruction_time.get()); _materialized = encode_variant_column(*_schema, *_physical); update_counter(_profile.variant_reconstructed_rows, static_cast(_physical->size())); @@ -738,10 +884,25 @@ class ParquetVariantShreddedState final : public VariantShreddedState { return *_materialized; } + const ColumnVariantV2& serialized_column() const override { + if (_complete) { + return materialized_column(); + } + std::lock_guard lock(_materialization_lock); + if (!_serialized) { + // Projected states intentionally omit root metadata, but an exchange buffer still + // needs self-contained bytes. Rebuild only retained paths; access-path planning is the + // invariant that prevents a downstream consumer from observing an omitted field. + _serialized = encode_variant_column(*_schema, *_physical, false); + } + return *_serialized; + } + private: - static void update_counter(RuntimeProfile::Counter* counter, int64_t value) { + static void update_counter(const std::shared_ptr& counter, + int64_t value) { if (counter != nullptr) { - COUNTER_UPDATE(counter, value); + COUNTER_UPDATE(counter.get(), value); } } @@ -751,6 +912,7 @@ class ParquetVariantShreddedState final : public VariantShreddedState { ParquetColumnReaderProfile _profile; mutable std::mutex _materialization_lock; mutable ColumnVariantV2::MutablePtr _materialized; + mutable ColumnVariantV2::MutablePtr _serialized; }; MutableColumnPtr build_variant_column(std::shared_ptr schema, diff --git a/be/src/format_v2/table/paimon_reader.cpp b/be/src/format_v2/table/paimon_reader.cpp index 93183af178babb..7c7a3fa7f3a9be 100644 --- a/be/src/format_v2/table/paimon_reader.cpp +++ b/be/src/format_v2/table/paimon_reader.cpp @@ -19,9 +19,16 @@ #include +#include +#include #include #include +#include "core/data_type/data_type_array.h" +#include "core/data_type/data_type_map.h" +#include "core/data_type/data_type_nullable.h" +#include "core/data_type/data_type_struct.h" +#include "core/data_type/data_type_variant_v2.h" #include "exprs/vexpr_context.h" #include "format/table/deletion_vector_reader.h" #include "format/table/paimon_reader.h" @@ -31,6 +38,148 @@ #include "gen_cpp/PlanNodes_types.h" namespace doris::format::paimon { +namespace { + +ColumnDefinition* find_file_column(const ColumnDefinition& table_column, + std::vector* file_schema, + TableColumnMappingMode mode) { + DORIS_CHECK(file_schema != nullptr); + if (mode == TableColumnMappingMode::BY_FIELD_ID) { + if (!table_column.has_identifier_field_id()) { + return nullptr; + } + const auto field_id = table_column.get_identifier_field_id(); + const auto it = std::ranges::find_if(*file_schema, [&](const auto& file_column) { + return file_column.has_identifier_field_id() && + file_column.get_identifier_field_id() == field_id; + }); + return it == file_schema->end() ? nullptr : &*it; + } + const auto* matched = format::find_column_by_name(table_column, *file_schema); + return matched == nullptr ? nullptr : &(*file_schema)[matched - file_schema->data()]; +} + +void rebuild_complex_type(ColumnDefinition* column) { + DORIS_CHECK(column != nullptr); + const bool nullable = column->type->is_nullable(); + const auto primitive = remove_nullable(column->type)->get_primitive_type(); + DataTypePtr rebuilt; + if (primitive == TYPE_ARRAY && column->children.size() == 1) { + rebuilt = std::make_shared(column->children[0].type); + } else if (primitive == TYPE_MAP && column->children.size() == 2) { + rebuilt = std::make_shared(column->children[0].type, column->children[1].type); + } else if (primitive == TYPE_STRUCT) { + DataTypes child_types; + Strings child_names; + child_types.reserve(column->children.size()); + child_names.reserve(column->children.size()); + for (const auto& child : column->children) { + child_types.push_back(child.type); + child_names.push_back(child.name); + } + rebuilt = std::make_shared(std::move(child_types), std::move(child_names)); + } + if (rebuilt != nullptr) { + column->type = nullable ? make_nullable(std::move(rebuilt)) : std::move(rebuilt); + } +} + +Status add_variant_schema_override(const std::vector& path, + std::vector* overrides) { + DORIS_CHECK(!path.empty()); + DORIS_CHECK(overrides != nullptr); + auto projection = LocalColumnIndex::local(path.back()); + for (size_t path_idx = path.size() - 1; path_idx > 0; --path_idx) { + auto parent = LocalColumnIndex::partial_local(path[path_idx - 1]); + parent.children.push_back(std::move(projection)); + projection = std::move(parent); + } + const auto existing = std::ranges::find_if(*overrides, [&](const auto& override) { + return override.local_id() == projection.local_id(); + }); + if (existing == overrides->end()) { + overrides->push_back(std::move(projection)); + } else { + RETURN_IF_ERROR(merge_local_column_index(&*existing, projection)); + } + return Status::OK(); +} + +bool contains_variant_type(const ColumnDefinition& column) { + if (column.type != nullptr && + remove_nullable(column.type)->get_primitive_type() == TYPE_VARIANT) { + return true; + } + return std::ranges::any_of(column.children, contains_variant_type); +} + +Status annotate_matched_paimon_variant(const ColumnDefinition& table_column, + ColumnDefinition* file_column, TableColumnMappingMode mode, + const std::vector& prefix, + std::vector* overrides) { + DORIS_CHECK(file_column != nullptr); + if (!contains_variant_type(table_column) || table_column.type == nullptr || + file_column->type == nullptr) { + return Status::OK(); + } + auto path = prefix; + path.push_back(file_column->local_id); + const auto table_primitive = remove_nullable(table_column.type)->get_primitive_type(); + const auto file_primitive = remove_nullable(file_column->type)->get_primitive_type(); + if (table_primitive == TYPE_VARIANT) { + if (file_primitive == TYPE_STRUCT) { + // Paimon omits the Parquet VARIANT annotation, so only a matched table Variant may + // reinterpret this carrier; ordinary STRUCT must stay a STRUCT. + DataTypePtr variant = std::make_shared(); + file_column->type = file_column->type->is_nullable() ? make_nullable(std::move(variant)) + : std::move(variant); + RETURN_IF_ERROR(add_variant_schema_override(path, overrides)); + } + return Status::OK(); + } + if (table_column.children.empty() || file_column->children.empty() || + table_primitive != file_primitive) { + return Status::OK(); + } + if (table_primitive == TYPE_ARRAY || table_primitive == TYPE_MAP) { + const auto child_count = + std::min(table_column.children.size(), file_column->children.size()); + for (size_t child_idx = 0; child_idx < child_count; ++child_idx) { + // ARRAY/MAP child names are writer-specific structural labels, so match these nodes by + // position and reserve name/field-id matching for actual STRUCT members. + RETURN_IF_ERROR(annotate_matched_paimon_variant(table_column.children[child_idx], + &file_column->children[child_idx], mode, + path, overrides)); + } + } else if (table_primitive == TYPE_STRUCT) { + for (const auto& table_child : table_column.children) { + auto* file_child = find_file_column(table_child, &file_column->children, mode); + if (file_child != nullptr) { + RETURN_IF_ERROR(annotate_matched_paimon_variant(table_child, file_child, mode, path, + overrides)); + } + } + } + rebuild_complex_type(file_column); + return Status::OK(); +} + +Status annotate_paimon_variants(const std::vector& table_schema, + std::vector* file_schema, + TableColumnMappingMode mode, + std::vector* overrides) { + DORIS_CHECK(file_schema != nullptr); + for (const auto& table_column : table_schema) { + auto* file_column = find_file_column(table_column, file_schema, mode); + if (file_column != nullptr) { + RETURN_IF_ERROR(annotate_matched_paimon_variant(table_column, file_column, mode, {}, + overrides)); + } + } + return Status::OK(); +} + +} // namespace Status PaimonReader::prepare_split(const format::SplitReadOptions& options) { { @@ -66,10 +215,25 @@ format::TableColumnMappingMode PaimonReader::mapping_mode() const { Status PaimonReader::annotate_file_schema(std::vector* file_schema) { DORIS_CHECK(file_schema != nullptr); - if (mapping_mode() != format::TableColumnMappingMode::BY_FIELD_ID) { - return Status::OK(); + _variant_schema_overrides.clear(); + const auto mode = mapping_mode(); + if (mode == format::TableColumnMappingMode::BY_FIELD_ID) { + RETURN_IF_ERROR(format::annotate_file_schema_from_history(_scan_params, _split_schema_id, + file_schema)); } - return format::annotate_file_schema_from_history(_scan_params, _split_schema_id, file_schema); + const bool projects_variant = std::ranges::any_of(_projected_columns, contains_variant_type); + if (projects_variant && _format == format::FileFormat::PARQUET) { + RETURN_IF_ERROR(annotate_paimon_variants(_projected_columns, file_schema, mode, + &_variant_schema_overrides)); + } + return Status::OK(); +} + +Status PaimonReader::customize_file_scan_request(format::FileScanRequest* file_request) { + DORIS_CHECK(file_request != nullptr); + RETURN_IF_ERROR(format::TableReader::customize_file_scan_request(file_request)); + file_request->variant_schema_overrides = _variant_schema_overrides; + return Status::OK(); } Status PaimonReader::_parse_deletion_vector_file(const TTableFormatFileDesc& t_desc, diff --git a/be/src/format_v2/table/paimon_reader.h b/be/src/format_v2/table/paimon_reader.h index 8570f2efba624e..ed2b9e75c1c722 100644 --- a/be/src/format_v2/table/paimon_reader.h +++ b/be/src/format_v2/table/paimon_reader.h @@ -35,10 +35,17 @@ class PaimonReader final : public format::TableReader { #ifdef BE_TEST void TEST_set_scan_params(TFileScanRangeParams* params) { _scan_params = params; } + void TEST_set_projected_columns(std::vector columns) { + _projected_columns = std::move(columns); + } + void TEST_set_format(format::FileFormat format) { _format = format; } format::TableColumnMappingMode TEST_mapping_mode() const { return mapping_mode(); } Status TEST_annotate_file_schema(std::vector* file_schema) { return annotate_file_schema(file_schema); } + Status TEST_customize_file_scan_request(format::FileScanRequest* request) { + return customize_file_scan_request(request); + } Status TEST_parse_deletion_vector_file(const TTableFormatFileDesc& t_desc, DeleteFileDesc* desc, bool* has_delete_file) { return _parse_deletion_vector_file(t_desc, desc, has_delete_file); @@ -48,12 +55,14 @@ class PaimonReader final : public format::TableReader { protected: format::TableColumnMappingMode mapping_mode() const override; Status annotate_file_schema(std::vector* file_schema) override; + Status customize_file_scan_request(format::FileScanRequest* file_request) override; Status _parse_deletion_vector_file(const TTableFormatFileDesc& t_desc, DeleteFileDesc* desc, bool* has_delete_file) override; private: int64_t _split_schema_id = -1; + std::vector _variant_schema_overrides; }; // Paimon scans can contain both native data-file splits and serialized JNI splits in the same diff --git a/be/src/runtime/runtime_profile.cpp b/be/src/runtime/runtime_profile.cpp index 8c7e9170758c63..311d434f6bcff9 100644 --- a/be/src/runtime/runtime_profile.cpp +++ b/be/src/runtime/runtime_profile.cpp @@ -521,6 +521,28 @@ RuntimeProfile::Counter* RuntimeProfile::add_counter(const std::string& name, TU return counter; } +std::shared_ptr RuntimeProfile::add_shared_counter( + const std::string& name, TUnit::type type, const std::string& parent_counter_name, + int64_t level) { + std::lock_guard l(_counter_map_lock); + + if (auto it = _shared_counter_pool.find(name); it != _shared_counter_pool.end()) { + DCHECK_EQ(it->second->type(), type); + return it->second; + } + + // A raw counter with the same name cannot be safely upgraded because external users may + // already hold its profile-owned address. + DCHECK(_counter_map.find(name) == _counter_map.end()); + DCHECK(parent_counter_name == ROOT_COUNTER || + _counter_map.find(parent_counter_name) != _counter_map.end()); + auto counter = std::make_shared(type, 0, level); + _shared_counter_pool.emplace(name, counter); + _counter_map[name] = counter.get(); + _child_counter_map[parent_counter_name].insert(name); + return counter; +} + RuntimeProfile::NonZeroCounter* RuntimeProfile::add_nonzero_counter( const std::string& name, TUnit::type type, const std::string& parent_counter_name, int64_t level) { diff --git a/be/src/runtime/runtime_profile.h b/be/src/runtime/runtime_profile.h index d7948a5a3f9c8d..e63480711ac397 100644 --- a/be/src/runtime/runtime_profile.h +++ b/be/src/runtime/runtime_profile.h @@ -536,6 +536,13 @@ class RuntimeProfile { return add_counter(name, type, RuntimeProfile::ROOT_COUNTER, level); } + // Add a counter whose storage may outlive this profile. Repeated registration returns the same + // shared counter, matching add_counter() semantics for reused scanner profiles. + std::shared_ptr add_shared_counter( + const std::string& name, TUnit::type type, + const std::string& parent_counter_name = RuntimeProfile::ROOT_COUNTER, + int64_t level = 2); + NonZeroCounter* add_nonzero_counter( const std::string& name, TUnit::type type, const std::string& parent_counter_name = RuntimeProfile::ROOT_COUNTER, @@ -659,7 +666,7 @@ class RuntimeProfile { std::unique_ptr _pool; // Pool for allocated counters. These counters are shared with some other objects. - std::map> _shared_counter_pool; + std::map> _shared_counter_pool; // Name for this runtime profile. std::string _name; diff --git a/be/test/core/column/column_variant_v2_test.cpp b/be/test/core/column/column_variant_v2_test.cpp index efca4e0ca999db..c9482107b46fd4 100644 --- a/be/test/core/column/column_variant_v2_test.cpp +++ b/be/test/core/column/column_variant_v2_test.cpp @@ -220,6 +220,56 @@ struct OwnedEncodedData { } }; +class CountingShreddedState final : public VariantShreddedState { +public: + CountingShreddedState(size_t rows, std::shared_ptr size_calls) + : _rows(rows), + _size_calls(std::move(size_calls)), + _serialized(ColumnVariantV2::create()) { + _serialized->insert_many_defaults(rows); + } + + size_t size() const override { + ++*_size_calls; + return _rows; + } + size_t byte_size() const override { return 0; } + size_t allocated_bytes() const override { return 0; } + void sanity_check() const override {} + void for_each_subcolumn(const IColumn::ImutableColumnCallback&) const override {} + std::shared_ptr filter(const IColumn::Filter& filter, + ssize_t) const override { + return std::make_shared( + std::count(filter.begin(), filter.end(), UInt8 {1}), _size_calls); + } + std::shared_ptr select_range(size_t, size_t length) const override { + return std::make_shared(length, _size_calls); + } + std::shared_ptr select_indices( + const uint32_t* indices_begin, const uint32_t* indices_end) const override { + return std::make_shared(indices_end - indices_begin, _size_calls); + } + bool can_materialize() const override { return false; } + bool try_append(const VariantShreddedState&) override { return false; } + std::optional find_typed_value( + std::span) const override { + return std::nullopt; + } + std::optional find_normalized_value( + std::span) const override { + return std::nullopt; + } + const ColumnVariantV2& materialized_column() const override { + throw Exception(ErrorCode::INTERNAL_ERROR, "counting shredded state cannot materialize"); + } + const ColumnVariantV2& serialized_column() const override { return *_serialized; } + +private: + size_t _rows; + std::shared_ptr _size_calls; + ColumnVariantV2::MutablePtr _serialized; +}; + template void expect_not_implemented(Function&& function, std::string_view marker) { try { @@ -1249,6 +1299,23 @@ TEST(ColumnVariantV2Test, PermuteMatchesColumnStringAndRejectsInvalidInputs) { expect_values_match(*source, *reference); } +TEST(ColumnVariantV2Test, CompositeShreddedSizeDoesNotRecountSegments) { + auto size_calls = std::make_shared(0); + auto first = ColumnVariantV2::create_shredded( + std::make_shared(2, size_calls)); + auto second = ColumnVariantV2::create_shredded( + std::make_shared(3, size_calls)); + auto composite = ColumnVariantV2::create(); + composite->insert_range_from(*first, 0, first->size()); + composite->insert_range_from(*second, 0, second->size()); + + *size_calls = 0; + for (size_t iteration = 0; iteration < 32; ++iteration) { + EXPECT_EQ(composite->size(), 5); + } + EXPECT_EQ(*size_calls, 0); +} + TEST(ColumnVariantV2Test, PopBackAndResizeCoverBoundsShrinkAndGrowth) { auto pop_column = ColumnVariantV2::create(); auto pop_reference = ColumnString::create(); diff --git a/be/test/format_v2/parquet/parquet_schema_test.cpp b/be/test/format_v2/parquet/parquet_schema_test.cpp index acd5e8600789ec..a92a4dee22842e 100644 --- a/be/test/format_v2/parquet/parquet_schema_test.cpp +++ b/be/test/format_v2/parquet/parquet_schema_test.cpp @@ -227,6 +227,97 @@ TEST(ParquetSchemaTest, NativeSchemaRecognizesVariantLogicalGroup) { } } +TEST(ParquetSchemaTest, AppliesTableFormatVariantOverrideToUnannotatedGroup) { + auto schema = unshredded_variant_schema(); + schema[1].__isset.logicalType = false; + NativeFieldDescriptor descriptor; + ASSERT_TRUE(descriptor.parse_from_thrift(schema).ok()); + + std::vector> fields; + ASSERT_TRUE(build_parquet_column_schema(descriptor, &fields).ok()); + ASSERT_EQ(fields.size(), 1); + ASSERT_EQ(fields[0]->kind, ParquetColumnSchemaKind::STRUCT); + + const std::vector overrides {format::LocalColumnIndex::top_level(format::LocalColumnId(0))}; + const auto status = apply_variant_schema_overrides(descriptor, overrides, &fields); + ASSERT_TRUE(status.ok()) << status; + EXPECT_EQ(fields[0]->kind, ParquetColumnSchemaKind::VARIANT); + EXPECT_TRUE(fields[0]->contains_variant); + EXPECT_EQ(remove_nullable(fields[0]->type)->get_primitive_type(), TYPE_VARIANT); + EXPECT_NE(fields[0]->variant_physical_type, nullptr); +} + +TEST(ParquetSchemaTest, AppliesPaimonShreddedVariantOverrideWithOptionalMetadata) { + auto schema = shredded_object_variant_schema(); + schema[1].__isset.logicalType = false; + schema[2].__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL); + NativeFieldDescriptor descriptor; + ASSERT_TRUE(descriptor.parse_from_thrift(schema).ok()); + + std::vector> fields; + ASSERT_TRUE(build_parquet_column_schema(descriptor, &fields).ok()); + const std::vector overrides {format::LocalColumnIndex::top_level(format::LocalColumnId(0))}; + const auto status = apply_variant_schema_overrides(descriptor, overrides, &fields); + ASSERT_TRUE(status.ok()) << status; + EXPECT_EQ(fields[0]->kind, ParquetColumnSchemaKind::VARIANT); +} + +TEST(ParquetSchemaTest, RejectsMalformedUnannotatedVariantOverride) { + auto schema = unshredded_variant_schema(); + schema[1].__isset.logicalType = false; + schema[2].__set_name("unexpected"); + NativeFieldDescriptor descriptor; + ASSERT_TRUE(descriptor.parse_from_thrift(schema).ok()); + + std::vector> fields; + ASSERT_TRUE(build_parquet_column_schema(descriptor, &fields).ok()); + const std::vector overrides {format::LocalColumnIndex::top_level(format::LocalColumnId(0))}; + const auto status = apply_variant_schema_overrides(descriptor, overrides, &fields); + EXPECT_TRUE(status.is()) << status; + EXPECT_NE(status.to_string().find("unexpected child"), std::string::npos); +} + +TEST(ParquetSchemaTest, RejectsOptionalMetadataOutsidePaimonShreddedOverride) { + auto annotated_shredded = shredded_object_variant_schema(); + annotated_shredded[2].__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL); + NativeFieldDescriptor descriptor; + const auto annotated_status = descriptor.parse_from_thrift(annotated_shredded); + EXPECT_TRUE(annotated_status.is()) << annotated_status; + + auto unannotated_unshredded = unshredded_variant_schema(); + unannotated_unshredded[1].__isset.logicalType = false; + unannotated_unshredded[2].__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL); + ASSERT_TRUE(descriptor.parse_from_thrift(unannotated_unshredded).ok()); + std::vector> fields; + ASSERT_TRUE(build_parquet_column_schema(descriptor, &fields).ok()); + const std::vector overrides {format::LocalColumnIndex::top_level(format::LocalColumnId(0))}; + const auto override_status = apply_variant_schema_overrides(descriptor, overrides, &fields); + EXPECT_TRUE(override_status.is()) << override_status; +} + +TEST(ParquetSchemaTest, AppliesNestedTableFormatVariantOverride) { + auto schema = struct_with_variant_schema(); + schema[3].__isset.logicalType = false; + NativeFieldDescriptor descriptor; + ASSERT_TRUE(descriptor.parse_from_thrift(schema).ok()); + + std::vector> fields; + ASSERT_TRUE(build_parquet_column_schema(descriptor, &fields).ok()); + ASSERT_EQ(fields.size(), 1); + ASSERT_EQ(fields[0]->kind, ParquetColumnSchemaKind::STRUCT); + ASSERT_EQ(fields[0]->children.size(), 2); + ASSERT_EQ(fields[0]->children[1]->kind, ParquetColumnSchemaKind::STRUCT); + + auto root_override = format::LocalColumnIndex::partial_local(0); + root_override.children.push_back(format::LocalColumnIndex::local(1)); + const auto status = apply_variant_schema_overrides(descriptor, {root_override}, &fields); + ASSERT_TRUE(status.ok()) << status; + EXPECT_TRUE(fields[0]->contains_variant); + EXPECT_EQ(fields[0]->children[1]->kind, ParquetColumnSchemaKind::VARIANT); + const auto& struct_type = assert_cast(*remove_nullable(fields[0]->type)); + EXPECT_EQ(remove_nullable(struct_type.get_element(1))->get_primitive_type(), TYPE_VARIANT); +} + TEST(ParquetSchemaTest, NativeSchemaAcceptsRequiredAndOptionalVariantGroups) { for (const auto repetition : {tparquet::FieldRepetitionType::REQUIRED, tparquet::FieldRepetitionType::OPTIONAL}) { @@ -358,6 +449,24 @@ TEST(ParquetSchemaTest, NativeVariantRejectsUnsupportedPrimitiveTypePairs) { mismatched_integer.logicalType.INTEGER.__set_isSigned(true); invalid_typed_values.push_back(mismatched_integer); + tparquet::SchemaElement mismatched_int32_annotation; + mismatched_int32_annotation.__set_type(tparquet::Type::INT64); + mismatched_int32_annotation.__set_logicalType(tparquet::LogicalType()); + mismatched_int32_annotation.logicalType.__set_INTEGER(tparquet::IntType()); + mismatched_int32_annotation.logicalType.INTEGER.__set_bitWidth(32); + mismatched_int32_annotation.logicalType.INTEGER.__set_isSigned(true); + mismatched_int32_annotation.__set_converted_type(tparquet::ConvertedType::INT_32); + invalid_typed_values.push_back(mismatched_int32_annotation); + + tparquet::SchemaElement mismatched_int64_annotation; + mismatched_int64_annotation.__set_type(tparquet::Type::INT32); + mismatched_int64_annotation.__set_logicalType(tparquet::LogicalType()); + mismatched_int64_annotation.logicalType.__set_INTEGER(tparquet::IntType()); + mismatched_int64_annotation.logicalType.INTEGER.__set_bitWidth(64); + mismatched_int64_annotation.logicalType.INTEGER.__set_isSigned(true); + mismatched_int64_annotation.__set_converted_type(tparquet::ConvertedType::INT_64); + invalid_typed_values.push_back(mismatched_int64_annotation); + tparquet::SchemaElement mismatched_decimal; mismatched_decimal.__set_type(tparquet::Type::INT32); mismatched_decimal.__set_logicalType(tparquet::LogicalType()); @@ -381,6 +490,25 @@ TEST(ParquetSchemaTest, NativeVariantRejectsUnsupportedPrimitiveTypePairs) { } } +TEST(ParquetSchemaTest, NativeVariantAcceptsIcebergFullWidthSignedIntegerAnnotations) { + for (const auto [physical_type, bit_width] : {std::pair {tparquet::Type::INT32, int8_t {32}}, + std::pair {tparquet::Type::INT64, int8_t {64}}}) { + tparquet::SchemaElement typed_value; + typed_value.__set_type(physical_type); + typed_value.__set_logicalType(tparquet::LogicalType()); + typed_value.logicalType.__set_INTEGER(tparquet::IntType()); + typed_value.logicalType.INTEGER.__set_bitWidth(bit_width); + typed_value.logicalType.INTEGER.__set_isSigned(true); + typed_value.__set_converted_type(bit_width == 32 ? tparquet::ConvertedType::INT_32 + : tparquet::ConvertedType::INT_64); + + NativeFieldDescriptor descriptor; + const auto status = descriptor.parse_from_thrift( + shredded_primitive_variant_schema(std::move(typed_value))); + EXPECT_TRUE(status.ok()) << status; + } +} + TEST(ParquetSchemaTest, NativeVariantRejectsRepeatedOuterGroup) { auto schema = unshredded_variant_schema(); schema[1].__set_repetition_type(tparquet::FieldRepetitionType::REPEATED); diff --git a/be/test/format_v2/parquet/variant_column_reader_test.cpp b/be/test/format_v2/parquet/variant_column_reader_test.cpp index 51a50319285c5b..9e2fc0774baf72 100644 --- a/be/test/format_v2/parquet/variant_column_reader_test.cpp +++ b/be/test/format_v2/parquet/variant_column_reader_test.cpp @@ -51,6 +51,8 @@ #include "core/value/variant/variant_parquet_encoding.h" #include "exprs/function/function_variant_element_v2.h" #include "format_v2/parquet/parquet_column_schema.h" +#include "format_v2/parquet/parquet_profile.h" +#include "runtime/runtime_profile.h" namespace doris::format::parquet { namespace { @@ -237,17 +239,57 @@ MutableColumnPtr projected_shredded_object_physical(const std::vector& return ColumnNullable::create(std::move(root), ColumnUInt8::create(values.size(), 0)); } +MutableColumnPtr projected_shredded_int32_object_physical(const std::vector& values) { + auto integers = ColumnInt32::create(); + integers->get_data().assign(values.begin(), values.end()); + MutableColumns wrapper_fields; + wrapper_fields.push_back( + ColumnNullable::create(std::move(integers), ColumnUInt8::create(values.size(), 0))); + auto wrapper = ColumnStruct::create(std::move(wrapper_fields)); + MutableColumns object_fields; + object_fields.push_back( + ColumnNullable::create(std::move(wrapper), ColumnUInt8::create(values.size(), 0))); + auto object = ColumnStruct::create(std::move(object_fields)); + MutableColumns root_fields; + root_fields.push_back( + ColumnNullable::create(std::move(object), ColumnUInt8::create(values.size(), 0))); + auto root = ColumnStruct::create(std::move(root_fields)); + return ColumnNullable::create(std::move(root), ColumnUInt8::create(values.size(), 0)); +} + +MutableColumnPtr projected_shredded_binary_object_physical( + const std::vector& values) { + std::vector refs; + refs.reserve(values.size()); + for (const auto value : values) { + refs.emplace_back(value.data(), value.size()); + } + MutableColumns wrapper_fields; + wrapper_fields.push_back(nullable_strings(refs, std::vector(values.size(), 0))); + auto wrapper = ColumnStruct::create(std::move(wrapper_fields)); + MutableColumns object_fields; + object_fields.push_back( + ColumnNullable::create(std::move(wrapper), ColumnUInt8::create(values.size(), 0))); + auto object = ColumnStruct::create(std::move(object_fields)); + MutableColumns root_fields; + root_fields.push_back( + ColumnNullable::create(std::move(object), ColumnUInt8::create(values.size(), 0))); + auto root = ColumnStruct::create(std::move(root_fields)); + return ColumnNullable::create(std::move(root), ColumnUInt8::create(values.size(), 0)); +} + MutableColumnPtr root_wrapper(MutableColumns fields, NullMap root_nulls = {0}); MutableColumnPtr nullable_int64(const std::vector& values, const std::vector& nulls); MutableColumnPtr complete_shredded_object_physical(std::string_view residual_key, - int64_t residual_value, int64_t typed_value) { + int64_t residual_value, int64_t typed_value, + uint8_t residual_width = 0) { VariantBatchBuilder builder; auto row = builder.begin_row(); auto object = row.start_object(); object.add_key(StringRef(residual_key.data(), residual_key.size())); - row.add_int(residual_value); + row.add_scalar(VariantScalarRef::integer(residual_value, residual_width)); object.finish(); row.finish(); VariantBatchBuilder batch = builder.finish_batch(); @@ -268,6 +310,51 @@ MutableColumnPtr complete_shredded_object_physical(std::string_view residual_key return root_wrapper(std::move(root_fields)); } +MutableColumnPtr complete_shredded_decimal_object_physical(std::string_view residual_key, + __int128 residual_value, + uint8_t residual_scale, + uint8_t residual_width, + int64_t typed_value) { + VariantBatchBuilder builder; + auto row = builder.begin_row(); + auto object = row.start_object(); + object.add_key(StringRef(residual_key.data(), residual_key.size())); + row.add_scalar(VariantScalarRef::decimal(residual_value, residual_scale, residual_width)); + object.finish(); + row.finish(); + VariantBatchBuilder batch = builder.finish_batch(); + const VariantRef residual = batch.value_at(0); + + MutableColumns wrapper_fields; + wrapper_fields.push_back(nullable_int64({typed_value}, {0})); + MutableColumns object_fields; + object_fields.push_back(ColumnNullable::create(ColumnStruct::create(std::move(wrapper_fields)), + ColumnUInt8::create(1, 0))); + MutableColumns root_fields; + root_fields.push_back( + nullable_strings({StringRef(residual.metadata.data, residual.metadata.size)}, {0})); + root_fields.push_back( + nullable_strings({StringRef(residual.value.data, residual.value.size)}, {0})); + root_fields.push_back(ColumnNullable::create(ColumnStruct::create(std::move(object_fields)), + ColumnUInt8::create(1, 0))); + return root_wrapper(std::move(root_fields)); +} + +MutableColumnPtr projected_shredded_decimal_object_physical(__int128 value, uint32_t scale) { + auto decimals = ColumnDecimal128V3::create(0, scale); + decimals->insert_value(Decimal128V3 {value}); + MutableColumns wrapper_fields; + wrapper_fields.push_back( + ColumnNullable::create(std::move(decimals), ColumnUInt8::create(1, 0))); + MutableColumns object_fields; + object_fields.push_back(ColumnNullable::create(ColumnStruct::create(std::move(wrapper_fields)), + ColumnUInt8::create(1, 0))); + MutableColumns root_fields; + root_fields.push_back(ColumnNullable::create(ColumnStruct::create(std::move(object_fields)), + ColumnUInt8::create(1, 0))); + return root_wrapper(std::move(root_fields)); +} + MutableColumnPtr projected_two_field_object_physical(const std::vector& first, const std::vector& second) { DORIS_CHECK(first.size() == second.size()); @@ -343,6 +430,17 @@ MutableColumnPtr nullable_int64(const std::vector& values, return ColumnNullable::create(std::move(data), std::move(null_map)); } +MutableColumnPtr binary_round_trip(const ColumnVariantV2& source) { + DataTypeVariantV2 type; + const int64_t maximum_size = type.get_uncompressed_serialized_bytes(source, 10); + std::vector bytes(maximum_size); + char* end = type.serialize(source, bytes.data(), 10); + bytes.resize(end - bytes.data()); + MutableColumnPtr destination = type.create_column(); + EXPECT_EQ(type.deserialize(bytes.data(), &destination, 10), bytes.data() + bytes.size()); + return destination; +} + template MutableColumnPtr nullable_fixed(std::initializer_list values, std::initializer_list nulls) { @@ -840,6 +938,506 @@ TEST(VariantColumnReaderTest, AppendsProjectedShreddedBatchesWithoutMaterializin EXPECT_EQ(plan.variant_state_schema.use_count(), 3); } +TEST(VariantColumnReaderTest, GathersConsecutiveProjectedShreddedBatches) { + auto schema = shredded_object_schema(); + schema.local_id = 0; + schema.children[2]->local_id = 2; + schema.children[2]->children[0]->local_id = 0; + schema.children[2]->children[0]->children[0]->local_id = 0; + auto projection = format::LocalColumnIndex::partial_local(0); + projection.children.push_back(format::LocalColumnIndex::partial_local(2)); + projection.children.back().children.push_back(format::LocalColumnIndex::partial_local(0)); + projection.children.back().children.back().children.push_back( + format::LocalColumnIndex::local(0)); + VariantMaterializationNode plan; + plan.schema = &schema; + plan.contains_variant = true; + plan.variant_projection = std::move(projection); + plan.variant_state_schema = create_variant_state_schema(schema, &*plan.variant_projection); + + auto first = make_nullable(std::make_shared())->create_column(); + auto second = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE( + materialize_variant_columns(plan, projected_shredded_object_physical({10, 20}), first) + .ok()); + ASSERT_TRUE(materialize_variant_columns(plan, projected_shredded_object_physical({30}), second) + .ok()); + + auto gathered = make_nullable(std::make_shared())->create_column(); + const std::array first_indices {1, 0}; + const std::array second_indices {0}; + gathered->insert_indices_from(*first, first_indices.begin(), first_indices.end()); + gathered->insert_indices_from(*second, second_indices.begin(), second_indices.end()); + + const auto& variants = assert_cast( + assert_cast(*gathered).get_nested_column()); + const std::array path {VariantShreddedPathSegment { + .kind = VariantShreddedPathSegment::Kind::OBJECT_KEY, .key = StringRef("a")}}; + const auto match = variants.find_shredded_typed_value(path); + ASSERT_TRUE(match.has_value()); + const auto& values = assert_cast( + assert_cast(*match->column).get_nested_column()); + EXPECT_EQ(values.get_data(), ColumnInt64::Container({20, 10, 30})); + + auto restored = binary_round_trip(variants); + const std::array path_segments {VariantElementV2PathSegment::object_key(StringRef("a"))}; + std::unique_ptr resolved_path; + ASSERT_TRUE(resolve_variant_element_v2_path(path_segments, &resolved_path).ok()); + ColumnPtr extracted; + ASSERT_TRUE(extract_variant_element_v2(assert_cast(*restored), + *resolved_path, {}, &extracted) + .ok()); + const auto& restored_values = assert_cast( + assert_cast(*extracted).get_nested_column()); + EXPECT_EQ(restored_values.get_value_ref(0).get_int(), 20); + EXPECT_EQ(restored_values.get_value_ref(1).get_int(), 10); + EXPECT_EQ(restored_values.get_value_ref(2).get_int(), 30); +} + +TEST(VariantColumnReaderTest, SelectsProjectedShreddedRowsWithoutMaterializing) { + auto schema = shredded_object_schema(); + schema.local_id = 0; + schema.children[2]->local_id = 2; + schema.children[2]->children[0]->local_id = 0; + schema.children[2]->children[0]->children[0]->local_id = 0; + auto projection = format::LocalColumnIndex::partial_local(0); + projection.children.push_back(format::LocalColumnIndex::partial_local(2)); + projection.children.back().children.push_back(format::LocalColumnIndex::partial_local(0)); + projection.children.back().children.back().children.push_back( + format::LocalColumnIndex::local(0)); + VariantMaterializationNode plan; + plan.schema = &schema; + plan.contains_variant = true; + plan.variant_projection = std::move(projection); + plan.variant_state_schema = create_variant_state_schema(schema, &*plan.variant_projection); + + auto output = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE(materialize_variant_columns(plan, projected_shredded_object_physical({10, 20, 30}), + output) + .ok()); + const auto& variants = assert_cast( + assert_cast(*output).get_nested_column()); + const IColumn::Permutation permutation {2, 0, 1}; + MutableColumnPtr permuted = variants.permute(permutation, 2); + MutableColumnPtr truncated = variants.clone_resized(2); + const std::array path {VariantShreddedPathSegment { + .kind = VariantShreddedPathSegment::Kind::OBJECT_KEY, .key = StringRef("a")}}; + auto verify = [&](const IColumn& column, const ColumnInt64::Container& expected) { + const auto& selected = assert_cast(column); + ASSERT_TRUE(selected.is_shredded()); + const auto match = selected.find_shredded_typed_value(path); + ASSERT_TRUE(match.has_value()); + EXPECT_EQ(assert_cast( + assert_cast(*match->column).get_nested_column()) + .get_data(), + expected); + }; + verify(*permuted, ColumnInt64::Container({30, 10})); + verify(*truncated, ColumnInt64::Container({10, 20})); +} + +TEST(VariantColumnReaderTest, GathersLocalProjectedAndRemoteSerializedRowsInEitherOrder) { + auto schema = shredded_object_schema(); + schema.local_id = 0; + schema.children[2]->local_id = 2; + schema.children[2]->children[0]->local_id = 0; + schema.children[2]->children[0]->children[0]->local_id = 0; + auto projection = format::LocalColumnIndex::partial_local(0); + projection.children.push_back(format::LocalColumnIndex::partial_local(2)); + projection.children.back().children.push_back(format::LocalColumnIndex::partial_local(0)); + projection.children.back().children.back().children.push_back( + format::LocalColumnIndex::local(0)); + VariantMaterializationNode plan; + plan.schema = &schema; + plan.contains_variant = true; + plan.variant_projection = std::move(projection); + plan.variant_state_schema = create_variant_state_schema(schema, &*plan.variant_projection); + + auto output = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE( + materialize_variant_columns(plan, projected_shredded_object_physical({10, 20}), output) + .ok()); + const auto& local = assert_cast( + assert_cast(*output).get_nested_column()); + ASSERT_TRUE(local.is_shredded()); + MutableColumnPtr remote = binary_round_trip(local); + ASSERT_FALSE(assert_cast(*remote).is_shredded()); + + const std::array path_segments {VariantElementV2PathSegment::object_key(StringRef("a"))}; + std::unique_ptr path; + ASSERT_TRUE(resolve_variant_element_v2_path(path_segments, &path).ok()); + auto verify_result = [&](const ColumnVariantV2& gathered, + const std::array& expected) { + ASSERT_EQ(gathered.size(), expected.size()); + + ColumnPtr extracted; + ASSERT_TRUE(extract_variant_element_v2(gathered, *path, {}, &extracted).ok()); + const auto& values = assert_cast( + assert_cast(*extracted).get_nested_column()); + for (size_t row = 0; row < expected.size(); ++row) { + EXPECT_EQ(values.get_value_ref(row).get_int(), expected[row]); + } + }; + auto verify = [&](const std::vector& sources, + const std::vector& positions, + const std::array& expected) { + auto gathered = ColumnVariantV2::create(); + gathered->insert_from_multi_column(sources, positions); + verify_result(*gathered, expected); + }; + + verify({&local, remote.get()}, {0, 1}, {10, 20}); + verify({remote.get(), &local}, {1, 0}, {20, 10}); + + const std::array first_row {0}; + const std::array second_row {1}; + auto indexed = ColumnVariantV2::create(); + indexed->insert_indices_from(local, first_row.begin(), first_row.end()); + indexed->insert_indices_from(*remote, second_row.begin(), second_row.end()); + verify_result(*indexed, {10, 20}); + + indexed = ColumnVariantV2::create(); + indexed->insert_indices_from(*remote, second_row.begin(), second_row.end()); + indexed->insert_indices_from(local, first_row.begin(), first_row.end()); + verify_result(*indexed, {20, 10}); +} + +TEST(VariantColumnReaderTest, ShrinksProjectedShreddedStateWithoutMaterializing) { + auto schema = shredded_object_schema(); + schema.local_id = 0; + schema.children[2]->local_id = 2; + schema.children[2]->children[0]->local_id = 0; + schema.children[2]->children[0]->children[0]->local_id = 0; + auto projection = format::LocalColumnIndex::partial_local(0); + projection.children.push_back(format::LocalColumnIndex::partial_local(2)); + projection.children.back().children.push_back(format::LocalColumnIndex::partial_local(0)); + projection.children.back().children.back().children.push_back( + format::LocalColumnIndex::local(0)); + VariantMaterializationNode plan; + plan.schema = &schema; + plan.contains_variant = true; + plan.variant_projection = std::move(projection); + plan.variant_state_schema = create_variant_state_schema(schema, &*plan.variant_projection); + + auto output = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE(materialize_variant_columns(plan, projected_shredded_object_physical({10, 20, 30}), + output) + .ok()); + const auto& variants = assert_cast( + assert_cast(*output).get_nested_column()); + ASSERT_TRUE(variants.is_shredded()); + + ColumnPtr shrink_source = variants.clone_resized(variants.size()); + ColumnPtr shrunk = shrink_source->shrink(2); + const auto& shrunk_variants = assert_cast(*shrunk); + ASSERT_TRUE(shrunk_variants.is_shredded()); + const std::array path {VariantShreddedPathSegment { + .kind = VariantShreddedPathSegment::Kind::OBJECT_KEY, .key = StringRef("a")}}; + const auto match = shrunk_variants.find_shredded_typed_value(path); + ASSERT_TRUE(match.has_value()); + EXPECT_EQ(assert_cast( + assert_cast(*match->column).get_nested_column()) + .get_data(), + ColumnInt64::Container({10, 20})); + + ColumnPtr empty_source = variants.clone_resized(variants.size()); + ColumnPtr empty = empty_source->shrink(0); + EXPECT_EQ(empty->size(), 0); + EXPECT_FALSE(assert_cast(*empty).is_shredded()); +} + +TEST(VariantColumnReaderTest, GathersCompleteAndProjectedShreddedBatches) { + auto projected_schema = shredded_object_schema(); + projected_schema.local_id = 0; + projected_schema.children[0]->local_id = 0; + projected_schema.children[1]->local_id = 1; + projected_schema.children[2]->local_id = 2; + projected_schema.children[2]->children[0]->local_id = 0; + projected_schema.children[2]->children[0]->children[0]->local_id = 0; + + auto projection = format::LocalColumnIndex::partial_local(0); + projection.children.push_back(format::LocalColumnIndex::partial_local(2)); + projection.children.back().children.push_back(format::LocalColumnIndex::partial_local(0)); + projection.children.back().children.back().children.push_back( + format::LocalColumnIndex::local(0)); + VariantMaterializationNode projected_plan; + projected_plan.schema = &projected_schema; + projected_plan.contains_variant = true; + projected_plan.variant_projection = std::move(projection); + projected_plan.variant_state_schema = + create_variant_state_schema(projected_schema, &*projected_plan.variant_projection); + + auto complete_schema = shredded_named_object_schema("b"); + VariantMaterializationNode complete_plan; + complete_plan.schema = &complete_schema; + complete_plan.contains_variant = true; + complete_plan.variant_state_schema = create_variant_state_schema(complete_schema, nullptr); + + auto projected = make_nullable(std::make_shared())->create_column(); + auto complete = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE(materialize_variant_columns(projected_plan, projected_shredded_object_physical({7}), + projected) + .ok()); + ASSERT_TRUE(materialize_variant_columns( + complete_plan, complete_shredded_object_physical("other", 9, 8), complete) + .ok()); + + const std::array selected {0}; + const std::array path_segments {VariantElementV2PathSegment::object_key(StringRef("a"))}; + std::unique_ptr path; + ASSERT_TRUE(resolve_variant_element_v2_path(path_segments, &path).ok()); + auto verify_order = [&](const IColumn& first, const IColumn& second, + const NullMap& expected_nulls, size_t value_row) { + auto gathered = make_nullable(std::make_shared())->create_column(); + gathered->insert_indices_from(first, selected.begin(), selected.end()); + gathered->insert_indices_from(second, selected.begin(), selected.end()); + + const auto& nullable = assert_cast(*gathered); + const auto& variants = assert_cast(nullable.get_nested_column()); + ColumnPtr extracted; + ASSERT_TRUE(extract_variant_element_v2(variants, *path, nullable.get_null_map_data(), + &extracted) + .ok()); + const auto& extracted_nullable = assert_cast(*extracted); + EXPECT_EQ(extracted_nullable.get_null_map_data(), expected_nulls); + const auto& extracted_values = + assert_cast(extracted_nullable.get_nested_column()); + EXPECT_EQ(extracted_values.get_value_ref(value_row).get_int(), 7); + }; + + // Complete and projected files can alternate in either order; a field absent from the + // complete file must contribute NULL without forcing the projected file to materialize. + verify_order(*projected, *complete, NullMap({0, 1}), 0); + verify_order(*complete, *projected, NullMap({1, 0}), 1); +} + +TEST(VariantColumnReaderTest, PreservesPrimitiveWidthsAcrossProjectedFiles) { + auto int64_schema = shredded_object_schema(); + auto int32_schema = shredded_object_schema(); + int32_schema.children[2]->children[0]->children[0]->type = + make_nullable(std::make_shared()); + int32_schema.children[2]->children[0]->children[0]->type_descriptor.integer_bit_width = 32; + for (auto* schema : {&int64_schema, &int32_schema}) { + schema->local_id = 0; + schema->children[2]->local_id = 2; + schema->children[2]->children[0]->local_id = 0; + schema->children[2]->children[0]->children[0]->local_id = 0; + } + auto make_plan = [](const ParquetColumnSchema& schema) { + auto projection = format::LocalColumnIndex::partial_local(0); + projection.children.push_back(format::LocalColumnIndex::partial_local(2)); + projection.children.back().children.push_back(format::LocalColumnIndex::partial_local(0)); + projection.children.back().children.back().children.push_back( + format::LocalColumnIndex::local(0)); + VariantMaterializationNode plan; + plan.schema = &schema; + plan.contains_variant = true; + plan.variant_projection = std::move(projection); + plan.variant_state_schema = create_variant_state_schema(schema, &*plan.variant_projection); + return plan; + }; + auto int64_plan = make_plan(int64_schema); + auto int32_plan = make_plan(int32_schema); + auto int64_rows = make_nullable(std::make_shared())->create_column(); + auto int32_rows = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE(materialize_variant_columns(int64_plan, projected_shredded_object_physical({7}), + int64_rows) + .ok()); + ASSERT_TRUE(materialize_variant_columns( + int32_plan, projected_shredded_int32_object_physical({8}), int32_rows) + .ok()); + + auto gathered = make_nullable(std::make_shared())->create_column(); + const std::array selected {0}; + gathered->insert_indices_from(*int64_rows, selected.begin(), selected.end()); + gathered->insert_indices_from(*int32_rows, selected.begin(), selected.end()); + + const auto& nullable = assert_cast(*gathered); + const auto& variants = assert_cast(nullable.get_nested_column()); + const std::array path_segments {VariantElementV2PathSegment::object_key(StringRef("a"))}; + std::unique_ptr path; + ASSERT_TRUE(resolve_variant_element_v2_path(path_segments, &path).ok()); + ColumnPtr extracted; + ASSERT_TRUE( + extract_variant_element_v2(variants, *path, nullable.get_null_map_data(), &extracted) + .ok()); + const auto& values = assert_cast( + assert_cast(*extracted).get_nested_column()); + EXPECT_EQ(values.get_value_ref(0).primitive_id(), VariantPrimitiveId::INT64); + EXPECT_EQ(values.get_value_ref(1).primitive_id(), VariantPrimitiveId::INT32); +} + +TEST(VariantColumnReaderTest, PreservesWidthsAcrossMaterializedPathFallback) { + auto projected_int_schema = shredded_object_schema(); + auto projected_decimal_schema = shredded_object_schema(); + auto* decimal_leaf = projected_decimal_schema.children[2]->children[0]->children[0].get(); + decimal_leaf->type = make_nullable(std::make_shared(38, 2)); + decimal_leaf->type_descriptor.decimal_precision = 38; + decimal_leaf->type_descriptor.decimal_scale = 2; + auto prepare_projected = [](ParquetColumnSchema& schema) { + schema.local_id = 0; + schema.children[2]->local_id = 2; + schema.children[2]->children[0]->local_id = 0; + schema.children[2]->children[0]->children[0]->local_id = 0; + auto projection = format::LocalColumnIndex::partial_local(0); + projection.children.push_back(format::LocalColumnIndex::partial_local(2)); + projection.children.back().children.push_back(format::LocalColumnIndex::partial_local(0)); + projection.children.back().children.back().children.push_back( + format::LocalColumnIndex::local(0)); + VariantMaterializationNode plan; + plan.schema = &schema; + plan.contains_variant = true; + plan.variant_projection = std::move(projection); + plan.variant_state_schema = create_variant_state_schema(schema, &*plan.variant_projection); + return plan; + }; + auto projected_int_plan = prepare_projected(projected_int_schema); + auto projected_decimal_plan = prepare_projected(projected_decimal_schema); + auto complete_schema = shredded_named_object_schema("b"); + VariantMaterializationNode complete_plan; + complete_plan.schema = &complete_schema; + complete_plan.contains_variant = true; + complete_plan.variant_state_schema = create_variant_state_schema(complete_schema, nullptr); + + auto verify = [&](VariantMaterializationNode& projected_plan, + MutableColumnPtr projected_physical, MutableColumnPtr complete_physical, + VariantPrimitiveId expected_id) { + auto projected = make_nullable(std::make_shared())->create_column(); + auto complete = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE(materialize_variant_columns(projected_plan, std::move(projected_physical), + projected) + .ok()); + ASSERT_TRUE( + materialize_variant_columns(complete_plan, std::move(complete_physical), complete) + .ok()); + const std::array selected {0}; + const std::array path_segments {VariantElementV2PathSegment::object_key(StringRef("a"))}; + std::unique_ptr path; + ASSERT_TRUE(resolve_variant_element_v2_path(path_segments, &path).ok()); + for (const auto order : {std::array {projected.get(), complete.get()}, + std::array {complete.get(), projected.get()}}) { + auto gathered = make_nullable(std::make_shared())->create_column(); + for (const IColumn* source : order) { + gathered->insert_indices_from(*source, selected.begin(), selected.end()); + } + const auto& nullable = assert_cast(*gathered); + const auto& variants = + assert_cast(nullable.get_nested_column()); + ColumnPtr extracted; + ASSERT_TRUE(extract_variant_element_v2(variants, *path, nullable.get_null_map_data(), + &extracted) + .ok()); + const auto& values = assert_cast( + assert_cast(*extracted).get_nested_column()); + EXPECT_EQ(values.get_value_ref(0).primitive_id(), expected_id); + EXPECT_EQ(values.get_value_ref(1).primitive_id(), expected_id); + } + }; + + verify(projected_int_plan, projected_shredded_object_physical({7}), + complete_shredded_object_physical("a", 8, 9, 8), VariantPrimitiveId::INT64); + verify(projected_decimal_plan, projected_shredded_decimal_object_physical(7, 2), + complete_shredded_decimal_object_physical("a", 8, 2, 16, 9), + VariantPrimitiveId::DECIMAL16); +} + +TEST(VariantColumnReaderTest, GathersProjectedShreddedBatchesWithDifferentLeafTypes) { + auto integer_schema = shredded_object_schema(); + auto string_schema = shredded_binary_object_schema(); + for (auto* schema : {&integer_schema, &string_schema}) { + schema->local_id = 0; + schema->children[2]->local_id = 2; + schema->children[2]->children[0]->local_id = 0; + schema->children[2]->children[0]->children[0]->local_id = 0; + } + auto make_plan = [](const ParquetColumnSchema& schema) { + auto projection = format::LocalColumnIndex::partial_local(0); + projection.children.push_back(format::LocalColumnIndex::partial_local(2)); + projection.children.back().children.push_back(format::LocalColumnIndex::partial_local(0)); + projection.children.back().children.back().children.push_back( + format::LocalColumnIndex::local(0)); + VariantMaterializationNode plan; + plan.schema = &schema; + plan.contains_variant = true; + plan.variant_projection = std::move(projection); + plan.variant_state_schema = create_variant_state_schema(schema, &*plan.variant_projection); + return plan; + }; + auto integer_plan = make_plan(integer_schema); + auto string_plan = make_plan(string_schema); + auto integers = make_nullable(std::make_shared())->create_column(); + auto strings = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE(materialize_variant_columns(integer_plan, + projected_shredded_object_physical({7, 8}), integers) + .ok()); + ASSERT_TRUE(materialize_variant_columns( + string_plan, projected_shredded_binary_object_physical({"seven"}), strings) + .ok()); + + auto gathered = make_nullable(std::make_shared())->create_column(); + const std::array integer_rows {0, 1}; + const std::array string_rows {0}; + gathered->insert_indices_from(*integers, integer_rows.begin(), integer_rows.end()); + gathered->insert_indices_from(*strings, string_rows.begin(), string_rows.end()); + + const auto& nullable = assert_cast(*gathered); + const auto& variants = assert_cast(nullable.get_nested_column()); + const std::array path_segments {VariantElementV2PathSegment::object_key(StringRef("a"))}; + std::unique_ptr path; + ASSERT_TRUE(resolve_variant_element_v2_path(path_segments, &path).ok()); + ColumnPtr extracted; + ASSERT_TRUE( + extract_variant_element_v2(variants, *path, nullable.get_null_map_data(), &extracted) + .ok()); + const auto& extracted_variants = assert_cast( + assert_cast(*extracted).get_nested_column()); + EXPECT_EQ(extracted_variants.get_value_ref(0).get_int(), 7); + EXPECT_EQ(extracted_variants.get_value_ref(1).get_int(), 8); + EXPECT_EQ(extracted_variants.get_value_ref(2).get_binary(), StringRef("seven")); + + variants.sanity_check(); + EXPECT_GT(variants.byte_size(), 0); + EXPECT_GE(variants.allocated_bytes(), variants.byte_size()); + + auto ranged = ColumnVariantV2::create(); + ranged->insert_range_from(variants, 1, 2); + const auto ranged_match = + ranged->find_shredded_typed_value(std::array {VariantShreddedPathSegment { + .kind = VariantShreddedPathSegment::Kind::OBJECT_KEY, .key = StringRef("a")}}); + ASSERT_TRUE(ranged_match.has_value()); + ASSERT_TRUE(ranged_match->normalized); + const auto& ranged_values = assert_cast( + assert_cast(*ranged_match->normalized).get_nested_column()); + EXPECT_EQ(ranged_values.get_value_ref(0).get_int(), 8); + EXPECT_EQ(ranged_values.get_value_ref(1).get_binary(), StringRef("seven")); + + const std::array shredded_path {VariantShreddedPathSegment { + .kind = VariantShreddedPathSegment::Kind::OBJECT_KEY, .key = StringRef("a")}}; + auto reordered = ColumnVariantV2::create(); + const std::array reversed {2, 1, 0}; + reordered->insert_indices_from(variants, reversed.begin(), reversed.end()); + const auto reordered_match = reordered->find_shredded_typed_value(shredded_path); + ASSERT_TRUE(reordered_match.has_value()); + ASSERT_TRUE(reordered_match->normalized); + const auto& reordered_values = assert_cast( + assert_cast(*reordered_match->normalized).get_nested_column()); + EXPECT_EQ(reordered_values.get_value_ref(0).get_binary(), StringRef("seven")); + EXPECT_EQ(reordered_values.get_value_ref(1).get_int(), 8); + EXPECT_EQ(reordered_values.get_value_ref(2).get_int(), 7); + + IColumn::Filter keep_integer {1, 0, 0}; + const auto filtered = variants.filter(keep_integer, 1); + const auto filtered_match = + assert_cast(*filtered).find_shredded_typed_value(shredded_path); + ASSERT_TRUE(filtered_match.has_value()); + ASSERT_TRUE(filtered_match->column); + EXPECT_EQ( + assert_cast( + assert_cast(*filtered_match->column).get_nested_column()) + .get_data()[0], + 7); +} + TEST(VariantColumnReaderTest, WideProjectionSharesSchemaAcrossBatchesAndSelections) { constexpr size_t width = 64; constexpr size_t batch_count = 16; @@ -1007,6 +1605,35 @@ TEST(VariantColumnReaderTest, MaterializedCacheParticipatesInMemoryAccounting) { EXPECT_GT(variants.allocated_bytes(), physical_allocated); } +TEST(VariantColumnReaderTest, ShreddedStateOutlivesScannerProfile) { + auto runtime_profile = std::make_unique("variant-reader-test"); + ParquetProfile parquet_profile; + parquet_profile.init(runtime_profile.get()); + ParquetProfile reused_profile; + reused_profile.init(runtime_profile.get()); + EXPECT_EQ(parquet_profile.variant_reconstructed_rows, + reused_profile.variant_reconstructed_rows); + + auto visible_output = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE(materialize_variant_rows(shredded_int64_schema(), shredded_int64_physical({7}), + visible_output, parquet_profile.column_reader_profile()) + .ok()); + const auto& visible_variants = assert_cast( + assert_cast(*visible_output).get_nested_column()); + EXPECT_EQ(visible_variants.get_value_ref(0).get_int(), 7); + EXPECT_EQ(runtime_profile->get_counter("VariantReconstructedRows")->value(), 1); + + auto output = make_nullable(std::make_shared())->create_column(); + ASSERT_TRUE(materialize_variant_rows(shredded_int64_schema(), shredded_int64_physical({42}), + output, parquet_profile.column_reader_profile()) + .ok()); + runtime_profile.reset(); + + const auto& nullable = assert_cast(*output); + const auto& variants = assert_cast(nullable.get_nested_column()); + EXPECT_EQ(variants.get_value_ref(0).get_int(), 42); +} + TEST(VariantColumnReaderTest, MaterializedShreddedCopiesDetachBeforeMutation) { auto first_output = make_nullable(std::make_shared())->create_column(); ASSERT_TRUE(materialize_variant_rows(shredded_int64_schema(), shredded_int64_physical({10, 20}), diff --git a/be/test/format_v2/table/paimon_variant_reader_test.cpp b/be/test/format_v2/table/paimon_variant_reader_test.cpp new file mode 100644 index 00000000000000..2cb8be17b80c8c --- /dev/null +++ b/be/test/format_v2/table/paimon_variant_reader_test.cpp @@ -0,0 +1,548 @@ +// 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 +#include +#include +#include + +#include +#include +#include +#include + +#include "core/assert_cast.h" +#include "core/block/block.h" +#include "core/column/column_nullable.h" +#include "core/column/variant_v2/column_variant_v2.h" +#include "core/data_type/data_type_array.h" +#include "core/data_type/data_type_nullable.h" +#include "core/data_type/data_type_number.h" +#include "core/data_type/data_type_string.h" +#include "core/data_type/data_type_struct.h" +#include "core/data_type/data_type_variant_v2.h" +#include "core/value/variant/variant_batch_builder.h" +#include "format_v2/parquet/parquet_reader.h" +#include "format_v2/table/paimon_reader.h" +#include "io/io_common.h" +#include "runtime/runtime_state.h" + +namespace doris::format { +namespace { + +ColumnDefinition table_column(std::string name, DataTypePtr type) { + ColumnDefinition column; + column.name = std::move(name); + column.type = make_nullable(std::move(type)); + return column; +} + +ColumnDefinition file_column(int32_t local_id, std::string name, DataTypePtr type) { + ColumnDefinition column; + column.local_id = local_id; + column.name = std::move(name); + column.type = make_nullable(std::move(type)); + return column; +} + +std::shared_ptr binary_array(const std::vector& values) { + arrow::BinaryBuilder builder; + for (const auto value : values) { + EXPECT_TRUE(builder.Append(reinterpret_cast(value.data), + static_cast(value.size)) + .ok()); + } + std::shared_ptr result; + EXPECT_TRUE(builder.Finish(&result).ok()); + return result; +} + +std::shared_ptr null_binary_array(size_t rows) { + arrow::BinaryBuilder builder; + for (size_t row = 0; row < rows; ++row) { + EXPECT_TRUE(builder.AppendNull().ok()); + } + std::shared_ptr result; + EXPECT_TRUE(builder.Finish(&result).ok()); + return result; +} + +std::shared_ptr int32_array(const std::vector& values) { + arrow::Int32Builder builder; + EXPECT_TRUE(builder.AppendValues(values).ok()); + std::shared_ptr result; + EXPECT_TRUE(builder.Finish(&result).ok()); + return result; +} + +void write_unannotated_paimon_variant_file(const std::string& path, + const std::vector& values) { + VariantBatchBuilder builder; + for (const auto value : values) { + auto row = builder.begin_row(); + auto object = row.start_object(); + object.add_key(StringRef("n")); + row.add_int(value); + object.finish(); + row.finish(); + } + auto batch = builder.finish_batch(); + std::vector value_rows; + std::vector metadata_rows; + for (size_t row = 0; row < values.size(); ++row) { + const auto encoded = batch.value_at(row); + value_rows.push_back(encoded.value); + metadata_rows.emplace_back(encoded.metadata.data, encoded.metadata.size); + } + + const auto payload_type = arrow::struct_({arrow::field("value", arrow::binary(), false), + arrow::field("metadata", arrow::binary(), false)}); + auto payload_result = arrow::StructArray::Make( + {binary_array(value_rows), binary_array(metadata_rows)}, payload_type->fields()); + ASSERT_TRUE(payload_result.ok()) << payload_result.status(); + auto table = arrow::Table::Make(arrow::schema({arrow::field("payload", payload_type)}), + {*payload_result}); + + auto file_result = arrow::io::FileOutputStream::Open(path); + ASSERT_TRUE(file_result.ok()) << file_result.status(); + ::parquet::WriterProperties::Builder properties; + properties.version(::parquet::ParquetVersion::PARQUET_2_6); + properties.compression(::parquet::Compression::UNCOMPRESSED); + PARQUET_THROW_NOT_OK( + ::parquet::arrow::WriteTable(*table, arrow::default_memory_pool(), *file_result, + static_cast(values.size()), properties.build())); +} + +void write_unannotated_shredded_paimon_variant_file(const std::string& path, + const std::vector& ages) { + VariantBatchBuilder builder; + for (const auto age : ages) { + auto row = builder.begin_row(); + auto object = row.start_object(); + object.add_key(StringRef("age")); + row.add_int(age); + object.finish(); + row.finish(); + } + const auto batch = builder.finish_batch(); + std::vector metadata_rows; + for (size_t row = 0; row < ages.size(); ++row) { + const auto encoded = batch.value_at(row); + metadata_rows.emplace_back(encoded.metadata.data, encoded.metadata.size); + } + + const auto age_wrapper_type = arrow::struct_( + {arrow::field("value", arrow::binary()), arrow::field("typed_value", arrow::int32())}); + auto age_result = arrow::StructArray::Make({null_binary_array(ages.size()), int32_array(ages)}, + age_wrapper_type->fields()); + ASSERT_TRUE(age_result.ok()) << age_result.status(); + const auto typed_value_type = arrow::struct_({arrow::field("age", age_wrapper_type, false)}); + auto typed_value_result = arrow::StructArray::Make({*age_result}, typed_value_type->fields()); + ASSERT_TRUE(typed_value_result.ok()) << typed_value_result.status(); + const auto payload_type = arrow::struct_({arrow::field("metadata", arrow::binary(), false), + arrow::field("value", arrow::binary()), + arrow::field("typed_value", typed_value_type)}); + auto payload_result = arrow::StructArray::Make( + {binary_array(metadata_rows), null_binary_array(ages.size()), *typed_value_result}, + payload_type->fields()); + ASSERT_TRUE(payload_result.ok()) << payload_result.status(); + auto table = arrow::Table::Make(arrow::schema({arrow::field("payload", payload_type)}), + {*payload_result}); + + auto file_result = arrow::io::FileOutputStream::Open(path); + ASSERT_TRUE(file_result.ok()) << file_result.status(); + ::parquet::WriterProperties::Builder properties; + properties.version(::parquet::ParquetVersion::PARQUET_2_6); + properties.compression(::parquet::Compression::UNCOMPRESSED); + PARQUET_THROW_NOT_OK( + ::parquet::arrow::WriteTable(*table, arrow::default_memory_pool(), *file_result, + static_cast(ages.size()), properties.build())); +} + +// Scenario: Paimon 1.3/1.4 writes Variant as an unannotated Parquet group. Only the Paimon table +// schema can distinguish that carrier from an ordinary STRUCT, so the table reader must expose the +// matched file node as Variant while retaining its physical children for native decoding. +TEST(PaimonVariantReaderTest, AnnotatesUnmarkedParquetVariantFromTableSchema) { + const auto binary = std::make_shared(); + auto physical_type = std::make_shared(DataTypes {binary, binary}, + Strings {"value", "metadata"}); + auto payload = file_column(0, "payload", physical_type); + payload.children = {file_column(0, "value", binary), file_column(1, "metadata", binary)}; + std::vector file_schema {std::move(payload)}; + + paimon::PaimonReader reader; + reader.TEST_set_format(FileFormat::PARQUET); + reader.TEST_set_projected_columns( + {table_column("payload", std::make_shared())}); + + ASSERT_TRUE(reader.TEST_annotate_file_schema(&file_schema).ok()); + ASSERT_EQ(file_schema.size(), 1); + EXPECT_EQ(remove_nullable(file_schema[0].type)->get_primitive_type(), TYPE_VARIANT); + ASSERT_EQ(file_schema[0].children.size(), 2); + EXPECT_EQ(file_schema[0].children[0].name, "value"); + EXPECT_EQ(file_schema[0].children[1].name, "metadata"); + + FileScanRequest request; + ASSERT_TRUE(reader.TEST_customize_file_scan_request(&request).ok()); + ASSERT_EQ(request.variant_schema_overrides.size(), 1); + EXPECT_EQ(request.variant_schema_overrides[0].local_id(), 0); + EXPECT_TRUE(request.variant_schema_overrides[0].project_all_children); +} + +TEST(PaimonVariantReaderTest, DoesNotGuessOrdinaryStructWithVariantCarrierNames) { + const auto binary = std::make_shared(); + auto struct_type = std::make_shared(DataTypes {binary, binary}, + Strings {"value", "metadata"}); + auto payload = file_column(0, "payload", struct_type); + payload.children = {file_column(0, "value", binary), file_column(1, "metadata", binary)}; + std::vector file_schema {payload}; + + auto projected = table_column("payload", struct_type); + projected.children = payload.children; + paimon::PaimonReader reader; + reader.TEST_set_format(FileFormat::PARQUET); + reader.TEST_set_projected_columns({std::move(projected)}); + + ASSERT_TRUE(reader.TEST_annotate_file_schema(&file_schema).ok()); + EXPECT_EQ(remove_nullable(file_schema[0].type)->get_primitive_type(), TYPE_STRUCT); + FileScanRequest request; + ASSERT_TRUE(reader.TEST_customize_file_scan_request(&request).ok()); + EXPECT_TRUE(request.variant_schema_overrides.empty()); +} + +TEST(PaimonVariantReaderTest, AnnotatesNestedArrayVariantByStructuralPosition) { + const auto binary = std::make_shared(); + auto carrier_type = std::make_shared(DataTypes {binary, binary}, + Strings {"value", "metadata"}); + auto element = file_column(0, "element", carrier_type); + element.children = {file_column(0, "value", binary), file_column(1, "metadata", binary)}; + auto values = file_column(0, "values", std::make_shared(element.type)); + values.children = {std::move(element)}; + std::vector file_schema {std::move(values)}; + + auto item = table_column("item", std::make_shared()); + auto projected = table_column("values", std::make_shared(item.type)); + projected.children = {std::move(item)}; + paimon::PaimonReader reader; + reader.TEST_set_format(FileFormat::PARQUET); + reader.TEST_set_projected_columns({std::move(projected)}); + + ASSERT_TRUE(reader.TEST_annotate_file_schema(&file_schema).ok()); + const auto& array_type = + assert_cast(*remove_nullable(file_schema[0].type)); + EXPECT_EQ(remove_nullable(array_type.get_nested_type())->get_primitive_type(), TYPE_VARIANT); + ASSERT_EQ(file_schema[0].children.size(), 1); + EXPECT_EQ(remove_nullable(file_schema[0].children[0].type)->get_primitive_type(), TYPE_VARIANT); + + FileScanRequest request; + ASSERT_TRUE(reader.TEST_customize_file_scan_request(&request).ok()); + ASSERT_EQ(request.variant_schema_overrides.size(), 1); + EXPECT_FALSE(request.variant_schema_overrides[0].project_all_children); + ASSERT_EQ(request.variant_schema_overrides[0].children.size(), 1); + EXPECT_TRUE(request.variant_schema_overrides[0].children[0].project_all_children); +} + +TEST(PaimonVariantReaderTest, MergesSiblingNestedVariantOverrides) { + const auto binary = std::make_shared(); + auto carrier_type = std::make_shared(DataTypes {binary, binary}, + Strings {"value", "metadata"}); + auto carrier = [&](int32_t local_id, std::string name) { + auto column = file_column(local_id, std::move(name), carrier_type); + column.children = {file_column(0, "value", binary), file_column(1, "metadata", binary)}; + return column; + }; + auto row = file_column(0, "row", + std::make_shared(DataTypes {carrier_type, carrier_type}, + Strings {"left", "right"})); + row.children = {carrier(0, "left"), carrier(1, "right")}; + std::vector file_schema {std::move(row)}; + + auto projected = table_column( + "row", std::make_shared( + DataTypes {make_nullable(std::make_shared()), + make_nullable(std::make_shared())}, + Strings {"left", "right"})); + projected.children = {table_column("left", std::make_shared()), + table_column("right", std::make_shared())}; + paimon::PaimonReader reader; + reader.TEST_set_format(FileFormat::PARQUET); + reader.TEST_set_projected_columns({std::move(projected)}); + + ASSERT_TRUE(reader.TEST_annotate_file_schema(&file_schema).ok()); + FileScanRequest request; + ASSERT_TRUE(reader.TEST_customize_file_scan_request(&request).ok()); + ASSERT_EQ(request.variant_schema_overrides.size(), 1); + ASSERT_EQ(request.variant_schema_overrides[0].children.size(), 2); + EXPECT_TRUE(request.variant_schema_overrides[0].children[0].project_all_children); + EXPECT_TRUE(request.variant_schema_overrides[0].children[1].project_all_children); +} + +TEST(PaimonVariantReaderTest, ReadsUnannotatedPaimonVariantWithNativeParquetReader) { + const auto test_dir = + std::filesystem::temp_directory_path() / "doris_paimon_native_unannotated_variant_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + const auto file_path = (test_dir / "data.parquet").string(); + write_unannotated_paimon_variant_file(file_path, {1, 2, 3}); + + std::vector projected_columns {table_column("payload", std::make_shared())}; + TFileScanRangeParams scan_params; + scan_params.__set_file_type(TFileType::FILE_LOCAL); + scan_params.__set_format_type(TFileFormatType::FORMAT_PARQUET); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + io::FileReaderStats file_reader_stats; + io::FileCacheStatistics file_cache_stats; + auto io_ctx = std::make_shared(); + io_ctx->file_reader_stats = &file_reader_stats; + io_ctx->file_cache_stats = &file_cache_stats; + + paimon::PaimonReader reader; + ASSERT_TRUE(reader.init({.projected_columns = projected_columns, + .conjuncts = {}, + .format = FileFormat::PARQUET, + .scan_params = &scan_params, + .io_ctx = io_ctx, + .runtime_state = &state, + .scanner_profile = nullptr}) + .ok()); + SplitReadOptions split; + split.current_range.__set_path(file_path); + split.current_range.__set_file_size( + static_cast(std::filesystem::file_size(file_path))); + TTableFormatFileDesc table_format; + table_format.__set_table_format_type("paimon"); + table_format.__set_paimon_params(TPaimonFileDesc {}); + split.current_range.__set_table_format_params(std::move(table_format)); + ASSERT_TRUE(reader.prepare_split(split).ok()); + + std::vector actual; + bool eos = false; + while (!eos) { + Block block; + block.insert({projected_columns[0].type->create_column(), projected_columns[0].type, + projected_columns[0].name}); + ASSERT_TRUE(reader.get_block(&block, &eos).ok()); + if (block.rows() == 0) { + continue; + } + const auto& nullable = assert_cast(*block.get_by_position(0).column); + const auto& variants = assert_cast(nullable.get_nested_column()); + for (size_t row = 0; row < variants.size(); ++row) { + VariantRef n; + ASSERT_TRUE(variants.get_value_ref(row).object_find(StringRef("n"), &n)); + actual.push_back(n.get_int()); + } + } + EXPECT_EQ(actual, std::vector({1, 2, 3})); + ASSERT_TRUE(reader.close().ok()); + std::filesystem::remove_all(test_dir); +} + +TEST(PaimonVariantReaderTest, ParquetReaderAppliesExplicitVariantSchemaOverride) { + const auto test_dir = + std::filesystem::temp_directory_path() / "doris_paimon_parquet_variant_override_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + const auto file_path = (test_dir / "data.parquet").string(); + write_unannotated_paimon_variant_file(file_path, {7, 8}); + + auto system_properties = std::make_shared(); + system_properties->system_type = TFileType::FILE_LOCAL; + auto file_description = std::make_unique(); + file_description->path = file_path; + file_description->file_size = static_cast(std::filesystem::file_size(file_path)); + file_description->range_start_offset = 0; + file_description->range_size = -1; + auto reader = std::make_unique( + system_properties, file_description, std::shared_ptr {}, nullptr); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + ASSERT_TRUE(reader->init(&state).ok()); + std::vector file_schema; + ASSERT_TRUE(reader->get_schema(&file_schema).ok()); + ASSERT_EQ(file_schema.size(), 1); + ASSERT_EQ(remove_nullable(file_schema[0].type)->get_primitive_type(), TYPE_STRUCT); + + auto request = std::make_shared(); + request->non_predicate_columns.push_back( + LocalColumnIndex::top_level(LocalColumnId(file_schema[0].local_id))); + request->local_positions.emplace(LocalColumnId(file_schema[0].local_id), LocalIndex(0)); + request->variant_schema_overrides.push_back( + LocalColumnIndex::top_level(LocalColumnId(file_schema[0].local_id))); + ASSERT_TRUE(reader->open(request).ok()); + + const auto variant_type = make_nullable(std::make_shared()); + std::vector actual; + bool eof = false; + while (!eof) { + Block block; + block.insert({variant_type->create_column(), variant_type, "payload"}); + size_t rows = 0; + ASSERT_TRUE(reader->get_block(&block, &rows, &eof).ok()); + const auto& nullable = assert_cast(*block.get_by_position(0).column); + const auto& variants = assert_cast(nullable.get_nested_column()); + for (size_t row = 0; row < variants.size(); ++row) { + VariantRef n; + ASSERT_TRUE(variants.get_value_ref(row).object_find(StringRef("n"), &n)); + actual.push_back(n.get_int()); + } + } + EXPECT_EQ(actual, std::vector({7, 8})); + ASSERT_TRUE(reader->close().ok()); + std::filesystem::remove_all(test_dir); +} + +TEST(PaimonVariantReaderTest, ReadsProjectedLeafFromUnannotatedShreddedVariant) { + const auto test_dir = std::filesystem::temp_directory_path() / + "doris_paimon_parquet_shredded_variant_override_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + const auto file_path = (test_dir / "data.parquet").string(); + write_unannotated_shredded_paimon_variant_file(file_path, {27, 42}); + + auto system_properties = std::make_shared(); + system_properties->system_type = TFileType::FILE_LOCAL; + auto file_description = std::make_unique(); + file_description->path = file_path; + file_description->file_size = static_cast(std::filesystem::file_size(file_path)); + file_description->range_start_offset = 0; + file_description->range_size = -1; + auto reader = std::make_unique( + system_properties, file_description, std::shared_ptr {}, nullptr); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + ASSERT_TRUE(reader->init(&state).ok()); + std::vector file_schema; + ASSERT_TRUE(reader->get_schema(&file_schema).ok()); + ASSERT_EQ(file_schema.size(), 1); + ASSERT_EQ(remove_nullable(file_schema[0].type)->get_primitive_type(), TYPE_STRUCT); + ASSERT_EQ(file_schema[0].children.size(), 3); + + auto find_child = [](const std::vector& children, + std::string_view name) -> const ColumnDefinition* { + const auto it = std::ranges::find_if( + children, [name](const auto& child) { return child.name == name; }); + return it == children.end() ? nullptr : &*it; + }; + const auto* typed_value = find_child(file_schema[0].children, "typed_value"); + ASSERT_NE(typed_value, nullptr); + const auto* age = find_child(typed_value->children, "age"); + ASSERT_NE(age, nullptr); + const auto* age_typed_value = find_child(age->children, "typed_value"); + ASSERT_NE(age_typed_value, nullptr); + + auto projection = LocalColumnIndex::partial_local(file_schema[0].local_id); + projection.children.push_back(LocalColumnIndex::partial_local(typed_value->local_id)); + projection.children.back().children.push_back(LocalColumnIndex::partial_local(age->local_id)); + projection.children.back().children.back().children.push_back( + LocalColumnIndex::local(age_typed_value->local_id)); + auto request = std::make_shared(); + request->non_predicate_columns.push_back(std::move(projection)); + request->local_positions.emplace(LocalColumnId(file_schema[0].local_id), LocalIndex(0)); + request->variant_schema_overrides.push_back( + LocalColumnIndex::top_level(LocalColumnId(file_schema[0].local_id))); + ASSERT_TRUE(reader->open(request).ok()); + + const auto variant_type = make_nullable(std::make_shared()); + Block block; + block.insert({variant_type->create_column(), variant_type, "payload"}); + size_t rows = 0; + bool eof = false; + while (!eof) { + size_t batch_rows = 0; + ASSERT_TRUE(reader->get_block(&block, &batch_rows, &eof).ok()); + rows += batch_rows; + } + EXPECT_EQ(rows, 2); + const auto& nullable = assert_cast(*block.get_by_position(0).column); + const auto& variants = assert_cast(nullable.get_nested_column()); + const std::array path {VariantShreddedPathSegment { + .kind = VariantShreddedPathSegment::Kind::OBJECT_KEY, .key = StringRef("age")}}; + const auto match = variants.find_shredded_typed_value(path); + ASSERT_TRUE(match.has_value()); + const auto& typed = assert_cast(*match->column); + const auto& values = assert_cast(typed.get_nested_column()); + ASSERT_EQ(values.size(), 2); + EXPECT_EQ(values.get_data()[0], 27); + EXPECT_EQ(values.get_data()[1], 42); + ASSERT_TRUE(reader->close().ok()); + std::filesystem::remove_all(test_dir); +} + +TEST(PaimonVariantReaderTest, AppendsUnshreddedAndShreddedPaimonFilesInEitherOrder) { + const auto test_dir = std::filesystem::temp_directory_path() / + "doris_paimon_parquet_mixed_variant_override_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + const auto unshredded_path = (test_dir / "unshredded.parquet").string(); + const auto shredded_path = (test_dir / "shredded.parquet").string(); + write_unannotated_paimon_variant_file(unshredded_path, {5}); + write_unannotated_shredded_paimon_variant_file(shredded_path, {27}); + + const auto assert_file_order = [&](const std::array& paths, + const std::array& fields, + const std::array& values) { + const auto variant_type = make_nullable(std::make_shared()); + Block block; + block.insert({variant_type->create_column(), variant_type, "payload"}); + const auto append_file = [&](const std::string& path) -> Status { + auto system_properties = std::make_shared(); + system_properties->system_type = TFileType::FILE_LOCAL; + auto file_description = std::make_unique(); + file_description->path = path; + file_description->file_size = static_cast(std::filesystem::file_size(path)); + file_description->range_start_offset = 0; + file_description->range_size = -1; + auto reader = std::make_unique( + system_properties, file_description, std::shared_ptr {}, + nullptr); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + RETURN_IF_ERROR(reader->init(&state)); + auto request = std::make_shared(); + request->non_predicate_columns.push_back(LocalColumnIndex::top_level(LocalColumnId(0))); + request->local_positions.emplace(LocalColumnId(0), LocalIndex(0)); + request->variant_schema_overrides.push_back( + LocalColumnIndex::top_level(LocalColumnId(0))); + RETURN_IF_ERROR(reader->open(request)); + bool eof = false; + while (!eof) { + size_t rows = 0; + RETURN_IF_ERROR(reader->get_block(&block, &rows, &eof)); + } + return reader->close(); + }; + + ASSERT_TRUE(append_file(paths[0]).ok()); + ASSERT_TRUE(append_file(paths[1]).ok()); + ASSERT_EQ(block.rows(), 2); + const auto& nullable = assert_cast(*block.get_by_position(0).column); + const auto& variants = assert_cast(nullable.get_nested_column()); + for (size_t row = 0; row < 2; ++row) { + VariantRef field; + ASSERT_TRUE(variants.get_value_ref(row).object_find(fields[row], &field)); + EXPECT_EQ(field.get_int(), values[row]); + } + }; + + assert_file_order({unshredded_path, shredded_path}, {StringRef("n"), StringRef("age")}, + {5, 27}); + assert_file_order({shredded_path, unshredded_path}, {StringRef("age"), StringRef("n")}, + {27, 5}); + + std::filesystem::remove_all(test_dir); +} + +} // namespace +} // namespace doris::format diff --git a/be/test/util/variant/variant_batch_builder_test.cpp b/be/test/util/variant/variant_batch_builder_test.cpp index f7e55ed1b3ad55..4e1f5817ab6515 100644 --- a/be/test/util/variant/variant_batch_builder_test.cpp +++ b/be/test/util/variant/variant_batch_builder_test.cpp @@ -567,6 +567,23 @@ TEST(VariantBatchBuilderTest, IntegerAndDecimalWidthsAreMinimal) { EXPECT_EQ(value.array_at(integers.size() + decimals.size()).get_decimal().width, 16); } +TEST(VariantBatchBuilderTest, AddValuePreservesExplicitPrimitiveWidths) { + const OwnedBuilderValue source = build_owned_value([](VariantBatchBuilder::Row& row) { + auto array = row.start_array(); + row.add_scalar(VariantScalarRef::integer(7, 8)); + row.add_scalar(VariantScalarRef::decimal(8, 2, 16)); + array.finish(); + }); + VariantBatchBuilder builder; + auto row = builder.begin_row(); + row.add_value(source.ref()); + row.finish(); + VariantBatchBuilder imported = builder.finish_batch(); + + EXPECT_EQ(imported.value_at(0).array_at(0).primitive_id(), VariantPrimitiveId::INT64); + EXPECT_EQ(imported.value_at(0).array_at(1).primitive_id(), VariantPrimitiveId::DECIMAL16); +} + TEST(VariantBatchBuilderTest, DecimalValidationLargeIntFallbackAndExplicitWidths) { VariantBatchBuilder builder; auto row = builder.begin_row(); diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql index 51ad9973ce970a..0f5804f610b6fe 100644 --- a/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql +++ b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql @@ -15,3 +15,56 @@ insert into variant_smoke values (1, parse_json('{"name":"alice","age":18,"active":true,"score":98.5,"tags":["flink","paimon"],"profile":{"city":"beijing","zip":100000},"missing":null}')), (2, parse_json('{"name":"bob","age":30,"active":false,"tags":["doris"],"profile":{"city":"shanghai"},"extra":{"levels":[1,2,3]}}')), (3, parse_json('[1,"mixed",false,null,{"k":"v"}]')); + +drop table if exists variant_shredded; +create table variant_shredded ( + id BIGINT, + event_date DATE, + payload VARIANT +) using paimon +partitioned by (event_date) +tblproperties ( + 'file.format' = 'parquet', + 'parquet.variant.shreddingSchema' = '{"type":"ROW","fields":[{"name":"payload","type":{"type":"ROW","fields":[{"name":"name","type":"STRING"},{"name":"age","type":"INT"}]}}]}' +); + +insert into variant_shredded values + (1, date '2026-06-01', parse_json('{"name":"alice","age":18,"extra":"shredded"}')), + (2, date '2026-06-01', parse_json('{"name":"bob","age":30}')); + +drop table if exists variant_mixed_us; +create table variant_mixed_us ( + id BIGINT, + event_date DATE, + payload VARIANT +) using paimon +partitioned by (event_date) +tblproperties ('file.format' = 'parquet'); + +insert into variant_mixed_us values + (1, date '2026-06-01', parse_json('{"name":"alice","age":18,"layout":"unshredded"}')); +alter table variant_mixed_us set tblproperties ( + 'parquet.variant.shreddingSchema' = '{"type":"ROW","fields":[{"name":"payload","type":{"type":"ROW","fields":[{"name":"name","type":"STRING"},{"name":"age","type":"INT"}]}}]}' +); +insert into variant_mixed_us values + (2, date '2026-07-01', parse_json('{"name":"bob","age":30,"layout":"shredded"}')); + +drop table if exists variant_mixed_su; +create table variant_mixed_su ( + id BIGINT, + event_date DATE, + payload VARIANT +) using paimon +partitioned by (event_date) +tblproperties ( + 'file.format' = 'parquet', + 'parquet.variant.shreddingSchema' = '{"type":"ROW","fields":[{"name":"payload","type":{"type":"ROW","fields":[{"name":"name","type":"STRING"},{"name":"age","type":"INT"}]}}]}' +); + +insert into variant_mixed_su values + (1, date '2026-06-01', parse_json('{"name":"alice","age":18,"layout":"shredded"}')); +alter table variant_mixed_su set tblproperties ( + 'parquet.variant.shreddingSchema' = '' +); +insert into variant_mixed_su values + (2, date '2026-07-01', parse_json('{"name":"bob","age":30,"layout":"unshredded"}')); diff --git a/regression-test/data/external_table_p0/iceberg/test_iceberg_variant_read.out b/regression-test/data/external_table_p0/iceberg/test_iceberg_variant_read.out index 88408fc6dea795..dc5e8bc3d42ccb 100644 --- a/regression-test/data/external_table_p0/iceberg/test_iceberg_variant_read.out +++ b/regression-test/data/external_table_p0/iceberg/test_iceberg_variant_read.out @@ -66,6 +66,12 @@ 4 40 \N \N \N {"c":4,"shared":40} 5 50 \N 5 500 {"b":5,"new_field":{"k":500},"shared":50} +-- !variant_projected_remote_gather -- +4093 4093 +4094 4094 +4095 4095 +5000 5000 + -- !variant_type_matrix -- true -128 -32768 2147483647 -9223372036854775808 true true -1234567890.1234 1970-01-02 1970-01-01T00:00:01.234567 "YmluYXJ5" false diff --git a/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out b/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out index f8966577392f4c..830e466af48556 100644 --- a/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out +++ b/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out @@ -24,3 +24,45 @@ payload variant= 20 ORDER BY id """ - String parallelScanToken = - "iceberg_variant_parallel_scan_" + UUID.randomUUID().toString() - List> parallelScanRows = sql """ - SELECT '${parallelScanToken}', id, - CAST(v['shared'] AS INT), - CAST(v['a'] AS INT), - CAST(v['b'] AS INT), - CAST(v['new_field']['k'] AS INT), - CAST(v AS STRING) - FROM variant_multi_file - WHERE v['shared'] >= 20 + // The stable snapshot contributes a genuinely shredded file, while the appended file uses + // the unshredded fallback. More than four rows qualify, forcing local TopN overshoot to be + // truncated after the merge exchange while the mapper-eligible projected path crosses the wire. + explain { + sql """ + SELECT id, CAST(projected['n'] AS INT) + FROM ( + SELECT id, v AS projected + FROM variant_page_pruning FOR VERSION AS OF ${mixedBeforeDeleteSnapshot} + WHERE CAST(v['n'] AS INT) > 3000 + ORDER BY id DESC + LIMIT 4 + ) gathered + """ + contains "VMERGING-EXCHANGE" + contains "inputSplitNum=2" + contains "all access paths: [v(2).n]" + } + String projectedGatherToken = + "iceberg_variant_projected_remote_gather_" + UUID.randomUUID().toString() + List> projectedGatherRows = sql """ + SELECT '${projectedGatherToken}', id, CAST(projected['n'] AS INT) + FROM ( + SELECT id, v AS projected + FROM variant_page_pruning FOR VERSION AS OF ${mixedBeforeDeleteSnapshot} + WHERE CAST(v['n'] AS INT) > 3000 + ORDER BY id DESC + LIMIT 4 + ) gathered + ORDER BY id + """ + assertEquals(4, projectedGatherRows.size()) + String projectedGatherProfile = getProfileByToken(projectedGatherToken, + ["VariantLeafProjections", "VariantDirectLeafPathMisses"]).toString() + assertTrue(counterSum(projectedGatherProfile, "VariantLeafProjections") > 0, + "The projected TopN did not read a physical shredded Variant leaf") + assertTrue(counterSum(projectedGatherProfile, "VariantDirectLeafPathMisses") > 0, + "The projected TopN did not combine the unshredded fallback file") + order_qt_variant_projected_remote_gather """ + SELECT id, + CAST(projected['n'] AS INT) + FROM ( + SELECT id, v AS projected + FROM variant_page_pruning FOR VERSION AS OF ${mixedBeforeDeleteSnapshot} + WHERE CAST(v['n'] AS INT) > 3000 + ORDER BY id DESC + LIMIT 4 + ) gathered ORDER BY id """ - assertEquals(4, parallelScanRows.size(), - "The parallel Variant query must read rows from multiple data files") - String parallelScanProfile = profileAction.getProfileBySql( - parallelScanToken, ["PerScannerRowsRead"]) - if (profileInfoValues(parallelScanProfile, "PerScannerRowsRead") - .count { long rows -> rows > 0 } <= 1) { - parallelScanProfile = profileAction.waitProfile({ - String profile = profileAction.getProfileBySql( - parallelScanToken, ["PerScannerRowsRead"]) - return profileInfoValues(profile, "PerScannerRowsRead") - .count { long rows -> rows > 0 } > 1 ? profile : "" - }, [], "Completed parallel Variant profile with multiple non-empty scanners") - } - assertTrue(profileInfoValues(parallelScanProfile, "PerScannerRowsRead") - .count { long rows -> rows > 0 } > 1, - "The parallel Variant query did not use multiple non-empty scanners") sql "set min_file_scanners_concurrency=1" order_qt_variant_type_matrix """ diff --git a/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy b/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy index e5a2cf79decee0..203d5788310aa2 100644 --- a/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy +++ b/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy @@ -34,6 +34,7 @@ suite("test_paimon_catalog_variant", "p0,external,doris,external_docker,external "s3.path.style.access" = "true" );""" sql """use `${catalogName}`.`test_paimon_spark`""" + sql """set enable_variant_v2 = true""" sql """set force_jni_scanner = true""" explain { @@ -87,5 +88,137 @@ suite("test_paimon_catalog_variant", "p0,external,doris,external_docker,external """ sql """set force_jni_scanner = false""" + + explain { + sql "select * from variant_smoke order by id" + check { explainString -> + def nativeSplits = explainString =~ /paimonNativeReadSplits=(\d+)\/(\d+)/ + // Paimon can change the physical split count; every planned split must stay native. + return nativeSplits.find() + && nativeSplits.group(1).toInteger() > 0 + && nativeSplits.group(1) == nativeSplits.group(2) + } + } + + order_qt_native_full_variant """ + select id, payload + from variant_smoke + order by id + """ + + order_qt_native_object_subpaths """ + select id, + cast(payload['name'] as string), + cast(payload['age'] as int), + cast(payload['profile']['city'] as string), + cast(payload['active'] as boolean) + from variant_smoke + order by id + """ + + order_qt_native_null_and_missing """ + select id, + payload['missing'] is null, + payload['not_exist'] is null + from variant_smoke + order by id + """ + + order_qt_native_root_array """ + select id, + cast(payload[1] as int), + cast(payload[2] as string), + cast(payload[3] as boolean), + cast(payload[4] as string), + cast(payload[5]['k'] as string) + from variant_smoke + where id = 3 + order by id + """ + + order_qt_native_subpath_predicate """ + select id, cast(payload['name'] as string) + from variant_smoke + where cast(payload['age'] as int) >= 20 + order by id + """ + + ["variant_shredded", "variant_mixed_us", "variant_mixed_su"].each { tableName -> + explain { + sql "select * from ${tableName} order by id" + check { explainString -> + def nativeSplits = explainString =~ /paimonNativeReadSplits=(\d+)\/(\d+)/ + return nativeSplits.find() + && nativeSplits.group(1).toInteger() > 0 + && nativeSplits.group(1) == nativeSplits.group(2) + } + } + } + + order_qt_native_shredded_projection """ + select id, + cast(payload['name'] as string), + cast(payload['age'] as int), + cast(payload['extra'] as string) + from variant_shredded + where cast(payload['age'] as int) >= 20 + order by id + """ + + order_qt_native_mixed_us_partitions """ + select id, + cast(event_date as string), + cast(payload['name'] as string), + cast(payload['age'] as int), + cast(payload['layout'] as string) + from variant_mixed_us + order by id + """ + + order_qt_native_mixed_us_root """ + select id, cast(payload as string) + from variant_mixed_us + order by id + """ + + order_qt_native_mixed_su_partitions """ + select id, + cast(event_date as string), + cast(payload['name'] as string), + cast(payload['age'] as int), + cast(payload['layout'] as string) + from variant_mixed_su + order by id + """ + + order_qt_native_mixed_su_root """ + select id, cast(payload as string) + from variant_mixed_su + order by id + """ + + String internalDb = context.config.getDbNameByFile(context.file) + String mvName = "paimon_variant_mixed_mv" + sql """switch internal""" + sql """use `${internalDb}`""" + sql """drop materialized view if exists ${mvName}""" + try { + sql """ + create materialized view ${mvName} + build deferred refresh complete on manual + distributed by random buckets 1 + properties ('replication_num' = '1') + as + select cast(payload['name'] as string) as name, count(*) as row_count + from ${catalogName}.`test_paimon_spark`.variant_mixed_us + where event_date = '2026-06-01' + group by cast(payload['name'] as string) + """ + sql """refresh materialized view ${mvName} complete""" + waitingMTMVTaskFinishedByMvName(mvName) + order_qt_native_mixed_us_mtmv """select * from ${mvName} order by name""" + } finally { + sql """drop materialized view if exists ${mvName}""" + } } }