This is an automated email from the ASF dual-hosted git repository.

Gabriel39 pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new 7d9cf5dbe2f [fix](parquet) Isolate Variant planning from ordinary 
scans (#66441)
7d9cf5dbe2f is described below

commit 7d9cf5dbe2f6f2eebd6f87b418089044f1678daa
Author: Gabriel <[email protected]>
AuthorDate: Tue Aug 4 22:36:01 2026 +0800

    [fix](parquet) Isolate Variant planning from ordinary scans (#66441)
    
    ### What
    
    - Cache whether a Parquet schema subtree contains Variant data.
    - Build Variant physical types and materialization plans only for
    Variant-containing columns.
    - Keep ordinary Parquet scans on their existing logical-type path.
    
    ### Why
    
    Variant planning was performed for every projected column and row group,
    even when the schema contained no Variant data. This added repeated
    schema traversal and allocation overhead to ordinary scans.
    
    ### Validation
    
    - Cherry-picked cleanly onto the latest `branch-4.1`.
    - The resulting patch is equivalent to the original isolated fix.
    - Local tests were not run for this cherry-pick-only PR.
---
 be/src/format_v2/parquet/parquet_column_schema.cpp | 14 +++---------
 be/src/format_v2/parquet/parquet_column_schema.h   |  4 ++++
 be/src/format_v2/parquet/parquet_file_context.cpp  |  2 ++
 be/src/format_v2/parquet/parquet_file_context.h    |  3 +++
 be/src/format_v2/parquet/parquet_reader.cpp        | 25 ++++++++++++++++------
 be/src/format_v2/parquet/parquet_statistics.cpp    | 10 +++++++--
 .../parquet/reader/native_column_reader.cpp        | 21 +++++++++++-------
 .../parquet/reader/variant_column_reader.cpp       |  1 +
 be/test/format_v2/parquet/parquet_schema_test.cpp  |  8 +++++++
 .../format_v2/parquet/parquet_statistics_test.cpp  |  1 +
 .../parquet/variant_column_reader_test.cpp         |  1 +
 11 files changed, 62 insertions(+), 28 deletions(-)

diff --git a/be/src/format_v2/parquet/parquet_column_schema.cpp 
b/be/src/format_v2/parquet/parquet_column_schema.cpp
index 623f1ece3dd..71416e17dc9 100644
--- a/be/src/format_v2/parquet/parquet_column_schema.cpp
+++ b/be/src/format_v2/parquet/parquet_column_schema.cpp
@@ -178,16 +178,6 @@ void propagate_native_max_levels(ParquetColumnSchema* 
schema) {
     }
 }
 
-bool contains_variant_node(const ParquetColumnSchema& schema) {
-    if (schema.kind == ParquetColumnSchemaKind::VARIANT) {
-        return true;
-    }
-    return std::ranges::any_of(schema.children, [](const auto& child) {
-        DORIS_CHECK(child != nullptr);
-        return contains_variant_node(*child);
-    });
-}
-
 std::unique_ptr<ParquetColumnSchema> build_native_node_schema(const 
NativeFieldSchema& field,
                                                               int32_t 
local_id) {
     auto result = std::make_unique<ParquetColumnSchema>();
@@ -222,6 +212,7 @@ std::unique_ptr<ParquetColumnSchema> 
build_native_node_schema(const NativeFieldS
     }
     if (field.variant_physical_type != nullptr) {
         result->kind = ParquetColumnSchemaKind::VARIANT;
+        result->contains_variant = true;
     } else if (primitive_type == TYPE_ARRAY) {
         result->kind = ParquetColumnSchemaKind::LIST;
     } else if (primitive_type == TYPE_MAP) {
@@ -233,10 +224,11 @@ std::unique_ptr<ParquetColumnSchema> 
build_native_node_schema(const NativeFieldS
     for (size_t child_idx = 0; child_idx < field.children.size(); ++child_idx) 
{
         result->children.push_back(
                 build_native_node_schema(field.children[child_idx], 
cast_set<int32_t>(child_idx)));
+        result->contains_variant |= result->children.back()->contains_variant;
     }
     // 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 && 
contains_variant_node(*result)) {
+    if (result->kind != ParquetColumnSchemaKind::VARIANT && 
result->contains_variant) {
         DataTypePtr logical_type;
         if (result->kind == ParquetColumnSchemaKind::LIST) {
             DORIS_CHECK(result->children.size() == 1);
diff --git a/be/src/format_v2/parquet/parquet_column_schema.h 
b/be/src/format_v2/parquet/parquet_column_schema.h
index 4cb29501833..3103b97cfc2 100644
--- a/be/src/format_v2/parquet/parquet_column_schema.h
+++ b/be/src/format_v2/parquet/parquet_column_schema.h
@@ -57,6 +57,10 @@ struct ParquetColumnSchema {
 
     ParquetColumnSchemaKind kind = ParquetColumnSchemaKind::PRIMITIVE;
 
+    // Cached during schema construction so readers created per row group do 
not repeatedly walk
+    // ordinary nested schemas to discover whether Variant-specific planning 
is needed.
+    bool contains_variant = false;
+
     // ======== Dremel Levels ========
 
     int16_t max_definition_level = 0;
diff --git a/be/src/format_v2/parquet/parquet_file_context.cpp 
b/be/src/format_v2/parquet/parquet_file_context.cpp
index c178e491b8a..8ba8cf94662 100644
--- a/be/src/format_v2/parquet/parquet_file_context.cpp
+++ b/be/src/format_v2/parquet/parquet_file_context.cpp
@@ -273,6 +273,7 @@ Status ParquetFileContext::open(io::FileReaderSPtr 
input_file_reader, io::IOCont
                                 bool enable_page_cache, const 
io::FileDescription& file_description,
                                 bool enable_mapping_timestamp_tz, bool 
enable_mapping_varbinary) {
     DORIS_CHECK(input_file_reader != nullptr);
+    contains_variant = false;
     if 
(detail::should_stage_small_http_file(input_file_reader->path().native(),
                                              input_file_reader->size(),
                                              config::in_memory_file_size)) {
@@ -602,6 +603,7 @@ Status ParquetFileContext::close() {
     native_io_ctx = nullptr;
     native_page_cache_enabled = false;
     native_page_cache_file_key.clear();
+    contains_variant = false;
     return Status::OK();
 }
 
diff --git a/be/src/format_v2/parquet/parquet_file_context.h 
b/be/src/format_v2/parquet/parquet_file_context.h
index 0cd413e1055..38e78438e39 100644
--- a/be/src/format_v2/parquet/parquet_file_context.h
+++ b/be/src/format_v2/parquet/parquet_file_context.h
@@ -139,6 +139,9 @@ struct ParquetFileContext {
     int64_t native_footer_cache_hits = 0;
     bool native_page_cache_enabled = false;
     std::string native_page_cache_file_key;
+    // Set once after the logical file schema is built. Per-request planning 
uses this guard so
+    // ordinary files never enter Variant projection or shredded-statistics 
paths.
+    bool contains_variant = false;
 
     Status open(io::FileReaderSPtr input_file_reader, io::IOContext* io_ctx, 
bool enable_page_cache,
                 const io::FileDescription& file_description,
diff --git a/be/src/format_v2/parquet/parquet_reader.cpp 
b/be/src/format_v2/parquet/parquet_reader.cpp
index bd44316463f..8c6efe27558 100644
--- a/be/src/format_v2/parquet/parquet_reader.cpp
+++ b/be/src/format_v2/parquet/parquet_reader.cpp
@@ -163,6 +163,9 @@ size_t finalize_variant_leaf_projections(
         if (local_id < 0 || local_id >= 
static_cast<int32_t>(file_schema.size())) {
             continue;
         }
+        if (!file_schema[local_id]->contains_variant) {
+            continue;
+        }
         retained += 
detail::finalize_variant_leaf_projection(metadata.to_thrift(),
                                                              
*file_schema[local_id], &projection);
     }
@@ -467,6 +470,11 @@ Status ParquetReader::init(RuntimeState* state) {
         SCOPED_TIMER(_parquet_profile.parse_meta_time);
         
RETURN_IF_ERROR(build_parquet_column_schema(_state->file_context.native_metadata->schema(),
                                                     &_state->file_schema));
+        _state->file_context.contains_variant =
+                std::ranges::any_of(_state->file_schema, [](const auto& 
column) {
+                    DORIS_CHECK(column != nullptr);
+                    return column->contains_variant;
+                });
         if (_enable_mapping_timestamp_tz) {
             for (auto& column_schema : _state->file_schema) {
                 apply_timestamp_tz_mapping(column_schema.get());
@@ -518,13 +526,16 @@ Status 
ParquetReader::open(std::shared_ptr<format::FileScanRequest> request) {
     }
     auto request_snapshot = request;
     DORIS_CHECK(request_snapshot != nullptr);
-    const size_t retained_variant_leaf_projections =
-            
finalize_variant_leaf_projections(*_state->file_context.native_metadata,
-                                              _state->file_schema,
-                                              
&request_snapshot->predicate_columns) +
-            
finalize_variant_leaf_projections(*_state->file_context.native_metadata,
-                                              _state->file_schema,
-                                              
&request_snapshot->non_predicate_columns);
+    size_t retained_variant_leaf_projections = 0;
+    if (_state->file_context.contains_variant) {
+        retained_variant_leaf_projections =
+                
finalize_variant_leaf_projections(*_state->file_context.native_metadata,
+                                                  _state->file_schema,
+                                                  
&request_snapshot->predicate_columns) +
+                
finalize_variant_leaf_projections(*_state->file_context.native_metadata,
+                                                  _state->file_schema,
+                                                  
&request_snapshot->non_predicate_columns);
+    }
     if (_parquet_profile.variant_leaf_projections != nullptr) {
         COUNTER_UPDATE(_parquet_profile.variant_leaf_projections,
                        retained_variant_leaf_projections);
diff --git a/be/src/format_v2/parquet/parquet_statistics.cpp 
b/be/src/format_v2/parquet/parquet_statistics.cpp
index b4f03677357..75d159d670d 100644
--- a/be/src/format_v2/parquet/parquet_statistics.cpp
+++ b/be/src/format_v2/parquet/parquet_statistics.cpp
@@ -1277,6 +1277,12 @@ Status select_row_groups_by_metadata(
     if (pruning_stats != nullptr) {
         pruning_stats->total_row_groups = cast_set<int64_t>(candidate_size);
     }
+    const bool contains_variant =
+            file_context != nullptr ? file_context->contains_variant
+                                    : std::ranges::any_of(file_schema, 
[](const auto& column) {
+                                          DORIS_CHECK(column != nullptr);
+                                          return column->contains_variant;
+                                      });
     selected_row_groups->reserve(candidate_size);
     for (size_t candidate_idx = 0; candidate_idx < candidate_size; 
++candidate_idx) {
         const int row_group_idx = candidate_row_groups == nullptr
@@ -1303,8 +1309,8 @@ Status select_row_groups_by_metadata(
             has_expr_zonemap_filter(request, runtime_state) &&
             (check_native_statistics(metadata, row_group, file_schema, 
request, pruning_stats,
                                      timezone) ||
-             check_shredded_variant_statistics(metadata, row_group, 
file_schema, request,
-                                               timezone))) {
+             (contains_variant && check_shredded_variant_statistics(
+                                          metadata, row_group, file_schema, 
request, timezone)))) {
             prune_reason = ParquetRowGroupPruneReason::STATISTICS;
         }
         if (probe_mode != ParquetMetadataProbeMode::FOOTER_ONLY &&
diff --git a/be/src/format_v2/parquet/reader/native_column_reader.cpp 
b/be/src/format_v2/parquet/reader/native_column_reader.cpp
index 618eca6f6a5..273964e5d1f 100644
--- a/be/src/format_v2/parquet/reader/native_column_reader.cpp
+++ b/be/src/format_v2/parquet/reader/native_column_reader.cpp
@@ -292,13 +292,18 @@ Status NativeColumnReader::create(
     }
 
     auto logical_type = projected_type(column_schema, projection, false);
-    auto native_type = projected_type(column_schema, projection, true);
+    auto native_type = logical_type;
+    std::unique_ptr<VariantMaterializationNode> variant_plan;
+    if (column_schema.contains_variant) {
+        // Native readers are instantiated per projected column and row group. 
Keep Variant tree
+        // construction out of ordinary scans instead of charging that setup 
cost at every split.
+        native_type = projected_type(column_schema, projection, true);
+        variant_plan = build_variant_plan(column_schema, projection);
+    }
     std::shared_ptr<NativeSchemaNode> schema_node;
     RETURN_IF_ERROR(build_native_schema_node(native_type, column_schema, 
&schema_node));
     std::set<uint64_t> projected_ids;
     collect_projected_ids(column_schema, projection, *field, &projected_ids);
-    auto variant_plan = build_variant_plan(column_schema, projection);
-
     auto native_reader = std::unique_ptr<NativeColumnReader>(
             new NativeColumnReader(column_schema, std::move(logical_type), 
std::move(native_type),
                                    std::move(variant_plan), profile));
@@ -364,7 +369,7 @@ Status NativeColumnReader::init(
             runtime_state != nullptr && runtime_state->enable_strict_mode()));
     DORIS_CHECK(_native_reader != nullptr);
     _skip_column = _native_type->create_column();
-    if (_variant_plan->contains_variant) {
+    if (_variant_plan != nullptr) {
         _variant_physical_column = _native_type->create_column();
     }
     return Status::OK();
@@ -387,7 +392,7 @@ Status NativeColumnReader::read_with_filter(int64_t rows, 
const uint8_t* filter_
     RETURN_IF_ERROR(filter.init(filter_data, static_cast<size_t>(rows), 
filter_all));
     _native_reader->reset_filter_map_index();
     const bool materialize_variant =
-            !dictionary_ids && _variant_plan->contains_variant && 
output_type->equals(*_type);
+            !dictionary_ids && _variant_plan != nullptr && 
output_type->equals(*_type);
     if (materialize_variant) {
         _variant_physical_column->clear();
     }
@@ -715,7 +720,7 @@ Status NativeColumnReader::select_with_dictionary_filter(
     DORIS_CHECK(row_filter != nullptr);
     DORIS_CHECK(survivor_count != nullptr);
     DORIS_CHECK(used_filter != nullptr);
-    if (_variant_plan->contains_variant) {
+    if (_variant_plan != nullptr) {
         row_filter->clear();
         *used_filter = false;
         return Status::OK();
@@ -871,7 +876,7 @@ Status NativeColumnReader::select_with_fixed_width_filter(
     DORIS_CHECK(row_filter != nullptr);
     DORIS_CHECK(used_filter != nullptr);
     DORIS_CHECK(execution_kind != nullptr);
-    if (_variant_plan->contains_variant) {
+    if (_variant_plan != nullptr) {
         // Direct fixed-width evaluation cannot preserve a Variant physical 
subtree's row shape.
         row_filter->clear();
         *used_filter = false;
@@ -1008,7 +1013,7 @@ bool NativeColumnReader::crossed_page_since_last_batch() {
 
 Result<MutableColumnPtr> NativeColumnReader::dictionary_values() {
     DORIS_CHECK(_native_reader != nullptr);
-    if (_variant_plan->contains_variant) {
+    if (_variant_plan != nullptr) {
         return ResultError(
                 Status::NotSupported("Parquet Variant columns do not expose 
dictionary values"));
     }
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 43ff808cf5b..ba9f63b9356 100644
--- a/be/src/format_v2/parquet/reader/variant_column_reader.cpp
+++ b/be/src/format_v2/parquet/reader/variant_column_reader.cpp
@@ -470,6 +470,7 @@ std::unique_ptr<ParquetColumnSchema> clone_schema(
     result->leaf_column_id = source.leaf_column_id;
     result->type_descriptor = source.type_descriptor;
     result->kind = source.kind;
+    result->contains_variant = source.contains_variant;
     result->max_definition_level = source.max_definition_level;
     result->max_repetition_level = source.max_repetition_level;
     result->nullable_definition_level = source.nullable_definition_level;
diff --git a/be/test/format_v2/parquet/parquet_schema_test.cpp 
b/be/test/format_v2/parquet/parquet_schema_test.cpp
index 9ecc1edf6fc..acd5e860078 100644
--- a/be/test/format_v2/parquet/parquet_schema_test.cpp
+++ b/be/test/format_v2/parquet/parquet_schema_test.cpp
@@ -217,6 +217,7 @@ TEST(ParquetSchemaTest, 
NativeSchemaRecognizesVariantLogicalGroup) {
         ASSERT_TRUE(status.ok()) << status;
         ASSERT_EQ(fields.size(), 1);
         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(typeid_cast<const 
DataTypeVariantV2*>(remove_nullable(fields[0]->type).get()),
                   nullptr);
@@ -247,6 +248,10 @@ TEST(ParquetSchemaTest, 
NestedVariantPropagatesIntoParentLogicalType) {
     const auto status = build_parquet_column_schema(descriptor, &fields);
     ASSERT_TRUE(status.ok()) << status;
     ASSERT_EQ(fields.size(), 1);
+    EXPECT_TRUE(fields[0]->contains_variant);
+    ASSERT_EQ(fields[0]->children.size(), 2);
+    EXPECT_FALSE(fields[0]->children[0]->contains_variant);
+    EXPECT_TRUE(fields[0]->children[1]->contains_variant);
     const auto* info_type =
             assert_cast<const 
DataTypeStruct*>(remove_nullable(fields[0]->type).get());
     ASSERT_EQ(info_type->get_elements().size(), 2);
@@ -481,9 +486,12 @@ TEST(ParquetSchemaTest, 
NativeMetadataTreePreservesNestedFieldNamesAndIds) {
     std::vector<std::unique_ptr<ParquetColumnSchema>> fields;
     ASSERT_TRUE(build_parquet_column_schema(native_schema, &fields).ok());
     ASSERT_EQ(fields.size(), 1);
+    EXPECT_FALSE(fields[0]->contains_variant);
     EXPECT_EQ(fields[0]->name, "protocol");
     EXPECT_EQ(fields[0]->parquet_field_id, 10);
     ASSERT_EQ(fields[0]->children.size(), 2);
+    EXPECT_FALSE(fields[0]->children[0]->contains_variant);
+    EXPECT_FALSE(fields[0]->children[1]->contains_variant);
     EXPECT_EQ(fields[0]->children[0]->name, "minReaderVersion");
     EXPECT_EQ(fields[0]->children[0]->leaf_column_id, 0);
     EXPECT_EQ(fields[0]->children[1]->name, "minWriterVersion");
diff --git a/be/test/format_v2/parquet/parquet_statistics_test.cpp 
b/be/test/format_v2/parquet/parquet_statistics_test.cpp
index 1a13e8ecea3..4a6dc62775e 100644
--- a/be/test/format_v2/parquet/parquet_statistics_test.cpp
+++ b/be/test/format_v2/parquet/parquet_statistics_test.cpp
@@ -924,6 +924,7 @@ TEST(NativeParquetStatisticsTest, 
ShreddedVariantTypedValueDrivesPageFiltering)
     variant->name = "v";
     variant->local_id = 0;
     variant->kind = format::parquet::ParquetColumnSchemaKind::VARIANT;
+    variant->contains_variant = true;
     variant->type = make_nullable(std::make_shared<DataTypeVariantV2>());
     variant->children.push_back(bytes("metadata", 0, 0));
     variant->children.push_back(bytes("value", 1, 1));
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 24ace262fab..51a50319285 100644
--- a/be/test/format_v2/parquet/variant_column_reader_test.cpp
+++ b/be/test/format_v2/parquet/variant_column_reader_test.cpp
@@ -70,6 +70,7 @@ ParquetColumnSchema unshredded_schema() {
     ParquetColumnSchema schema;
     schema.name = "payload";
     schema.kind = ParquetColumnSchemaKind::VARIANT;
+    schema.contains_variant = true;
     schema.type = make_nullable(std::make_shared<DataTypeVariantV2>());
     const auto binary = make_nullable(std::make_shared<DataTypeString>());
     schema.variant_physical_type = 
make_nullable(std::make_shared<DataTypeStruct>(


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

Reply via email to