Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions be/src/agent/heartbeat_server.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ void HeartbeatServer::heartbeat(THeartbeatResult& heartbeat_result,
heartbeat_result.backend_info.__set_be_rpc_port(-1);
heartbeat_result.backend_info.__set_brpc_port(config::brpc_port);
heartbeat_result.backend_info.__set_arrow_flight_sql_port(config::arrow_flight_sql_port);
heartbeat_result.backend_info.__set_arrow_flight_native_variant_supported(true);
heartbeat_result.backend_info.__set_version(get_short_version());
heartbeat_result.backend_info.__set_be_start_time(_be_epoch);
heartbeat_result.backend_info.__set_be_node_role(config::be_node_role);
Expand Down
5 changes: 3 additions & 2 deletions be/src/core/column/column_variant.h
Original file line number Diff line number Diff line change
Expand Up @@ -366,6 +366,9 @@ class ColumnVariant final : public COWHelper<IColumn, ColumnVariant> {
// Only single scalar root column
bool is_scalar_variant() const;

// Output adapters must use the same root/document precedence as legacy JSON serialization.
bool is_visible_root_value(size_t nrow) const;

ColumnPtr get_root() const { return subcolumns.get_root()->data.get_finalized_column_ptr(); }

bool has_subcolumn(const PathInData& key) const;
Expand Down Expand Up @@ -680,8 +683,6 @@ class ColumnVariant final : public COWHelper<IColumn, ColumnVariant> {
size_t start, size_t length);

bool try_add_new_subcolumn(const PathInData& path);

bool is_visible_root_value(size_t nrow) const;
};

} // namespace doris
338 changes: 338 additions & 0 deletions be/src/core/data_type_serde/data_type_variant_serde.cpp

Large diffs are not rendered by default.

80 changes: 73 additions & 7 deletions be/src/core/data_type_serde/data_type_variant_v2_serde.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
#include <algorithm>
#include <cstring>
#include <limits>
#include <optional>
#include <orc/Vector.hh>
#include <span>
#include <utility>
Expand All @@ -42,6 +43,7 @@
#include "core/data_type/data_type_nullable.h"
#include "core/data_type/data_type_string.h"
#include "core/data_type_serde/data_type_string_serde.h"
#include "core/data_type_serde/variant_arrow_utils.h"
#include "core/types.h"
#include "core/value/jsonb_value.h"
#include "core/value/variant/variant_batch_builder.h"
Expand All @@ -51,6 +53,49 @@
#include "util/mysql_row_buffer.h"

