github-actions[bot] commented on code in PR #68667:
URL: https://github.com/apache/doris/pull/68667#discussion_r4141016458
##########
fe/fe-core/src/main/java/org/apache/doris/planner/ResultSink.java:
##########
@@ -34,15 +35,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 && context != null
Review Comment:
[P1] Gate native VARIANT on every BE that may serve the result. During a
rolling upgrade, this bit is sent to all result sinks but an older BE ignores
it and still publishes UTF8. With parallel sinks on both versions,
`fetchArrowFlightSchema` rejects the differing schemas before returning
FlightInfo. Even with one upgraded sink, a ticket routed through an older
Flight BE fails when its remote reader tries to write the native struct using
the old plain Variant writer. Negotiate capability or keep the UTF8 format
until all result and proxy BEs support this representation.
##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -157,6 +163,64 @@ 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) {
+ // Legacy documents need JSON conversion; typed scalar roots can keep
their type.
+ // 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);
+ }
+ if (var->is_scalar_variant()) {
+ auto scalar_type = remove_nullable(var->get_root_type());
+ if (scalar_type->get_primitive_type() == TYPE_DECIMAL256) {
+ return Status::NotSupported(
+ "Native Arrow Variant does not support Decimal256
roots");
+ }
+ if
(is_supported_variant_typed_identity(scalar_type->get_primitive_type())) {
+ // Avoid a JSON round trip that would turn exact decimal roots
into doubles.
+ auto typed =
+
ColumnVariantV2::create_typed(make_nullable(var->get_root()), scalar_type);
+ return DataTypeVariantV2SerDe().write_column_to_arrow(
+ *typed, null_map, array_builder, start, end, ctz);
+ }
+ }
+ JsonToVariantOptions parse_options;
+ parse_options.throw_on_invalid_json = true;
+ // Stored keys were already accepted at ingestion; mutable parse
limits must not reject reads.
+ parse_options.max_json_key_length =
std::numeric_limits<uint32_t>::max();
+ parse_options.check_duplicate_json_path = false;
+ JsonStringToVariantEncoder encoder(parse_options);
Review Comment:
[P2] Preserve readable legacy documents past the V2 depth guard. Legacy
ingestion and UTF8 Flight can handle an object with 129 nested levels, but
constructing this V2 encoder imposes a 128-level limit on every legacy row
reparse. Opting into native Flight therefore makes an existing valid row fail
at read time. Support the stored legacy depth in this conversion, or explicitly
constrain and document this native-mode compatibility limit, and add a
129-level regression case.
##########
be/src/format/arrow/arrow_block_convertor.cpp:
##########
@@ -472,6 +472,26 @@ Status ArrowBlockConvertor::init() {
return Status::OK();
}
+Status ArrowFlightArrowBlockConvertor::write_column(const
std::shared_ptr<const IDataType>& type,
+ const DataTypeSerDe& serde,
+ const IColumn& column,
const NullMap* null_map,
+ const
std::shared_ptr<arrow::Field>& field,
+ arrow::ArrayBuilder*
array_builder,
+ int64_t start, int64_t end,
+ const cctz::time_zone&
ctz) const {
+ if (contains_extension_type(field->type())) {
+ std::shared_ptr<arrow::DataType> native_type;
+ RETURN_IF_ERROR(convert_to_arrow_type(type, &native_type, ctz.name(),
true, true));
Review Comment:
[P2] Compare native nested types using the published timezone label. With
`time_zone='+08:00'` and a result `STRUCT<v:VARIANT,t:TIMESTAMPTZ>`, the sink
schema labels the timestamp `+08:00`, but this conversion rebuilds it with
`ctz.name()` (`Fixed/UTC+08:00:00`). `Equals` fails on the timestamp child,
then the plain fallback rejects the native VARIANT storage, so the Flight query
fails. Preserve the sink's timezone spelling here or allow equivalent timestamp
zones in the nested binding check, and cover this mixed struct under an offset
session.
##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -157,6 +163,64 @@ 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) {
+ // Legacy documents need JSON conversion; typed scalar roots can keep
their type.
+ // 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);
+ }
+ if (var->is_scalar_variant()) {
+ auto scalar_type = remove_nullable(var->get_root_type());
+ if (scalar_type->get_primitive_type() == TYPE_DECIMAL256) {
+ return Status::NotSupported(
+ "Native Arrow Variant does not support Decimal256
roots");
+ }
+ if
(is_supported_variant_typed_identity(scalar_type->get_primitive_type())) {
+ // Avoid a JSON round trip that would turn exact decimal roots
into doubles.
+ auto typed =
+
ColumnVariantV2::create_typed(make_nullable(var->get_root()), scalar_type);
+ return DataTypeVariantV2SerDe().write_column_to_arrow(
+ *typed, null_map, array_builder, start, end, ctz);
+ }
+ }
+ JsonToVariantOptions parse_options;
+ parse_options.throw_on_invalid_json = true;
+ // Stored keys were already accepted at ingestion; mutable parse
limits must not reject reads.
+ parse_options.max_json_key_length =
std::numeric_limits<uint32_t>::max();
+ parse_options.check_duplicate_json_path = false;
+ JsonStringToVariantEncoder encoder(parse_options);
+ FormatOptions options;
+ options.timezone = &ctz;
+ NullMap selected_nulls;
+ if (null_map != nullptr) {
+ selected_nulls.assign(null_map->begin() + start, null_map->begin()
+ end);
+ }
+ for (int64_t row = start; row < end; ++row) {
+ std::string json;
+ if (null_map != nullptr && (*null_map)[row]) {
+ json = "null";
+ } else {
+ var->serialize_one_row_to_string(row, &json, options);
+ if (var->get_root_type()->get_primitive_type() == TYPE_STRING
&&
+ !var->get_root()->is_null_at(row)) {
+ // The legacy root string serializer emits raw text, not a
JSON string literal.
+ auto quoted = ColumnString::create();
+ VectorBufferWriter writer(*quoted);
+ writer.write_json_string(json);
+ writer.commit();
+ json = quoted->get_data_at(0).to_string();
+ }
+ }
+ encoder.add_json({json.data(), json.size()});
Review Comment:
[P1] Preserve typed values in the legacy JSON fallback. A VARIANT made by
casting an array containing DECIMAL(20,2) value 9007199254740993.01 takes this
path; the decimal SerDe emits the exact number, but `JsonTreeCollector`
reparses it as `double`, changing the value and dropping its DECIMAL type.
Arrays containing NaN/Infinity and mixed-path batches with a DATE root can
instead emit invalid JSON and fail the Flight query. Encode typed values
directly (or reject unsupported ones explicitly) and cover composite and
mixed-path roots; the scalar decimal test takes the shortcut above.
##########
be/src/core/data_type_serde/data_type_variant_serde.cpp:
##########
@@ -157,6 +163,64 @@ 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) {
+ // Legacy documents need JSON conversion; typed scalar roots can keep
their type.
+ // 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);
+ }
+ if (var->is_scalar_variant()) {
+ auto scalar_type = remove_nullable(var->get_root_type());
+ if (scalar_type->get_primitive_type() == TYPE_DECIMAL256) {
+ return Status::NotSupported(
+ "Native Arrow Variant does not support Decimal256
roots");
+ }
+ if
(is_supported_variant_typed_identity(scalar_type->get_primitive_type())) {
+ // Avoid a JSON round trip that would turn exact decimal roots
into doubles.
+ auto typed =
+
ColumnVariantV2::create_typed(make_nullable(var->get_root()), scalar_type);
+ return DataTypeVariantV2SerDe().write_column_to_arrow(
+ *typed, null_map, array_builder, start, end, ctz);
+ }
+ }
+ JsonToVariantOptions parse_options;
+ parse_options.throw_on_invalid_json = true;
+ // Stored keys were already accepted at ingestion; mutable parse
limits must not reject reads.
+ parse_options.max_json_key_length =
std::numeric_limits<uint32_t>::max();
+ parse_options.check_duplicate_json_path = false;
+ JsonStringToVariantEncoder encoder(parse_options);
+ FormatOptions options;
+ options.timezone = &ctz;
+ NullMap selected_nulls;
+ if (null_map != nullptr) {
+ selected_nulls.assign(null_map->begin() + start, null_map->begin()
+ end);
+ }
+ for (int64_t row = start; row < end; ++row) {
+ std::string json;
+ if (null_map != nullptr && (*null_map)[row]) {
+ json = "null";
+ } else {
+ var->serialize_one_row_to_string(row, &json, options);
+ if (var->get_root_type()->get_primitive_type() == TYPE_STRING
&&
Review Comment:
[P1] Quote all string-family roots in mixed legacy VARIANT columns. When an
IF result merges a VARCHAR-root VARIANT row with an object-path VARIANT row,
`is_scalar_variant()` is false for the combined column but the visible VARCHAR
root still serializes as raw text. This guard handles only TYPE_STRING, so
`hello` fails JSON parsing and `true` becomes a Boolean instead of a string.
Apply the root quoting rule to CHAR/VARCHAR as well, and cover mixed
scalar/object rows with ordinary and JSON-looking text.
--
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]