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

yiguolei 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 042429a7cfa branch-4.1: [improvement](be) Optimize Variant V2 
ingestion and STRING casts (#66744)
042429a7cfa is described below

commit 042429a7cfaa7bfb611f0482a009e8cb557b8060
Author: lihangyu <[email protected]>
AuthorDate: Tue Aug 18 09:42:40 2026 +0800

    branch-4.1: [improvement](be) Optimize Variant V2 ingestion and STRING 
casts (#66744)
    
    cherry-pick #66709
---
 be/src/core/value/variant/variant_value.cpp        |  23 +-
 be/src/core/value/variant/variant_value.h          |  23 +-
 .../cast/variant_v2/cast_variant_to_string.cpp     |   9 +-
 .../segment/variant/v2/variant_path_builder.cpp    | 180 ++++++--
 .../segment/variant/v2/variant_path_builder.h      |   5 +
 .../segment/variant/v2/variant_shredder.cpp        | 145 ++++--
 .../storage/segment/variant/v2/variant_shredder.h  |   8 +
 .../function/cast/cast_variant_v2_from_test.cpp    |  62 +++
 .../variant/variant_column_writer_reader_test.cpp  | 485 ++++++++++++++++++++-
 be/test/util/variant/variant_value_test.cpp        |  40 ++
 10 files changed, 895 insertions(+), 85 deletions(-)

diff --git a/be/src/core/value/variant/variant_value.cpp 
b/be/src/core/value/variant/variant_value.cpp
index 1f8101af30c..fac75c70096 100644
--- a/be/src/core/value/variant/variant_value.cpp
+++ b/be/src/core/value/variant/variant_value.cpp
@@ -368,7 +368,14 @@ uint32_t VariantRef::num_elements() const {
     return _container_layout(type).count;
 }
 
-uint32_t VariantRef::_object_field_id(const ContainerLayout& layout, uint32_t 
index) const {
+VariantRef::ObjectView VariantRef::object_view() const {
+    ContainerLayout layout = _container_layout(VariantBasicType::OBJECT);
+    const uint32_t dictionary_size = layout.count == 0 ? 0 : 
metadata.dict_size();
+    return ObjectView(*this, layout, dictionary_size);
+}
+
+uint32_t VariantRef::_object_field_id(const ContainerLayout& layout, uint32_t 
index,
+                                      const uint32_t* dictionary_size) const {
     if (index >= layout.count) {
         throw Exception(ErrorCode::INVALID_ARGUMENT,
                         "Variant object index {} is out of range [0, {})", 
index, layout.count);
@@ -376,10 +383,12 @@ uint32_t VariantRef::_object_field_id(const 
ContainerLayout& layout, uint32_t in
     const auto field_id = static_cast<uint32_t>(read_unsigned(
             value.data + layout.ids_offset + static_cast<size_t>(index) * 
layout.id_width,
             layout.id_width));
-    if (field_id >= metadata.dict_size()) {
+    const uint32_t metadata_dictionary_size =
+            dictionary_size != nullptr ? *dictionary_size : 
metadata.dict_size();
+    if (field_id >= metadata_dictionary_size) {
         throw Exception(ErrorCode::CORRUPTION,
                         "Variant object field id {} is outside metadata 
dictionary of size {}",
-                        field_id, metadata.dict_size());
+                        field_id, metadata_dictionary_size);
     }
     return field_id;
 }
@@ -427,6 +436,14 @@ VariantRef VariantRef::object_value_at(uint32_t index, 
uint32_t* field_id_out) c
     return _container_value_at(layout, index, false);
 }
 
+VariantRef VariantRef::ObjectView::value_at(uint32_t index, uint32_t* 
field_id_out) const {
+    const uint32_t field_id = _value._object_field_id(_layout, index, 
&_dictionary_size);
+    if (field_id_out != nullptr) {
+        *field_id_out = field_id;
+    }
+    return _value._container_value_at(_layout, index, false);
+}
+
 VariantRef VariantRef::array_at(uint32_t index) const {
     return _container_value_at(_container_layout(VariantBasicType::ARRAY), 
index, true);
 }
diff --git a/be/src/core/value/variant/variant_value.h 
b/be/src/core/value/variant/variant_value.h
index 7dd14b24a6b..c09794fcd8a 100644
--- a/be/src/core/value/variant/variant_value.h
+++ b/be/src/core/value/variant/variant_value.h
@@ -28,6 +28,8 @@
 namespace doris {
 
 struct VariantRef {
+    class ObjectView;
+
     VariantMetadataRef metadata;
     StringRef value;
 
@@ -52,6 +54,7 @@ struct VariantRef {
     std::array<uint8_t, 16> get_uuid() const;
 
     uint32_t num_elements() const;
+    ObjectView object_view() const;
     bool object_find(StringRef key, VariantRef* out) const;
     bool object_find_by_id(uint32_t field_id, VariantRef* out) const;
     VariantRef object_value_at(uint32_t index, uint32_t* field_id_out) const;
@@ -69,11 +72,29 @@ private:
     };
 
     ContainerLayout _container_layout(VariantBasicType expected_type) const;
-    uint32_t _object_field_id(const ContainerLayout& layout, uint32_t index) 
const;
+    uint32_t _object_field_id(const ContainerLayout& layout, uint32_t index,
+                              const uint32_t* dictionary_size = nullptr) const;
     bool _object_find_by_id(const ContainerLayout& layout, uint32_t field_id,
                             VariantRef* out) const;
     VariantRef _container_value_at(const ContainerLayout& layout, uint32_t 
index,
                                    bool require_array_boundary) const;
 };
 
+// Parses and validates an object's physical layout and metadata dictionary 
size once, then reuses
+// them while iterating its children. The referenced metadata and value bytes 
must outlive the view.
+class VariantRef::ObjectView {
+public:
+    uint32_t size() const { return _layout.count; }
+    VariantRef value_at(uint32_t index, uint32_t* field_id_out = nullptr) 
const;
+
+private:
+    friend struct VariantRef;
+    ObjectView(VariantRef value, ContainerLayout layout, uint32_t 
dictionary_size)
+            : _value(value), _layout(layout), 
_dictionary_size(dictionary_size) {}
+
+    VariantRef _value;
+    ContainerLayout _layout;
+    uint32_t _dictionary_size;
+};
+
 } // namespace doris
diff --git a/be/src/exprs/function/cast/variant_v2/cast_variant_to_string.cpp 
b/be/src/exprs/function/cast/variant_v2/cast_variant_to_string.cpp
index 56f085c81c1..73ba82537ae 100644
--- a/be/src/exprs/function/cast/variant_v2/cast_variant_to_string.cpp
+++ b/be/src/exprs/function/cast/variant_v2/cast_variant_to_string.cpp
@@ -164,8 +164,15 @@ Status cast_values_to_string(FunctionContext* context, 
size_t rows, ForcedNulls
 Status cast_typed_variant_to_string(FunctionContext* context, const 
ColumnVariantV2& source,
                                     size_t rows, ForcedNulls forced_nulls, 
ColumnPtr* output) {
     const auto& typed = assert_cast<const 
ColumnNullable&>(source.typed_column());
-    const DataTypePtr string_type = std::make_shared<DataTypeString>();
+    const PrimitiveType source_primitive = 
source.typed_type()->get_primitive_type();
     const NullMap& inner_nulls = typed.get_null_map_data();
+    if (is_string_type(source_primitive) && !typed.has_null()) {
+        // The physical payload already has the requested representation. An 
inner null is a
+        // Variant null and must still stringify as literal "null".
+        return apply_forced_nulls(typed.get_ptr(), forced_nulls, output);
+    }
+
+    const DataTypePtr string_type = std::make_shared<DataTypeString>();
     size_t concrete_rows = 0;
     for (size_t row = 0; row < rows; ++row) {
         if (inner_nulls[row] == 0 && (forced_nulls.empty() || 
forced_nulls[row] == 0)) {
diff --git a/be/src/storage/segment/variant/v2/variant_path_builder.cpp 
b/be/src/storage/segment/variant/v2/variant_path_builder.cpp
index c7542669fe4..ea609d3bdac 100644
--- a/be/src/storage/segment/variant/v2/variant_path_builder.cpp
+++ b/be/src/storage/segment/variant/v2/variant_path_builder.cpp
@@ -74,6 +74,34 @@ enum class ValueKind : uint8_t {
     ARRAY,
 };
 
+enum class ScalarPhysicalKind : uint8_t { OTHER, NULL_VALUE, SHORT_STRING, 
PRIMITIVE };
+
+struct ScalarPhysical {
+    ScalarPhysicalKind kind;
+    VariantPrimitiveId primitive_id = VariantPrimitiveId::NULL_VALUE;
+};
+
+ScalarPhysical scalar_physical(VariantRef value) {
+    switch (value.basic_type()) {
+    case VariantBasicType::SHORT_STRING:
+        return {.kind = ScalarPhysicalKind::SHORT_STRING};
+    case VariantBasicType::PRIMITIVE: {
+        const VariantPrimitiveId primitive_id = value.primitive_id();
+        return {.kind = primitive_id == VariantPrimitiveId::NULL_VALUE
+                                ? ScalarPhysicalKind::NULL_VALUE
+                                : ScalarPhysicalKind::PRIMITIVE,
+                .primitive_id = primitive_id};
+    }
+    case VariantBasicType::OBJECT:
+    case VariantBasicType::ARRAY:
+        return {.kind = ScalarPhysicalKind::OTHER};
+    }
+    // Match VariantRef::is_null(): an unknown basic type is not null. The
+    // authoritative slow path reports the corrupt type after preserving
+    // row-validation ordering.
+    return {.kind = ScalarPhysicalKind::OTHER};
+}
+
 const DataTypePtr& jsonb_type() {
     static const DataTypePtr type = std::make_shared<DataTypeJsonb>();
     return type;
@@ -197,37 +225,8 @@ DataTypePtr infer_type(VariantRef value, const 
DataTypePtr& reusable_type = null
         return type;
     }
     case ValueKind::INT64: {
-        PrimitiveType primitive = TYPE_BIGINT;
-        switch (value.primitive_id()) {
-        case VariantPrimitiveId::INT8:
-            primitive = TYPE_TINYINT;
-            break;
-        case VariantPrimitiveId::INT16:
-            primitive = TYPE_SMALLINT;
-            break;
-        case VariantPrimitiveId::INT32:
-            primitive = TYPE_INT;
-            break;
-        case VariantPrimitiveId::INT64:
-            break;
-        default:
-            throw Exception(ErrorCode::CORRUPTION, "Invalid Variant integer 
primitive id");
-        }
-        static const std::array<DataTypePtr, 4> types {
-                std::make_shared<DataTypeInt8>(), 
std::make_shared<DataTypeInt16>(),
-                std::make_shared<DataTypeInt32>(), 
std::make_shared<DataTypeInt64>()};
-        switch (primitive) {
-        case TYPE_TINYINT:
-            return types[0];
-        case TYPE_SMALLINT:
-            return types[1];
-        case TYPE_INT:
-            return types[2];
-        case TYPE_BIGINT:
-            return types[3];
-        default:
-            throw Exception(ErrorCode::CORRUPTION, "Invalid Variant integer 
type {}", primitive);
-        }
+        static const DataTypePtr type = std::make_shared<DataTypeInt64>();
+        return type;
     }
     case ValueKind::LARGEINT: {
         static const DataTypePtr type = std::make_shared<DataTypeInt128>();
@@ -694,6 +693,72 @@ void append_timestamp(VariantRef value, PrimitiveType 
target_type, IColumn* targ
     assert_cast<ColumnTimeStampTz&>(*target).insert_value(converted);
 }
 
+template <typename Value>
+bool stable_scalar_matches_type(const Value& value, const ScalarPhysical& 
physical,
+                                const DataTypePtr& target_type) {
+    const PrimitiveType target_primitive = target_type->get_primitive_type();
+    if (target_primitive == TYPE_JSONB) {
+        return physical.kind == ScalarPhysicalKind::SHORT_STRING ||
+               physical.kind == ScalarPhysicalKind::PRIMITIVE;
+    }
+    if (physical.kind == ScalarPhysicalKind::SHORT_STRING) {
+        return target_primitive == TYPE_STRING;
+    }
+    if (physical.kind != ScalarPhysicalKind::PRIMITIVE) {
+        return false;
+    }
+
+    switch (physical.primitive_id) {
+    case VariantPrimitiveId::TRUE_VALUE:
+    case VariantPrimitiveId::FALSE_VALUE:
+        return target_primitive == TYPE_BOOLEAN;
+    case VariantPrimitiveId::INT8:
+    case VariantPrimitiveId::INT16:
+    case VariantPrimitiveId::INT32:
+    case VariantPrimitiveId::INT64:
+        // Encoded integer widths are a physical detail. Integer paths use 
BIGINT from their first
+        // value, avoiding repeated inference and column rewrites when later 
values use a wider or
+        // narrower physical width. LARGEINT remains valid after a real 
128-bit promotion.
+        return target_primitive == TYPE_BIGINT || target_primitive == 
TYPE_LARGEINT;
+    case VariantPrimitiveId::FLOAT:
+        return target_primitive == TYPE_FLOAT;
+    case VariantPrimitiveId::DOUBLE:
+        return target_primitive == TYPE_DOUBLE;
+    case VariantPrimitiveId::DECIMAL4:
+    case VariantPrimitiveId::DECIMAL8:
+    case VariantPrimitiveId::DECIMAL16: {
+        const VariantDecimal decimal = value.get_decimal();
+        if (physical.primitive_id == VariantPrimitiveId::DECIMAL16 && 
decimal.scale == 0) {
+            return target_primitive == TYPE_LARGEINT;
+        }
+        if (target_primitive != TYPE_DECIMAL128I || 
target_type->get_precision() != 38 ||
+            target_type->get_scale() != decimal.scale) {
+            return false;
+        }
+        __int128 converted = 0;
+        return try_rescale_decimal_value(value, target_type, &converted);
+    }
+    case VariantPrimitiveId::DATE:
+        return target_primitive == TYPE_DATEV2 && 
date_fits_doris_range(value.get_date());
+    case VariantPrimitiveId::TIMESTAMP_MICROS:
+        return target_primitive == TYPE_TIMESTAMPTZ &&
+               timestamp_fits_doris_range(value.get_timestamp_micros());
+    case VariantPrimitiveId::TIMESTAMP_NTZ_MICROS:
+        return target_primitive == TYPE_DATETIMEV2 &&
+               timestamp_fits_doris_range(value.get_timestamp_ntz_micros());
+    case VariantPrimitiveId::STRING:
+        return target_primitive == TYPE_STRING;
+    case VariantPrimitiveId::NULL_VALUE:
+    case VariantPrimitiveId::BINARY:
+    case VariantPrimitiveId::TIME_NTZ_MICROS:
+    case VariantPrimitiveId::TIMESTAMP_NANOS:
+    case VariantPrimitiveId::TIMESTAMP_NTZ_NANOS:
+    case VariantPrimitiveId::UUID:
+        return false;
+    }
+    throw Exception(ErrorCode::CORRUPTION, "Unknown Variant primitive id");
+}
+
 void append_value(VariantRef value, const DataTypePtr& target_type, IColumn* 
target);
 
 void append_array(VariantRef value, const DataTypePtr& target_type, IColumn* 
target) {
@@ -934,6 +999,7 @@ struct VariantPathBuilder::Impl {
         type = remove_nullable(initial_type);
         nullable_type = make_nullable(type);
         column = nullable_type->create_column();
+        binary_serde.reset();
         return Status::OK();
     }
 
@@ -991,6 +1057,7 @@ struct VariantPathBuilder::Impl {
         column = IColumn::mutate(std::move(promoted));
         type = std::move(target);
         nullable_type = make_nullable(type);
+        binary_serde.reset();
 #ifdef BE_TEST
         ++promotions;
 #endif
@@ -1000,12 +1067,16 @@ struct VariantPathBuilder::Impl {
     PathInData path;
     DataTypePtr type;
     DataTypePtr nullable_type;
+    DataTypeSerDeSPtr binary_serde;
     MutableColumnPtr column;
     DorisVector<uint32_t> rowids;
     size_t logical_rows = 0;
 #ifdef BE_TEST
     size_t promotions = 0;
 #endif
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+    size_t stable_scalar_appends = 0;
+#endif
 };
 
 VariantPathBuilder::VariantPathBuilder(PathInData path, size_t prefix_rows)
@@ -1016,7 +1087,8 @@ VariantPathBuilder& 
VariantPathBuilder::operator=(VariantPathBuilder&&) noexcept
 
 Status VariantPathBuilder::append(VariantRef value, size_t row) {
     try {
-        if (value.is_null()) {
+        const ScalarPhysical physical = scalar_physical(value);
+        if (physical.kind == ScalarPhysicalKind::NULL_VALUE) {
             return Status::InvalidArgument("Variant path builder {} must not 
append JSON null",
                                            _impl->path.get_path());
         }
@@ -1028,16 +1100,26 @@ Status VariantPathBuilder::append(VariantRef value, 
size_t row) {
             return Status::InvalidArgument("Variant path builder {} row {} 
exceeds uint32 limit",
                                            _impl->path.get_path(), row);
         }
-        RETURN_IF_ERROR(complete_rows(row));
-
-        if (!_impl->column) {
-            RETURN_IF_ERROR(_impl->initialize(infer_type(value)));
-        } else if (_impl->type->get_primitive_type() != TYPE_JSONB) {
-            DataTypePtr inferred_type = infer_type(value, _impl->type);
-            DataTypePtr common_type = path_least_common_type(_impl->type, 
inferred_type);
-            RETURN_IF_ERROR(_impl->promote(common_type, false));
+        _impl->logical_rows = row;
+
+        const bool stable_scalar =
+                _impl->column && stable_scalar_matches_type(value, physical, 
_impl->type);
+        if (!stable_scalar) {
+            if (!_impl->column) {
+                RETURN_IF_ERROR(_impl->initialize(infer_type(value)));
+            } else if (_impl->type->get_primitive_type() != TYPE_JSONB) {
+                DataTypePtr inferred_type = infer_type(value, _impl->type);
+                if (_impl->type.get() != inferred_type.get() &&
+                    !_impl->type->equals(*inferred_type)) {
+                    DataTypePtr common_type = 
path_least_common_type(_impl->type, inferred_type);
+                    if (_impl->type.get() != common_type.get() &&
+                        !_impl->type->equals(*common_type)) {
+                        RETURN_IF_ERROR(_impl->promote(common_type, false));
+                    }
+                }
+            }
         }
-        const bool is_array = value.basic_type() == VariantBasicType::ARRAY;
+        const bool is_array = !stable_scalar && value_kind(value) == 
ValueKind::ARRAY;
         if (is_array && !value_is_representable(value, _impl->type)) {
             RETURN_IF_ERROR(_impl->promote(jsonb_type(), false));
         }
@@ -1057,6 +1139,9 @@ Status VariantPathBuilder::append(VariantRef value, 
size_t row) {
         }
         _impl->rowids.push_back(static_cast<uint32_t>(row));
         _impl->logical_rows = row + 1;
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+        _impl->stable_scalar_appends += stable_scalar;
+#endif
         return Status::OK();
     } catch (const Exception& exception) {
         return exception.to_status();
@@ -1110,6 +1195,12 @@ size_t VariantPathBuilder::promotion_count() const {
 }
 #endif
 
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+size_t VariantPathBuilder::stable_scalar_append_count() const {
+    return _impl->stable_scalar_appends;
+}
+#endif
+
 size_t VariantPathBuilder::byte_size() const {
     return sizeof(Impl) + path_allocated_bytes(_impl->path) +
            _impl->rowids.capacity() * sizeof(uint32_t) +
@@ -1181,8 +1272,11 @@ Status VariantPathBuilder::write_sparse_cell(size_t 
value_index, ColumnString::C
                                      _impl->path.get_path());
     }
     try {
-        
_impl->type->get_serde(2)->write_one_cell_to_binary(nullable.get_nested_column(),
 *chars,
-                                                            value_index);
+        if (!_impl->binary_serde) {
+            _impl->binary_serde = _impl->type->get_serde(2);
+        }
+        
_impl->binary_serde->write_one_cell_to_binary(nullable.get_nested_column(), 
*chars,
+                                                      value_index);
         return Status::OK();
     } catch (const Exception& exception) {
         return exception.to_status();
diff --git a/be/src/storage/segment/variant/v2/variant_path_builder.h 
b/be/src/storage/segment/variant/v2/variant_path_builder.h
index 707aa3f762e..7b467a17ae7 100644
--- a/be/src/storage/segment/variant/v2/variant_path_builder.h
+++ b/be/src/storage/segment/variant/v2/variant_path_builder.h
@@ -74,6 +74,11 @@ public:
 #ifdef BE_TEST
     size_t rows() const;
     size_t promotion_count() const;
+#endif
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+    size_t stable_scalar_append_count() const;
+#endif
+#ifdef BE_TEST
     bool is_null_at(size_t row) const;
     Status materialize(ColumnPtr* result) const;
 #endif
diff --git a/be/src/storage/segment/variant/v2/variant_shredder.cpp 
b/be/src/storage/segment/variant/v2/variant_shredder.cpp
index d61b51aed96..f197d96d30a 100644
--- a/be/src/storage/segment/variant/v2/variant_shredder.cpp
+++ b/be/src/storage/segment/variant/v2/variant_shredder.cpp
@@ -64,6 +64,8 @@ struct VariantShredder::Impl {
     using PathIndex = uint32_t;
     using ParentFieldKey = uint64_t;
     using ChildPathCache = doris::flat_hash_map<ParentFieldKey, PathIndex>;
+    static constexpr PathIndex UNRESOLVED_PATH = 
std::numeric_limits<PathIndex>::max();
+    static constexpr size_t MAX_BINARY_CELLS_PER_CHUNK = 1U << 20;
 
     // Metadata bytes belong to the input ReadView. This cache never escapes 
one append call, so
     // it can borrow the dictionary and retain only parent+field transitions 
observed in that
@@ -197,7 +199,7 @@ struct VariantShredder::Impl {
         if (const auto found = path_indices.find(child); found != 
path_indices.end()) {
             child_index = found->second;
         } else {
-            if (paths.size() > std::numeric_limits<PathIndex>::max()) {
+            if (paths.size() >= UNRESOLVED_PATH) {
                 throw Exception(ErrorCode::INVALID_ARGUMENT,
                                 "Variant path count exceeds uint32 limit");
             }
@@ -218,10 +220,10 @@ struct VariantShredder::Impl {
         if (value.basic_type() != VariantBasicType::OBJECT) {
             return append_leaf(value, path_index, row);
         }
-        const uint32_t children = value.num_elements();
-        for (uint32_t index = 0; index < children; ++index) {
+        const VariantRef::ObjectView object = value.object_view();
+        for (uint32_t index = 0; index < object.size(); ++index) {
             uint32_t field = 0;
-            VariantRef child = value.object_value_at(index, &field);
+            VariantRef child = object.value_at(index, &field);
             const PathIndex child_path = resolve_child_path(metadata_cache, 
path_index, field);
             RETURN_IF_ERROR(validate_doc_path(child_path));
             RETURN_IF_ERROR(visit(child, metadata_cache, child_path, row));
@@ -384,16 +386,22 @@ struct VariantShredder::Impl {
     template <typename BinaryPlan>
     Status append_binary_rows(const DorisVector<BinaryPlan>& binary_plan,
                               const DorisVector<ColumnMap*>& maps) const {
-        struct BinaryCell {
-            size_t plan_index = 0;
-            size_t value_index = 0;
-        };
-
-        // Build a compact row index in two passes. Each path contributes only 
its present values,
-        // and paths are visited in publication order so cells within one row 
preserve path order.
+        DorisVector<size_t> bucket_cells(maps.size(), 0);
+        DorisVector<size_t> bucket_key_bytes(maps.size(), 0);
+        DorisVector<size_t> bucket_value_bytes(maps.size(), 0);
+        // Build a compact row index in two passes. Each path contributes only 
its
+        // present values, and paths are visited in publication order so cells
+        // within one row preserve path order.
         DorisVector<size_t> row_offsets(rows + 1, 0);
         for (const BinaryPlan& plan : binary_plan) {
-            for (uint32_t row : plan.builder->rowids()) {
+            DORIS_CHECK_LT(plan.bucket, maps.size());
+            const std::span<const uint32_t> rowids = plan.builder->rowids();
+            bucket_cells[plan.bucket] += rowids.size();
+            bucket_key_bytes[plan.bucket] += rowids.size() * plan.path->size();
+            const ColumnPtr column = plan.builder->column();
+            DORIS_CHECK(column);
+            bucket_value_bytes[plan.bucket] += column->byte_size();
+            for (uint32_t row : rowids) {
                 if (row >= rows) {
                     return Status::InternalError("Variant path {} row {} 
exceeds {} rows",
                                                  *plan.path, row, rows);
@@ -401,32 +409,88 @@ struct VariantShredder::Impl {
                 ++row_offsets[row + 1];
             }
         }
+        for (size_t bucket = 0; bucket < maps.size(); ++bucket) {
+            auto& keys = assert_cast<ColumnString&>(maps[bucket]->get_keys());
+            auto& values = 
assert_cast<ColumnString&>(maps[bucket]->get_values());
+            maps[bucket]->get_offsets().reserve(rows);
+            keys.reserve(bucket_cells[bucket]);
+            keys.get_chars().reserve(bucket_key_bytes[bucket]);
+            values.reserve(bucket_cells[bucket]);
+            values.get_chars().reserve(bucket_value_bytes[bucket]);
+        }
+        if (binary_plan.size() > std::numeric_limits<uint32_t>::max()) {
+            return Status::InternalError("Variant binary path count {} exceeds 
uint32 limit",
+                                         binary_plan.size());
+        }
         std::partial_sum(row_offsets.begin(), row_offsets.end(), 
row_offsets.begin());
-        DorisVector<BinaryCell> cells(row_offsets.back());
-        DorisVector<size_t> next_cell = row_offsets;
-        for (size_t plan_index = 0; plan_index < binary_plan.size(); 
++plan_index) {
-            const auto rowids = binary_plan[plan_index].builder->rowids();
-            for (size_t value_index = 0; value_index < rowids.size(); 
++value_index) {
-                cells[next_cell[rowids[value_index]]++] = {.plan_index = 
plan_index,
-                                                           .value_index = 
value_index};
+
+        // Transpose path-major builders into row-major maps in bounded 
chunks. A
+        // chunk always ends at a row boundary, so the path-sorted plan order
+        // remains the canonical key order within each row and bucket. The 
per-plan
+        // cursor also recovers value_index without storing it in every cell.
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+        const size_t max_binary_cells_per_chunk = binary_cells_per_chunk;
+#else
+        constexpr size_t max_binary_cells_per_chunk = 
MAX_BINARY_CELLS_PER_CHUNK;
+#endif
+        DorisVector<size_t> value_indices(binary_plan.size(), 0);
+        DorisVector<uint32_t> cells;
+        DorisVector<size_t> next_cell;
+        size_t row_begin = 0;
+        while (row_begin < rows) {
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+            ++binary_chunk_count;
+#endif
+            const size_t first_cell = row_offsets[row_begin];
+            size_t row_end = rows;
+            if (row_offsets.back() - first_cell > max_binary_cells_per_chunk) {
+                const auto first_too_large =
+                        std::upper_bound(row_offsets.begin() + row_begin + 1, 
row_offsets.end(),
+                                         first_cell + 
max_binary_cells_per_chunk);
+                row_end = static_cast<size_t>(first_too_large - 
row_offsets.begin() - 1);
+                // One exceptionally wide row may exceed the bound, but must 
stay
+                // intact.
+                row_end = std::max(row_end, row_begin + 1);
             }
-        }
 
-        for (size_t row = 0; row < rows; ++row) {
-            for (size_t cell_index = row_offsets[row]; cell_index < 
row_offsets[row + 1];
-                 ++cell_index) {
-                const BinaryCell& cell = cells[cell_index];
-                const BinaryPlan& plan = binary_plan[cell.plan_index];
-                auto& keys = 
assert_cast<ColumnString&>(maps[plan.bucket]->get_keys());
-                auto& values = 
assert_cast<ColumnString&>(maps[plan.bucket]->get_values());
-                keys.insert_data(plan.path->data(), plan.path->size());
-                RETURN_IF_ERROR(
-                        plan.builder->write_sparse_cell(cell.value_index, 
&values.get_chars()));
-                values.get_offsets().push_back(values.get_chars().size());
+            cells.resize(row_offsets[row_end] - first_cell);
+            next_cell.resize(row_end - row_begin);
+            for (size_t row = row_begin; row < row_end; ++row) {
+                next_cell[row - row_begin] = row_offsets[row] - first_cell;
+            }
+            for (size_t plan_index = 0; plan_index < binary_plan.size(); 
++plan_index) {
+                const std::span<const uint32_t> rowids = 
binary_plan[plan_index].builder->rowids();
+                size_t value_index = value_indices[plan_index];
+                DORIS_CHECK(value_index == rowids.size() || 
rowids[value_index] >= row_begin);
+                while (value_index < rowids.size() && rowids[value_index] < 
row_end) {
+                    const uint32_t row = rowids[value_index++];
+                    cells[next_cell[row - row_begin]++] = 
static_cast<uint32_t>(plan_index);
+                }
             }
-            for (ColumnMap* map : maps) {
-                map->get_offsets().push_back(map->get_keys().size());
+
+            for (size_t row = row_begin; row < row_end; ++row) {
+                DORIS_CHECK_EQ(next_cell[row - row_begin], row_offsets[row + 
1] - first_cell);
+                for (size_t cell_index = row_offsets[row] - first_cell;
+                     cell_index < row_offsets[row + 1] - first_cell; 
++cell_index) {
+                    const uint32_t plan_index = cells[cell_index];
+                    const BinaryPlan& plan = binary_plan[plan_index];
+                    const size_t value_index = value_indices[plan_index]++;
+                    auto& keys = 
assert_cast<ColumnString&>(maps[plan.bucket]->get_keys());
+                    auto& values = 
assert_cast<ColumnString&>(maps[plan.bucket]->get_values());
+                    keys.insert_data(plan.path->data(), plan.path->size());
+                    RETURN_IF_ERROR(
+                            plan.builder->write_sparse_cell(value_index, 
&values.get_chars()));
+                    values.get_offsets().push_back(values.get_chars().size());
+                }
+                for (ColumnMap* map : maps) {
+                    map->get_offsets().push_back(map->get_keys().size());
+                }
             }
+            row_begin = row_end;
+        }
+        for (size_t plan_index = 0; plan_index < binary_plan.size(); 
++plan_index) {
+            DORIS_CHECK_EQ(value_indices[plan_index],
+                           binary_plan[plan_index].builder->rowids().size());
         }
         return Status::OK();
     }
@@ -587,6 +651,10 @@ struct VariantShredder::Impl {
     DorisVector<PathState> paths;
     ColumnString::MutablePtr root_values = ColumnString::create();
     JsonbWriter root_writer;
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+    size_t binary_cells_per_chunk = MAX_BINARY_CELLS_PER_CHUNK;
+    mutable size_t binary_chunk_count = 0;
+#endif
 };
 
 VariantShredder::VariantShredder(VariantShredderOptions options)
@@ -700,4 +768,15 @@ size_t VariantShredder::byte_size() const {
     return size;
 }
 
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+size_t VariantShredder::TestAccess::binary_chunk_count(const VariantShredder& 
shredder) {
+    return shredder._impl->binary_chunk_count;
+}
+
+void VariantShredder::TestAccess::set_binary_cells_per_chunk(VariantShredder& 
shredder,
+                                                             size_t cells) {
+    DORIS_CHECK(cells > 0);
+    shredder._impl->binary_cells_per_chunk = cells;
+}
+#endif
 } // namespace doris::segment_v2
diff --git a/be/src/storage/segment/variant/v2/variant_shredder.h 
b/be/src/storage/segment/variant/v2/variant_shredder.h
index ce9535241ef..91b8c8b4944 100644
--- a/be/src/storage/segment/variant/v2/variant_shredder.h
+++ b/be/src/storage/segment/variant/v2/variant_shredder.h
@@ -87,6 +87,14 @@ public:
     Status finish(VariantShreddedColumns* output);
     size_t byte_size() const;
 
+#ifdef BE_TEST
+    struct TestAccess {
+#if !defined(BE_BENCHMARK)
+        static size_t binary_chunk_count(const VariantShredder& shredder);
+        static void set_binary_cells_per_chunk(VariantShredder& shredder, 
size_t cells);
+#endif
+    };
+#endif
 private:
     struct Impl;
     std::unique_ptr<Impl> _impl;
diff --git a/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp 
b/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp
index cc594acfc9c..a1e20984965 100644
--- a/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp
+++ b/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp
@@ -121,6 +121,16 @@ ColumnVariantV2::MutablePtr typed_ints() {
             std::make_shared<DataTypeInt32>());
 }
 
+ColumnVariantV2::MutablePtr typed_strings() {
+    auto values = ColumnString::create();
+    values->insert_data("alice", 5);
+    values->insert_data("bob", 3);
+    values->insert_data("carol", 5);
+    return ColumnVariantV2::create_typed(
+            ColumnNullable::create(std::move(values), ColumnUInt8::create(3, 
0)),
+            std::make_shared<DataTypeString>());
+}
+
 const ColumnNullable& nullable_result(const ColumnPtr& column) {
     return assert_cast<const ColumnNullable&>(*column);
 }
@@ -290,6 +300,58 @@ TEST(CastVariantV2FromTest, 
TypedInnerNullStringifiesAsLiteralNull) {
     EXPECT_FALSE(nullable.has_null());
 }
 
+TEST(CastVariantV2FromTest, TypedStringIdentityReusesPayload) {
+    ColumnPtr source = typed_strings();
+    const auto& typed = assert_cast<const ColumnVariantV2&>(*source);
+    const auto& source_nullable = assert_cast<const 
ColumnNullable&>(typed.typed_column());
+
+    CastResult identity = execute_from_variant(source, 
std::make_shared<DataTypeString>());
+    ASSERT_TRUE(identity.status.ok()) << identity.status;
+    EXPECT_EQ(identity.column.get(), &typed.typed_column());
+
+    constexpr std::array<NullMap::value_type, 3> FORCED_NULLS {0, 1, 0};
+    CastResult masked =
+            execute_from_variant(source, std::make_shared<DataTypeString>(), 
FORCED_NULLS.data());
+    ASSERT_TRUE(masked.status.ok()) << masked.status;
+    const auto& masked_nullable = nullable_result(masked.column);
+    EXPECT_EQ(masked_nullable.get_nested_column_ptr().get(),
+              source_nullable.get_nested_column_ptr().get());
+    EXPECT_EQ(masked_nullable.get_null_map_data(), (NullMap {0, 1, 0}));
+}
+
+TEST(CastVariantV2FromTest, TypedStringPreservesForcedInnerNullSemantics) {
+    auto values = ColumnString::create();
+    values->insert_data("alice", 5);
+    values->insert_default();
+    values->insert_data("carol", 5);
+    auto inner_nulls = ColumnUInt8::create();
+    inner_nulls->insert_value(0);
+    inner_nulls->insert_value(1);
+    inner_nulls->insert_value(0);
+    ColumnPtr source = ColumnVariantV2::create_typed(
+            ColumnNullable::create(std::move(values), std::move(inner_nulls)),
+            std::make_shared<DataTypeString>());
+
+    CastResult visible_null = execute_from_variant(source, 
std::make_shared<DataTypeString>());
+    ASSERT_TRUE(visible_null.status.ok()) << visible_null.status;
+    const auto& visible_nullable = nullable_result(visible_null.column);
+    const auto& visible_strings =
+            assert_cast<const 
ColumnString&>(visible_nullable.get_nested_column());
+    EXPECT_EQ(visible_strings.get_data_at(1), StringRef("null"));
+    EXPECT_FALSE(visible_nullable.has_null());
+
+    constexpr std::array<NullMap::value_type, 3> FORCED_NULLS {0, 1, 0};
+    CastResult masked =
+            execute_from_variant(source, std::make_shared<DataTypeString>(), 
FORCED_NULLS.data());
+    ASSERT_TRUE(masked.status.ok()) << masked.status;
+    const auto& masked_nullable = nullable_result(masked.column);
+    const auto& masked_strings =
+            assert_cast<const 
ColumnString&>(masked_nullable.get_nested_column());
+    EXPECT_EQ(masked_strings.get_data_at(0), StringRef("alice"));
+    EXPECT_EQ(masked_strings.get_data_at(2), StringRef("carol"));
+    EXPECT_EQ(masked_nullable.get_null_map_data(), (NullMap {0, 1, 0}));
+}
+
 TEST(CastVariantV2FromTest, TypedStringUsesCanonicalTimestampScale) {
     DateV2Value<DateTimeV2ValueType> value;
     value.unchecked_set_time(2024, 1, 2, 3, 4, 5, 123000);
diff --git a/be/test/storage/variant/variant_column_writer_reader_test.cpp 
b/be/test/storage/variant/variant_column_writer_reader_test.cpp
index c19f1eda764..fd9cf208648 100644
--- a/be/test/storage/variant/variant_column_writer_reader_test.cpp
+++ b/be/test/storage/variant/variant_column_writer_reader_test.cpp
@@ -240,7 +240,7 @@ static std::string variant_json_at(const IColumn& column, 
size_t row) {
 TEST(VariantPathBuilderTest, PromotesValuesAndMaterializesMissingRows) {
     VariantBatchBuilder value_builder;
     auto integer_row = value_builder.begin_row();
-    integer_row.add_int(1);
+    integer_row.add_float(1.0F);
     integer_row.finish();
     auto double_row = value_builder.begin_row();
     double_row.add_double(2.5);
@@ -283,7 +283,7 @@ TEST(VariantPathBuilderTest, 
PromotesValuesAndMaterializesMissingRows) {
     }
 }
 
-TEST(VariantPathBuilderTest, PreservesIntegerWidthAcrossPromotion) {
+TEST(VariantPathBuilderTest, UsesBigintForEncodedIntegersWithoutPromotion) {
     VariantBatchBuilder value_builder;
     auto tiny_row = value_builder.begin_row();
     tiny_row.add_int(1);
@@ -295,9 +295,98 @@ TEST(VariantPathBuilderTest, 
PreservesIntegerWidthAcrossPromotion) {
 
     segment_v2::VariantPathBuilder builder(PathInData("metric"));
     ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
-    EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(), 
TYPE_TINYINT);
+    EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(), 
TYPE_BIGINT);
     ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
-    EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(), TYPE_INT);
+    EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(), 
TYPE_BIGINT);
+    EXPECT_EQ(builder.promotion_count(), 0);
+}
+
+TEST(VariantPathBuilderTest, 
MatchesLegacyMixedIntegerDoubleInferenceInEitherOrder) {
+    for (const bool reverse : {false, true}) {
+        SCOPED_TRACE(testing::Message() << "reverse=" << reverse);
+        const std::vector<std::string> jsons =
+                reverse ? std::vector<std::string> {R"({"metric":1.5})", 
R"({"metric":1})"}
+                        : std::vector<std::string> {R"({"metric":1})", 
R"({"metric":1.5})"};
+
+        auto legacy = ColumnVariant::create(0, false);
+        auto json = ColumnString::create();
+        for (const std::string& value : jsons) {
+            json->insert_data(value.data(), value.size());
+        }
+        ParseConfig parse_config;
+        parse_config.parse_to = ParseConfig::ParseTo::OnlySubcolumns;
+        variant_util::parse_json_to_variant(*legacy, *json, parse_config);
+        legacy->finalize();
+        auto* legacy_metric = legacy->get_subcolumn(PathInData("metric"));
+        ASSERT_NE(legacy_metric, nullptr);
+        
EXPECT_EQ(remove_nullable(legacy_metric->get_least_common_type())->get_primitive_type(),
+                  TYPE_JSONB);
+
+        VariantBatchBuilder value_builder;
+        if (reverse) {
+            auto double_row = value_builder.begin_row();
+            double_row.add_double(1.5);
+            double_row.finish();
+            auto integer_row = value_builder.begin_row();
+            integer_row.add_int(1);
+            integer_row.finish();
+        } else {
+            auto integer_row = value_builder.begin_row();
+            integer_row.add_int(1);
+            integer_row.finish();
+            auto double_row = value_builder.begin_row();
+            double_row.add_double(1.5);
+            double_row.finish();
+        }
+        VariantBatchBuilder values = value_builder.finish_batch();
+        segment_v2::VariantPathBuilder builder(PathInData("metric"));
+        ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+        ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+        ASSERT_EQ(remove_nullable(builder.type())->get_primitive_type(), 
TYPE_JSONB);
+
+        for (size_t row = 0; row < jsons.size(); ++row) {
+            auto legacy_keys = ColumnString::create();
+            auto legacy_values = ColumnString::create();
+            legacy_metric->serialize_to_binary_column(legacy_keys.get(), 
"metric",
+                                                      legacy_values.get(), 
row);
+            ASSERT_EQ(legacy_values->size(), 1);
+            ColumnString::Chars v2_cell;
+            ASSERT_TRUE(builder.write_sparse_cell(row, &v2_cell).ok());
+            const StringRef legacy_cell = legacy_values->get_data_at(0);
+            ASSERT_EQ(v2_cell.size(), legacy_cell.size);
+            EXPECT_EQ(std::memcmp(v2_cell.data(), legacy_cell.data, 
legacy_cell.size), 0);
+        }
+    }
+}
+
+TEST(VariantPathBuilderTest, CachedBinarySerdeFollowsFloatingPromotion) {
+    VariantBatchBuilder value_builder;
+    auto tiny_row = value_builder.begin_row();
+    tiny_row.add_float(1.0F);
+    tiny_row.finish();
+    auto int_row = value_builder.begin_row();
+    int_row.add_double(2.0);
+    int_row.finish();
+    VariantBatchBuilder values = value_builder.finish_batch();
+
+    segment_v2::VariantPathBuilder builder(PathInData("metric"));
+    ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+    ColumnString::Chars binary;
+    ASSERT_TRUE(builder.write_sparse_cell(0, &binary).ok());
+    ASSERT_FALSE(binary.empty());
+    EXPECT_EQ(static_cast<FieldType>(binary.front()), 
FieldType::OLAP_FIELD_TYPE_FLOAT);
+
+    ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+    binary.clear();
+    ASSERT_TRUE(builder.write_sparse_cell(1, &binary).ok());
+    ASSERT_FALSE(binary.empty());
+    EXPECT_EQ(static_cast<FieldType>(binary.front()), 
FieldType::OLAP_FIELD_TYPE_DOUBLE);
+
+    ASSERT_TRUE(builder.convert_to(std::make_shared<DataTypeInt64>()).ok());
+    binary.clear();
+    ASSERT_TRUE(builder.write_sparse_cell(1, &binary).ok());
+    ASSERT_FALSE(binary.empty());
+    EXPECT_EQ(static_cast<FieldType>(binary.front()), 
FieldType::OLAP_FIELD_TYPE_BIGINT);
 }
 
 TEST(VariantPathBuilderTest, 
PromotesCompatibleDecimalScalesInEitherOrderAndInsideArrays) {
@@ -346,6 +435,201 @@ TEST(VariantPathBuilderTest, 
PromotesCompatibleDecimalScalesInEitherOrderAndInsi
     }
 }
 
+TEST(VariantPathBuilderTest, 
EqualFullTypeDoesNotPromoteButDifferentDecimalScaleDoes) {
+    VariantBatchBuilder value_builder;
+    for (const uint8_t scale : {2, 2, 4}) {
+        auto row = value_builder.begin_row();
+        row.add_decimal(123, scale, 16);
+        row.finish();
+    }
+    VariantBatchBuilder values = value_builder.finish_batch();
+
+    segment_v2::VariantPathBuilder builder(PathInData("metric"));
+    ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+    ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+    EXPECT_EQ(builder.promotion_count(), 0);
+    ASSERT_TRUE(builder.append(values.value_at(2), 2).ok());
+    EXPECT_EQ(builder.promotion_count(), 1);
+    EXPECT_EQ(builder.stable_scalar_append_count(), 1);
+    EXPECT_EQ(remove_nullable(builder.type())->get_scale(), 4);
+}
+
+TEST(VariantPathBuilderTest, StableScalarFastPathPreservesInferenceBoundaries) 
{
+    const auto verify_encoded = []<typename Append>(Append append) {
+        VariantBatchBuilder value_builder;
+        for (size_t row = 0; row < 2; ++row) {
+            auto value_row = value_builder.begin_row();
+            append(value_row);
+            value_row.finish();
+        }
+        VariantBatchBuilder values = value_builder.finish_batch();
+        segment_v2::VariantPathBuilder builder(PathInData("metric"));
+        ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+        ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+        EXPECT_EQ(builder.promotion_count(), 0);
+        EXPECT_EQ(builder.stable_scalar_append_count(), 1);
+    };
+
+    verify_encoded([](auto& row) { row.add_bool(true); });
+    verify_encoded([](auto& row) { row.add_int(1 << 20); });
+    verify_encoded([](auto& row) { row.add_largeint(static_cast<__int128>(1) 
<< 80); });
+    verify_encoded([](auto& row) { row.add_float(1.25F); });
+    verify_encoded([](auto& row) { row.add_double(2.5); });
+    verify_encoded([](auto& row) { row.add_string(StringRef("stable")); });
+    const std::string long_string(80, 'x');
+    verify_encoded([&](auto& row) { row.add_string(StringRef(long_string)); });
+    verify_encoded([](auto& row) { row.add_date(0); });
+    verify_encoded([](auto& row) { row.add_timestamp_micros(1, true); });
+    verify_encoded([](auto& row) { row.add_timestamp_micros(1, false); });
+    verify_encoded([](auto& row) { row.add_decimal(12300, 4, 8); });
+
+    VariantBatchBuilder widening_values_builder;
+    for (const int64_t value : std::array<int64_t, 4> {1, 1LL << 20, 1LL << 
40, 1}) {
+        auto row = widening_values_builder.begin_row();
+        row.add_int(value);
+        row.finish();
+    }
+    VariantBatchBuilder widening_values = 
widening_values_builder.finish_batch();
+    segment_v2::VariantPathBuilder widening(PathInData("widening"));
+    ASSERT_TRUE(widening.append(widening_values.value_at(0), 0).ok());
+    ASSERT_TRUE(widening.append(widening_values.value_at(1), 1).ok());
+    EXPECT_EQ(widening.stable_scalar_append_count(), 1);
+    EXPECT_EQ(widening.promotion_count(), 0);
+    ASSERT_TRUE(widening.append(widening_values.value_at(2), 2).ok());
+    EXPECT_EQ(widening.stable_scalar_append_count(), 2);
+    ASSERT_TRUE(widening.append(widening_values.value_at(3), 3).ok());
+    EXPECT_EQ(widening.stable_scalar_append_count(), 3);
+    EXPECT_EQ(widening.promotion_count(), 0);
+    EXPECT_EQ(widening.type()->to_string(*widening.column(), 3), "1");
+
+    VariantBatchBuilder conflict_values_builder;
+    auto integer_row = conflict_values_builder.begin_row();
+    integer_row.add_int(1);
+    integer_row.finish();
+    for (size_t row = 0; row < 2; ++row) {
+        auto bool_row = conflict_values_builder.begin_row();
+        bool_row.add_bool(true);
+        bool_row.finish();
+    }
+    VariantBatchBuilder conflict_values = 
conflict_values_builder.finish_batch();
+    segment_v2::VariantPathBuilder conflict(PathInData("conflict"));
+    for (size_t row = 0; row < conflict_values.num_rows(); ++row) {
+        ASSERT_TRUE(conflict.append(conflict_values.value_at(row), row).ok());
+    }
+    EXPECT_EQ(remove_nullable(conflict.type())->get_primitive_type(), 
TYPE_JSONB);
+    EXPECT_EQ(conflict.stable_scalar_append_count(), 1);
+
+    VariantBatchBuilder binary_values_builder;
+    for (size_t row = 0; row < 2; ++row) {
+        auto binary_row = binary_values_builder.begin_row();
+        binary_row.add_binary(StringRef("binary"));
+        binary_row.finish();
+    }
+    VariantBatchBuilder binary_values = binary_values_builder.finish_batch();
+    segment_v2::VariantPathBuilder binary(PathInData("binary"));
+    ASSERT_TRUE(binary.append(binary_values.value_at(0), 0).ok());
+    ASSERT_TRUE(binary.append(binary_values.value_at(1), 1).ok());
+    EXPECT_EQ(remove_nullable(binary.type())->get_primitive_type(), 
TYPE_JSONB);
+    EXPECT_EQ(binary.stable_scalar_append_count(), 1);
+
+    VariantBatchBuilder array_values_builder;
+    for (size_t row = 0; row < 2; ++row) {
+        auto value_row = array_values_builder.begin_row();
+        auto array = value_row.start_array();
+        value_row.add_int(1);
+        array.finish();
+        value_row.finish();
+    }
+    VariantBatchBuilder array_values = array_values_builder.finish_batch();
+    segment_v2::VariantPathBuilder arrays(PathInData("arrays"));
+    ASSERT_TRUE(arrays.append(array_values.value_at(0), 0).ok());
+    ASSERT_TRUE(arrays.append(array_values.value_at(1), 1).ok());
+    EXPECT_EQ(arrays.stable_scalar_append_count(), 0);
+
+    const auto verify_out_of_range_falls_back = []<typename Append>(Append 
append) {
+        VariantBatchBuilder value_builder;
+        auto valid_row = value_builder.begin_row();
+        append(valid_row, false);
+        valid_row.finish();
+        auto invalid_row = value_builder.begin_row();
+        append(invalid_row, true);
+        invalid_row.finish();
+        VariantBatchBuilder values = value_builder.finish_batch();
+
+        segment_v2::VariantPathBuilder builder(PathInData("range"));
+        ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+        ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+        EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(), 
TYPE_JSONB);
+        EXPECT_EQ(builder.stable_scalar_append_count(), 0);
+    };
+    verify_out_of_range_falls_back(
+            [](auto& row, bool invalid) { row.add_date(invalid ? 3'000'000 : 
0); });
+    verify_out_of_range_falls_back([](auto& row, bool invalid) {
+        row.add_timestamp_micros(invalid ? 253'402'300'800'000'000LL : 0, 
true);
+    });
+}
+
+TEST(VariantPathBuilderTest, 
StableScalarGuardRetainsDecimalAndAppendFailureFallbacks) {
+    VariantBatchBuilder value_builder;
+    auto valid_row = value_builder.begin_row();
+    valid_row.add_decimal(1, 1, 16);
+    valid_row.finish();
+    VariantBatchBuilder values = value_builder.finish_batch();
+
+    segment_v2::VariantPathBuilder decimal(PathInData("decimal"));
+    ASSERT_TRUE(decimal.append(values.value_at(0), 0).ok());
+
+    std::array<char, 18> overflow_decimal {};
+    overflow_decimal[0] = 
static_cast<char>(static_cast<uint8_t>(VariantPrimitiveId::DECIMAL16)
+                                            << VARIANT_VALUE_HEADER_SHIFT);
+    overflow_decimal[1] = 1;
+    unsigned __int128 unscaled = VARIANT_DECIMAL16_MAX + 1;
+    for (size_t byte = 0; byte < 16; ++byte) {
+        overflow_decimal[byte + 2] = static_cast<char>(unscaled >> (byte * 8));
+    }
+    ASSERT_TRUE(
+            decimal.append(VariantRef {.metadata = {},
+                                       .value = {overflow_decimal.data(), 
overflow_decimal.size()}},
+                           1)
+                    .ok());
+    EXPECT_EQ(remove_nullable(decimal.type())->get_primitive_type(), 
TYPE_JSONB);
+    EXPECT_EQ(decimal.promotion_count(), 1);
+    EXPECT_EQ(decimal.stable_scalar_append_count(), 0);
+    EXPECT_EQ(decimal.non_null_rows(), 2);
+
+    VariantBatchBuilder narrow_values_builder;
+    auto narrow_row = narrow_values_builder.begin_row();
+    narrow_row.add_decimal(9999, 1, 4);
+    narrow_row.finish();
+    VariantBatchBuilder narrow_values = narrow_values_builder.finish_batch();
+    segment_v2::VariantPathBuilder 
narrow_decimal(PathInData("narrow_decimal"));
+    ASSERT_TRUE(narrow_decimal.append(values.value_at(0), 0).ok());
+    
ASSERT_TRUE(narrow_decimal.convert_to(std::make_shared<DataTypeDecimal128>(3, 
1)).ok());
+    ASSERT_TRUE(narrow_decimal.append(narrow_values.value_at(0), 1).ok());
+    EXPECT_EQ(remove_nullable(narrow_decimal.type())->get_primitive_type(), 
TYPE_JSONB);
+    EXPECT_EQ(narrow_decimal.promotion_count(), 2);
+    EXPECT_EQ(narrow_decimal.stable_scalar_append_count(), 0);
+    EXPECT_EQ(narrow_decimal.non_null_rows(), 2);
+
+    VariantBatchBuilder integer_builder;
+    auto integer_row = integer_builder.begin_row();
+    integer_row.add_int(1);
+    integer_row.finish();
+    VariantBatchBuilder integer_values = integer_builder.finish_batch();
+
+    segment_v2::VariantPathBuilder truncated(PathInData("truncated"));
+    ASSERT_TRUE(truncated.append(integer_values.value_at(0), 0).ok());
+    const char truncated_int8 = 
static_cast<char>(static_cast<uint8_t>(VariantPrimitiveId::INT8)
+                                                  << 
VARIANT_VALUE_HEADER_SHIFT);
+    const Status truncated_status =
+            truncated.append(VariantRef {.metadata = {}, .value = 
{&truncated_int8, 1}}, 1);
+    EXPECT_FALSE(truncated_status.ok());
+    EXPECT_EQ(remove_nullable(truncated.type())->get_primitive_type(), 
TYPE_JSONB);
+    EXPECT_EQ(truncated.promotion_count(), 1);
+    EXPECT_EQ(truncated.stable_scalar_append_count(), 0);
+    EXPECT_EQ(truncated.non_null_rows(), 1);
+}
+
 TEST(VariantPathBuilderTest, JsonbFallbackPreservesCanonicalNumericBytes) {
     VariantBatchBuilder value_builder;
     auto row = value_builder.begin_row();
@@ -637,6 +921,199 @@ TEST(VariantPathBuilderTest, 
ShredderReusesCanonicalPathsAcrossAppends) {
     EXPECT_EQ(selected.type->to_string(*selected.column, 1), "2");
 }
 
+TEST(VariantShredderTest, ReusesRootAndNestedPathsForSharedMetadata) {
+    DataTypeVariantV2SerDe serde;
+    DataTypeSerDe::FormatOptions format_options;
+    auto values = ColumnVariantV2::create();
+    for (const std::string_view json :
+         {R"({"a":1,"nested":{"x":2,"y":3}})", 
R"({"a":4,"nested":{"x":5,"y":6}})"}) {
+        Slice slice(json.data(), json.size());
+        ASSERT_TRUE(serde.deserialize_one_cell_from_json(*values, slice, 
format_options).ok());
+    }
+    ASSERT_EQ(values->read_view().metadata_count(), 1);
+
+    segment_v2::VariantShredderOptions options;
+    options.max_subcolumns_count = 0;
+    options.sparse_bucket_count = 1;
+    segment_v2::VariantShredder shredder(std::move(options));
+    ASSERT_TRUE(shredder.append(values->read_view(), 0, values->size()).ok());
+
+    segment_v2::VariantShreddedColumns shredded;
+    ASSERT_TRUE(shredder.finish(&shredded).ok());
+    ASSERT_EQ(shredded.materialized.size(), 3);
+    const auto expect_path = [&](std::string_view path, std::string_view first,
+                                 std::string_view second) {
+        const auto found = std::ranges::find_if(shredded.materialized, 
[&](const auto& column) {
+            return column.path.get_path() == path;
+        });
+        ASSERT_NE(found, shredded.materialized.end()) << path;
+        ASSERT_TRUE(found->column);
+        EXPECT_EQ(found->rowids, (DorisVector<uint32_t> {0, 1}));
+        EXPECT_EQ(found->type->to_string(*found->column, 0), first);
+        EXPECT_EQ(found->type->to_string(*found->column, 1), second);
+    };
+    expect_path("a", "1", "4");
+    expect_path("nested.x", "2", "5");
+    expect_path("nested.y", "3", "6");
+}
+
+static void expect_variant_statistics_equal(const 
segment_v2::VariantStatistics& actual,
+                                            const 
segment_v2::VariantStatistics& expected) {
+    EXPECT_EQ(actual.subcolumns_non_null_size, 
expected.subcolumns_non_null_size);
+    EXPECT_EQ(actual.sparse_column_non_null_size, 
expected.sparse_column_non_null_size);
+    EXPECT_EQ(actual.doc_value_column_non_null_size, 
expected.doc_value_column_non_null_size);
+    EXPECT_EQ(actual.has_nested_group, expected.has_nested_group);
+}
+
+static void expect_physical_column_data_equal(const IColumn& actual, const 
IColumn& expected) {
+    ASSERT_EQ(actual.get_name(), expected.get_name());
+    ASSERT_EQ(actual.size(), expected.size());
+    if (const auto* actual_string = 
check_and_get_column<ColumnString>(actual)) {
+        const auto& expected_string = assert_cast<const 
ColumnString&>(expected);
+        EXPECT_EQ(actual_string->get_chars(), expected_string.get_chars());
+        EXPECT_EQ(actual_string->get_offsets(), expected_string.get_offsets());
+        return;
+    }
+    if (const auto* actual_nullable = 
check_and_get_column<ColumnNullable>(actual)) {
+        const auto& expected_nullable = assert_cast<const 
ColumnNullable&>(expected);
+        EXPECT_EQ(actual_nullable->get_null_map_data(), 
expected_nullable.get_null_map_data());
+        expect_physical_column_data_equal(actual_nullable->get_nested_column(),
+                                          
expected_nullable.get_nested_column());
+        return;
+    }
+    if (const auto* actual_map = check_and_get_column<ColumnMap>(actual)) {
+        const auto& expected_map = assert_cast<const ColumnMap&>(expected);
+        EXPECT_EQ(actual_map->get_offsets(), expected_map.get_offsets());
+        expect_physical_column_data_equal(actual_map->get_keys(), 
expected_map.get_keys());
+        expect_physical_column_data_equal(actual_map->get_values(), 
expected_map.get_values());
+        return;
+    }
+    if (const auto* actual_array = check_and_get_column<ColumnArray>(actual)) {
+        const auto& expected_array = assert_cast<const ColumnArray&>(expected);
+        EXPECT_EQ(actual_array->get_offsets(), expected_array.get_offsets());
+        expect_physical_column_data_equal(actual_array->get_data(), 
expected_array.get_data());
+        return;
+    }
+    for (size_t row = 0; row < actual.size(); ++row) {
+        EXPECT_EQ(actual.compare_at(row, row, expected, -1), 0) << "row=" << 
row;
+    }
+}
+
+static void expect_physical_columns_equal(const ColumnPtr& actual, const 
ColumnPtr& expected) {
+    ASSERT_TRUE(actual);
+    ASSERT_TRUE(expected);
+    expect_physical_column_data_equal(*actual, *expected);
+}
+
+static void expect_shredded_columns_equal(const 
segment_v2::VariantShreddedColumns& actual,
+                                          const 
segment_v2::VariantShreddedColumns& expected) {
+    ASSERT_EQ(actual.num_rows, expected.num_rows);
+    expect_physical_columns_equal(actual.root_jsonb, expected.root_jsonb);
+
+    ASSERT_EQ(actual.materialized.size(), expected.materialized.size());
+    for (size_t index = 0; index < actual.materialized.size(); ++index) {
+        const auto& actual_path = actual.materialized[index];
+        const auto& expected_path = expected.materialized[index];
+        EXPECT_EQ(actual_path.path, expected_path.path) << "index=" << index;
+        ASSERT_TRUE(actual_path.type);
+        ASSERT_TRUE(expected_path.type);
+        EXPECT_TRUE(actual_path.type->equals(*expected_path.type)) << "index=" 
<< index;
+        EXPECT_EQ(actual_path.rowids, expected_path.rowids) << "index=" << 
index;
+        expect_physical_columns_equal(actual_path.column, 
expected_path.column);
+    }
+
+    ASSERT_EQ(actual.binary_buckets.size(), expected.binary_buckets.size());
+    for (size_t bucket = 0; bucket < actual.binary_buckets.size(); ++bucket) {
+        expect_physical_columns_equal(actual.binary_buckets[bucket].column,
+                                      expected.binary_buckets[bucket].column);
+        
expect_variant_statistics_equal(actual.binary_buckets[bucket].statistics,
+                                        
expected.binary_buckets[bucket].statistics);
+    }
+    expect_variant_statistics_equal(actual.statistics, expected.statistics);
+}
+
+TEST(VariantShredderTest, 
ChunkedBinaryTransposeMatchesSingleChunkForOrdinaryAndDoc) {
+    const auto path_for_bucket = [](uint32_t bucket, std::string_view prefix) {
+        for (uint32_t suffix = 0; suffix < 1024; ++suffix) {
+            std::string path(prefix);
+            path += std::to_string(suffix);
+            if (variant_util::variant_binary_shard_of({path.data(), 
path.size()}, 2) == bucket) {
+                return path;
+            }
+        }
+        DORIS_CHECK(false) << "failed to find path for bucket " << bucket;
+        return std::string {};
+    };
+    const std::array<std::string, 3> sparse_paths {
+            path_for_bucket(0, "chunk_left_"),
+            path_for_bucket(0, "chunk_middle_"),
+            path_for_bucket(1, "chunk_right_"),
+    };
+    // Ordinary sparse cells per row are 3,1,0,2,3,1,0,2. A four-cell limit
+    // crosses exact boundaries and empty rows; a two-cell limit also forces 
the
+    // over-wide-row branch. DOC additionally publishes hot in the binary map, 
so
+    // both layouts exercise multiple chunks.
+    const std::array<uint8_t, 8> sparse_masks {0b111, 0b001, 0b000, 0b110,
+                                               0b111, 0b100, 0b000, 0b011};
+    auto values = ColumnVariantV2::create();
+    DataTypeVariantV2SerDe serde;
+    DataTypeSerDe::FormatOptions format_options;
+    for (size_t row = 0; row < sparse_masks.size(); ++row) {
+        std::string json = "{\"hot\":" + std::to_string(row);
+        for (size_t path = 0; path < sparse_paths.size(); ++path) {
+            if ((sparse_masks[row] & (1U << path)) != 0) {
+                json += ",\"" + sparse_paths[path] + "\":" + 
std::to_string(row * 10 + path);
+            }
+        }
+        json += "}";
+        Slice slice(json.data(), json.size());
+        ASSERT_TRUE(serde.deserialize_one_cell_from_json(*values, slice, 
format_options).ok());
+    }
+
+    for (const auto physical_layout : 
{segment_v2::VariantShredderPhysicalLayout::ORDINARY,
+                                       
segment_v2::VariantShredderPhysicalLayout::DOC}) {
+        SCOPED_TRACE(physical_layout == 
segment_v2::VariantShredderPhysicalLayout::ORDINARY
+                             ? "ordinary"
+                             : "doc");
+        const segment_v2::VariantShredderOptions options {
+                .physical_layout = physical_layout,
+                .max_subcolumns_count = 1,
+                .sparse_bucket_count = 2,
+                .doc_bucket_count = 2,
+                .doc_materialization_min_rows = sparse_masks.size() + 1,
+        };
+        segment_v2::VariantShredder single_chunk(options);
+        ASSERT_TRUE(single_chunk.append(values->read_view(), 0, 
values->size()).ok());
+        segment_v2::VariantShreddedColumns expected;
+        ASSERT_TRUE(single_chunk.finish(&expected).ok());
+
+        for (const size_t chunk_limit : {size_t {4}, size_t {2}}) {
+            SCOPED_TRACE(testing::Message() << "chunk_limit=" << chunk_limit);
+            segment_v2::VariantShredder chunked(options);
+            
segment_v2::VariantShredder::TestAccess::set_binary_cells_per_chunk(chunked,
+                                                                               
 chunk_limit);
+            ASSERT_TRUE(chunked.append(values->read_view(), 0, 
values->size()).ok());
+            segment_v2::VariantShreddedColumns actual;
+            ASSERT_TRUE(chunked.finish(&actual).ok());
+
+            const size_t expected_chunks =
+                    physical_layout == 
segment_v2::VariantShredderPhysicalLayout::ORDINARY
+                            ? (chunk_limit == 4 ? 4 : 6)
+                            : (chunk_limit == 4 ? 6 : 8);
+            
EXPECT_EQ(segment_v2::VariantShredder::TestAccess::binary_chunk_count(chunked),
+                      expected_chunks);
+            expect_shredded_columns_equal(actual, expected);
+            ASSERT_EQ(actual.binary_buckets.size(), 2);
+            for (size_t bucket = 0; bucket < actual.binary_buckets.size(); 
++bucket) {
+                const auto& map =
+                        assert_cast<const 
ColumnMap&>(*actual.binary_buckets[bucket].column);
+                EXPECT_EQ(map.size(), sparse_masks.size());
+                EXPECT_GT(map.get_keys().size(), 0) << "bucket=" << bucket;
+            }
+        }
+    }
+}
+
 static void construct_column(ColumnPB* column_pb, int32_t col_unique_id,
                              const std::string& column_type, const 
std::string& column_name,
                              int variant_max_subcolumns_count = 3, bool is_key 
= false,
diff --git a/be/test/util/variant/variant_value_test.cpp 
b/be/test/util/variant/variant_value_test.cpp
index 88b80f53af8..0d8cc8e6cdb 100644
--- a/be/test/util/variant/variant_value_test.cpp
+++ b/be/test/util/variant/variant_value_test.cpp
@@ -405,6 +405,46 @@ TEST(VariantValueTest, 
ObjectLookupSortedAndUnsortedMetadata) {
     EXPECT_EQ(id, 1);
 }
 
+TEST(VariantValueTest, ObjectViewMatchesRandomAccessAndRetainsBoundsChecks) {
+    const std::string sorted_metadata = metadata({"a", "b"}, true);
+    const std::string false_value = primitive(VariantPrimitiveId::FALSE_VALUE);
+    const std::string true_value = primitive(VariantPrimitiveId::TRUE_VALUE);
+    const std::string encoded = object_value({0, 1}, {1, 0}, {false_value, 
true_value});
+    const VariantRef ref = value_ref(sorted_metadata, encoded);
+
+    const VariantRef::ObjectView object = ref.object_view();
+    ASSERT_EQ(object.size(), 2);
+    for (uint32_t index = 0; index < object.size(); ++index) {
+        uint32_t view_field = std::numeric_limits<uint32_t>::max();
+        uint32_t direct_field = std::numeric_limits<uint32_t>::max();
+        const VariantRef view_value = object.value_at(index, &view_field);
+        const VariantRef direct_value = ref.object_value_at(index, 
&direct_field);
+        EXPECT_EQ(view_field, direct_field);
+        EXPECT_EQ(view_value.value.data, direct_value.value.data);
+        EXPECT_EQ(view_value.value.size, direct_value.value.size);
+    }
+    EXPECT_TRUE(object.value_at(0).get_bool());
+    EXPECT_FALSE(object.value_at(1).get_bool());
+    EXPECT_THROW(object.value_at(object.size()), Exception);
+
+    const std::string invalid_id_object =
+            object_value({2}, {0}, 
{primitive(VariantPrimitiveId::NULL_VALUE)});
+    const VariantRef::ObjectView invalid_id =
+            value_ref(sorted_metadata, invalid_id_object).object_view();
+    EXPECT_THROW(invalid_id.value_at(0), Exception);
+
+    const std::string truncated_object(1, 
static_cast<char>(VariantBasicType::OBJECT));
+    EXPECT_THROW(value_ref(sorted_metadata, truncated_object).object_view(), 
Exception);
+    const std::string truncated_metadata = sorted_metadata.substr(0, 2);
+    EXPECT_THROW(value_ref(truncated_metadata, encoded).object_view(), 
Exception);
+
+    // An empty object has no field ids, so iterating it must not inspect 
otherwise unused
+    // metadata. This preserves the random-access API's validation boundary.
+    const std::string empty_object = object_value({}, {}, {});
+    const VariantRef::ObjectView empty = value_ref(truncated_metadata, 
empty_object).object_view();
+    EXPECT_EQ(empty.size(), 0);
+}
+
 TEST(VariantValueTest, ObjectFindRejectsInvalidReceivers) {
     const std::string empty_metadata = metadata({}, true);
     VariantRef found;


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

Reply via email to