This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 042429a7cfa branch-4.1: [improvement](be) Optimize Variant V2
ingestion and STRING casts (#66744)
042429a7cfa is described below
commit 042429a7cfaa7bfb611f0482a009e8cb557b8060
Author: lihangyu <[email protected]>
AuthorDate: Tue Aug 18 09:42:40 2026 +0800
branch-4.1: [improvement](be) Optimize Variant V2 ingestion and STRING
casts (#66744)
cherry-pick #66709
---
be/src/core/value/variant/variant_value.cpp | 23 +-
be/src/core/value/variant/variant_value.h | 23 +-
.../cast/variant_v2/cast_variant_to_string.cpp | 9 +-
.../segment/variant/v2/variant_path_builder.cpp | 180 ++++++--
.../segment/variant/v2/variant_path_builder.h | 5 +
.../segment/variant/v2/variant_shredder.cpp | 145 ++++--
.../storage/segment/variant/v2/variant_shredder.h | 8 +
.../function/cast/cast_variant_v2_from_test.cpp | 62 +++
.../variant/variant_column_writer_reader_test.cpp | 485 ++++++++++++++++++++-
be/test/util/variant/variant_value_test.cpp | 40 ++
10 files changed, 895 insertions(+), 85 deletions(-)
diff --git a/be/src/core/value/variant/variant_value.cpp
b/be/src/core/value/variant/variant_value.cpp
index 1f8101af30c..fac75c70096 100644
--- a/be/src/core/value/variant/variant_value.cpp
+++ b/be/src/core/value/variant/variant_value.cpp
@@ -368,7 +368,14 @@ uint32_t VariantRef::num_elements() const {
return _container_layout(type).count;
}
-uint32_t VariantRef::_object_field_id(const ContainerLayout& layout, uint32_t
index) const {
+VariantRef::ObjectView VariantRef::object_view() const {
+ ContainerLayout layout = _container_layout(VariantBasicType::OBJECT);
+ const uint32_t dictionary_size = layout.count == 0 ? 0 :
metadata.dict_size();
+ return ObjectView(*this, layout, dictionary_size);
+}
+
+uint32_t VariantRef::_object_field_id(const ContainerLayout& layout, uint32_t
index,
+ const uint32_t* dictionary_size) const {
if (index >= layout.count) {
throw Exception(ErrorCode::INVALID_ARGUMENT,
"Variant object index {} is out of range [0, {})",
index, layout.count);
@@ -376,10 +383,12 @@ uint32_t VariantRef::_object_field_id(const
ContainerLayout& layout, uint32_t in
const auto field_id = static_cast<uint32_t>(read_unsigned(
value.data + layout.ids_offset + static_cast<size_t>(index) *
layout.id_width,
layout.id_width));
- if (field_id >= metadata.dict_size()) {
+ const uint32_t metadata_dictionary_size =
+ dictionary_size != nullptr ? *dictionary_size :
metadata.dict_size();
+ if (field_id >= metadata_dictionary_size) {
throw Exception(ErrorCode::CORRUPTION,
"Variant object field id {} is outside metadata
dictionary of size {}",
- field_id, metadata.dict_size());
+ field_id, metadata_dictionary_size);
}
return field_id;
}
@@ -427,6 +436,14 @@ VariantRef VariantRef::object_value_at(uint32_t index,
uint32_t* field_id_out) c
return _container_value_at(layout, index, false);
}
+VariantRef VariantRef::ObjectView::value_at(uint32_t index, uint32_t*
field_id_out) const {
+ const uint32_t field_id = _value._object_field_id(_layout, index,
&_dictionary_size);
+ if (field_id_out != nullptr) {
+ *field_id_out = field_id;
+ }
+ return _value._container_value_at(_layout, index, false);
+}
+
VariantRef VariantRef::array_at(uint32_t index) const {
return _container_value_at(_container_layout(VariantBasicType::ARRAY),
index, true);
}
diff --git a/be/src/core/value/variant/variant_value.h
b/be/src/core/value/variant/variant_value.h
index 7dd14b24a6b..c09794fcd8a 100644
--- a/be/src/core/value/variant/variant_value.h
+++ b/be/src/core/value/variant/variant_value.h
@@ -28,6 +28,8 @@
namespace doris {
struct VariantRef {
+ class ObjectView;
+
VariantMetadataRef metadata;
StringRef value;
@@ -52,6 +54,7 @@ struct VariantRef {
std::array<uint8_t, 16> get_uuid() const;
uint32_t num_elements() const;
+ ObjectView object_view() const;
bool object_find(StringRef key, VariantRef* out) const;
bool object_find_by_id(uint32_t field_id, VariantRef* out) const;
VariantRef object_value_at(uint32_t index, uint32_t* field_id_out) const;
@@ -69,11 +72,29 @@ private:
};
ContainerLayout _container_layout(VariantBasicType expected_type) const;
- uint32_t _object_field_id(const ContainerLayout& layout, uint32_t index)
const;
+ uint32_t _object_field_id(const ContainerLayout& layout, uint32_t index,
+ const uint32_t* dictionary_size = nullptr) const;
bool _object_find_by_id(const ContainerLayout& layout, uint32_t field_id,
VariantRef* out) const;
VariantRef _container_value_at(const ContainerLayout& layout, uint32_t
index,
bool require_array_boundary) const;
};
+// Parses and validates an object's physical layout and metadata dictionary
size once, then reuses
+// them while iterating its children. The referenced metadata and value bytes
must outlive the view.
+class VariantRef::ObjectView {
+public:
+ uint32_t size() const { return _layout.count; }
+ VariantRef value_at(uint32_t index, uint32_t* field_id_out = nullptr)
const;
+
+private:
+ friend struct VariantRef;
+ ObjectView(VariantRef value, ContainerLayout layout, uint32_t
dictionary_size)
+ : _value(value), _layout(layout),
_dictionary_size(dictionary_size) {}
+
+ VariantRef _value;
+ ContainerLayout _layout;
+ uint32_t _dictionary_size;
+};
+
} // namespace doris
diff --git a/be/src/exprs/function/cast/variant_v2/cast_variant_to_string.cpp
b/be/src/exprs/function/cast/variant_v2/cast_variant_to_string.cpp
index 56f085c81c1..73ba82537ae 100644
--- a/be/src/exprs/function/cast/variant_v2/cast_variant_to_string.cpp
+++ b/be/src/exprs/function/cast/variant_v2/cast_variant_to_string.cpp
@@ -164,8 +164,15 @@ Status cast_values_to_string(FunctionContext* context,
size_t rows, ForcedNulls
Status cast_typed_variant_to_string(FunctionContext* context, const
ColumnVariantV2& source,
size_t rows, ForcedNulls forced_nulls,
ColumnPtr* output) {
const auto& typed = assert_cast<const
ColumnNullable&>(source.typed_column());
- const DataTypePtr string_type = std::make_shared<DataTypeString>();
+ const PrimitiveType source_primitive =
source.typed_type()->get_primitive_type();
const NullMap& inner_nulls = typed.get_null_map_data();
+ if (is_string_type(source_primitive) && !typed.has_null()) {
+ // The physical payload already has the requested representation. An
inner null is a
+ // Variant null and must still stringify as literal "null".
+ return apply_forced_nulls(typed.get_ptr(), forced_nulls, output);
+ }
+
+ const DataTypePtr string_type = std::make_shared<DataTypeString>();
size_t concrete_rows = 0;
for (size_t row = 0; row < rows; ++row) {
if (inner_nulls[row] == 0 && (forced_nulls.empty() ||
forced_nulls[row] == 0)) {
diff --git a/be/src/storage/segment/variant/v2/variant_path_builder.cpp
b/be/src/storage/segment/variant/v2/variant_path_builder.cpp
index c7542669fe4..ea609d3bdac 100644
--- a/be/src/storage/segment/variant/v2/variant_path_builder.cpp
+++ b/be/src/storage/segment/variant/v2/variant_path_builder.cpp
@@ -74,6 +74,34 @@ enum class ValueKind : uint8_t {
ARRAY,
};
+enum class ScalarPhysicalKind : uint8_t { OTHER, NULL_VALUE, SHORT_STRING,
PRIMITIVE };
+
+struct ScalarPhysical {
+ ScalarPhysicalKind kind;
+ VariantPrimitiveId primitive_id = VariantPrimitiveId::NULL_VALUE;
+};
+
+ScalarPhysical scalar_physical(VariantRef value) {
+ switch (value.basic_type()) {
+ case VariantBasicType::SHORT_STRING:
+ return {.kind = ScalarPhysicalKind::SHORT_STRING};
+ case VariantBasicType::PRIMITIVE: {
+ const VariantPrimitiveId primitive_id = value.primitive_id();
+ return {.kind = primitive_id == VariantPrimitiveId::NULL_VALUE
+ ? ScalarPhysicalKind::NULL_VALUE
+ : ScalarPhysicalKind::PRIMITIVE,
+ .primitive_id = primitive_id};
+ }
+ case VariantBasicType::OBJECT:
+ case VariantBasicType::ARRAY:
+ return {.kind = ScalarPhysicalKind::OTHER};
+ }
+ // Match VariantRef::is_null(): an unknown basic type is not null. The
+ // authoritative slow path reports the corrupt type after preserving
+ // row-validation ordering.
+ return {.kind = ScalarPhysicalKind::OTHER};
+}
+
const DataTypePtr& jsonb_type() {
static const DataTypePtr type = std::make_shared<DataTypeJsonb>();
return type;
@@ -197,37 +225,8 @@ DataTypePtr infer_type(VariantRef value, const
DataTypePtr& reusable_type = null
return type;
}
case ValueKind::INT64: {
- PrimitiveType primitive = TYPE_BIGINT;
- switch (value.primitive_id()) {
- case VariantPrimitiveId::INT8:
- primitive = TYPE_TINYINT;
- break;
- case VariantPrimitiveId::INT16:
- primitive = TYPE_SMALLINT;
- break;
- case VariantPrimitiveId::INT32:
- primitive = TYPE_INT;
- break;
- case VariantPrimitiveId::INT64:
- break;
- default:
- throw Exception(ErrorCode::CORRUPTION, "Invalid Variant integer
primitive id");
- }
- static const std::array<DataTypePtr, 4> types {
- std::make_shared<DataTypeInt8>(),
std::make_shared<DataTypeInt16>(),
- std::make_shared<DataTypeInt32>(),
std::make_shared<DataTypeInt64>()};
- switch (primitive) {
- case TYPE_TINYINT:
- return types[0];
- case TYPE_SMALLINT:
- return types[1];
- case TYPE_INT:
- return types[2];
- case TYPE_BIGINT:
- return types[3];
- default:
- throw Exception(ErrorCode::CORRUPTION, "Invalid Variant integer
type {}", primitive);
- }
+ static const DataTypePtr type = std::make_shared<DataTypeInt64>();
+ return type;
}
case ValueKind::LARGEINT: {
static const DataTypePtr type = std::make_shared<DataTypeInt128>();
@@ -694,6 +693,72 @@ void append_timestamp(VariantRef value, PrimitiveType
target_type, IColumn* targ
assert_cast<ColumnTimeStampTz&>(*target).insert_value(converted);
}
+template <typename Value>
+bool stable_scalar_matches_type(const Value& value, const ScalarPhysical&
physical,
+ const DataTypePtr& target_type) {
+ const PrimitiveType target_primitive = target_type->get_primitive_type();
+ if (target_primitive == TYPE_JSONB) {
+ return physical.kind == ScalarPhysicalKind::SHORT_STRING ||
+ physical.kind == ScalarPhysicalKind::PRIMITIVE;
+ }
+ if (physical.kind == ScalarPhysicalKind::SHORT_STRING) {
+ return target_primitive == TYPE_STRING;
+ }
+ if (physical.kind != ScalarPhysicalKind::PRIMITIVE) {
+ return false;
+ }
+
+ switch (physical.primitive_id) {
+ case VariantPrimitiveId::TRUE_VALUE:
+ case VariantPrimitiveId::FALSE_VALUE:
+ return target_primitive == TYPE_BOOLEAN;
+ case VariantPrimitiveId::INT8:
+ case VariantPrimitiveId::INT16:
+ case VariantPrimitiveId::INT32:
+ case VariantPrimitiveId::INT64:
+ // Encoded integer widths are a physical detail. Integer paths use
BIGINT from their first
+ // value, avoiding repeated inference and column rewrites when later
values use a wider or
+ // narrower physical width. LARGEINT remains valid after a real
128-bit promotion.
+ return target_primitive == TYPE_BIGINT || target_primitive ==
TYPE_LARGEINT;
+ case VariantPrimitiveId::FLOAT:
+ return target_primitive == TYPE_FLOAT;
+ case VariantPrimitiveId::DOUBLE:
+ return target_primitive == TYPE_DOUBLE;
+ case VariantPrimitiveId::DECIMAL4:
+ case VariantPrimitiveId::DECIMAL8:
+ case VariantPrimitiveId::DECIMAL16: {
+ const VariantDecimal decimal = value.get_decimal();
+ if (physical.primitive_id == VariantPrimitiveId::DECIMAL16 &&
decimal.scale == 0) {
+ return target_primitive == TYPE_LARGEINT;
+ }
+ if (target_primitive != TYPE_DECIMAL128I ||
target_type->get_precision() != 38 ||
+ target_type->get_scale() != decimal.scale) {
+ return false;
+ }
+ __int128 converted = 0;
+ return try_rescale_decimal_value(value, target_type, &converted);
+ }
+ case VariantPrimitiveId::DATE:
+ return target_primitive == TYPE_DATEV2 &&
date_fits_doris_range(value.get_date());
+ case VariantPrimitiveId::TIMESTAMP_MICROS:
+ return target_primitive == TYPE_TIMESTAMPTZ &&
+ timestamp_fits_doris_range(value.get_timestamp_micros());
+ case VariantPrimitiveId::TIMESTAMP_NTZ_MICROS:
+ return target_primitive == TYPE_DATETIMEV2 &&
+ timestamp_fits_doris_range(value.get_timestamp_ntz_micros());
+ case VariantPrimitiveId::STRING:
+ return target_primitive == TYPE_STRING;
+ case VariantPrimitiveId::NULL_VALUE:
+ case VariantPrimitiveId::BINARY:
+ case VariantPrimitiveId::TIME_NTZ_MICROS:
+ case VariantPrimitiveId::TIMESTAMP_NANOS:
+ case VariantPrimitiveId::TIMESTAMP_NTZ_NANOS:
+ case VariantPrimitiveId::UUID:
+ return false;
+ }
+ throw Exception(ErrorCode::CORRUPTION, "Unknown Variant primitive id");
+}
+
void append_value(VariantRef value, const DataTypePtr& target_type, IColumn*
target);
void append_array(VariantRef value, const DataTypePtr& target_type, IColumn*
target) {
@@ -934,6 +999,7 @@ struct VariantPathBuilder::Impl {
type = remove_nullable(initial_type);
nullable_type = make_nullable(type);
column = nullable_type->create_column();
+ binary_serde.reset();
return Status::OK();
}
@@ -991,6 +1057,7 @@ struct VariantPathBuilder::Impl {
column = IColumn::mutate(std::move(promoted));
type = std::move(target);
nullable_type = make_nullable(type);
+ binary_serde.reset();
#ifdef BE_TEST
++promotions;
#endif
@@ -1000,12 +1067,16 @@ struct VariantPathBuilder::Impl {
PathInData path;
DataTypePtr type;
DataTypePtr nullable_type;
+ DataTypeSerDeSPtr binary_serde;
MutableColumnPtr column;
DorisVector<uint32_t> rowids;
size_t logical_rows = 0;
#ifdef BE_TEST
size_t promotions = 0;
#endif
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+ size_t stable_scalar_appends = 0;
+#endif
};
VariantPathBuilder::VariantPathBuilder(PathInData path, size_t prefix_rows)
@@ -1016,7 +1087,8 @@ VariantPathBuilder&
VariantPathBuilder::operator=(VariantPathBuilder&&) noexcept
Status VariantPathBuilder::append(VariantRef value, size_t row) {
try {
- if (value.is_null()) {
+ const ScalarPhysical physical = scalar_physical(value);
+ if (physical.kind == ScalarPhysicalKind::NULL_VALUE) {
return Status::InvalidArgument("Variant path builder {} must not
append JSON null",
_impl->path.get_path());
}
@@ -1028,16 +1100,26 @@ Status VariantPathBuilder::append(VariantRef value,
size_t row) {
return Status::InvalidArgument("Variant path builder {} row {}
exceeds uint32 limit",
_impl->path.get_path(), row);
}
- RETURN_IF_ERROR(complete_rows(row));
-
- if (!_impl->column) {
- RETURN_IF_ERROR(_impl->initialize(infer_type(value)));
- } else if (_impl->type->get_primitive_type() != TYPE_JSONB) {
- DataTypePtr inferred_type = infer_type(value, _impl->type);
- DataTypePtr common_type = path_least_common_type(_impl->type,
inferred_type);
- RETURN_IF_ERROR(_impl->promote(common_type, false));
+ _impl->logical_rows = row;
+
+ const bool stable_scalar =
+ _impl->column && stable_scalar_matches_type(value, physical,
_impl->type);
+ if (!stable_scalar) {
+ if (!_impl->column) {
+ RETURN_IF_ERROR(_impl->initialize(infer_type(value)));
+ } else if (_impl->type->get_primitive_type() != TYPE_JSONB) {
+ DataTypePtr inferred_type = infer_type(value, _impl->type);
+ if (_impl->type.get() != inferred_type.get() &&
+ !_impl->type->equals(*inferred_type)) {
+ DataTypePtr common_type =
path_least_common_type(_impl->type, inferred_type);
+ if (_impl->type.get() != common_type.get() &&
+ !_impl->type->equals(*common_type)) {
+ RETURN_IF_ERROR(_impl->promote(common_type, false));
+ }
+ }
+ }
}
- const bool is_array = value.basic_type() == VariantBasicType::ARRAY;
+ const bool is_array = !stable_scalar && value_kind(value) ==
ValueKind::ARRAY;
if (is_array && !value_is_representable(value, _impl->type)) {
RETURN_IF_ERROR(_impl->promote(jsonb_type(), false));
}
@@ -1057,6 +1139,9 @@ Status VariantPathBuilder::append(VariantRef value,
size_t row) {
}
_impl->rowids.push_back(static_cast<uint32_t>(row));
_impl->logical_rows = row + 1;
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+ _impl->stable_scalar_appends += stable_scalar;
+#endif
return Status::OK();
} catch (const Exception& exception) {
return exception.to_status();
@@ -1110,6 +1195,12 @@ size_t VariantPathBuilder::promotion_count() const {
}
#endif
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+size_t VariantPathBuilder::stable_scalar_append_count() const {
+ return _impl->stable_scalar_appends;
+}
+#endif
+
size_t VariantPathBuilder::byte_size() const {
return sizeof(Impl) + path_allocated_bytes(_impl->path) +
_impl->rowids.capacity() * sizeof(uint32_t) +
@@ -1181,8 +1272,11 @@ Status VariantPathBuilder::write_sparse_cell(size_t
value_index, ColumnString::C
_impl->path.get_path());
}
try {
-
_impl->type->get_serde(2)->write_one_cell_to_binary(nullable.get_nested_column(),
*chars,
- value_index);
+ if (!_impl->binary_serde) {
+ _impl->binary_serde = _impl->type->get_serde(2);
+ }
+
_impl->binary_serde->write_one_cell_to_binary(nullable.get_nested_column(),
*chars,
+ value_index);
return Status::OK();
} catch (const Exception& exception) {
return exception.to_status();
diff --git a/be/src/storage/segment/variant/v2/variant_path_builder.h
b/be/src/storage/segment/variant/v2/variant_path_builder.h
index 707aa3f762e..7b467a17ae7 100644
--- a/be/src/storage/segment/variant/v2/variant_path_builder.h
+++ b/be/src/storage/segment/variant/v2/variant_path_builder.h
@@ -74,6 +74,11 @@ public:
#ifdef BE_TEST
size_t rows() const;
size_t promotion_count() const;
+#endif
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+ size_t stable_scalar_append_count() const;
+#endif
+#ifdef BE_TEST
bool is_null_at(size_t row) const;
Status materialize(ColumnPtr* result) const;
#endif
diff --git a/be/src/storage/segment/variant/v2/variant_shredder.cpp
b/be/src/storage/segment/variant/v2/variant_shredder.cpp
index d61b51aed96..f197d96d30a 100644
--- a/be/src/storage/segment/variant/v2/variant_shredder.cpp
+++ b/be/src/storage/segment/variant/v2/variant_shredder.cpp
@@ -64,6 +64,8 @@ struct VariantShredder::Impl {
using PathIndex = uint32_t;
using ParentFieldKey = uint64_t;
using ChildPathCache = doris::flat_hash_map<ParentFieldKey, PathIndex>;
+ static constexpr PathIndex UNRESOLVED_PATH =
std::numeric_limits<PathIndex>::max();
+ static constexpr size_t MAX_BINARY_CELLS_PER_CHUNK = 1U << 20;
// Metadata bytes belong to the input ReadView. This cache never escapes
one append call, so
// it can borrow the dictionary and retain only parent+field transitions
observed in that
@@ -197,7 +199,7 @@ struct VariantShredder::Impl {
if (const auto found = path_indices.find(child); found !=
path_indices.end()) {
child_index = found->second;
} else {
- if (paths.size() > std::numeric_limits<PathIndex>::max()) {
+ if (paths.size() >= UNRESOLVED_PATH) {
throw Exception(ErrorCode::INVALID_ARGUMENT,
"Variant path count exceeds uint32 limit");
}
@@ -218,10 +220,10 @@ struct VariantShredder::Impl {
if (value.basic_type() != VariantBasicType::OBJECT) {
return append_leaf(value, path_index, row);
}
- const uint32_t children = value.num_elements();
- for (uint32_t index = 0; index < children; ++index) {
+ const VariantRef::ObjectView object = value.object_view();
+ for (uint32_t index = 0; index < object.size(); ++index) {
uint32_t field = 0;
- VariantRef child = value.object_value_at(index, &field);
+ VariantRef child = object.value_at(index, &field);
const PathIndex child_path = resolve_child_path(metadata_cache,
path_index, field);
RETURN_IF_ERROR(validate_doc_path(child_path));
RETURN_IF_ERROR(visit(child, metadata_cache, child_path, row));
@@ -384,16 +386,22 @@ struct VariantShredder::Impl {
template <typename BinaryPlan>
Status append_binary_rows(const DorisVector<BinaryPlan>& binary_plan,
const DorisVector<ColumnMap*>& maps) const {
- struct BinaryCell {
- size_t plan_index = 0;
- size_t value_index = 0;
- };
-
- // Build a compact row index in two passes. Each path contributes only
its present values,
- // and paths are visited in publication order so cells within one row
preserve path order.
+ DorisVector<size_t> bucket_cells(maps.size(), 0);
+ DorisVector<size_t> bucket_key_bytes(maps.size(), 0);
+ DorisVector<size_t> bucket_value_bytes(maps.size(), 0);
+ // Build a compact row index in two passes. Each path contributes only
its
+ // present values, and paths are visited in publication order so cells
+ // within one row preserve path order.
DorisVector<size_t> row_offsets(rows + 1, 0);
for (const BinaryPlan& plan : binary_plan) {
- for (uint32_t row : plan.builder->rowids()) {
+ DORIS_CHECK_LT(plan.bucket, maps.size());
+ const std::span<const uint32_t> rowids = plan.builder->rowids();
+ bucket_cells[plan.bucket] += rowids.size();
+ bucket_key_bytes[plan.bucket] += rowids.size() * plan.path->size();
+ const ColumnPtr column = plan.builder->column();
+ DORIS_CHECK(column);
+ bucket_value_bytes[plan.bucket] += column->byte_size();
+ for (uint32_t row : rowids) {
if (row >= rows) {
return Status::InternalError("Variant path {} row {}
exceeds {} rows",
*plan.path, row, rows);
@@ -401,32 +409,88 @@ struct VariantShredder::Impl {
++row_offsets[row + 1];
}
}
+ for (size_t bucket = 0; bucket < maps.size(); ++bucket) {
+ auto& keys = assert_cast<ColumnString&>(maps[bucket]->get_keys());
+ auto& values =
assert_cast<ColumnString&>(maps[bucket]->get_values());
+ maps[bucket]->get_offsets().reserve(rows);
+ keys.reserve(bucket_cells[bucket]);
+ keys.get_chars().reserve(bucket_key_bytes[bucket]);
+ values.reserve(bucket_cells[bucket]);
+ values.get_chars().reserve(bucket_value_bytes[bucket]);
+ }
+ if (binary_plan.size() > std::numeric_limits<uint32_t>::max()) {
+ return Status::InternalError("Variant binary path count {} exceeds
uint32 limit",
+ binary_plan.size());
+ }
std::partial_sum(row_offsets.begin(), row_offsets.end(),
row_offsets.begin());
- DorisVector<BinaryCell> cells(row_offsets.back());
- DorisVector<size_t> next_cell = row_offsets;
- for (size_t plan_index = 0; plan_index < binary_plan.size();
++plan_index) {
- const auto rowids = binary_plan[plan_index].builder->rowids();
- for (size_t value_index = 0; value_index < rowids.size();
++value_index) {
- cells[next_cell[rowids[value_index]]++] = {.plan_index =
plan_index,
- .value_index =
value_index};
+
+ // Transpose path-major builders into row-major maps in bounded
chunks. A
+ // chunk always ends at a row boundary, so the path-sorted plan order
+ // remains the canonical key order within each row and bucket. The
per-plan
+ // cursor also recovers value_index without storing it in every cell.
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+ const size_t max_binary_cells_per_chunk = binary_cells_per_chunk;
+#else
+ constexpr size_t max_binary_cells_per_chunk =
MAX_BINARY_CELLS_PER_CHUNK;
+#endif
+ DorisVector<size_t> value_indices(binary_plan.size(), 0);
+ DorisVector<uint32_t> cells;
+ DorisVector<size_t> next_cell;
+ size_t row_begin = 0;
+ while (row_begin < rows) {
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+ ++binary_chunk_count;
+#endif
+ const size_t first_cell = row_offsets[row_begin];
+ size_t row_end = rows;
+ if (row_offsets.back() - first_cell > max_binary_cells_per_chunk) {
+ const auto first_too_large =
+ std::upper_bound(row_offsets.begin() + row_begin + 1,
row_offsets.end(),
+ first_cell +
max_binary_cells_per_chunk);
+ row_end = static_cast<size_t>(first_too_large -
row_offsets.begin() - 1);
+ // One exceptionally wide row may exceed the bound, but must
stay
+ // intact.
+ row_end = std::max(row_end, row_begin + 1);
}
- }
- for (size_t row = 0; row < rows; ++row) {
- for (size_t cell_index = row_offsets[row]; cell_index <
row_offsets[row + 1];
- ++cell_index) {
- const BinaryCell& cell = cells[cell_index];
- const BinaryPlan& plan = binary_plan[cell.plan_index];
- auto& keys =
assert_cast<ColumnString&>(maps[plan.bucket]->get_keys());
- auto& values =
assert_cast<ColumnString&>(maps[plan.bucket]->get_values());
- keys.insert_data(plan.path->data(), plan.path->size());
- RETURN_IF_ERROR(
- plan.builder->write_sparse_cell(cell.value_index,
&values.get_chars()));
- values.get_offsets().push_back(values.get_chars().size());
+ cells.resize(row_offsets[row_end] - first_cell);
+ next_cell.resize(row_end - row_begin);
+ for (size_t row = row_begin; row < row_end; ++row) {
+ next_cell[row - row_begin] = row_offsets[row] - first_cell;
+ }
+ for (size_t plan_index = 0; plan_index < binary_plan.size();
++plan_index) {
+ const std::span<const uint32_t> rowids =
binary_plan[plan_index].builder->rowids();
+ size_t value_index = value_indices[plan_index];
+ DORIS_CHECK(value_index == rowids.size() ||
rowids[value_index] >= row_begin);
+ while (value_index < rowids.size() && rowids[value_index] <
row_end) {
+ const uint32_t row = rowids[value_index++];
+ cells[next_cell[row - row_begin]++] =
static_cast<uint32_t>(plan_index);
+ }
}
- for (ColumnMap* map : maps) {
- map->get_offsets().push_back(map->get_keys().size());
+
+ for (size_t row = row_begin; row < row_end; ++row) {
+ DORIS_CHECK_EQ(next_cell[row - row_begin], row_offsets[row +
1] - first_cell);
+ for (size_t cell_index = row_offsets[row] - first_cell;
+ cell_index < row_offsets[row + 1] - first_cell;
++cell_index) {
+ const uint32_t plan_index = cells[cell_index];
+ const BinaryPlan& plan = binary_plan[plan_index];
+ const size_t value_index = value_indices[plan_index]++;
+ auto& keys =
assert_cast<ColumnString&>(maps[plan.bucket]->get_keys());
+ auto& values =
assert_cast<ColumnString&>(maps[plan.bucket]->get_values());
+ keys.insert_data(plan.path->data(), plan.path->size());
+ RETURN_IF_ERROR(
+ plan.builder->write_sparse_cell(value_index,
&values.get_chars()));
+ values.get_offsets().push_back(values.get_chars().size());
+ }
+ for (ColumnMap* map : maps) {
+ map->get_offsets().push_back(map->get_keys().size());
+ }
}
+ row_begin = row_end;
+ }
+ for (size_t plan_index = 0; plan_index < binary_plan.size();
++plan_index) {
+ DORIS_CHECK_EQ(value_indices[plan_index],
+ binary_plan[plan_index].builder->rowids().size());
}
return Status::OK();
}
@@ -587,6 +651,10 @@ struct VariantShredder::Impl {
DorisVector<PathState> paths;
ColumnString::MutablePtr root_values = ColumnString::create();
JsonbWriter root_writer;
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+ size_t binary_cells_per_chunk = MAX_BINARY_CELLS_PER_CHUNK;
+ mutable size_t binary_chunk_count = 0;
+#endif
};
VariantShredder::VariantShredder(VariantShredderOptions options)
@@ -700,4 +768,15 @@ size_t VariantShredder::byte_size() const {
return size;
}
+#if defined(BE_TEST) && !defined(BE_BENCHMARK)
+size_t VariantShredder::TestAccess::binary_chunk_count(const VariantShredder&
shredder) {
+ return shredder._impl->binary_chunk_count;
+}
+
+void VariantShredder::TestAccess::set_binary_cells_per_chunk(VariantShredder&
shredder,
+ size_t cells) {
+ DORIS_CHECK(cells > 0);
+ shredder._impl->binary_cells_per_chunk = cells;
+}
+#endif
} // namespace doris::segment_v2
diff --git a/be/src/storage/segment/variant/v2/variant_shredder.h
b/be/src/storage/segment/variant/v2/variant_shredder.h
index ce9535241ef..91b8c8b4944 100644
--- a/be/src/storage/segment/variant/v2/variant_shredder.h
+++ b/be/src/storage/segment/variant/v2/variant_shredder.h
@@ -87,6 +87,14 @@ public:
Status finish(VariantShreddedColumns* output);
size_t byte_size() const;
+#ifdef BE_TEST
+ struct TestAccess {
+#if !defined(BE_BENCHMARK)
+ static size_t binary_chunk_count(const VariantShredder& shredder);
+ static void set_binary_cells_per_chunk(VariantShredder& shredder,
size_t cells);
+#endif
+ };
+#endif
private:
struct Impl;
std::unique_ptr<Impl> _impl;
diff --git a/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp
b/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp
index cc594acfc9c..a1e20984965 100644
--- a/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp
+++ b/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp
@@ -121,6 +121,16 @@ ColumnVariantV2::MutablePtr typed_ints() {
std::make_shared<DataTypeInt32>());
}
+ColumnVariantV2::MutablePtr typed_strings() {
+ auto values = ColumnString::create();
+ values->insert_data("alice", 5);
+ values->insert_data("bob", 3);
+ values->insert_data("carol", 5);
+ return ColumnVariantV2::create_typed(
+ ColumnNullable::create(std::move(values), ColumnUInt8::create(3,
0)),
+ std::make_shared<DataTypeString>());
+}
+
const ColumnNullable& nullable_result(const ColumnPtr& column) {
return assert_cast<const ColumnNullable&>(*column);
}
@@ -290,6 +300,58 @@ TEST(CastVariantV2FromTest,
TypedInnerNullStringifiesAsLiteralNull) {
EXPECT_FALSE(nullable.has_null());
}
+TEST(CastVariantV2FromTest, TypedStringIdentityReusesPayload) {
+ ColumnPtr source = typed_strings();
+ const auto& typed = assert_cast<const ColumnVariantV2&>(*source);
+ const auto& source_nullable = assert_cast<const
ColumnNullable&>(typed.typed_column());
+
+ CastResult identity = execute_from_variant(source,
std::make_shared<DataTypeString>());
+ ASSERT_TRUE(identity.status.ok()) << identity.status;
+ EXPECT_EQ(identity.column.get(), &typed.typed_column());
+
+ constexpr std::array<NullMap::value_type, 3> FORCED_NULLS {0, 1, 0};
+ CastResult masked =
+ execute_from_variant(source, std::make_shared<DataTypeString>(),
FORCED_NULLS.data());
+ ASSERT_TRUE(masked.status.ok()) << masked.status;
+ const auto& masked_nullable = nullable_result(masked.column);
+ EXPECT_EQ(masked_nullable.get_nested_column_ptr().get(),
+ source_nullable.get_nested_column_ptr().get());
+ EXPECT_EQ(masked_nullable.get_null_map_data(), (NullMap {0, 1, 0}));
+}
+
+TEST(CastVariantV2FromTest, TypedStringPreservesForcedInnerNullSemantics) {
+ auto values = ColumnString::create();
+ values->insert_data("alice", 5);
+ values->insert_default();
+ values->insert_data("carol", 5);
+ auto inner_nulls = ColumnUInt8::create();
+ inner_nulls->insert_value(0);
+ inner_nulls->insert_value(1);
+ inner_nulls->insert_value(0);
+ ColumnPtr source = ColumnVariantV2::create_typed(
+ ColumnNullable::create(std::move(values), std::move(inner_nulls)),
+ std::make_shared<DataTypeString>());
+
+ CastResult visible_null = execute_from_variant(source,
std::make_shared<DataTypeString>());
+ ASSERT_TRUE(visible_null.status.ok()) << visible_null.status;
+ const auto& visible_nullable = nullable_result(visible_null.column);
+ const auto& visible_strings =
+ assert_cast<const
ColumnString&>(visible_nullable.get_nested_column());
+ EXPECT_EQ(visible_strings.get_data_at(1), StringRef("null"));
+ EXPECT_FALSE(visible_nullable.has_null());
+
+ constexpr std::array<NullMap::value_type, 3> FORCED_NULLS {0, 1, 0};
+ CastResult masked =
+ execute_from_variant(source, std::make_shared<DataTypeString>(),
FORCED_NULLS.data());
+ ASSERT_TRUE(masked.status.ok()) << masked.status;
+ const auto& masked_nullable = nullable_result(masked.column);
+ const auto& masked_strings =
+ assert_cast<const
ColumnString&>(masked_nullable.get_nested_column());
+ EXPECT_EQ(masked_strings.get_data_at(0), StringRef("alice"));
+ EXPECT_EQ(masked_strings.get_data_at(2), StringRef("carol"));
+ EXPECT_EQ(masked_nullable.get_null_map_data(), (NullMap {0, 1, 0}));
+}
+
TEST(CastVariantV2FromTest, TypedStringUsesCanonicalTimestampScale) {
DateV2Value<DateTimeV2ValueType> value;
value.unchecked_set_time(2024, 1, 2, 3, 4, 5, 123000);
diff --git a/be/test/storage/variant/variant_column_writer_reader_test.cpp
b/be/test/storage/variant/variant_column_writer_reader_test.cpp
index c19f1eda764..fd9cf208648 100644
--- a/be/test/storage/variant/variant_column_writer_reader_test.cpp
+++ b/be/test/storage/variant/variant_column_writer_reader_test.cpp
@@ -240,7 +240,7 @@ static std::string variant_json_at(const IColumn& column,
size_t row) {
TEST(VariantPathBuilderTest, PromotesValuesAndMaterializesMissingRows) {
VariantBatchBuilder value_builder;
auto integer_row = value_builder.begin_row();
- integer_row.add_int(1);
+ integer_row.add_float(1.0F);
integer_row.finish();
auto double_row = value_builder.begin_row();
double_row.add_double(2.5);
@@ -283,7 +283,7 @@ TEST(VariantPathBuilderTest,
PromotesValuesAndMaterializesMissingRows) {
}
}
-TEST(VariantPathBuilderTest, PreservesIntegerWidthAcrossPromotion) {
+TEST(VariantPathBuilderTest, UsesBigintForEncodedIntegersWithoutPromotion) {
VariantBatchBuilder value_builder;
auto tiny_row = value_builder.begin_row();
tiny_row.add_int(1);
@@ -295,9 +295,98 @@ TEST(VariantPathBuilderTest,
PreservesIntegerWidthAcrossPromotion) {
segment_v2::VariantPathBuilder builder(PathInData("metric"));
ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
- EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(),
TYPE_TINYINT);
+ EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(),
TYPE_BIGINT);
ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
- EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(), TYPE_INT);
+ EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(),
TYPE_BIGINT);
+ EXPECT_EQ(builder.promotion_count(), 0);
+}
+
+TEST(VariantPathBuilderTest,
MatchesLegacyMixedIntegerDoubleInferenceInEitherOrder) {
+ for (const bool reverse : {false, true}) {
+ SCOPED_TRACE(testing::Message() << "reverse=" << reverse);
+ const std::vector<std::string> jsons =
+ reverse ? std::vector<std::string> {R"({"metric":1.5})",
R"({"metric":1})"}
+ : std::vector<std::string> {R"({"metric":1})",
R"({"metric":1.5})"};
+
+ auto legacy = ColumnVariant::create(0, false);
+ auto json = ColumnString::create();
+ for (const std::string& value : jsons) {
+ json->insert_data(value.data(), value.size());
+ }
+ ParseConfig parse_config;
+ parse_config.parse_to = ParseConfig::ParseTo::OnlySubcolumns;
+ variant_util::parse_json_to_variant(*legacy, *json, parse_config);
+ legacy->finalize();
+ auto* legacy_metric = legacy->get_subcolumn(PathInData("metric"));
+ ASSERT_NE(legacy_metric, nullptr);
+
EXPECT_EQ(remove_nullable(legacy_metric->get_least_common_type())->get_primitive_type(),
+ TYPE_JSONB);
+
+ VariantBatchBuilder value_builder;
+ if (reverse) {
+ auto double_row = value_builder.begin_row();
+ double_row.add_double(1.5);
+ double_row.finish();
+ auto integer_row = value_builder.begin_row();
+ integer_row.add_int(1);
+ integer_row.finish();
+ } else {
+ auto integer_row = value_builder.begin_row();
+ integer_row.add_int(1);
+ integer_row.finish();
+ auto double_row = value_builder.begin_row();
+ double_row.add_double(1.5);
+ double_row.finish();
+ }
+ VariantBatchBuilder values = value_builder.finish_batch();
+ segment_v2::VariantPathBuilder builder(PathInData("metric"));
+ ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+ ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+ ASSERT_EQ(remove_nullable(builder.type())->get_primitive_type(),
TYPE_JSONB);
+
+ for (size_t row = 0; row < jsons.size(); ++row) {
+ auto legacy_keys = ColumnString::create();
+ auto legacy_values = ColumnString::create();
+ legacy_metric->serialize_to_binary_column(legacy_keys.get(),
"metric",
+ legacy_values.get(),
row);
+ ASSERT_EQ(legacy_values->size(), 1);
+ ColumnString::Chars v2_cell;
+ ASSERT_TRUE(builder.write_sparse_cell(row, &v2_cell).ok());
+ const StringRef legacy_cell = legacy_values->get_data_at(0);
+ ASSERT_EQ(v2_cell.size(), legacy_cell.size);
+ EXPECT_EQ(std::memcmp(v2_cell.data(), legacy_cell.data,
legacy_cell.size), 0);
+ }
+ }
+}
+
+TEST(VariantPathBuilderTest, CachedBinarySerdeFollowsFloatingPromotion) {
+ VariantBatchBuilder value_builder;
+ auto tiny_row = value_builder.begin_row();
+ tiny_row.add_float(1.0F);
+ tiny_row.finish();
+ auto int_row = value_builder.begin_row();
+ int_row.add_double(2.0);
+ int_row.finish();
+ VariantBatchBuilder values = value_builder.finish_batch();
+
+ segment_v2::VariantPathBuilder builder(PathInData("metric"));
+ ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+ ColumnString::Chars binary;
+ ASSERT_TRUE(builder.write_sparse_cell(0, &binary).ok());
+ ASSERT_FALSE(binary.empty());
+ EXPECT_EQ(static_cast<FieldType>(binary.front()),
FieldType::OLAP_FIELD_TYPE_FLOAT);
+
+ ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+ binary.clear();
+ ASSERT_TRUE(builder.write_sparse_cell(1, &binary).ok());
+ ASSERT_FALSE(binary.empty());
+ EXPECT_EQ(static_cast<FieldType>(binary.front()),
FieldType::OLAP_FIELD_TYPE_DOUBLE);
+
+ ASSERT_TRUE(builder.convert_to(std::make_shared<DataTypeInt64>()).ok());
+ binary.clear();
+ ASSERT_TRUE(builder.write_sparse_cell(1, &binary).ok());
+ ASSERT_FALSE(binary.empty());
+ EXPECT_EQ(static_cast<FieldType>(binary.front()),
FieldType::OLAP_FIELD_TYPE_BIGINT);
}
TEST(VariantPathBuilderTest,
PromotesCompatibleDecimalScalesInEitherOrderAndInsideArrays) {
@@ -346,6 +435,201 @@ TEST(VariantPathBuilderTest,
PromotesCompatibleDecimalScalesInEitherOrderAndInsi
}
}
+TEST(VariantPathBuilderTest,
EqualFullTypeDoesNotPromoteButDifferentDecimalScaleDoes) {
+ VariantBatchBuilder value_builder;
+ for (const uint8_t scale : {2, 2, 4}) {
+ auto row = value_builder.begin_row();
+ row.add_decimal(123, scale, 16);
+ row.finish();
+ }
+ VariantBatchBuilder values = value_builder.finish_batch();
+
+ segment_v2::VariantPathBuilder builder(PathInData("metric"));
+ ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+ ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+ EXPECT_EQ(builder.promotion_count(), 0);
+ ASSERT_TRUE(builder.append(values.value_at(2), 2).ok());
+ EXPECT_EQ(builder.promotion_count(), 1);
+ EXPECT_EQ(builder.stable_scalar_append_count(), 1);
+ EXPECT_EQ(remove_nullable(builder.type())->get_scale(), 4);
+}
+
+TEST(VariantPathBuilderTest, StableScalarFastPathPreservesInferenceBoundaries)
{
+ const auto verify_encoded = []<typename Append>(Append append) {
+ VariantBatchBuilder value_builder;
+ for (size_t row = 0; row < 2; ++row) {
+ auto value_row = value_builder.begin_row();
+ append(value_row);
+ value_row.finish();
+ }
+ VariantBatchBuilder values = value_builder.finish_batch();
+ segment_v2::VariantPathBuilder builder(PathInData("metric"));
+ ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+ ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+ EXPECT_EQ(builder.promotion_count(), 0);
+ EXPECT_EQ(builder.stable_scalar_append_count(), 1);
+ };
+
+ verify_encoded([](auto& row) { row.add_bool(true); });
+ verify_encoded([](auto& row) { row.add_int(1 << 20); });
+ verify_encoded([](auto& row) { row.add_largeint(static_cast<__int128>(1)
<< 80); });
+ verify_encoded([](auto& row) { row.add_float(1.25F); });
+ verify_encoded([](auto& row) { row.add_double(2.5); });
+ verify_encoded([](auto& row) { row.add_string(StringRef("stable")); });
+ const std::string long_string(80, 'x');
+ verify_encoded([&](auto& row) { row.add_string(StringRef(long_string)); });
+ verify_encoded([](auto& row) { row.add_date(0); });
+ verify_encoded([](auto& row) { row.add_timestamp_micros(1, true); });
+ verify_encoded([](auto& row) { row.add_timestamp_micros(1, false); });
+ verify_encoded([](auto& row) { row.add_decimal(12300, 4, 8); });
+
+ VariantBatchBuilder widening_values_builder;
+ for (const int64_t value : std::array<int64_t, 4> {1, 1LL << 20, 1LL <<
40, 1}) {
+ auto row = widening_values_builder.begin_row();
+ row.add_int(value);
+ row.finish();
+ }
+ VariantBatchBuilder widening_values =
widening_values_builder.finish_batch();
+ segment_v2::VariantPathBuilder widening(PathInData("widening"));
+ ASSERT_TRUE(widening.append(widening_values.value_at(0), 0).ok());
+ ASSERT_TRUE(widening.append(widening_values.value_at(1), 1).ok());
+ EXPECT_EQ(widening.stable_scalar_append_count(), 1);
+ EXPECT_EQ(widening.promotion_count(), 0);
+ ASSERT_TRUE(widening.append(widening_values.value_at(2), 2).ok());
+ EXPECT_EQ(widening.stable_scalar_append_count(), 2);
+ ASSERT_TRUE(widening.append(widening_values.value_at(3), 3).ok());
+ EXPECT_EQ(widening.stable_scalar_append_count(), 3);
+ EXPECT_EQ(widening.promotion_count(), 0);
+ EXPECT_EQ(widening.type()->to_string(*widening.column(), 3), "1");
+
+ VariantBatchBuilder conflict_values_builder;
+ auto integer_row = conflict_values_builder.begin_row();
+ integer_row.add_int(1);
+ integer_row.finish();
+ for (size_t row = 0; row < 2; ++row) {
+ auto bool_row = conflict_values_builder.begin_row();
+ bool_row.add_bool(true);
+ bool_row.finish();
+ }
+ VariantBatchBuilder conflict_values =
conflict_values_builder.finish_batch();
+ segment_v2::VariantPathBuilder conflict(PathInData("conflict"));
+ for (size_t row = 0; row < conflict_values.num_rows(); ++row) {
+ ASSERT_TRUE(conflict.append(conflict_values.value_at(row), row).ok());
+ }
+ EXPECT_EQ(remove_nullable(conflict.type())->get_primitive_type(),
TYPE_JSONB);
+ EXPECT_EQ(conflict.stable_scalar_append_count(), 1);
+
+ VariantBatchBuilder binary_values_builder;
+ for (size_t row = 0; row < 2; ++row) {
+ auto binary_row = binary_values_builder.begin_row();
+ binary_row.add_binary(StringRef("binary"));
+ binary_row.finish();
+ }
+ VariantBatchBuilder binary_values = binary_values_builder.finish_batch();
+ segment_v2::VariantPathBuilder binary(PathInData("binary"));
+ ASSERT_TRUE(binary.append(binary_values.value_at(0), 0).ok());
+ ASSERT_TRUE(binary.append(binary_values.value_at(1), 1).ok());
+ EXPECT_EQ(remove_nullable(binary.type())->get_primitive_type(),
TYPE_JSONB);
+ EXPECT_EQ(binary.stable_scalar_append_count(), 1);
+
+ VariantBatchBuilder array_values_builder;
+ for (size_t row = 0; row < 2; ++row) {
+ auto value_row = array_values_builder.begin_row();
+ auto array = value_row.start_array();
+ value_row.add_int(1);
+ array.finish();
+ value_row.finish();
+ }
+ VariantBatchBuilder array_values = array_values_builder.finish_batch();
+ segment_v2::VariantPathBuilder arrays(PathInData("arrays"));
+ ASSERT_TRUE(arrays.append(array_values.value_at(0), 0).ok());
+ ASSERT_TRUE(arrays.append(array_values.value_at(1), 1).ok());
+ EXPECT_EQ(arrays.stable_scalar_append_count(), 0);
+
+ const auto verify_out_of_range_falls_back = []<typename Append>(Append
append) {
+ VariantBatchBuilder value_builder;
+ auto valid_row = value_builder.begin_row();
+ append(valid_row, false);
+ valid_row.finish();
+ auto invalid_row = value_builder.begin_row();
+ append(invalid_row, true);
+ invalid_row.finish();
+ VariantBatchBuilder values = value_builder.finish_batch();
+
+ segment_v2::VariantPathBuilder builder(PathInData("range"));
+ ASSERT_TRUE(builder.append(values.value_at(0), 0).ok());
+ ASSERT_TRUE(builder.append(values.value_at(1), 1).ok());
+ EXPECT_EQ(remove_nullable(builder.type())->get_primitive_type(),
TYPE_JSONB);
+ EXPECT_EQ(builder.stable_scalar_append_count(), 0);
+ };
+ verify_out_of_range_falls_back(
+ [](auto& row, bool invalid) { row.add_date(invalid ? 3'000'000 :
0); });
+ verify_out_of_range_falls_back([](auto& row, bool invalid) {
+ row.add_timestamp_micros(invalid ? 253'402'300'800'000'000LL : 0,
true);
+ });
+}
+
+TEST(VariantPathBuilderTest,
StableScalarGuardRetainsDecimalAndAppendFailureFallbacks) {
+ VariantBatchBuilder value_builder;
+ auto valid_row = value_builder.begin_row();
+ valid_row.add_decimal(1, 1, 16);
+ valid_row.finish();
+ VariantBatchBuilder values = value_builder.finish_batch();
+
+ segment_v2::VariantPathBuilder decimal(PathInData("decimal"));
+ ASSERT_TRUE(decimal.append(values.value_at(0), 0).ok());
+
+ std::array<char, 18> overflow_decimal {};
+ overflow_decimal[0] =
static_cast<char>(static_cast<uint8_t>(VariantPrimitiveId::DECIMAL16)
+ << VARIANT_VALUE_HEADER_SHIFT);
+ overflow_decimal[1] = 1;
+ unsigned __int128 unscaled = VARIANT_DECIMAL16_MAX + 1;
+ for (size_t byte = 0; byte < 16; ++byte) {
+ overflow_decimal[byte + 2] = static_cast<char>(unscaled >> (byte * 8));
+ }
+ ASSERT_TRUE(
+ decimal.append(VariantRef {.metadata = {},
+ .value = {overflow_decimal.data(),
overflow_decimal.size()}},
+ 1)
+ .ok());
+ EXPECT_EQ(remove_nullable(decimal.type())->get_primitive_type(),
TYPE_JSONB);
+ EXPECT_EQ(decimal.promotion_count(), 1);
+ EXPECT_EQ(decimal.stable_scalar_append_count(), 0);
+ EXPECT_EQ(decimal.non_null_rows(), 2);
+
+ VariantBatchBuilder narrow_values_builder;
+ auto narrow_row = narrow_values_builder.begin_row();
+ narrow_row.add_decimal(9999, 1, 4);
+ narrow_row.finish();
+ VariantBatchBuilder narrow_values = narrow_values_builder.finish_batch();
+ segment_v2::VariantPathBuilder
narrow_decimal(PathInData("narrow_decimal"));
+ ASSERT_TRUE(narrow_decimal.append(values.value_at(0), 0).ok());
+
ASSERT_TRUE(narrow_decimal.convert_to(std::make_shared<DataTypeDecimal128>(3,
1)).ok());
+ ASSERT_TRUE(narrow_decimal.append(narrow_values.value_at(0), 1).ok());
+ EXPECT_EQ(remove_nullable(narrow_decimal.type())->get_primitive_type(),
TYPE_JSONB);
+ EXPECT_EQ(narrow_decimal.promotion_count(), 2);
+ EXPECT_EQ(narrow_decimal.stable_scalar_append_count(), 0);
+ EXPECT_EQ(narrow_decimal.non_null_rows(), 2);
+
+ VariantBatchBuilder integer_builder;
+ auto integer_row = integer_builder.begin_row();
+ integer_row.add_int(1);
+ integer_row.finish();
+ VariantBatchBuilder integer_values = integer_builder.finish_batch();
+
+ segment_v2::VariantPathBuilder truncated(PathInData("truncated"));
+ ASSERT_TRUE(truncated.append(integer_values.value_at(0), 0).ok());
+ const char truncated_int8 =
static_cast<char>(static_cast<uint8_t>(VariantPrimitiveId::INT8)
+ <<
VARIANT_VALUE_HEADER_SHIFT);
+ const Status truncated_status =
+ truncated.append(VariantRef {.metadata = {}, .value =
{&truncated_int8, 1}}, 1);
+ EXPECT_FALSE(truncated_status.ok());
+ EXPECT_EQ(remove_nullable(truncated.type())->get_primitive_type(),
TYPE_JSONB);
+ EXPECT_EQ(truncated.promotion_count(), 1);
+ EXPECT_EQ(truncated.stable_scalar_append_count(), 0);
+ EXPECT_EQ(truncated.non_null_rows(), 1);
+}
+
TEST(VariantPathBuilderTest, JsonbFallbackPreservesCanonicalNumericBytes) {
VariantBatchBuilder value_builder;
auto row = value_builder.begin_row();
@@ -637,6 +921,199 @@ TEST(VariantPathBuilderTest,
ShredderReusesCanonicalPathsAcrossAppends) {
EXPECT_EQ(selected.type->to_string(*selected.column, 1), "2");
}
+TEST(VariantShredderTest, ReusesRootAndNestedPathsForSharedMetadata) {
+ DataTypeVariantV2SerDe serde;
+ DataTypeSerDe::FormatOptions format_options;
+ auto values = ColumnVariantV2::create();
+ for (const std::string_view json :
+ {R"({"a":1,"nested":{"x":2,"y":3}})",
R"({"a":4,"nested":{"x":5,"y":6}})"}) {
+ Slice slice(json.data(), json.size());
+ ASSERT_TRUE(serde.deserialize_one_cell_from_json(*values, slice,
format_options).ok());
+ }
+ ASSERT_EQ(values->read_view().metadata_count(), 1);
+
+ segment_v2::VariantShredderOptions options;
+ options.max_subcolumns_count = 0;
+ options.sparse_bucket_count = 1;
+ segment_v2::VariantShredder shredder(std::move(options));
+ ASSERT_TRUE(shredder.append(values->read_view(), 0, values->size()).ok());
+
+ segment_v2::VariantShreddedColumns shredded;
+ ASSERT_TRUE(shredder.finish(&shredded).ok());
+ ASSERT_EQ(shredded.materialized.size(), 3);
+ const auto expect_path = [&](std::string_view path, std::string_view first,
+ std::string_view second) {
+ const auto found = std::ranges::find_if(shredded.materialized,
[&](const auto& column) {
+ return column.path.get_path() == path;
+ });
+ ASSERT_NE(found, shredded.materialized.end()) << path;
+ ASSERT_TRUE(found->column);
+ EXPECT_EQ(found->rowids, (DorisVector<uint32_t> {0, 1}));
+ EXPECT_EQ(found->type->to_string(*found->column, 0), first);
+ EXPECT_EQ(found->type->to_string(*found->column, 1), second);
+ };
+ expect_path("a", "1", "4");
+ expect_path("nested.x", "2", "5");
+ expect_path("nested.y", "3", "6");
+}
+
+static void expect_variant_statistics_equal(const
segment_v2::VariantStatistics& actual,
+ const
segment_v2::VariantStatistics& expected) {
+ EXPECT_EQ(actual.subcolumns_non_null_size,
expected.subcolumns_non_null_size);
+ EXPECT_EQ(actual.sparse_column_non_null_size,
expected.sparse_column_non_null_size);
+ EXPECT_EQ(actual.doc_value_column_non_null_size,
expected.doc_value_column_non_null_size);
+ EXPECT_EQ(actual.has_nested_group, expected.has_nested_group);
+}
+
+static void expect_physical_column_data_equal(const IColumn& actual, const
IColumn& expected) {
+ ASSERT_EQ(actual.get_name(), expected.get_name());
+ ASSERT_EQ(actual.size(), expected.size());
+ if (const auto* actual_string =
check_and_get_column<ColumnString>(actual)) {
+ const auto& expected_string = assert_cast<const
ColumnString&>(expected);
+ EXPECT_EQ(actual_string->get_chars(), expected_string.get_chars());
+ EXPECT_EQ(actual_string->get_offsets(), expected_string.get_offsets());
+ return;
+ }
+ if (const auto* actual_nullable =
check_and_get_column<ColumnNullable>(actual)) {
+ const auto& expected_nullable = assert_cast<const
ColumnNullable&>(expected);
+ EXPECT_EQ(actual_nullable->get_null_map_data(),
expected_nullable.get_null_map_data());
+ expect_physical_column_data_equal(actual_nullable->get_nested_column(),
+
expected_nullable.get_nested_column());
+ return;
+ }
+ if (const auto* actual_map = check_and_get_column<ColumnMap>(actual)) {
+ const auto& expected_map = assert_cast<const ColumnMap&>(expected);
+ EXPECT_EQ(actual_map->get_offsets(), expected_map.get_offsets());
+ expect_physical_column_data_equal(actual_map->get_keys(),
expected_map.get_keys());
+ expect_physical_column_data_equal(actual_map->get_values(),
expected_map.get_values());
+ return;
+ }
+ if (const auto* actual_array = check_and_get_column<ColumnArray>(actual)) {
+ const auto& expected_array = assert_cast<const ColumnArray&>(expected);
+ EXPECT_EQ(actual_array->get_offsets(), expected_array.get_offsets());
+ expect_physical_column_data_equal(actual_array->get_data(),
expected_array.get_data());
+ return;
+ }
+ for (size_t row = 0; row < actual.size(); ++row) {
+ EXPECT_EQ(actual.compare_at(row, row, expected, -1), 0) << "row=" <<
row;
+ }
+}
+
+static void expect_physical_columns_equal(const ColumnPtr& actual, const
ColumnPtr& expected) {
+ ASSERT_TRUE(actual);
+ ASSERT_TRUE(expected);
+ expect_physical_column_data_equal(*actual, *expected);
+}
+
+static void expect_shredded_columns_equal(const
segment_v2::VariantShreddedColumns& actual,
+ const
segment_v2::VariantShreddedColumns& expected) {
+ ASSERT_EQ(actual.num_rows, expected.num_rows);
+ expect_physical_columns_equal(actual.root_jsonb, expected.root_jsonb);
+
+ ASSERT_EQ(actual.materialized.size(), expected.materialized.size());
+ for (size_t index = 0; index < actual.materialized.size(); ++index) {
+ const auto& actual_path = actual.materialized[index];
+ const auto& expected_path = expected.materialized[index];
+ EXPECT_EQ(actual_path.path, expected_path.path) << "index=" << index;
+ ASSERT_TRUE(actual_path.type);
+ ASSERT_TRUE(expected_path.type);
+ EXPECT_TRUE(actual_path.type->equals(*expected_path.type)) << "index="
<< index;
+ EXPECT_EQ(actual_path.rowids, expected_path.rowids) << "index=" <<
index;
+ expect_physical_columns_equal(actual_path.column,
expected_path.column);
+ }
+
+ ASSERT_EQ(actual.binary_buckets.size(), expected.binary_buckets.size());
+ for (size_t bucket = 0; bucket < actual.binary_buckets.size(); ++bucket) {
+ expect_physical_columns_equal(actual.binary_buckets[bucket].column,
+ expected.binary_buckets[bucket].column);
+
expect_variant_statistics_equal(actual.binary_buckets[bucket].statistics,
+
expected.binary_buckets[bucket].statistics);
+ }
+ expect_variant_statistics_equal(actual.statistics, expected.statistics);
+}
+
+TEST(VariantShredderTest,
ChunkedBinaryTransposeMatchesSingleChunkForOrdinaryAndDoc) {
+ const auto path_for_bucket = [](uint32_t bucket, std::string_view prefix) {
+ for (uint32_t suffix = 0; suffix < 1024; ++suffix) {
+ std::string path(prefix);
+ path += std::to_string(suffix);
+ if (variant_util::variant_binary_shard_of({path.data(),
path.size()}, 2) == bucket) {
+ return path;
+ }
+ }
+ DORIS_CHECK(false) << "failed to find path for bucket " << bucket;
+ return std::string {};
+ };
+ const std::array<std::string, 3> sparse_paths {
+ path_for_bucket(0, "chunk_left_"),
+ path_for_bucket(0, "chunk_middle_"),
+ path_for_bucket(1, "chunk_right_"),
+ };
+ // Ordinary sparse cells per row are 3,1,0,2,3,1,0,2. A four-cell limit
+ // crosses exact boundaries and empty rows; a two-cell limit also forces
the
+ // over-wide-row branch. DOC additionally publishes hot in the binary map,
so
+ // both layouts exercise multiple chunks.
+ const std::array<uint8_t, 8> sparse_masks {0b111, 0b001, 0b000, 0b110,
+ 0b111, 0b100, 0b000, 0b011};
+ auto values = ColumnVariantV2::create();
+ DataTypeVariantV2SerDe serde;
+ DataTypeSerDe::FormatOptions format_options;
+ for (size_t row = 0; row < sparse_masks.size(); ++row) {
+ std::string json = "{\"hot\":" + std::to_string(row);
+ for (size_t path = 0; path < sparse_paths.size(); ++path) {
+ if ((sparse_masks[row] & (1U << path)) != 0) {
+ json += ",\"" + sparse_paths[path] + "\":" +
std::to_string(row * 10 + path);
+ }
+ }
+ json += "}";
+ Slice slice(json.data(), json.size());
+ ASSERT_TRUE(serde.deserialize_one_cell_from_json(*values, slice,
format_options).ok());
+ }
+
+ for (const auto physical_layout :
{segment_v2::VariantShredderPhysicalLayout::ORDINARY,
+
segment_v2::VariantShredderPhysicalLayout::DOC}) {
+ SCOPED_TRACE(physical_layout ==
segment_v2::VariantShredderPhysicalLayout::ORDINARY
+ ? "ordinary"
+ : "doc");
+ const segment_v2::VariantShredderOptions options {
+ .physical_layout = physical_layout,
+ .max_subcolumns_count = 1,
+ .sparse_bucket_count = 2,
+ .doc_bucket_count = 2,
+ .doc_materialization_min_rows = sparse_masks.size() + 1,
+ };
+ segment_v2::VariantShredder single_chunk(options);
+ ASSERT_TRUE(single_chunk.append(values->read_view(), 0,
values->size()).ok());
+ segment_v2::VariantShreddedColumns expected;
+ ASSERT_TRUE(single_chunk.finish(&expected).ok());
+
+ for (const size_t chunk_limit : {size_t {4}, size_t {2}}) {
+ SCOPED_TRACE(testing::Message() << "chunk_limit=" << chunk_limit);
+ segment_v2::VariantShredder chunked(options);
+
segment_v2::VariantShredder::TestAccess::set_binary_cells_per_chunk(chunked,
+
chunk_limit);
+ ASSERT_TRUE(chunked.append(values->read_view(), 0,
values->size()).ok());
+ segment_v2::VariantShreddedColumns actual;
+ ASSERT_TRUE(chunked.finish(&actual).ok());
+
+ const size_t expected_chunks =
+ physical_layout ==
segment_v2::VariantShredderPhysicalLayout::ORDINARY
+ ? (chunk_limit == 4 ? 4 : 6)
+ : (chunk_limit == 4 ? 6 : 8);
+
EXPECT_EQ(segment_v2::VariantShredder::TestAccess::binary_chunk_count(chunked),
+ expected_chunks);
+ expect_shredded_columns_equal(actual, expected);
+ ASSERT_EQ(actual.binary_buckets.size(), 2);
+ for (size_t bucket = 0; bucket < actual.binary_buckets.size();
++bucket) {
+ const auto& map =
+ assert_cast<const
ColumnMap&>(*actual.binary_buckets[bucket].column);
+ EXPECT_EQ(map.size(), sparse_masks.size());
+ EXPECT_GT(map.get_keys().size(), 0) << "bucket=" << bucket;
+ }
+ }
+ }
+}
+
static void construct_column(ColumnPB* column_pb, int32_t col_unique_id,
const std::string& column_type, const
std::string& column_name,
int variant_max_subcolumns_count = 3, bool is_key
= false,
diff --git a/be/test/util/variant/variant_value_test.cpp
b/be/test/util/variant/variant_value_test.cpp
index 88b80f53af8..0d8cc8e6cdb 100644
--- a/be/test/util/variant/variant_value_test.cpp
+++ b/be/test/util/variant/variant_value_test.cpp
@@ -405,6 +405,46 @@ TEST(VariantValueTest,
ObjectLookupSortedAndUnsortedMetadata) {
EXPECT_EQ(id, 1);
}
+TEST(VariantValueTest, ObjectViewMatchesRandomAccessAndRetainsBoundsChecks) {
+ const std::string sorted_metadata = metadata({"a", "b"}, true);
+ const std::string false_value = primitive(VariantPrimitiveId::FALSE_VALUE);
+ const std::string true_value = primitive(VariantPrimitiveId::TRUE_VALUE);
+ const std::string encoded = object_value({0, 1}, {1, 0}, {false_value,
true_value});
+ const VariantRef ref = value_ref(sorted_metadata, encoded);
+
+ const VariantRef::ObjectView object = ref.object_view();
+ ASSERT_EQ(object.size(), 2);
+ for (uint32_t index = 0; index < object.size(); ++index) {
+ uint32_t view_field = std::numeric_limits<uint32_t>::max();
+ uint32_t direct_field = std::numeric_limits<uint32_t>::max();
+ const VariantRef view_value = object.value_at(index, &view_field);
+ const VariantRef direct_value = ref.object_value_at(index,
&direct_field);
+ EXPECT_EQ(view_field, direct_field);
+ EXPECT_EQ(view_value.value.data, direct_value.value.data);
+ EXPECT_EQ(view_value.value.size, direct_value.value.size);
+ }
+ EXPECT_TRUE(object.value_at(0).get_bool());
+ EXPECT_FALSE(object.value_at(1).get_bool());
+ EXPECT_THROW(object.value_at(object.size()), Exception);
+
+ const std::string invalid_id_object =
+ object_value({2}, {0},
{primitive(VariantPrimitiveId::NULL_VALUE)});
+ const VariantRef::ObjectView invalid_id =
+ value_ref(sorted_metadata, invalid_id_object).object_view();
+ EXPECT_THROW(invalid_id.value_at(0), Exception);
+
+ const std::string truncated_object(1,
static_cast<char>(VariantBasicType::OBJECT));
+ EXPECT_THROW(value_ref(sorted_metadata, truncated_object).object_view(),
Exception);
+ const std::string truncated_metadata = sorted_metadata.substr(0, 2);
+ EXPECT_THROW(value_ref(truncated_metadata, encoded).object_view(),
Exception);
+
+ // An empty object has no field ids, so iterating it must not inspect
otherwise unused
+ // metadata. This preserves the random-access API's validation boundary.
+ const std::string empty_object = object_value({}, {}, {});
+ const VariantRef::ObjectView empty = value_ref(truncated_metadata,
empty_object).object_view();
+ EXPECT_EQ(empty.size(), 0);
+}
+
TEST(VariantValueTest, ObjectFindRejectsInvalidReceivers) {
const std::string empty_metadata = metadata({}, true);
VariantRef found;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]