Gabriel39 commented on code in PR #68667:
URL: https://github.com/apache/doris/pull/68667#discussion_r4154707977


##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -19,28 +19,268 @@
 
 #include <arrow/array/builder_binary.h>
 
+#include <algorithm>
+#include <cmath>
 #include <cstdint>
 #include <string>
+#include <string_view>
+#include <vector>
 
 #include "common/cast_set.h"
 #include "common/config.h"
 #include "common/exception.h"
 #include "common/status.h"
 #include "core/assert_cast.h"
 #include "core/column/column.h"
+#include "core/column/column_array.h"
+#include "core/column/column_map.h"
+#include "core/column/column_struct.h"
 #include "core/column/column_variant.h"
+#include "core/column/variant_v2/column_variant_v2.h"
+#include "core/column/variant_v2/column_variant_v2_typed_column.h"
+#include "core/data_type/data_type_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_serde/data_type_serde.h"
+#include "core/data_type_serde/data_type_variant_v2_serde.h"
 #include "core/field.h"
 #include "core/string_ref.h"
 #include "core/types.h"
 #include "core/value/jsonb_value.h"
 #include "exec/common/variant_util.h"
+#include "exprs/function/parse/variant_jsonb_parse.h"
+#include "exprs/function/parse/variant_string_parse.h"
 #include "util/json/json_parser.h"
 #include "util/jsonb_writer.h"
 
 namespace doris {
 namespace {
 
+Status append_legacy_arrow_document(const ColumnVariant& column, size_t index,
+                                    VariantBatchBuilder::Row& output,
+                                    const DataTypeSerDe::FormatOptions& 
options, size_t depth);
+
+// Legacy CAST accepts more root families than V2 CAST. Encode their structure 
here so
+// Flight output does not reject valid roots or lose typed leaves through JSON 
reparsing.
+Status append_legacy_arrow_value(const IColumn& column, const DataTypePtr& 
type, size_t index,
+                                 VariantBatchBuilder::Row& output,
+                                 const DataTypeSerDe::FormatOptions& options, 
size_t depth = 0) {
+    if (depth > VARIANT_MAX_NESTING_DEPTH) {
+        return Status::NotSupported(
+                "Native Arrow Variant nesting exceeds {}; "
+                "use enable_arrow_flight_sql_native_variant=false for UTF8 
output",
+                VARIANT_MAX_NESTING_DEPTH);
+    }
+    if (const auto* constant = check_and_get_column<ColumnConst>(column)) {
+        return append_legacy_arrow_value(constant->get_data_column(), type, 0, 
output, options,
+                                         depth);
+    }
+    if (const auto* nullable = check_and_get_column<ColumnNullable>(column)) {
+        if (nullable->is_null_at(index)) {
+            output.add_null();
+            return Status::OK();
+        }
+        return append_legacy_arrow_value(nullable->get_nested_column(), 
remove_nullable(type),
+                                         index, output, options, depth);
+    }
+    const auto primitive = type->get_primitive_type();
+    if (is_supported_variant_typed_identity(primitive)) {
+        dispatch_variant_typed_column(
+                column, primitive, [&]<PrimitiveType Type>(const auto& scalar) 
{
+                    with_variant_typed_scalar<Type>(
+                            scalar, index, 
cast_set<uint8_t>(type->get_scale()),
+                            [&](const VariantScalarRef& value) { 
output.add_scalar(value); });
+                });
+    } else if (primitive == TYPE_TIMEV2) {
+        // TIMEV2 already stores microseconds; treating its physical double as 
a number loses its type.
+        const double micros = assert_cast<const 
ColumnTimeV2&>(column).get_data()[index];
+        // Parquet TIME is a time of day, whereas Doris TIME also represents 
signed durations.
+        // Reject unrepresentable durations instead of wrapping them or 
emitting invalid TIME values.
+        constexpr int64_t micros_per_day = 86400000000;
+        if (!std::isfinite(micros) || micros < 0 || micros >= micros_per_day ||
+            std::llround(micros) >= micros_per_day) {
+            return Status::NotSupported(
+                    "Native Arrow Variant TIMEV2 requires a time in [00:00:00, 
24:00:00); "
+                    "use enable_arrow_flight_sql_native_variant=false for UTF8 
output");
+        }
+        output.add_time_ntz_micros(std::llround(micros));
+    } else if (primitive == TYPE_VARBINARY) {
+        // Binary leaves must retain arbitrary bytes, including NUL and 
non-UTF8 data.
+        output.add_binary(column.get_data_at(index));
+    } else if (primitive == TYPE_JSONB) {
+        // JSONB leaves retain their own depth limit, but also consume the 
enclosing Variant depth.
+        try {
+            jsonb_to_variant(column.get_data_at(index), output, 
cast_set<uint32_t>(depth));
+        } catch (const Exception& e) {
+            if (e.code() != ErrorCode::INVALID_ARGUMENT) {
+                return e.to_status();
+            }
+            return Status::NotSupported(
+                    "Native Arrow Variant cannot encode JSONB leaf: {}; "
+                    "use enable_arrow_flight_sql_native_variant=false for UTF8 
output",
+                    e.what());
+        }
+    } else if (primitive == TYPE_ARRAY) {
+        const auto& array = assert_cast<const ColumnArray&>(column);
+        const auto& array_type = assert_cast<const DataTypeArray&>(*type);
+        auto scope = output.start_array();
+        for (size_t element = array.offset_at(index); element < 
array.get_offsets()[index];
+             ++element) {
+            RETURN_IF_ERROR(append_legacy_arrow_value(array.get_data(),
+                                                      
array_type.get_nested_type(), element, output,
+                                                      options, depth + 1));
+        }
+        scope.finish();
+    } else if (primitive == TYPE_MAP) {
+        const auto& map = assert_cast<const ColumnMap&>(column);
+        const auto& map_type = assert_cast<const DataTypeMap&>(*type);
+        auto scope = output.start_object();
+        for (size_t element = map.get_offsets()[static_cast<ssize_t>(index) - 
1];
+             element < map.get_offsets()[index]; ++element) {
+            // Variant object keys cannot distinguish SQL NULL from the 
literal string "null".
+            if (map.get_keys().is_null_at(element)) {
+                return Status::NotSupported(
+                        "Native Arrow Variant cannot represent MAP with NULL 
keys; "
+                        "use enable_arrow_flight_sql_native_variant=false for 
UTF8 output");
+            }
+            auto key = map_type.get_key_type()->to_string(map.get_keys(), 
element, options);
+            scope.add_key({key.data(), key.size()});
+            RETURN_IF_ERROR(append_legacy_arrow_value(map.get_values(), 
map_type.get_value_type(),
+                                                      element, output, 
options, depth + 1));
+        }
+        scope.finish();
+    } else if (primitive == TYPE_STRUCT) {
+        const auto& structure = assert_cast<const ColumnStruct&>(column);
+        const auto& struct_type = assert_cast<const DataTypeStruct&>(*type);
+        auto scope = output.start_object();
+        for (size_t field = 0; field < struct_type.get_elements().size(); 
++field) {
+            const auto& name = struct_type.get_element_names()[field];
+            scope.add_key({name.data(), name.size()});
+            
RETURN_IF_ERROR(append_legacy_arrow_value(structure.get_column(field),
+                                                      
struct_type.get_element(field), index, output,
+                                                      options, depth + 1));
+        }
+        scope.finish();
+    } else if (primitive == TYPE_VARIANT) {
+        if (const auto* legacy = check_and_get_column<ColumnVariant>(column)) {
+            const bool visible = legacy->is_scalar_variant()
+                                         ? 
!legacy->get_root()->is_null_at(index)
+                                         : 
legacy->is_visible_root_value(index);
+            if (visible) {
+                return append_legacy_arrow_value(*legacy->get_root(), 
legacy->get_root_type(),
+                                                 index, output, options, 
depth);
+            }
+            RETURN_IF_ERROR(append_legacy_arrow_document(*legacy, index, 
output, options, depth));
+        } else {
+            visit_variant_v2_values(
+                    column, index, index + 1, {}, [&](size_t) { 
output.add_null(); },
+                    [&](size_t, VariantRef value) { output.add_value(value); 
});

Review Comment:
   Fixed in 70746f4eb77. Nested V2 leaves now use the same selected-value 
importer as native V2 output and pass the enclosing legacy depth. Exceeding the 
limit returns NotSupported with `enable_arrow_flight_sql_native_variant=false`, 
rather than throwing an import error that becomes InternalError.
   
   Added a regression covering scalar, empty-array and empty-object leaves at 
the boundary, including successful UTF8 fallback. It failed on the previous 
code and passes with this fix. All 53 focused BE tests passed; the complete 
native Flight regression suite and Python ADBC test also passed locally with 
both V1 and V2 using the newly compiled BE.
   



##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -19,28 +19,268 @@
 
 #include <arrow/array/builder_binary.h>
 
+#include <algorithm>
+#include <cmath>
 #include <cstdint>
 #include <string>
+#include <string_view>
+#include <vector>
 
 #include "common/cast_set.h"
 #include "common/config.h"
 #include "common/exception.h"
 #include "common/status.h"
 #include "core/assert_cast.h"
 #include "core/column/column.h"
+#include "core/column/column_array.h"
+#include "core/column/column_map.h"
+#include "core/column/column_struct.h"
 #include "core/column/column_variant.h"
+#include "core/column/variant_v2/column_variant_v2.h"
+#include "core/column/variant_v2/column_variant_v2_typed_column.h"
+#include "core/data_type/data_type_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_serde/data_type_serde.h"
+#include "core/data_type_serde/data_type_variant_v2_serde.h"
 #include "core/field.h"
 #include "core/string_ref.h"
 #include "core/types.h"
 #include "core/value/jsonb_value.h"
 #include "exec/common/variant_util.h"
+#include "exprs/function/parse/variant_jsonb_parse.h"
+#include "exprs/function/parse/variant_string_parse.h"
 #include "util/json/json_parser.h"
 #include "util/jsonb_writer.h"
 
 namespace doris {
 namespace {
 
+Status append_legacy_arrow_document(const ColumnVariant& column, size_t index,
+                                    VariantBatchBuilder::Row& output,
+                                    const DataTypeSerDe::FormatOptions& 
options, size_t depth);
+
+// Legacy CAST accepts more root families than V2 CAST. Encode their structure 
here so
+// Flight output does not reject valid roots or lose typed leaves through JSON 
reparsing.
+Status append_legacy_arrow_value(const IColumn& column, const DataTypePtr& 
type, size_t index,
+                                 VariantBatchBuilder::Row& output,
+                                 const DataTypeSerDe::FormatOptions& options, 
size_t depth = 0) {
+    if (depth > VARIANT_MAX_NESTING_DEPTH) {
+        return Status::NotSupported(
+                "Native Arrow Variant nesting exceeds {}; "
+                "use enable_arrow_flight_sql_native_variant=false for UTF8 
output",
+                VARIANT_MAX_NESTING_DEPTH);
+    }
+    if (const auto* constant = check_and_get_column<ColumnConst>(column)) {
+        return append_legacy_arrow_value(constant->get_data_column(), type, 0, 
output, options,
+                                         depth);
+    }
+    if (const auto* nullable = check_and_get_column<ColumnNullable>(column)) {
+        if (nullable->is_null_at(index)) {
+            output.add_null();
+            return Status::OK();
+        }
+        return append_legacy_arrow_value(nullable->get_nested_column(), 
remove_nullable(type),
+                                         index, output, options, depth);
+    }
+    const auto primitive = type->get_primitive_type();
+    if (is_supported_variant_typed_identity(primitive)) {
+        dispatch_variant_typed_column(
+                column, primitive, [&]<PrimitiveType Type>(const auto& scalar) 
{
+                    with_variant_typed_scalar<Type>(
+                            scalar, index, 
cast_set<uint8_t>(type->get_scale()),
+                            [&](const VariantScalarRef& value) { 
output.add_scalar(value); });
+                });
+    } else if (primitive == TYPE_TIMEV2) {
+        // TIMEV2 already stores microseconds; treating its physical double as 
a number loses its type.
+        const double micros = assert_cast<const 
ColumnTimeV2&>(column).get_data()[index];
+        // Parquet TIME is a time of day, whereas Doris TIME also represents 
signed durations.
+        // Reject unrepresentable durations instead of wrapping them or 
emitting invalid TIME values.
+        constexpr int64_t micros_per_day = 86400000000;
+        if (!std::isfinite(micros) || micros < 0 || micros >= micros_per_day ||
+            std::llround(micros) >= micros_per_day) {
+            return Status::NotSupported(
+                    "Native Arrow Variant TIMEV2 requires a time in [00:00:00, 
24:00:00); "
+                    "use enable_arrow_flight_sql_native_variant=false for UTF8 
output");
+        }
+        output.add_time_ntz_micros(std::llround(micros));
+    } else if (primitive == TYPE_VARBINARY) {
+        // Binary leaves must retain arbitrary bytes, including NUL and 
non-UTF8 data.
+        output.add_binary(column.get_data_at(index));
+    } else if (primitive == TYPE_JSONB) {
+        // JSONB leaves retain their own depth limit, but also consume the 
enclosing Variant depth.
+        try {
+            jsonb_to_variant(column.get_data_at(index), output, 
cast_set<uint32_t>(depth));
+        } catch (const Exception& e) {
+            if (e.code() != ErrorCode::INVALID_ARGUMENT) {
+                return e.to_status();
+            }
+            return Status::NotSupported(
+                    "Native Arrow Variant cannot encode JSONB leaf: {}; "
+                    "use enable_arrow_flight_sql_native_variant=false for UTF8 
output",
+                    e.what());
+        }
+    } else if (primitive == TYPE_ARRAY) {
+        const auto& array = assert_cast<const ColumnArray&>(column);
+        const auto& array_type = assert_cast<const DataTypeArray&>(*type);
+        auto scope = output.start_array();
+        for (size_t element = array.offset_at(index); element < 
array.get_offsets()[index];
+             ++element) {
+            RETURN_IF_ERROR(append_legacy_arrow_value(array.get_data(),
+                                                      
array_type.get_nested_type(), element, output,
+                                                      options, depth + 1));
+        }
+        scope.finish();
+    } else if (primitive == TYPE_MAP) {
+        const auto& map = assert_cast<const ColumnMap&>(column);
+        const auto& map_type = assert_cast<const DataTypeMap&>(*type);
+        auto scope = output.start_object();
+        for (size_t element = map.get_offsets()[static_cast<ssize_t>(index) - 
1];
+             element < map.get_offsets()[index]; ++element) {
+            // Variant object keys cannot distinguish SQL NULL from the 
literal string "null".
+            if (map.get_keys().is_null_at(element)) {
+                return Status::NotSupported(
+                        "Native Arrow Variant cannot represent MAP with NULL 
keys; "
+                        "use enable_arrow_flight_sql_native_variant=false for 
UTF8 output");
+            }
+            auto key = map_type.get_key_type()->to_string(map.get_keys(), 
element, options);
+            scope.add_key({key.data(), key.size()});
+            RETURN_IF_ERROR(append_legacy_arrow_value(map.get_values(), 
map_type.get_value_type(),
+                                                      element, output, 
options, depth + 1));
+        }
+        scope.finish();
+    } else if (primitive == TYPE_STRUCT) {
+        const auto& structure = assert_cast<const ColumnStruct&>(column);
+        const auto& struct_type = assert_cast<const DataTypeStruct&>(*type);
+        auto scope = output.start_object();
+        for (size_t field = 0; field < struct_type.get_elements().size(); 
++field) {
+            const auto& name = struct_type.get_element_names()[field];
+            scope.add_key({name.data(), name.size()});
+            
RETURN_IF_ERROR(append_legacy_arrow_value(structure.get_column(field),
+                                                      
struct_type.get_element(field), index, output,
+                                                      options, depth + 1));
+        }
+        scope.finish();
+    } else if (primitive == TYPE_VARIANT) {
+        if (const auto* legacy = check_and_get_column<ColumnVariant>(column)) {
+            const bool visible = legacy->is_scalar_variant()
+                                         ? 
!legacy->get_root()->is_null_at(index)
+                                         : 
legacy->is_visible_root_value(index);
+            if (visible) {
+                return append_legacy_arrow_value(*legacy->get_root(), 
legacy->get_root_type(),
+                                                 index, output, options, 
depth);
+            }
+            RETURN_IF_ERROR(append_legacy_arrow_document(*legacy, index, 
output, options, depth));
+        } else {
+            visit_variant_v2_values(
+                    column, index, index + 1, {}, [&](size_t) { 
output.add_null(); },
+                    [&](size_t, VariantRef value) { output.add_value(value); 
});

Review Comment:
   Fixed in 70746f4eb77. Legacy containers now reuse the existing trusted 
selected-value importer, which visits only keys referenced by each validated V2 
value and preserves primitive physical types. It no longer validates the whole 
shared dictionary for every element.
   
   Added a regression with 4096 V2 objects sharing one dictionary inside a 
legacy ARRAY root, checking every key and value. A local ASAN scaling check 
with 1024/2048/4096/8192 elements improved from 308/985/3521/13695 ms to 
55/111/240/465 ms. The regression has no timing assertions. All 53 focused BE 
tests and the local V1/V2 Flight and ADBC checks passed before pushing.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to