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]