namespace doris {

// ColumnVariantV2 already validates its encoded dictionaries and value structure. Only walk
// the selected row here: validating every unused dictionary key per row is quadratic for
// shared dictionaries, including when a nested ARRAY invokes this writer on one-row slices.
Status append_flight_variant_value(VariantRef value, VariantBatchBuilder::Row& output,
size_t depth) {
const auto basic_type = value.basic_type();
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 (value.value_size() != value.value.size) {
throw Exception(ErrorCode::CORRUPTION,
"Native Arrow Variant contains trailing value bytes");
}
if (basic_type == VariantBasicType::OBJECT) {
auto object = output.start_object();
auto fields = value.object_view();
for (uint32_t i = 0; i < fields.size(); ++i) {
uint32_t field_id;
auto child = fields.value_at(i, &field_id);
object.add_key(value.metadata.key_at(field_id));
RETURN_IF_ERROR(append_flight_variant_value(child, output, depth + 1));
}
object.finish();
} else if (basic_type == VariantBasicType::ARRAY) {
auto array = output.start_array();
for (uint32_t i = 0; i < value.num_elements(); ++i) {
RETURN_IF_ERROR(append_flight_variant_value(value.array_at(i), output, depth + 1));
}
array.finish();
} else {
// Primitives never reference dictionary keys. Reuse physical import to retain widths,
// decimal scales and non-JSON types; canonical equality encoding normalizes those away.
static constexpr char empty_metadata[] = {0x11, 0, 0};
value.metadata = {empty_metadata, sizeof(empty_metadata)};
output.add_value(value);
}
return Status::OK();
}

namespace {

using MetaIdsColumn = ColumnVector<TYPE_UINT32>;
Expand Down Expand Up @@ -625,13 +670,14 @@ Status write_paimon_variant(const IColumn& column, const NullMap* null_map,
return status;
}

Status write_iceberg_variant(const IColumn& column, const NullMap* null_map,
arrow::ArrayBuilder* array_builder, int64_t start, int64_t end) {
Status write_parquet_variant_arrow(const IColumn& column, const NullMap* null_map,
arrow::ArrayBuilder* array_builder, int64_t start, int64_t end,
bool compact_metadata) {
if (start < 0 || end < start) {
return Status::InvalidArgument("Invalid Iceberg Variant row range [{}, {})", start, end);
return Status::InvalidArgument("Invalid Variant Arrow row range [{}, {})", start, end);
}
if (array_builder->type()->id() != arrow::Type::STRUCT) {
return Status::InvalidArgument("Iceberg Variant writer requires a struct builder, got {}",
return Status::InvalidArgument("Variant Arrow writer requires a struct builder, got {}",
array_builder->type()->ToString());
}
auto& builder = assert_cast<arrow::StructBuilder&>(*array_builder);
Expand All @@ -640,7 +686,7 @@ Status write_iceberg_variant(const IColumn& column, const NullMap* null_map,
type->field(1)->name() != "value" || type->field(0)->type()->id() != arrow::Type::BINARY ||
type->field(1)->type()->id() != arrow::Type::BINARY) {
return Status::InvalidArgument(
"Iceberg Variant writer requires struct<metadata: binary, value: binary>, got {}",
"Variant Arrow writer requires struct<metadata: binary, value: binary>, got {}",
type->ToString());
}
auto& metadata_builder = assert_cast<arrow::BinaryBuilder&>(*builder.field_builder(0));
Expand All @@ -657,10 +703,27 @@ Status write_iceberg_variant(const IColumn& column, const NullMap* null_map,
if (!status.ok()) {
return;
}
std::optional<VariantBatchBuilder> compacted;
if (compact_metadata) {
const auto keys = value.metadata.dict_size();
// An empty dictionary or an object using every key already has row-local metadata.
if (keys != 0 && (value.basic_type() != VariantBasicType::OBJECT ||
value.num_elements() != keys)) {
VariantBatchBuilder encoder;
auto row = encoder.begin_row();
status = append_flight_variant_value(value, row);
if (!status.ok()) {
return;
}
row.finish();
compacted.emplace(encoder.finish_batch());
value = compacted->value_at(0);
}
}
if (value.metadata.size > std::numeric_limits<int32_t>::max() ||
value.value.size > std::numeric_limits<int32_t>::max()) {
status = Status::InvalidArgument(
"Iceberg Variant metadata/value exceeds Arrow binary size limit");
"Variant Arrow metadata/value exceeds Arrow binary size limit");
return;
}
status = checkArrowStatus(builder.Append(), column, builder);
Expand Down Expand Up @@ -733,6 +796,9 @@ Status DataTypeVariantV2SerDe::write_column_to_arrow(const IColumn& column, cons
options.timezone = &ctz;
const size_t first = checked_row(start);
const size_t last = checked_row(end);
if (array_builder->type()->id() == arrow::Type::STRUCT) {
return write_parquet_variant_arrow(column, null_map, array_builder, start, end, true);
}
if (array_builder->type()->id() == arrow::Type::STRING) {
return write_arrow(column, null_map, assert_cast<arrow::StringBuilder&>(*array_builder),
first, last, options);
Expand All @@ -759,7 +825,7 @@ Status DataTypeVariantV2SerDe::write_column_to_iceberg_arrow(
const std::shared_ptr<const IDataType>&, const IColumn& column, const NullMap* null_map,
const std::shared_ptr<arrow::Field>&, arrow::ArrayBuilder* array_builder, int64_t start,
int64_t end, const cctz::time_zone&) const {
return write_iceberg_variant(column, null_map, array_builder, start, end);
return write_parquet_variant_arrow(column, null_map, array_builder, start, end, false);
}

Status DataTypeVariantV2SerDe::write_column_to_orc(const std::string&, const IColumn& column,
Expand Down
32 changes: 32 additions & 0 deletions be/src/core/data_type_serde/variant_arrow_utils.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
// 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.

#pragma once

#include <cstddef>

#include "common/status.h"
#include "core/value/variant/variant_batch_builder.h"

namespace doris {

// Import only values obtained from validated ColumnVariantV2 storage. Depth includes any
// enclosing legacy containers so all native Flight paths enforce the same nesting limit.
Status append_flight_variant_value(VariantRef value, VariantBatchBuilder::Row& output,
size_t depth = 0);

} // namespace doris
13 changes: 7 additions & 6 deletions be/src/exec/operator/result_sink_operator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -57,9 +57,9 @@ Status ResultSinkLocalState::init(RuntimeState* state, LocalSinkStateInfo& info)
} else {
std::shared_ptr<arrow::Schema> arrow_schema;
if (p._sink_type == TResultSinkType::ARROW_FLIGHT_PROTOCOL) {
RETURN_IF_ERROR(get_arrow_schema_from_expr_ctxs(_output_vexpr_ctxs, &arrow_schema,
state->timezone(),
/*datetime_naive=*/true));
RETURN_IF_ERROR(get_arrow_schema_from_expr_ctxs(
_output_vexpr_ctxs, &arrow_schema, state->timezone(),
/*datetime_naive=*/true, p._native_variant));
}
VLOG_DEBUG << "create sender in INIT with instance id " << fragment_instance_id;
RETURN_IF_ERROR(state->exec_env()->result_mgr()->create_sender(
Expand Down Expand Up @@ -103,6 +103,7 @@ ResultSinkOperatorX::ResultSinkOperatorX(int operator_id, int node_id,
_sink_type(!sink.__isset.type || sink.type == TResultSinkType::MYSQL_PROTOCOL
? TResultSinkType::MYSQL_PROTOCOL
: sink.type),
_native_variant(sink.native_variant),
_result_sink_buffer_size_rows(_sink_type == TResultSinkType::ARROW_FLIGHT_PROTOCOL
? config::arrow_flight_result_sink_buffer_size_rows
: RESULT_SINK_BUFFER_SIZE),
Expand All @@ -123,9 +124,9 @@ Status ResultSinkOperatorX::prepare(RuntimeState* state) {
if (state->query_options().enable_parallel_result_sink) {
std::shared_ptr<arrow::Schema> arrow_schema;
if (_sink_type == TResultSinkType::ARROW_FLIGHT_PROTOCOL) {
RETURN_IF_ERROR(get_arrow_schema_from_expr_ctxs(_output_vexpr_ctxs, &arrow_schema,
state->timezone(),
/*datetime_naive=*/true));
RETURN_IF_ERROR(get_arrow_schema_from_expr_ctxs(
_output_vexpr_ctxs, &arrow_schema, state->timezone(),
/*datetime_naive=*/true, _native_variant));
}
VLOG_DEBUG << "create sender in prepare with query id " << state->query_id();
RETURN_IF_ERROR(state->exec_env()->result_mgr()->create_sender(
Expand Down
1 change: 1 addition & 0 deletions be/src/exec/operator/result_sink_operator.h
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,7 @@ class ResultSinkOperatorX final : public DataSinkOperatorX<ResultSinkLocalState>

Status _second_phase_fetch_data(RuntimeState* state, Block* final_block);
const TResultSinkType::type _sink_type;
const bool _native_variant;
const int _result_sink_buffer_size_rows;
// set file options when sink type is FILE
std::unique_ptr<ResultFileOptions> _file_opts = nullptr;
Expand Down
21 changes: 21 additions & 0 deletions be/src/format/arrow/arrow_block_convertor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -472,6 +472,27 @@ 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));
Comment thread
Gabriel39 marked this conversation as resolved.
// Check the extension identity and its complete nested shape before allowing the
// Variant SerDe to write binary storage. An arbitrary STRUCT is not a Variant binding.
// Timestamp labels may differ for equivalent fixed offsets, including inside containers.
if (is_declared_plain_arrow_binding(type, native_type, field->type())) {
return serde.write_column_to_arrow(column, null_map, array_builder, start, end, ctz);
}
}
return DorisArrowBlockConvertor::write_column(type, serde, column, null_map, field,
array_builder, start, end, ctz);
}

Status ArrowFlightArrowBlockConvertor::convert_to_arrow(const Block& block, arrow::MemoryPool* pool,
std::shared_ptr<arrow::RecordBatch>* result,
size_t start_row, size_t end_row) const {
Expand Down
7 changes: 7 additions & 0 deletions be/src/format/arrow/arrow_block_convertor.h
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,13 @@ class ArrowFlightArrowBlockConvertor final : public DorisArrowBlockConvertor {
Status convert_to_arrow(const Block& block, arrow::MemoryPool* pool,
std::shared_ptr<arrow::RecordBatch>* result, size_t start_row = 0,
size_t end_row = 0) const override;

protected:
Status 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 override;
};

class PythonArrowBlockConvertor final : public DorisArrowBlockConvertor {
Expand Down
Loading
Loading