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


##########
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:
   Fixed in eac03d1d196. Query schema analysis now evaluates the same 
session/capability gate as the result sink and propagates the decision through 
ARRAY/MAP/STRUCT fields. GetTables and query analysis share the native Variant 
field builder. Added FE coverage for nested schemas, Groovy coverage for 
GetSchema, Prepare, prepared execution and empty results in UTF8/native modes, 
and Python ADBC ExecuteSchema/Prepare coverage. Local validation: 9 FE tests 
passed and the integration additions compile; full Doris cluster execution 
remains for CI.



##########
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:
   Fixed in eac03d1d196. JSONB conversion now receives the enclosing Variant 
depth. Unsupported conversions include UTF8 fallback guidance; other error 
codes are preserved. Added a reproduction using a 100-container JSONB leaf: 28 
enclosing object levels succeed and 29 are rejected. The test failed before the 
fix and passes now, along with all 47 related BE tests.



##########
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:
   Fixed in eac03d1d196. Unsupported scalar roots now go through the row path 
after the outer SQL null mask is applied. Added coverage for an all-null 
Decimal256 root batch, a null-only slice with an unmasked value outside the 
slice, and rejection of a selected unmasked Decimal256 value. The reproduction 
failed before the fix and passes now; all 47 related BE tests pass.



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