Gabriel39 commented on code in PR #68667:
URL: https://github.com/apache/doris/pull/68667#discussion_r4145786208
##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -28,19 +30,141 @@
#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 {
+// 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];
+ if (!std::isfinite(micros) ||
+ std::abs(micros) >=
static_cast<double>(std::numeric_limits<int64_t>::max())) {
+ return Status::InvalidArgument("Invalid native Arrow Variant
TIMEV2 value");
+ }
+ output.add_time_ntz_micros(std::llround(micros));
Review Comment:
Addressed and retained after the rebase. Native TIMEV2 output now requires a
time within one day; negative and 24-hour-or-longer durations return an
explicit unsupported-value error with the UTF8 setting. Regression tests cover
both boundary values and nested array leaves.
##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -28,19 +30,141 @@
#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 {
+// 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];
+ if (!std::isfinite(micros) ||
+ std::abs(micros) >=
static_cast<double>(std::numeric_limits<int64_t>::max())) {
+ return Status::InvalidArgument("Invalid native Arrow Variant
TIMEV2 value");
+ }
+ output.add_time_ntz_micros(std::llround(micros));
+ } else if (primitive == TYPE_JSONB) {
+ jsonb_to_variant(column.get_data_at(index), output);
+ } 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 objects have textual keys, matching the legacy document
representation.
+ auto key =
+ map.get_keys().is_null_at(element)
+ ? std::string("null")
Review Comment:
Addressed and retained after the rebase. A MAP containing SQL NULL keys now
returns an explicit unsupported-value error with the UTF8 setting, before
converting NULL to an object name. The literal string key "null" remains valid.
Root and nested MAP tests cover the collision and the valid string-key case.
##########
regression-test/suites/arrow_flight_sql_p0/test_flight_native_variant.groovy:
##########
@@ -0,0 +1,128 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+import org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.CallOptions
+import org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.FlightClient
+import org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.Location
+import
org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.sql.FlightSqlClient
+import
org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.memory.RootAllocator
+
+import java.util.concurrent.TimeUnit
+
+suite("test_flight_native_variant", "arrow_flight_sql") {
+ def frontend = jdbc_sql_return_maparray("SHOW FRONTENDS").find {
+ it.IsMaster.toString().equalsIgnoreCase("true") &&
it.Alive.toString().equalsIgnoreCase("true")
+ }
+ assertNotNull(frontend)
+ assertTrue(frontend.ArrowFlightSqlPort.toString().toInteger() > 0)
+ def database = jdbc_sql("SELECT DATABASE()")[0][0]
+ // Match ingestion to the configured Variant representation; legacy
expression roots
+ // cannot be cast to a table Variant with a different subcolumn limit.
+ def variantV2Function = getFeConfig("enable_variant_v2").toBoolean() ?
"parse_to_variant" : ""
+ def table = "${database}.flight_native_variant_input"
+ def allocator = new RootAllocator(Long.MAX_VALUE)
+ def feClient = FlightClient.builder(allocator,
+ Location.forGrpcInsecure(frontend.Host.toString(),
frontend.ArrowFlightSqlPort.toString().toInteger())).build()
+ def client = new FlightSqlClient(feClient)
+ def auth
+ def read = { String query, Closure inspect ->
+ int count = 0
+ client.execute(query, auth).endpoints.each { endpoint ->
+ FlightClient.builder(allocator,
endpoint.locations[0]).build().withCloseable { beClient ->
+ beClient.getStream(endpoint.ticket, auth,
CallOptions.timeout(30, TimeUnit.SECONDS)).withCloseable { stream ->
+ while (stream.next()) {
+ inspect(stream.root)
+ count += stream.root.rowCount
+ }
+ }
+ }
+ }
+ count
+ }
+ def executeSetting = { String query -> read(query, { root -> }) }
+ try {
+ auth =
feClient.authenticateBasicToken(context.config.otherConfigs.get("extArrowFlightSqlUser"),
+
context.config.otherConfigs.get("extArrowFlightSqlPassword")).get()
+ executeSetting("SET enable_sql_cache=false")
+ jdbc_sql("DROP TABLE IF EXISTS ${table}")
+ jdbc_sql("""CREATE TABLE ${table} (id INT, v VARIANT)
+ DUPLICATE KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 3
+ PROPERTIES("replication_num"="1")""")
+ jdbc_sql("""INSERT INTO ${table} VALUES
+ (1, ${variantV2Function}('42')), (2,
${variantV2Function}('"text"')),
+ (3, ${variantV2Function}('{"a":[1,null,"x"]}')), (4, NULL)""")
+ [false, true].each { parallel ->
Review Comment:
Addressed in 46625fe8881. The fixture now uses 60 buckets and 60 rows,
enables the distributed planner, and checks table placement. Parallel reads
require multiple endpoints when multiple result backends are available, reject
duplicate result BE addresses, and compare every stream schema with the
published schema. Local execution of the Groovy suite against Arrow IPC
fixtures covered V1/V2, UTF8/native and single/multiple endpoints; negative
controls verified that one endpoint, duplicate BE addresses and mismatched
schemas are rejected. This fixture validation does not replace a full cluster
run, which remains for CI.
--
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]