github-actions[bot] commented on code in PR #68667:
URL: https://github.com/apache/doris/pull/68667#discussion_r4146490552


##########
fe/fe-core/src/main/java/org/apache/doris/planner/ResultSink.java:
##########
@@ -34,15 +36,20 @@ public class ResultSink extends DataSink {
     // Two phase fetch option
     private TFetchOption fetchOption;
 
+    private final boolean nativeVariant;
     private TResultSinkType resultSinkType = TResultSinkType.MYSQL_PROTOCOL;
 
     public ResultSink(PlanNodeId exchNodeId) {
-        this.exchNodeId = exchNodeId;
+        this(exchNodeId, TResultSinkType.MYSQL_PROTOCOL);
     }
 
     public ResultSink(PlanNodeId exchNodeId, TResultSinkType resultSinkType) {
         this.exchNodeId = exchNodeId;
         this.resultSinkType = resultSinkType;
+        ConnectContext context = ConnectContext.get();
+        // The session may change before deferred result fetching; pin the 
format during planning.
+        nativeVariant = resultSinkType == TResultSinkType.ARROW_FLIGHT_PROTOCOL

Review Comment:
   [P1] Teach Flight schema analysis about native VARIANT before enabling this 
sink format. `FlightSqlQuerySchema.field()` still maps every VARIANT (including 
nested leaves) to UTF8 through `FlightSqlSchemaHelper.getArrowType`, whereas 
this sink publishes the `arrow.parquet.variant` struct. With the native session 
setting, preparing `SELECT CAST(42 AS VARIANT)` advertises UTF8, then 
`getFlightInfoPreparedStatement` compares it to the BE struct and rejects the 
query as "Prepared statement schema changed"; `GetSchemaStatement` and ADBC 
ExecuteSchema also return the wrong schema. Use the same capability/session 
decision in schema analysis and cover prepared and schema discovery paths.



##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -19,28 +19,257 @@
 
 #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_to_variant(column.get_data_at(index), output);

Review Comment:
   [P2] Carry the enclosing path depth into JSONB conversion. A legacy JSON row 
with 29 flattened object keys and an array containing 99 nested objects stores 
that array as a JSONB leaf, which is valid under JSONB's 100-container limit. 
This branch calls `jsonb_to_variant` with its default depth of zero, although 
`append_legacy_arrow_document` has already counted the 29 keys; the builder can 
therefore emit a 129-level native Variant that Doris later rejects as corrupt 
instead of returning the documented native-mode depth error. The earlier 
plain-document depth guard does not cover JSONB leaves. Pass `depth` as 
`initial_depth` and cover this combined-depth boundary.



##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -157,6 +386,100 @@ Status DataTypeVariantSerDe::write_column_to_arrow(const 
IColumn& column, const
                                                    int64_t start, int64_t end,
                                                    const cctz::time_zone& ctz) 
const {
     const auto* var = check_and_get_column<ColumnVariant>(column);
+    if (array_builder->type()->id() == arrow::Type::STRUCT) {
+        // Keep legacy scalar and document leaves in their original types.
+        // The outer null map must remain SQL NULL on the wire.
+        if (start < 0 || end < start || end > column.size() ||
+            (null_map != nullptr && end > null_map->size())) {
+            return Status::InvalidArgument("Invalid Variant Arrow row range 
[{}, {})", start, end);
+        }
+        // A legacy null root renders as {}, not Variant null, even in a 
scalar-only batch.
+        if (var->is_scalar_variant() && !var->get_root()->has_null(start, 
end)) {
+            auto scalar_type = remove_nullable(var->get_root_type());
+            if (scalar_type->get_primitive_type() == TYPE_DECIMAL256) {

Review Comment:
   [P2] Check the outer SQL null map before rejecting Decimal256 roots. With 
legacy VARIANT, casting a nullable DECIMAL(76,2) column keeps a nonnull default 
decimal payload under the result's outer null map on SQL NULL rows. A native 
Flight query selecting only those rows enters this scalar shortcut and returns 
`NotSupported` here before examining `null_map`, even though the correct output 
is null structs; UTF8 already skips masked rows. Reject Decimal256 only when a 
selected unmasked root is visible, and cover an all-null or sliced nullable 
Decimal256 cast.



-- 
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