This is an automated email from the ASF dual-hosted git repository.
eldenmoon pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 8104526f018 [improvement](be) Speed up unshredded Parquet Variant
extraction (#67016)
8104526f018 is described below
commit 8104526f018c2a78d3a43b514e21167d16eb7fd8
Author: lihangyu <[email protected]>
AuthorDate: Mon Sep 14 17:16:36 2026 +0800
[improvement](be) Speed up unshredded Parquet Variant extraction (#67016)
---
be/src/core/value/variant/variant_metadata.cpp | 16 +-
be/src/core/value/variant/variant_metadata.h | 1 +
be/src/core/value/variant/variant_value.cpp | 25 +-
be/src/format_v2/parquet/parquet_profile.cpp | 15 +
be/src/format_v2/parquet/parquet_profile.h | 11 +
.../parquet/reader/variant_column_reader.cpp | 752 ++++++++++++++++++++-
.../parquet/variant_column_reader_test.cpp | 546 +++++++++++++++
be/test/util/variant/variant_value_test.cpp | 7 +
8 files changed, 1341 insertions(+), 32 deletions(-)
diff --git a/be/src/core/value/variant/variant_metadata.cpp
b/be/src/core/value/variant/variant_metadata.cpp
index 50c1d94074a..5fbb09a715e 100644
--- a/be/src/core/value/variant/variant_metadata.cpp
+++ b/be/src/core/value/variant/variant_metadata.cpp
@@ -107,7 +107,10 @@ uint32_t VariantMetadataRef::dict_size() const {
}
StringRef VariantMetadataRef::key_at(uint32_t id) const {
- const Layout layout = _layout();
+ return _key_at(_layout(), id);
+}
+
+StringRef VariantMetadataRef::_key_at(const Layout& layout, uint32_t id) const
{
if (id >= layout.num_keys) {
throw Exception(ErrorCode::INVALID_ARGUMENT,
"Variant metadata dictionary id {} is out of range [0,
{})", id,
@@ -132,7 +135,7 @@ int64_t VariantMetadataRef::find_key(StringRef key) const {
const Layout layout = _layout();
if (!sorted_strings()) {
for (uint32_t id = 0; id < layout.num_keys; ++id) {
- if (key_at(id) == key) {
+ if (_key_at(layout, id) == key) {
return id;
}
}
@@ -143,14 +146,14 @@ int64_t VariantMetadataRef::find_key(StringRef key) const
{
uint32_t end = layout.num_keys;
while (begin < end) {
const uint32_t middle = begin + (end - begin) / 2;
- const int comparison = key_at(middle).compare(key);
+ const int comparison = _key_at(layout, middle).compare(key);
if (comparison < 0) {
begin = middle + 1;
} else {
end = middle;
}
}
- if (begin < layout.num_keys && key_at(begin) == key) {
+ if (begin < layout.num_keys && _key_at(layout, begin) == key) {
return begin;
}
return -1;
@@ -158,10 +161,11 @@ int64_t VariantMetadataRef::find_key(StringRef key) const
{
void VariantMetadataRef::validate() const {
const Layout layout = _layout();
+ const bool strings_are_sorted = sorted_strings();
StringRef previous;
for (uint32_t id = 0; id < layout.num_keys; ++id) {
- const StringRef current = key_at(id);
- if (sorted_strings() && id != 0 && previous.compare(current) >= 0) {
+ const StringRef current = _key_at(layout, id);
+ if (strings_are_sorted && id != 0 && previous.compare(current) >= 0) {
throw Exception(ErrorCode::CORRUPTION,
"Variant metadata dictionary is not sorted and
unique at id {}", id);
}
diff --git a/be/src/core/value/variant/variant_metadata.h
b/be/src/core/value/variant/variant_metadata.h
index d2a92c92ee7..09f35789c40 100644
--- a/be/src/core/value/variant/variant_metadata.h
+++ b/be/src/core/value/variant/variant_metadata.h
@@ -47,6 +47,7 @@ private:
};
Layout _layout() const;
+ StringRef _key_at(const Layout& layout, uint32_t id) const;
};
} // namespace doris
diff --git a/be/src/core/value/variant/variant_value.cpp
b/be/src/core/value/variant/variant_value.cpp
index fac75c70096..c19c52a057f 100644
--- a/be/src/core/value/variant/variant_value.cpp
+++ b/be/src/core/value/variant/variant_value.cpp
@@ -469,24 +469,35 @@ bool VariantRef::object_find_by_id(uint32_t field_id,
VariantRef* out) const {
bool VariantRef::_object_find_by_id(const ContainerLayout& layout, uint32_t
field_id,
VariantRef* out) const {
+ const uint32_t dictionary_size = metadata.dict_size();
+ if (field_id >= dictionary_size) {
+ throw Exception(ErrorCode::INVALID_ARGUMENT,
+ "Variant metadata dictionary id {} is out of range [0,
{})", field_id,
+ dictionary_size);
+ }
+ const bool strings_are_sorted = metadata.sorted_strings();
+ // Besides supplying the key for unsorted dictionaries, this preserves the
public lookup
+ // contract's validation of the requested dictionary entry for sorted
dictionaries.
const StringRef target_key = metadata.key_at(field_id);
uint32_t begin = 0;
uint32_t end = layout.count;
while (begin < end) {
const uint32_t middle = begin + (end - begin) / 2;
- const uint32_t middle_id = _object_field_id(layout, middle);
- const int comparison = metadata.sorted_strings()
- ? (middle_id > field_id) - (middle_id <
field_id)
- :
metadata.key_at(middle_id).compare(target_key);
+ const uint32_t middle_id = _object_field_id(layout, middle,
&dictionary_size);
+ const int comparison = strings_are_sorted ? (middle_id > field_id) -
(middle_id < field_id)
+ :
metadata.key_at(middle_id).compare(target_key);
if (comparison < 0) {
begin = middle + 1;
} else {
end = middle;
}
}
- if (begin == layout.count || _object_field_id(layout, begin) != field_id) {
- if (!metadata.sorted_strings() && begin < layout.count &&
- metadata.key_at(_object_field_id(layout, begin)) == target_key) {
+ if (begin == layout.count) {
+ return false;
+ }
+ const uint32_t found_id = _object_field_id(layout, begin,
&dictionary_size);
+ if (found_id != field_id) {
+ if (!strings_are_sorted && metadata.key_at(found_id) == target_key) {
*out = _container_value_at(layout, begin, false);
return true;
}
diff --git a/be/src/format_v2/parquet/parquet_profile.cpp
b/be/src/format_v2/parquet/parquet_profile.cpp
index ef167847c18..8da28ac7d19 100644
--- a/be/src/format_v2/parquet/parquet_profile.cpp
+++ b/be/src/format_v2/parquet/parquet_profile.cpp
@@ -107,6 +107,16 @@ void ParquetProfile::init(RuntimeProfile* profile) {
TUnit::TIME_NS,
parquet_profile);
variant_reconstructed_rows = add_persistent_counter(profile,
"VariantReconstructedRows",
TUnit::UNIT,
parquet_profile);
+ variant_unshredded_direct_seek_time = add_persistent_counter(
+ profile, "VariantUnshreddedDirectSeekTime", TUnit::TIME_NS,
parquet_profile);
+ variant_unshredded_direct_seek_rows = add_persistent_counter(
+ profile, "VariantUnshreddedDirectSeekRows", TUnit::UNIT,
parquet_profile);
+ variant_unshredded_direct_seek_bytes = add_persistent_counter(
+ profile, "VariantUnshreddedDirectSeekBytes", TUnit::BYTES,
parquet_profile);
+ variant_unshredded_prefix_reuse_rows = add_persistent_counter(
+ profile, "VariantUnshreddedPrefixReuseRows", TUnit::UNIT,
parquet_profile);
+ variant_direct_subtree_rows = add_persistent_counter(profile,
"VariantDirectSubtreeRows",
+ TUnit::UNIT,
parquet_profile);
variant_direct_leaf_rows =
add_persistent_counter(profile, "VariantDirectLeafRows",
TUnit::UNIT, parquet_profile);
variant_direct_leaf_path_misses = add_persistent_counter(profile,
"VariantDirectLeafPathMisses",
@@ -345,6 +355,11 @@ ParquetColumnReaderProfile
ParquetProfile::column_reader_profile() const {
.materialization_time = materialization_time,
.variant_reconstruction_time = variant_reconstruction_time,
.variant_reconstructed_rows = variant_reconstructed_rows,
+ .variant_unshredded_direct_seek_time =
variant_unshredded_direct_seek_time,
+ .variant_unshredded_direct_seek_rows =
variant_unshredded_direct_seek_rows,
+ .variant_unshredded_direct_seek_bytes =
variant_unshredded_direct_seek_bytes,
+ .variant_unshredded_prefix_reuse_rows =
variant_unshredded_prefix_reuse_rows,
+ .variant_direct_subtree_rows = variant_direct_subtree_rows,
.variant_direct_leaf_rows = variant_direct_leaf_rows,
.variant_direct_leaf_path_misses = variant_direct_leaf_path_misses,
.variant_direct_leaf_residual_fallbacks =
variant_direct_leaf_residual_fallbacks,
diff --git a/be/src/format_v2/parquet/parquet_profile.h
b/be/src/format_v2/parquet/parquet_profile.h
index cde385bdb72..7ec112e0a55 100644
--- a/be/src/format_v2/parquet/parquet_profile.h
+++ b/be/src/format_v2/parquet/parquet_profile.h
@@ -42,6 +42,12 @@ struct ParquetColumnReaderProfile {
RuntimeProfile::Counter* materialization_time = nullptr; // value
materialization time (ns)
std::shared_ptr<RuntimeProfile::Counter> variant_reconstruction_time;
std::shared_ptr<RuntimeProfile::Counter> variant_reconstructed_rows;
+ // Pure metadata+value roots can seek a requested path before constructing
a root column.
+ std::shared_ptr<RuntimeProfile::Counter>
variant_unshredded_direct_seek_time;
+ std::shared_ptr<RuntimeProfile::Counter>
variant_unshredded_direct_seek_rows;
+ std::shared_ptr<RuntimeProfile::Counter>
variant_unshredded_direct_seek_bytes;
+ std::shared_ptr<RuntimeProfile::Counter>
variant_unshredded_prefix_reuse_rows;
+ std::shared_ptr<RuntimeProfile::Counter> variant_direct_subtree_rows;
std::shared_ptr<RuntimeProfile::Counter> variant_direct_leaf_rows;
std::shared_ptr<RuntimeProfile::Counter> variant_direct_leaf_path_misses;
std::shared_ptr<RuntimeProfile::Counter>
variant_direct_leaf_residual_fallbacks;
@@ -178,6 +184,11 @@ struct ParquetProfile {
RuntimeProfile::Counter* materialization_time = nullptr;
std::shared_ptr<RuntimeProfile::Counter> variant_reconstruction_time;
std::shared_ptr<RuntimeProfile::Counter> variant_reconstructed_rows;
+ std::shared_ptr<RuntimeProfile::Counter>
variant_unshredded_direct_seek_time;
+ std::shared_ptr<RuntimeProfile::Counter>
variant_unshredded_direct_seek_rows;
+ std::shared_ptr<RuntimeProfile::Counter>
variant_unshredded_direct_seek_bytes;
+ std::shared_ptr<RuntimeProfile::Counter>
variant_unshredded_prefix_reuse_rows;
+ std::shared_ptr<RuntimeProfile::Counter> variant_direct_subtree_rows;
std::shared_ptr<RuntimeProfile::Counter> variant_direct_leaf_rows;
std::shared_ptr<RuntimeProfile::Counter> variant_direct_leaf_path_misses;
std::shared_ptr<RuntimeProfile::Counter>
variant_direct_leaf_residual_fallbacks;
diff --git a/be/src/format_v2/parquet/reader/variant_column_reader.cpp
b/be/src/format_v2/parquet/reader/variant_column_reader.cpp
index 3a6826c3d1b..23f9309f418 100644
--- a/be/src/format_v2/parquet/reader/variant_column_reader.cpp
+++ b/be/src/format_v2/parquet/reader/variant_column_reader.cpp
@@ -25,7 +25,10 @@
#include <limits>
#include <mutex>
#include <optional>
+#include <string>
#include <string_view>
+#include <unordered_map>
+#include <utility>
#include <vector>
#include "common/exception.h"
@@ -34,14 +37,19 @@
#include "core/column/column_decimal.h"
#include "core/column/column_map.h"
#include "core/column/column_nullable.h"
+#include "core/column/column_string.h"
#include "core/column/column_struct.h"
#include "core/column/column_vector.h"
#include "core/column/variant_v2/column_variant_v2.h"
#include "core/column/variant_v2/column_variant_v2_typed_column.h"
+#include "core/custom_allocator.h"
#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_number.h"
+#include "core/data_type/data_type_string.h"
#include "core/data_type/data_type_variant_v2.h"
#include "core/value/variant/variant_batch_builder.h"
#include "core/value/variant/variant_metadata.h"
+#include "core/value/variant/variant_scalar.h"
#include "format_v2/parquet/parquet_column_schema.h"
namespace doris::format::parquet {
@@ -52,6 +60,9 @@ struct Cell {
bool is_null = false;
};
+constexpr std::array<char, 1> VARIANT_NULL_VALUE {static_cast<char>(
+ static_cast<uint8_t>(VariantPrimitiveId::NULL_VALUE) <<
VARIANT_VALUE_HEADER_SHIFT)};
+
Cell cell_at(const IColumn& column, size_t row) {
if (row >= column.size()) {
throw Exception(ErrorCode::CORRUPTION, "Parquet Variant row {} exceeds
column size {}", row,
@@ -444,6 +455,20 @@ void encode_variant_range(const ParquetColumnSchema&
schema, const IColumn& wrap
}
}
+std::optional<std::pair<size_t, size_t>> unshredded_child_indices(
+ const ParquetColumnSchema& schema) {
+ if (schema.children.size() != 2 || find_child(schema, "typed_value",
nullptr) != nullptr) {
+ return std::nullopt;
+ }
+ size_t metadata_index = 0;
+ size_t value_index = 0;
+ if (find_child(schema, "metadata", &metadata_index) == nullptr ||
+ find_child(schema, "value", &value_index) == nullptr) {
+ return std::nullopt;
+ }
+ return std::pair {metadata_index, value_index};
+}
+
ColumnVariantV2::MutablePtr encode_variant_column(const ParquetColumnSchema&
schema,
const IColumn& physical,
bool require_metadata =
true) {
@@ -586,13 +611,105 @@ ColumnPtr normalize_projected_primitive_leaf(const
ParquetColumnSchema& schema,
return ColumnNullable::create(std::move(values), std::move(nulls));
}
-bool find_materialized_path(VariantRef current, std::span<const
VariantShreddedPathSegment> path,
- VariantRef* output) {
- DORIS_CHECK(output != nullptr);
- for (const auto& segment : path) {
+struct UnshreddedPathCacheEntry {
+ static constexpr uint32_t MISSING = std::numeric_limits<uint32_t>::max();
+
+ uint32_t value_offset = MISSING;
+ uint32_t value_size = 0;
+
+ bool present() const noexcept { return value_offset != MISSING; }
+};
+
+struct UnshreddedPathCache {
+ DorisVector<UnshreddedPathCacheEntry> entries;
+
+ size_t byte_size() const noexcept { return entries.size() *
sizeof(UnshreddedPathCacheEntry); }
+ size_t allocated_bytes() const noexcept {
+ return entries.capacity() * sizeof(UnshreddedPathCacheEntry);
+ }
+};
+
+struct UnshreddedMetadataIndex {
+ static constexpr uint32_t NULL_ROW = std::numeric_limits<uint32_t>::max();
+
+ const IColumn* physical_identity = nullptr;
+ DorisVector<VariantMetadataRef> dictionaries;
+ DorisVector<uint32_t> row_dictionary_ids;
+
+ size_t byte_size() const noexcept {
+ return dictionaries.size() * sizeof(VariantMetadataRef) +
+ row_dictionary_ids.size() * sizeof(uint32_t);
+ }
+ size_t allocated_bytes() const noexcept {
+ return dictionaries.capacity() * sizeof(VariantMetadataRef) +
+ row_dictionary_ids.capacity() * sizeof(uint32_t);
+ }
+};
+
+std::shared_ptr<const UnshreddedMetadataIndex> build_unshredded_metadata_index(
+ const ParquetColumnSchema& schema, const IColumn& physical,
+ std::pair<size_t, size_t> child_indices) {
+ const auto* outer_nullable =
check_and_get_column<ColumnNullable>(physical);
+ const IColumn& wrapper =
+ outer_nullable == nullptr ? physical :
outer_nullable->get_nested_column();
+ const auto& structure = assert_cast<const ColumnStruct&>(wrapper);
+ DORIS_CHECK_EQ(structure.tuple_size(), schema.children.size());
+
+ auto index = std::make_shared<UnshreddedMetadataIndex>();
+ index->physical_identity = &physical;
+ index->row_dictionary_ids.resize(physical.size(),
UnshreddedMetadataIndex::NULL_ROW);
+ using MetadataIdMap =
+ std::unordered_map<std::string_view, uint32_t,
std::hash<std::string_view>,
+ std::equal_to<std::string_view>,
+ CustomStdAllocator<std::pair<const
std::string_view, uint32_t>>>;
+ MetadataIdMap dictionary_ids;
+ for (size_t row = 0; row < physical.size(); ++row) {
+ if (outer_nullable != nullptr &&
outer_nullable->get_null_map_data()[row] != 0) {
+ continue;
+ }
+ const Cell metadata =
cell_at(structure.get_column(child_indices.first), row);
+ if (metadata.is_null) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} has
null metadata at row {}",
+ schema.name, row);
+ }
+ const StringRef bytes = metadata.column->get_data_at(row);
+ uint32_t dictionary_id = 0;
+ if (index->dictionaries.empty()) {
+ index->dictionaries.push_back({.data = bytes.data, .size =
bytes.size});
+ } else if (index->dictionaries.size() == 1 && dictionary_ids.empty() &&
+ StringRef(index->dictionaries.front().data,
index->dictionaries.front().size) ==
+ bytes) {
+ // Iceberg normally repeats one metadata dictionary throughout a
decoded block. Delay
+ // the hash table until a second dictionary is actually observed.
+ } else {
+ if (dictionary_ids.empty()) {
+ const VariantMetadataRef first = index->dictionaries.front();
+ dictionary_ids.emplace(std::string_view(first.data,
first.size), 0);
+ }
+ const std::string_view key(bytes.data, bytes.size);
+ if (const auto found = dictionary_ids.find(key); found !=
dictionary_ids.end()) {
+ dictionary_id = found->second;
+ } else {
+ dictionary_id =
static_cast<uint32_t>(index->dictionaries.size());
+ index->dictionaries.push_back({.data = bytes.data, .size =
bytes.size});
+ dictionary_ids.emplace(key, dictionary_id);
+ }
+ }
+ index->row_dictionary_ids[row] = dictionary_id;
+ }
+ return index;
+}
+
+template <typename ObjectFinder>
+bool find_materialized_path_impl(VariantRef current,
+ std::span<const VariantShreddedPathSegment>
path,
+ const ObjectFinder& find_object, VariantRef*
output) {
+ DCHECK(output != nullptr);
+ for (size_t position = 0; position < path.size(); ++position) {
+ const auto& segment = path[position];
if (segment.kind == VariantShreddedPathSegment::Kind::OBJECT_KEY) {
if (current.basic_type() != VariantBasicType::OBJECT ||
- !current.object_find(segment.key, ¤t)) {
+ !find_object(current, segment.key, position, ¤t)) {
return false;
}
continue;
@@ -611,6 +728,400 @@ bool find_materialized_path(VariantRef current,
std::span<const VariantShreddedP
return true;
}
+bool find_materialized_path(VariantRef current, std::span<const
VariantShreddedPathSegment> path,
+ VariantRef* output) {
+ return find_materialized_path_impl(
+ current, path,
+ [](VariantRef object, StringRef key, size_t, VariantRef* found) {
+ return object.object_find(key, found);
+ },
+ output);
+}
+
+bool find_materialized_path_with_index(VariantRef current, uint32_t
dictionary_id,
+ const UnshreddedMetadataIndex&
metadata_index,
+ std::span<const
VariantShreddedPathSegment> path,
+ DorisVector<int64_t>&
resolved_field_ids,
+ VariantRef* output) {
+ DCHECK_LT(dictionary_id, metadata_index.dictionaries.size());
+ DCHECK_EQ(resolved_field_ids.size(), metadata_index.dictionaries.size() *
path.size());
+ constexpr int64_t UNRESOLVED_FIELD_ID = -2;
+ return find_materialized_path_impl(
+ current, path,
+ [&](VariantRef object, StringRef key, size_t position, VariantRef*
found) {
+ int64_t& field_id = resolved_field_ids[dictionary_id *
path.size() + position];
+ bool layout_validated = false;
+ if (field_id == UNRESOLVED_FIELD_ID) {
+ // object_find() validates the object layout before
consulting metadata.
+ static_cast<void>(object.num_elements());
+ layout_validated = true;
+ field_id =
metadata_index.dictionaries[dictionary_id].find_key(key);
+ }
+ if (field_id < 0) {
+ // A cached metadata miss must not hide a corrupt object
in a later row.
+ if (!layout_validated) {
+ static_cast<void>(object.num_elements());
+ }
+ return false;
+ }
+ return
object.object_find_by_id(static_cast<uint32_t>(field_id), found);
+ },
+ output);
+}
+
+struct UnshreddedPathScan {
+ MutableColumnPtr outer_nulls;
+ MutableColumnPtr typed_values;
+ DataTypePtr typed_type;
+ std::shared_ptr<const UnshreddedPathCache> path_cache;
+ int64_t copied_bytes = 0;
+};
+
+enum class UnshreddedTypedKind : uint8_t { UNKNOWN, STRING, INTEGER,
UNSUPPORTED };
+
+class UnshreddedTypedValueBuilder {
+public:
+ UnshreddedTypedValueBuilder(size_t rows, const UnshreddedMetadataIndex&
metadata_index)
+ : _rows(rows),
+ _metadata_index(metadata_index),
+ _inner_nulls(ColumnUInt8::create()),
+ _result_nulls(ColumnUInt8::create()),
+ _validated_metadata(metadata_index.dictionaries.size(), 0) {
+ _inner_nulls->reserve(rows);
+ _result_nulls->reserve(rows);
+ }
+
+ void append_outer_null() { append_null(1); }
+
+ void append_json_null(uint32_t dictionary_id) {
+ if (_typed_kind == UnshreddedTypedKind::UNKNOWN) {
+ _pending_json_null_dictionaries.push_back(dictionary_id);
+ } else if (_typed_kind == UnshreddedTypedKind::INTEGER &&
+ !validate_integer_metadata(dictionary_id)) {
+ mark_unsupported();
+ }
+ append_null(0);
+ }
+
+ void append_scalar(const VariantRef& found, uint32_t dictionary_id, size_t
row) {
+ if (_typed_kind == UnshreddedTypedKind::UNSUPPORTED) {
+ append_null(0);
+ return;
+ }
+
+ const VariantBasicType basic_type = found.basic_type();
+ const bool is_string = basic_type == VariantBasicType::SHORT_STRING ||
+ (basic_type == VariantBasicType::PRIMITIVE &&
+ found.primitive_id() ==
VariantPrimitiveId::STRING);
+ if (is_string) {
+ if (!prepare(UnshreddedTypedKind::STRING, row)) {
+ append_null(0);
+ return;
+ }
+ const StringRef string = found.get_string();
+
assert_cast<ColumnString&>(*_typed_values).insert_data(string.data,
string.size);
+ _inner_nulls->insert_value(0);
+ _result_nulls->insert_value(0);
+ DCHECK_LE(string.size,
+ static_cast<size_t>(std::numeric_limits<int64_t>::max()
- _copied_bytes));
+ _copied_bytes += static_cast<int64_t>(string.size);
+ return;
+ }
+
+ const auto primitive_id = basic_type == VariantBasicType::PRIMITIVE
+ ? found.primitive_id()
+ : VariantPrimitiveId::NULL_VALUE;
+ const bool is_integer = primitive_id == VariantPrimitiveId::INT8 ||
+ primitive_id == VariantPrimitiveId::INT16 ||
+ primitive_id == VariantPrimitiveId::INT32 ||
+ primitive_id == VariantPrimitiveId::INT64;
+ const int64_t integer = is_integer ? found.get_int() : 0;
+ // Typed Variant integers are re-encoded using the narrowest width.
Keep explicitly widened
+ // source integers on the encoded path so observable physical types
remain unchanged.
+ const bool has_canonical_width =
+ is_integer &&
VariantScalarRef::integer(integer).encoded_size() == found.value.size;
+ if (!has_canonical_width || !validate_integer_metadata(dictionary_id)
||
+ !prepare(UnshreddedTypedKind::INTEGER, row)) {
+ mark_unsupported();
+ append_null(0);
+ return;
+ }
+ assert_cast<ColumnInt64&>(*_typed_values).insert_value(integer);
+ _inner_nulls->insert_value(0);
+ _result_nulls->insert_value(0);
+ DCHECK_LE(static_cast<int64_t>(sizeof(int64_t)),
+ std::numeric_limits<int64_t>::max() - _copied_bytes);
+ _copied_bytes += sizeof(int64_t);
+ }
+
+ UnshreddedPathScan finish(std::shared_ptr<const UnshreddedPathCache>
path_cache) && {
+ DataTypePtr typed_type;
+ if (_typed_values) {
+ typed_type = _typed_kind == UnshreddedTypedKind::STRING
+ ?
DataTypePtr(std::make_shared<DataTypeString>())
+ :
DataTypePtr(std::make_shared<DataTypeInt64>());
+ _typed_values =
+ ColumnNullable::create(std::move(_typed_values),
std::move(_inner_nulls));
+ }
+ return {.outer_nulls = std::move(_result_nulls),
+ .typed_values = std::move(_typed_values),
+ .typed_type = std::move(typed_type),
+ .path_cache = std::move(path_cache),
+ .copied_bytes = _copied_bytes};
+ }
+
+private:
+ bool validate_integer_metadata(uint32_t dictionary_id) {
+ DCHECK_LT(dictionary_id, _metadata_index.dictionaries.size());
+ try {
+ if (_validated_metadata[dictionary_id] == 0) {
+
validate_variant_metadata(_metadata_index.dictionaries[dictionary_id]);
+ _validated_metadata[dictionary_id] = 1;
+ }
+ return true;
+ } catch (const Exception&) {
+ return false;
+ }
+ }
+
+ bool prepare(UnshreddedTypedKind kind, size_t row) {
+ if (_typed_kind == UnshreddedTypedKind::UNKNOWN) {
+ if (kind == UnshreddedTypedKind::INTEGER) {
+ for (const uint32_t dictionary_id :
_pending_json_null_dictionaries) {
+ if (!validate_integer_metadata(dictionary_id)) {
+ mark_unsupported();
+ return false;
+ }
+ }
+ _typed_values = ColumnInt64::create();
+ } else {
+ DCHECK(kind == UnshreddedTypedKind::STRING);
+ _typed_values = ColumnString::create();
+ }
+ _pending_json_null_dictionaries.clear();
+ _typed_kind = kind;
+ _typed_values->reserve(_rows);
+ _typed_values->insert_many_defaults(row);
+ return true;
+ }
+ if (_typed_kind == kind) {
+ return true;
+ }
+ mark_unsupported();
+ return false;
+ }
+
+ void mark_unsupported() {
+ _typed_kind = UnshreddedTypedKind::UNSUPPORTED;
+ _typed_values.reset();
+ _pending_json_null_dictionaries.clear();
+ }
+
+ void append_null(uint8_t outer_null) {
+ if (_typed_values) {
+ _typed_values->insert_default();
+ }
+ _inner_nulls->insert_value(1);
+ _result_nulls->insert_value(outer_null);
+ }
+
+ size_t _rows;
+ const UnshreddedMetadataIndex& _metadata_index;
+ MutableColumnPtr _typed_values;
+ ColumnUInt8::MutablePtr _inner_nulls;
+ ColumnUInt8::MutablePtr _result_nulls;
+ UnshreddedTypedKind _typed_kind = UnshreddedTypedKind::UNKNOWN;
+ int64_t _copied_bytes = 0;
+ DorisVector<uint8_t> _validated_metadata;
+ DorisVector<uint32_t> _pending_json_null_dictionaries;
+};
+
+std::optional<VariantRef> unshredded_root_at(const ColumnNullable*
outer_nullable,
+ const ColumnStruct& structure,
+ std::pair<size_t, size_t>
child_indices,
+ const UnshreddedMetadataIndex&
metadata_index,
+ size_t row) {
+ if (outer_nullable != nullptr && outer_nullable->get_null_map_data()[row]
!= 0) {
+ return std::nullopt;
+ }
+
+ const uint32_t dictionary_id = metadata_index.row_dictionary_ids[row];
+ DCHECK_LT(dictionary_id, metadata_index.dictionaries.size());
+ const Cell value = cell_at(structure.get_column(child_indices.second),
row);
+ const StringRef value_bytes =
+ value.is_null ? StringRef(VARIANT_NULL_VALUE.data(),
VARIANT_NULL_VALUE.size())
+ : value.column->get_data_at(row);
+ return VariantRef {.metadata = metadata_index.dictionaries[dictionary_id],
+ .value = value_bytes};
+}
+
+UnshreddedPathCacheEntry make_unshredded_path_cache_entry(const VariantRef&
root,
+ const VariantRef&
found) {
+ const uintptr_t root_begin = reinterpret_cast<uintptr_t>(root.value.data);
+ const uintptr_t root_end = root_begin + root.value.size;
+ const uintptr_t found_begin =
reinterpret_cast<uintptr_t>(found.value.data);
+ const uintptr_t found_end = found_begin + found.value.size;
+ DCHECK_GE(found_begin, root_begin);
+ DCHECK_LE(found_end, root_end);
+ const size_t offset = found_begin - root_begin;
+ DCHECK_LT(offset, static_cast<size_t>(UnshreddedPathCacheEntry::MISSING));
+ DCHECK_LE(found.value.size,
static_cast<size_t>(std::numeric_limits<uint32_t>::max()));
+ return {.value_offset = static_cast<uint32_t>(offset),
+ .value_size = static_cast<uint32_t>(found.value.size)};
+}
+
+UnshreddedPathScan scan_unshredded_path(const ParquetColumnSchema& schema,
const IColumn& physical,
+ std::pair<size_t, size_t>
child_indices,
+ const UnshreddedMetadataIndex&
metadata_index,
+ std::span<const
VariantShreddedPathSegment> path,
+ const UnshreddedPathCache*
prefix_cache = nullptr) {
+ const auto* outer_nullable =
check_and_get_column<ColumnNullable>(physical);
+ const IColumn& wrapper =
+ outer_nullable == nullptr ? physical :
outer_nullable->get_nested_column();
+ const auto& structure = assert_cast<const ColumnStruct&>(wrapper);
+ DORIS_CHECK_EQ(structure.tuple_size(), schema.children.size());
+
+ auto path_cache = std::make_shared<UnshreddedPathCache>();
+ path_cache->entries.reserve(physical.size());
+ DORIS_CHECK(prefix_cache == nullptr || prefix_cache->entries.size() ==
physical.size());
+ DORIS_CHECK_EQ(metadata_index.physical_identity, &physical);
+ DORIS_CHECK_EQ(metadata_index.row_dictionary_ids.size(), physical.size());
+ constexpr int64_t UNRESOLVED_FIELD_ID = -2;
+ DorisVector<int64_t> resolved_field_ids(metadata_index.dictionaries.size()
* path.size(),
+ UNRESOLVED_FIELD_ID);
+ UnshreddedTypedValueBuilder typed_values(physical.size(), metadata_index);
+
+ for (size_t row = 0; row < physical.size(); ++row) {
+ const auto root =
+ unshredded_root_at(outer_nullable, structure, child_indices,
metadata_index, row);
+ if (!root.has_value()) {
+ DCHECK(prefix_cache == nullptr ||
!prefix_cache->entries[row].present());
+ path_cache->entries.emplace_back();
+ typed_values.append_outer_null();
+ continue;
+ }
+
+ const uint32_t dictionary_id = metadata_index.row_dictionary_ids[row];
+ VariantRef current = *root;
+ if (prefix_cache != nullptr) {
+ const auto& cached = prefix_cache->entries[row];
+ if (!cached.present()) {
+ path_cache->entries.emplace_back();
+ typed_values.append_outer_null();
+ continue;
+ }
+ DCHECK_LE(static_cast<size_t>(cached.value_offset) +
cached.value_size,
+ root->value.size);
+ current.value = {root->value.data + cached.value_offset,
cached.value_size};
+ }
+ VariantRef found;
+ if (!find_materialized_path_with_index(current, dictionary_id,
metadata_index, path,
+ resolved_field_ids, &found)) {
+ path_cache->entries.emplace_back();
+ typed_values.append_outer_null();
+ continue;
+ }
+ path_cache->entries.push_back(make_unshredded_path_cache_entry(*root,
found));
+ if (found.is_null()) {
+ typed_values.append_json_null(dictionary_id);
+ } else {
+ typed_values.append_scalar(found, dictionary_id, row);
+ }
+ }
+ return std::move(typed_values).finish(std::move(path_cache));
+}
+
+void append_variant_rows(ColumnVariantV2& output, std::span<const VariantRef>
rows) {
+ try {
+ VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows =
rows.size()});
+ for (const VariantRef value : rows) {
+ auto row = builder.begin_row();
+ row.add_value(value);
+ row.finish();
+ }
+ output.insert_encoded_batch(builder.finish_batch());
+ } catch (...) {
+ if (rows.size() <= 1) {
+ throw;
+ }
+ // A builder has one metadata dictionary. Split heterogeneous batches
while preserving the
+ // original validation error for a malformed single row.
+ const size_t middle = rows.size() / 2;
+ append_variant_rows(output, rows.first(middle));
+ append_variant_rows(output, rows.subspan(middle));
+ }
+}
+
+ColumnPtr normalize_unshredded_path(const ParquetColumnSchema& schema, const
IColumn& physical,
+ std::pair<size_t, size_t> child_indices,
+ std::span<const
VariantShreddedPathSegment> path,
+ int64_t* copied_bytes) {
+ DORIS_CHECK(copied_bytes != nullptr);
+ const auto* outer_nullable =
check_and_get_column<ColumnNullable>(physical);
+ const IColumn& wrapper =
+ outer_nullable == nullptr ? physical :
outer_nullable->get_nested_column();
+ const auto& structure = assert_cast<const ColumnStruct&>(wrapper);
+ DORIS_CHECK_EQ(structure.tuple_size(), schema.children.size());
+
+ auto values = ColumnVariantV2::create();
+ auto nulls = ColumnUInt8::create();
+ nulls->reserve(physical.size());
+ constexpr size_t MAX_DIRECT_SEEK_BATCH_ROWS = 4096;
+ DorisVector<VariantRef> encoded_rows;
+ encoded_rows.reserve(std::min(physical.size(),
MAX_DIRECT_SEEK_BATCH_ROWS));
+ const VariantRef null_value {.metadata = {.data =
VARIANT_EMPTY_METADATA.data(),
+ .size =
VARIANT_EMPTY_METADATA.size()},
+ .value = {VARIANT_NULL_VALUE.data(),
VARIANT_NULL_VALUE.size()}};
+ int64_t total_copied_bytes = 0;
+ auto count_copied_bytes = [&](const VariantRef value) {
+ for (const size_t bytes : {value.metadata.size, value.value.size}) {
+ DORIS_CHECK_LE(bytes,
static_cast<size_t>(std::numeric_limits<int64_t>::max() -
+ total_copied_bytes));
+ total_copied_bytes += static_cast<int64_t>(bytes);
+ }
+ };
+
+ for (size_t begin = 0; begin < physical.size(); begin +=
MAX_DIRECT_SEEK_BATCH_ROWS) {
+ const size_t end = std::min(physical.size(), begin +
MAX_DIRECT_SEEK_BATCH_ROWS);
+ encoded_rows.clear();
+ for (size_t row = begin; row < end; ++row) {
+ if (outer_nullable != nullptr &&
outer_nullable->get_null_map_data()[row] != 0) {
+ encoded_rows.push_back(null_value);
+ nulls->insert_value(1);
+ count_copied_bytes(null_value);
+ continue;
+ }
+
+ const Cell metadata =
cell_at(structure.get_column(child_indices.first), row);
+ if (metadata.is_null) {
+ throw Exception(ErrorCode::CORRUPTION,
+ "Parquet Variant {} has null metadata at row
{}", schema.name, row);
+ }
+ const StringRef metadata_bytes = metadata.column->get_data_at(row);
+ const Cell value =
cell_at(structure.get_column(child_indices.second), row);
+ const StringRef value_bytes =
+ value.is_null ? StringRef(VARIANT_NULL_VALUE.data(),
VARIANT_NULL_VALUE.size())
+ : value.column->get_data_at(row);
+ const VariantRef root {
+ .metadata = {.data = metadata_bytes.data, .size =
metadata_bytes.size},
+ .value = value_bytes};
+ VariantRef found;
+ if (find_materialized_path(root, path, &found)) {
+ encoded_rows.push_back(found);
+ nulls->insert_value(0);
+ count_copied_bytes(found);
+ } else {
+ encoded_rows.push_back(null_value);
+ nulls->insert_value(1);
+ count_copied_bytes(null_value);
+ }
+ }
+ append_variant_rows(*values, encoded_rows);
+ }
+ *copied_bytes = total_copied_bytes;
+ return ColumnNullable::create(std::move(values), std::move(nulls));
+}
+
ColumnPtr normalize_materialized_path(const ColumnVariantV2& materialized,
std::span<const
VariantShreddedPathSegment> path) {
VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows =
materialized.size()});
@@ -671,15 +1182,27 @@ bool same_shredded_schema(const ParquetColumnSchema&
left, const ParquetColumnSc
void append_compatible_column(IColumn& output, const IColumn& converted);
void validate_compatible_column(const IColumn& output, const IColumn&
converted);
+struct OwnedShreddedPathSegment {
+ VariantShreddedPathSegment::Kind kind =
VariantShreddedPathSegment::Kind::OBJECT_KEY;
+ std::string key;
+ int64_t index = 0;
+
+ bool operator==(const OwnedShreddedPathSegment&) const = default;
+};
+
class ParquetVariantShreddedState final : public VariantShreddedState {
public:
- ParquetVariantShreddedState(std::shared_ptr<const ParquetColumnSchema>
schema,
- ColumnPtr physical, bool complete,
- ParquetColumnReaderProfile profile = {})
+ ParquetVariantShreddedState(
+ std::shared_ptr<const ParquetColumnSchema> schema, ColumnPtr
physical, bool complete,
+ ParquetColumnReaderProfile profile = {},
+ std::vector<OwnedShreddedPathSegment> unshredded_prefix = {},
+ std::shared_ptr<const UnshreddedPathCache> unshredded_path_cache =
nullptr)
: _schema(std::move(schema)),
_physical(std::move(physical)),
_complete(complete),
- _profile(profile) {
+ _profile(profile),
+ _unshredded_prefix(std::move(unshredded_prefix)),
+ _unshredded_path_cache(std::move(unshredded_path_cache)) {
DORIS_CHECK(_schema != nullptr && static_cast<bool>(_physical));
const ColumnPtr wrapper = unwrap_nullable(_physical);
const auto* structure = check_and_get_column<ColumnStruct>(*wrapper);
@@ -687,21 +1210,49 @@ public:
throw Exception(ErrorCode::CORRUPTION,
"Parquet Variant {} physical field count
mismatch", _schema->name);
}
+ DORIS_CHECK(_unshredded_prefix.empty() ||
unshredded_child_indices(*_schema).has_value());
+ DORIS_CHECK(!_unshredded_path_cache ||
+ (!_unshredded_prefix.empty() &&
+ _unshredded_path_cache->entries.size() ==
_physical->size()));
}
size_t size() const override { return _physical->size(); }
size_t byte_size() const override {
std::lock_guard lock(_materialization_lock);
- return _physical->byte_size() + (_materialized ?
_materialized->byte_size() : 0) +
+ return _physical->byte_size() +
+ (_unshredded_path_cache ? _unshredded_path_cache->byte_size() :
0) +
+ (_unshredded_metadata_index ?
_unshredded_metadata_index->byte_size() : 0) +
+ (_normalized_prefix ? _normalized_prefix->byte_size() : 0) +
+ (_materialized ? _materialized->byte_size() : 0) +
(_serialized ? _serialized->byte_size() : 0);
}
size_t allocated_bytes() const override {
std::lock_guard lock(_materialization_lock);
return _physical->allocated_bytes() +
+ (_unshredded_path_cache ?
_unshredded_path_cache->allocated_bytes() : 0) +
+ (_unshredded_metadata_index ?
_unshredded_metadata_index->allocated_bytes() : 0) +
+ (_normalized_prefix ? _normalized_prefix->allocated_bytes() :
0) +
(_materialized ? _materialized->allocated_bytes() : 0) +
(_serialized ? _serialized->allocated_bytes() : 0);
}
- void sanity_check() const override { _physical->sanity_check(); }
+ void sanity_check() const override {
+ _physical->sanity_check();
+ DORIS_CHECK(!_unshredded_path_cache ||
+ _unshredded_path_cache->entries.size() ==
_physical->size());
+ std::lock_guard lock(_materialization_lock);
+ if (_unshredded_metadata_index) {
+ DORIS_CHECK_EQ(_unshredded_metadata_index->physical_identity,
_physical.get());
+
DORIS_CHECK_EQ(_unshredded_metadata_index->row_dictionary_ids.size(),
+ _physical->size());
+ for (const uint32_t dictionary_id :
_unshredded_metadata_index->row_dictionary_ids) {
+ DORIS_CHECK(dictionary_id == UnshreddedMetadataIndex::NULL_ROW
||
+ dictionary_id <
_unshredded_metadata_index->dictionaries.size());
+ }
+ }
+ if (_normalized_prefix) {
+ _normalized_prefix->sanity_check();
+ }
+ }
void for_each_subcolumn(IColumn::ColumnCallback callback) const override {
callback(*_physical);
@@ -713,21 +1264,52 @@ public:
// no metadata/value columns from which a canonical Variant could be
reconstructed.
// The projection schema is immutable and reader-scoped, so derived
selections share it
// instead of cloning the whole shredded tree for every filter
operation.
- return std::make_shared<ParquetVariantShreddedState>(
- _schema, _physical->filter(filter, result_size_hint),
_complete, _profile);
+ auto select_path_cache = [&](const UnshreddedPathCache& source) {
+ DORIS_CHECK_EQ(filter.size(), source.entries.size());
+ auto selected = std::make_shared<UnshreddedPathCache>();
+ if (result_size_hint > 0) {
+ selected->entries.reserve(cast_set<size_t>(result_size_hint));
+ }
+ for (size_t row = 0; row < source.entries.size(); ++row) {
+ if (filter[row] != 0) {
+ selected->entries.push_back(source.entries[row]);
+ }
+ }
+ return selected;
+ };
+ return select_state(_physical->filter(filter, result_size_hint),
select_path_cache);
}
std::shared_ptr<VariantShreddedState> select_range(size_t start, size_t
length) const override {
- return std::make_shared<ParquetVariantShreddedState>(_schema,
_physical->cut(start, length),
- _complete,
_profile);
+ auto select_path_cache = [&](const UnshreddedPathCache& source) {
+ DORIS_CHECK_LE(start, source.entries.size());
+ DORIS_CHECK_LE(length, source.entries.size() - start);
+ auto selected = std::make_shared<UnshreddedPathCache>();
+ selected->entries.insert(selected->entries.end(),
+ source.entries.begin() +
cast_set<ssize_t>(start),
+ source.entries.begin() +
cast_set<ssize_t>(start + length));
+ return selected;
+ };
+ return select_state(_physical->cut(start, length), select_path_cache);
}
std::shared_ptr<VariantShreddedState> select_indices(
const uint32_t* indices_begin, const uint32_t* indices_end) const
override {
- MutableColumnPtr selected = _physical->clone_empty();
- selected->insert_indices_from(*_physical, indices_begin, indices_end);
- return std::make_shared<ParquetVariantShreddedState>(_schema,
std::move(selected),
- _complete,
_profile);
+ auto select_column = [&](const ColumnPtr& column) {
+ MutableColumnPtr selected = column->clone_empty();
+ selected->insert_indices_from(*column, indices_begin, indices_end);
+ return ColumnPtr(std::move(selected));
+ };
+ auto select_path_cache = [&](const UnshreddedPathCache& source) {
+ auto selected = std::make_shared<UnshreddedPathCache>();
+ selected->entries.reserve(indices_end - indices_begin);
+ for (const uint32_t* index = indices_begin; index != indices_end;
++index) {
+ DORIS_CHECK_LT(*index, source.entries.size());
+ selected->entries.push_back(source.entries[*index]);
+ }
+ return selected;
+ };
+ return select_state(select_column(_physical), select_path_cache);
}
bool can_materialize() const override { return _complete; }
@@ -735,6 +1317,7 @@ public:
bool try_append(const VariantShreddedState& source) override {
const auto* parquet_source = dynamic_cast<const
ParquetVariantShreddedState*>(&source);
if (parquet_source == nullptr || _complete !=
parquet_source->_complete ||
+ _unshredded_prefix != parquet_source->_unshredded_prefix ||
!same_shredded_schema(*_schema, *parquet_source->_schema)) {
return false;
}
@@ -743,6 +1326,9 @@ public:
append_compatible_column(*mutable_physical,
*parquet_source->_physical);
_physical = std::move(mutable_physical);
std::lock_guard lock(_materialization_lock);
+ _unshredded_path_cache.reset();
+ _unshredded_metadata_index.reset();
+ _normalized_prefix.reset();
_materialized.reset();
_serialized.reset();
return true;
@@ -758,6 +1344,11 @@ public:
return path_miss();
}
+ if (unshredded_child_indices(*_schema).has_value()) {
+ return VariantShreddedTypedValue {
+ .column = nullptr, .type = nullptr, .normalized =
direct_unshredded_path(path)};
+ }
+
const ParquetColumnSchema* typed_schema = nullptr;
ColumnPtr typed = struct_child(*_schema, _physical, "typed_value",
&typed_schema);
if (!typed || typed_schema->kind != ParquetColumnSchemaKind::STRUCT) {
@@ -827,6 +1418,10 @@ public:
return std::nullopt;
}
+ if (unshredded_child_indices(*_schema).has_value()) {
+ return direct_unshredded_path(path);
+ }
+
const ParquetColumnSchema* typed_schema = nullptr;
ColumnPtr typed = struct_child(*_schema, _physical, "typed_value",
&typed_schema);
if (typed && typed_schema->kind == ParquetColumnSchemaKind::STRUCT) {
@@ -878,6 +1473,24 @@ public:
ErrorCode::INTERNAL_ERROR,
"A projected Parquet Variant can only serve its validated
shredded leaves");
}
+ if (!_unshredded_prefix.empty()) {
+ if (!_normalized_prefix) {
+ const auto child_indices = unshredded_child_indices(*_schema);
+ DORIS_CHECK(child_indices.has_value());
+ const auto borrowed =
borrow_unshredded_path(_unshredded_prefix);
+ int64_t copied_bytes = 0;
+ {
+
SCOPED_TIMER(_profile.variant_unshredded_direct_seek_time.get());
+ _normalized_prefix = normalize_unshredded_path(
+ *_schema, *_physical, *child_indices, borrowed,
&copied_bytes);
+ }
+ const auto rows = static_cast<int64_t>(_physical->size());
+ update_counter(_profile.variant_unshredded_direct_seek_rows,
rows);
+ update_counter(_profile.variant_unshredded_direct_seek_bytes,
copied_bytes);
+ }
+ const auto& nullable = assert_cast<const
ColumnNullable&>(*_normalized_prefix);
+ return assert_cast<const
ColumnVariantV2&>(nullable.get_nested_column());
+ }
if (!_materialized) {
SCOPED_TIMER(_profile.variant_reconstruction_time.get());
_materialized = encode_variant_column(*_schema, *_physical);
@@ -902,6 +1515,103 @@ public:
}
private:
+ template <typename PathCacheSelector>
+ std::shared_ptr<ParquetVariantShreddedState> select_state(
+ ColumnPtr selected_physical, const PathCacheSelector&
select_path_cache) const {
+ std::shared_ptr<const UnshreddedPathCache> path_cache;
+ {
+ std::lock_guard lock(_materialization_lock);
+ path_cache = _unshredded_path_cache;
+ }
+
+ std::shared_ptr<const UnshreddedPathCache> selected_path_cache;
+ if (path_cache) {
+ selected_path_cache = select_path_cache(*path_cache);
+ }
+ return std::make_shared<ParquetVariantShreddedState>(
+ _schema, std::move(selected_physical), _complete, _profile,
_unshredded_prefix,
+ std::move(selected_path_cache));
+ }
+
+ std::vector<OwnedShreddedPathSegment> combined_unshredded_path(
+ std::span<const VariantShreddedPathSegment> suffix) const {
+ std::vector<OwnedShreddedPathSegment> combined = _unshredded_prefix;
+ combined.reserve(combined.size() + suffix.size());
+ for (const auto& segment : suffix) {
+ combined.push_back(
+ {.kind = segment.kind,
+ .key = segment.kind ==
VariantShreddedPathSegment::Kind::OBJECT_KEY
+ ? (segment.key.size == 0
+ ? std::string()
+ : std::string(segment.key.data,
segment.key.size))
+ : std::string(),
+ .index = segment.index});
+ }
+ return combined;
+ }
+
+ static std::vector<VariantShreddedPathSegment> borrow_unshredded_path(
+ const std::vector<OwnedShreddedPathSegment>& owned) {
+ std::vector<VariantShreddedPathSegment> borrowed;
+ borrowed.reserve(owned.size());
+ for (const auto& segment : owned) {
+ borrowed.push_back({.kind = segment.kind,
+ .key = {segment.key.data(),
segment.key.size()},
+ .index = segment.index});
+ }
+ return borrowed;
+ }
+
+ std::shared_ptr<const UnshreddedMetadataIndex>
get_unshredded_metadata_index(
+ std::pair<size_t, size_t> child_indices) const {
+ std::lock_guard lock(_materialization_lock);
+ if (!_unshredded_metadata_index) {
+ _unshredded_metadata_index =
+ build_unshredded_metadata_index(*_schema, *_physical,
child_indices);
+ }
+ return _unshredded_metadata_index;
+ }
+
+ ColumnPtr direct_unshredded_path(std::span<const
VariantShreddedPathSegment> suffix) const {
+ const auto child_indices = unshredded_child_indices(*_schema);
+ DORIS_CHECK(child_indices.has_value());
+ std::vector<OwnedShreddedPathSegment> combined =
combined_unshredded_path(suffix);
+ const auto borrowed_combined = borrow_unshredded_path(combined);
+ std::shared_ptr<const UnshreddedMetadataIndex> metadata_index;
+ UnshreddedPathScan scan;
+ {
+ SCOPED_TIMER(_profile.variant_unshredded_direct_seek_time.get());
+ metadata_index = get_unshredded_metadata_index(*child_indices);
+ scan = scan_unshredded_path(
+ *_schema, *_physical, *child_indices, *metadata_index,
+ _unshredded_path_cache
+ ? suffix
+ : std::span<const
VariantShreddedPathSegment>(borrowed_combined),
+ _unshredded_path_cache.get());
+ }
+
+ const auto rows = static_cast<int64_t>(_physical->size());
+ update_counter(_profile.variant_unshredded_direct_seek_rows, rows);
+ if (_unshredded_path_cache) {
+ update_counter(_profile.variant_unshredded_prefix_reuse_rows,
rows);
+ }
+ MutableColumnPtr values;
+ std::shared_ptr<ParquetVariantShreddedState> subtree_state;
+ if (scan.typed_values) {
+ values =
ColumnVariantV2::create_typed(std::move(scan.typed_values),
+ std::move(scan.typed_type));
+ update_counter(_profile.variant_unshredded_direct_seek_bytes,
scan.copied_bytes);
+ update_counter(_profile.variant_direct_leaf_rows, rows);
+ } else {
+ subtree_state = std::make_shared<ParquetVariantShreddedState>(
+ _schema, _physical, _complete, _profile, combined,
std::move(scan.path_cache));
+ values = ColumnVariantV2::create_shredded(subtree_state);
+ update_counter(_profile.variant_direct_subtree_rows, rows);
+ }
+ ColumnPtr result = ColumnNullable::create(std::move(values),
std::move(scan.outer_nulls));
+ return result;
+ }
+
static void update_counter(const std::shared_ptr<RuntimeProfile::Counter>&
counter,
int64_t value) {
if (counter != nullptr) {
@@ -913,7 +1623,11 @@ private:
ColumnPtr _physical;
bool _complete = true;
ParquetColumnReaderProfile _profile;
+ std::vector<OwnedShreddedPathSegment> _unshredded_prefix;
+ std::shared_ptr<const UnshreddedPathCache> _unshredded_path_cache;
mutable std::mutex _materialization_lock;
+ mutable std::shared_ptr<const UnshreddedMetadataIndex>
_unshredded_metadata_index;
+ mutable ColumnPtr _normalized_prefix;
mutable ColumnVariantV2::MutablePtr _materialized;
mutable ColumnVariantV2::MutablePtr _serialized;
};
diff --git a/be/test/format_v2/parquet/variant_column_reader_test.cpp
b/be/test/format_v2/parquet/variant_column_reader_test.cpp
index 3a19c18d649..311094a90e3 100644
--- a/be/test/format_v2/parquet/variant_column_reader_test.cpp
+++ b/be/test/format_v2/parquet/variant_column_reader_test.cpp
@@ -68,6 +68,16 @@ MutableColumnPtr nullable_strings(const
std::vector<StringRef>& values,
return ColumnNullable::create(std::move(data), std::move(null_map));
}
+Status extract_object_key(const ColumnVariantV2& source, const ColumnNullable&
nullable,
+ StringRef key, ColumnPtr* result) {
+ const std::array segments {VariantElementV2PathSegment::object_key(key)};
+ std::unique_ptr<ResolvedVariantElementV2Path> path;
+ RETURN_IF_ERROR(resolve_variant_element_v2_path(segments, &path));
+ const auto& null_map = nullable.get_null_map_data();
+ return extract_variant_element_v2(
+ source, *path, std::span<const uint8_t>(null_map.data(),
null_map.size()), result);
+}
+
ParquetColumnSchema unshredded_schema() {
ParquetColumnSchema schema;
schema.name = "payload";
@@ -495,6 +505,542 @@ TEST(VariantColumnReaderTest,
UnshreddedRowsPreserveSqlNullAndVariantNull) {
EXPECT_TRUE(variants.get_value_ref(2).is_null());
}
+TEST(VariantColumnReaderTest,
UnshreddedElementChainSeeksWithoutRootReconstruction) {
+ VariantBatchBuilder builder;
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("commit"));
+ auto commit = row.start_object();
+ commit.add_key(StringRef("collection"));
+ row.add_string(StringRef("app.bsky.feed.post"));
+ commit.add_key(StringRef("operation"));
+ row.add_string(StringRef("create"));
+ commit.finish();
+ root.finish();
+ row.finish();
+ }
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("commit"));
+ auto commit = row.start_object();
+ commit.add_key(StringRef("collection"));
+ row.add_null();
+ commit.add_key(StringRef("operation"));
+ row.add_string(StringRef("create"));
+ commit.finish();
+ root.finish();
+ row.finish();
+ }
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("commit"));
+ auto commit = row.start_object();
+ commit.finish();
+ root.finish();
+ row.finish();
+ }
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("other"));
+ row.add_int(1);
+ root.finish();
+ row.finish();
+ }
+ {
+ auto row = builder.begin_row();
+ row.add_int(7);
+ row.finish();
+ }
+ VariantBatchBuilder batch = builder.finish_batch();
+ std::vector<StringRef> metadata;
+ std::vector<StringRef> values;
+ for (size_t row = 0; row < batch.num_rows(); ++row) {
+ const VariantRef value = batch.value_at(row);
+ metadata.emplace_back(value.metadata.data, value.metadata.size);
+ values.push_back(value.value);
+ }
+
+ MutableColumns fields;
+ fields.push_back(nullable_strings(metadata,
std::vector<uint8_t>(batch.num_rows(), 0)));
+ fields.push_back(nullable_strings(values,
std::vector<uint8_t>(batch.num_rows(), 0)));
+ auto physical = root_wrapper(std::move(fields), {0, 0, 0, 0, 1});
+
+ RuntimeProfile runtime_profile("unshredded-direct-path");
+ ParquetProfile parquet_profile;
+ parquet_profile.init(&runtime_profile);
+ auto output =
make_nullable(std::make_shared<DataTypeVariantV2>())->create_column();
+ const Status materialize_status = materialize_variant_rows(
+ unshredded_schema(), *physical, output,
parquet_profile.column_reader_profile());
+ ASSERT_TRUE(materialize_status.ok()) << materialize_status;
+ const auto& root_nullable = assert_cast<const ColumnNullable&>(*output);
+ const auto& root_variants =
+ assert_cast<const
ColumnVariantV2&>(root_nullable.get_nested_column());
+
+ auto extract_key = [](const ColumnVariantV2& source, const ColumnNullable&
nullable,
+ StringRef key) {
+ const std::array segments
{VariantElementV2PathSegment::object_key(key)};
+ std::unique_ptr<ResolvedVariantElementV2Path> path;
+ EXPECT_TRUE(resolve_variant_element_v2_path(segments, &path).ok());
+ ColumnPtr result;
+ const auto& null_map = nullable.get_null_map_data();
+ const Status status = extract_variant_element_v2(
+ source, *path, std::span<const uint8_t>(null_map.data(),
null_map.size()), &result);
+ EXPECT_TRUE(status.ok()) << status;
+ return result;
+ };
+
+ ColumnPtr commit_result = extract_key(root_variants, root_nullable,
StringRef("commit"));
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
5);
+ const auto& commit_nullable = assert_cast<const
ColumnNullable&>(*commit_result);
+ const auto& commit_variants =
+ assert_cast<const
ColumnVariantV2&>(commit_nullable.get_nested_column());
+ EXPECT_TRUE(commit_variants.is_shredded());
+ ColumnPtr collection_result =
+ extract_key(commit_variants, commit_nullable,
StringRef("collection"));
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
10);
+ const auto& collection_nullable = assert_cast<const
ColumnNullable&>(*collection_result);
+ const auto& collection_variants =
+ assert_cast<const
ColumnVariantV2&>(collection_nullable.get_nested_column());
+
+ EXPECT_EQ(collection_nullable.get_null_map_data(), (NullMap {0, 0, 1, 1,
1}));
+ ASSERT_TRUE(collection_variants.is_typed());
+ EXPECT_TRUE(collection_variants.typed_type()->equals(DataTypeString()));
+ const auto& typed_collection =
+ assert_cast<const
ColumnNullable&>(collection_variants.typed_column());
+ EXPECT_EQ(typed_collection.get_null_map_data(), (NullMap {0, 1, 1, 1, 1}));
+ const auto& collection_strings =
+ assert_cast<const
ColumnString&>(typed_collection.get_nested_column());
+ EXPECT_EQ(collection_strings.get_data_at(0),
StringRef("app.bsky.feed.post"));
+
+ ColumnPtr repeated_commit_result =
+ extract_key(root_variants, root_nullable, StringRef("commit"));
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
15);
+ const auto& repeated_commit_nullable =
+ assert_cast<const ColumnNullable&>(*repeated_commit_result);
+ const auto& repeated_commit_variants =
+ assert_cast<const
ColumnVariantV2&>(repeated_commit_nullable.get_nested_column());
+ ColumnPtr repeated_collection_result = extract_key(
+ repeated_commit_variants, repeated_commit_nullable,
StringRef("collection"));
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
20);
+ const auto& repeated_collection_nullable =
+ assert_cast<const ColumnNullable&>(*repeated_collection_result);
+ const auto& repeated_collection_variants =
+ assert_cast<const
ColumnVariantV2&>(repeated_collection_nullable.get_nested_column());
+ const auto& repeated_typed_collection =
+ assert_cast<const
ColumnNullable&>(repeated_collection_variants.typed_column());
+ const auto& repeated_collection_strings =
+ assert_cast<const
ColumnString&>(repeated_typed_collection.get_nested_column());
+ EXPECT_EQ(repeated_collection_nullable.get_null_map_data(), (NullMap {0,
0, 1, 1, 1}));
+ EXPECT_EQ(repeated_typed_collection.get_null_map_data(), (NullMap {0, 1,
1, 1, 1}));
+ EXPECT_EQ(repeated_collection_strings.get_data_at(0),
StringRef("app.bsky.feed.post"));
+
+ ASSERT_NE(runtime_profile.get_counter("VariantReconstructedRows"),
nullptr);
+
EXPECT_EQ(runtime_profile.get_counter("VariantReconstructedRows")->value(), 0);
+ ASSERT_NE(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows"),
nullptr);
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
20);
+ ASSERT_NE(runtime_profile.get_counter("VariantUnshreddedDirectSeekBytes"),
nullptr);
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekBytes")->value(),
36);
+ ASSERT_NE(runtime_profile.get_counter("VariantUnshreddedPrefixReuseRows"),
nullptr);
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedPrefixReuseRows")->value(),
10);
+ ASSERT_NE(runtime_profile.get_counter("VariantDirectLeafRows"), nullptr);
+ EXPECT_EQ(runtime_profile.get_counter("VariantDirectLeafRows")->value(),
10);
+ ASSERT_NE(runtime_profile.get_counter("VariantDirectSubtreeRows"),
nullptr);
+
EXPECT_EQ(runtime_profile.get_counter("VariantDirectSubtreeRows")->value(), 10);
+
+ IColumn::Filter keep_first_two {1, 1, 0, 0, 0};
+ ColumnPtr filtered_root = root_nullable.filter(keep_first_two, 2);
+ const auto& filtered_nullable = assert_cast<const
ColumnNullable&>(*filtered_root);
+ const auto& filtered_variants =
+ assert_cast<const
ColumnVariantV2&>(filtered_nullable.get_nested_column());
+ ColumnPtr filtered_commit_result =
+ extract_key(filtered_variants, filtered_nullable,
StringRef("commit"));
+ const auto& filtered_commit_nullable =
+ assert_cast<const ColumnNullable&>(*filtered_commit_result);
+ const auto& filtered_commit_variants =
+ assert_cast<const
ColumnVariantV2&>(filtered_commit_nullable.get_nested_column());
+ ColumnPtr filtered_collection_result = extract_key(
+ filtered_commit_variants, filtered_commit_nullable,
StringRef("collection"));
+ EXPECT_EQ(assert_cast<const
ColumnNullable&>(*filtered_collection_result).get_null_map_data(),
+ (NullMap {0, 0}));
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
24);
+
+ ColumnPtr filtered_operation_result =
+ extract_key(filtered_commit_variants, filtered_commit_nullable,
StringRef("operation"));
+ const auto& filtered_operation_variants = assert_cast<const
ColumnVariantV2&>(
+ assert_cast<const
ColumnNullable&>(*filtered_operation_result).get_nested_column());
+ const auto& filtered_operation_strings = assert_cast<const ColumnString&>(
+ assert_cast<const
ColumnNullable&>(filtered_operation_variants.typed_column())
+ .get_nested_column());
+ EXPECT_EQ(filtered_operation_strings.get_data_at(0), StringRef("create"));
+ EXPECT_EQ(filtered_operation_strings.get_data_at(1), StringRef("create"));
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
26);
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedPrefixReuseRows")->value(),
14);
+
+ ColumnPtr ranged_root = root_nullable.cut(0, 2);
+ const auto& ranged_nullable = assert_cast<const
ColumnNullable&>(*ranged_root);
+ const auto& ranged_variants =
+ assert_cast<const
ColumnVariantV2&>(ranged_nullable.get_nested_column());
+ ColumnPtr ranged_commit_result =
+ extract_key(ranged_variants, ranged_nullable, StringRef("commit"));
+ const auto& ranged_commit_nullable = assert_cast<const
ColumnNullable&>(*ranged_commit_result);
+ const auto& ranged_commit_variants =
+ assert_cast<const
ColumnVariantV2&>(ranged_commit_nullable.get_nested_column());
+ (void)extract_key(ranged_commit_variants, ranged_commit_nullable,
StringRef("collection"));
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
30);
+
+ auto gathered_root = root_nullable.clone_empty();
+ const std::array<uint32_t, 2> reversed_rows {1, 0};
+ gathered_root->insert_indices_from(root_nullable, reversed_rows.begin(),
reversed_rows.end());
+ const auto& gathered_nullable = assert_cast<const
ColumnNullable&>(*gathered_root);
+ const auto& gathered_variants =
+ assert_cast<const
ColumnVariantV2&>(gathered_nullable.get_nested_column());
+ ColumnPtr gathered_commit_result =
+ extract_key(gathered_variants, gathered_nullable,
StringRef("commit"));
+ const auto& gathered_commit_nullable =
+ assert_cast<const ColumnNullable&>(*gathered_commit_result);
+ const auto& gathered_commit_variants =
+ assert_cast<const
ColumnVariantV2&>(gathered_commit_nullable.get_nested_column());
+ (void)extract_key(gathered_commit_variants, gathered_commit_nullable,
StringRef("collection"));
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekRows")->value(),
34);
+}
+
+TEST(VariantColumnReaderTest, UnshreddedIntegerLeafSeeksAsTypedBigInt) {
+ VariantBatchBuilder builder;
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("time_us"));
+ row.add_null();
+ root.finish();
+ row.finish();
+ }
+ for (const int64_t value :
+ {int64_t {7}, int64_t {1} << 8, int64_t {1} << 20, int64_t {1} <<
40}) {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("time_us"));
+ row.add_int(value);
+ root.finish();
+ row.finish();
+ }
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("other"));
+ row.add_int(1);
+ root.finish();
+ row.finish();
+ }
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("time_us"));
+ row.add_int(9);
+ root.finish();
+ row.finish();
+ }
+ VariantBatchBuilder batch = builder.finish_batch();
+ std::vector<StringRef> metadata;
+ std::vector<StringRef> values;
+ for (size_t row = 0; row < batch.num_rows(); ++row) {
+ const VariantRef value = batch.value_at(row);
+ metadata.emplace_back(value.metadata.data, value.metadata.size);
+ values.push_back(value.value);
+ }
+
+ MutableColumns fields;
+ fields.push_back(nullable_strings(metadata,
std::vector<uint8_t>(batch.num_rows(), 0)));
+ fields.push_back(nullable_strings(values,
std::vector<uint8_t>(batch.num_rows(), 0)));
+ auto physical = root_wrapper(std::move(fields), {0, 0, 0, 0, 0, 0, 1});
+
+ RuntimeProfile runtime_profile("unshredded-integer-direct-path");
+ ParquetProfile parquet_profile;
+ parquet_profile.init(&runtime_profile);
+ auto output =
make_nullable(std::make_shared<DataTypeVariantV2>())->create_column();
+ ASSERT_TRUE(materialize_variant_rows(unshredded_schema(), *physical,
output,
+
parquet_profile.column_reader_profile())
+ .ok());
+ const auto& root_nullable = assert_cast<const ColumnNullable&>(*output);
+ const auto& root_variants =
+ assert_cast<const
ColumnVariantV2&>(root_nullable.get_nested_column());
+ const std::array segments
{VariantElementV2PathSegment::object_key(StringRef("time_us"))};
+ std::unique_ptr<ResolvedVariantElementV2Path> path;
+ ASSERT_TRUE(resolve_variant_element_v2_path(segments, &path).ok());
+ ColumnPtr result;
+ const Status extract_status = extract_variant_element_v2(
+ root_variants, *path,
+ std::span<const uint8_t>(root_nullable.get_null_map_data().data(),
+ root_nullable.get_null_map_data().size()),
+ &result);
+ ASSERT_TRUE(extract_status.ok()) << extract_status;
+
+ const auto& result_nullable = assert_cast<const ColumnNullable&>(*result);
+ EXPECT_EQ(result_nullable.get_null_map_data(), (NullMap {0, 0, 0, 0, 0, 1,
1}));
+ const auto& result_variants =
+ assert_cast<const
ColumnVariantV2&>(result_nullable.get_nested_column());
+ ASSERT_TRUE(result_variants.is_typed());
+ EXPECT_TRUE(result_variants.typed_type()->equals(DataTypeInt64()));
+ const auto& typed = assert_cast<const
ColumnNullable&>(result_variants.typed_column());
+ EXPECT_EQ(typed.get_null_map_data(), (NullMap {1, 0, 0, 0, 0, 1, 1}));
+ const auto& integers = assert_cast<const
ColumnInt64&>(typed.get_nested_column());
+ EXPECT_EQ(integers.get_data()[1], 7);
+ EXPECT_EQ(integers.get_data()[2], int64_t {1} << 8);
+ EXPECT_EQ(integers.get_data()[3], int64_t {1} << 20);
+ EXPECT_EQ(integers.get_data()[4], int64_t {1} << 40);
+
EXPECT_EQ(runtime_profile.get_counter("VariantReconstructedRows")->value(), 0);
+ EXPECT_EQ(runtime_profile.get_counter("VariantDirectLeafRows")->value(),
7);
+
EXPECT_EQ(runtime_profile.get_counter("VariantDirectSubtreeRows")->value(), 0);
+
EXPECT_EQ(runtime_profile.get_counter("VariantUnshreddedDirectSeekBytes")->value(),
32);
+}
+
+TEST(VariantColumnReaderTest, UnshreddedExplicitIntegerWidthStaysEncoded) {
+ VariantBatchBuilder builder;
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("value"));
+ row.add_scalar(VariantScalarRef::integer(7, 8));
+ root.finish();
+ row.finish();
+ VariantBatchBuilder batch = builder.finish_batch();
+ const VariantRef encoded = batch.value_at(0);
+
+ MutableColumns fields;
+ fields.push_back(
+ nullable_strings({StringRef(encoded.metadata.data,
encoded.metadata.size)}, {0}));
+ fields.push_back(nullable_strings({encoded.value}, {0}));
+ auto physical = root_wrapper(std::move(fields));
+ auto output =
make_nullable(std::make_shared<DataTypeVariantV2>())->create_column();
+ ASSERT_TRUE(materialize_variant_rows(unshredded_schema(), *physical,
output).ok());
+ const auto& root_nullable = assert_cast<const ColumnNullable&>(*output);
+ const auto& root_variants =
+ assert_cast<const
ColumnVariantV2&>(root_nullable.get_nested_column());
+ const std::array segments
{VariantElementV2PathSegment::object_key(StringRef("value"))};
+ std::unique_ptr<ResolvedVariantElementV2Path> path;
+ ASSERT_TRUE(resolve_variant_element_v2_path(segments, &path).ok());
+ ColumnPtr result;
+ ASSERT_TRUE(extract_variant_element_v2(
+ root_variants, *path,
+ std::span<const
uint8_t>(root_nullable.get_null_map_data().data(),
+
root_nullable.get_null_map_data().size()),
+ &result)
+ .ok());
+
+ const auto& result_variants = assert_cast<const ColumnVariantV2&>(
+ assert_cast<const ColumnNullable&>(*result).get_nested_column());
+ ASSERT_TRUE(result_variants.is_shredded());
+ EXPECT_EQ(result_variants.get_value_ref(0).primitive_id(),
VariantPrimitiveId::INT64);
+ EXPECT_EQ(result_variants.get_value_ref(0).get_int(), 7);
+}
+
+TEST(VariantColumnReaderTest,
UnshreddedDirectPathsSeparateMetadataDictionaries) {
+ VariantBatchBuilder first_builder;
+ auto first_row = first_builder.begin_row();
+ auto first_root = first_row.start_object();
+ first_root.add_key(StringRef("a"));
+ first_row.add_int(1);
+ first_root.add_key(StringRef("target"));
+ first_row.add_int(11);
+ first_root.finish();
+ first_row.finish();
+ VariantBatchBuilder first_batch = first_builder.finish_batch();
+
+ VariantBatchBuilder second_builder;
+ auto second_row = second_builder.begin_row();
+ auto second_root = second_row.start_object();
+ second_root.add_key(StringRef("target"));
+ second_row.add_int(22);
+ second_root.add_key(StringRef("z"));
+ second_row.add_int(2);
+ second_root.finish();
+ second_row.finish();
+ VariantBatchBuilder second_batch = second_builder.finish_batch();
+
+ VariantBatchBuilder missing_builder;
+ auto missing_row = missing_builder.begin_row();
+ auto missing_root = missing_row.start_object();
+ missing_root.add_key(StringRef("other"));
+ missing_row.add_int(3);
+ missing_root.finish();
+ missing_row.finish();
+ VariantBatchBuilder missing_batch = missing_builder.finish_batch();
+
+ const VariantRef first = first_batch.value_at(0);
+ const VariantRef second = second_batch.value_at(0);
+ const VariantRef missing = missing_batch.value_at(0);
+ const std::vector<VariantRef> rows {first, second, first, missing, second};
+ std::vector<StringRef> metadata;
+ std::vector<StringRef> values;
+ for (const VariantRef value : rows) {
+ metadata.emplace_back(value.metadata.data, value.metadata.size);
+ values.push_back(value.value);
+ }
+ MutableColumns fields;
+ fields.push_back(nullable_strings(metadata,
std::vector<uint8_t>(rows.size(), 0)));
+ fields.push_back(nullable_strings(values,
std::vector<uint8_t>(rows.size(), 0)));
+ auto physical = root_wrapper(std::move(fields), NullMap(rows.size(), 0));
+ auto output =
make_nullable(std::make_shared<DataTypeVariantV2>())->create_column();
+ ASSERT_TRUE(materialize_variant_rows(unshredded_schema(), *physical,
output).ok());
+ const auto& root_nullable = assert_cast<const ColumnNullable&>(*output);
+ const auto& root_variants =
+ assert_cast<const
ColumnVariantV2&>(root_nullable.get_nested_column());
+
+ ColumnPtr target_result;
+ ASSERT_TRUE(
+ extract_object_key(root_variants, root_nullable,
StringRef("target"), &target_result)
+ .ok());
+ const auto& target_nullable = assert_cast<const
ColumnNullable&>(*target_result);
+ EXPECT_EQ(target_nullable.get_null_map_data(), (NullMap {0, 0, 0, 1, 0}));
+ const auto& target_variants =
+ assert_cast<const
ColumnVariantV2&>(target_nullable.get_nested_column());
+ ASSERT_TRUE(target_variants.is_typed());
+ const auto& target_values = assert_cast<const ColumnInt64&>(
+ assert_cast<const
ColumnNullable&>(target_variants.typed_column()).get_nested_column());
+ EXPECT_EQ(target_values.get_data()[0], 11);
+ EXPECT_EQ(target_values.get_data()[1], 22);
+ EXPECT_EQ(target_values.get_data()[2], 11);
+ EXPECT_EQ(target_values.get_data()[4], 22);
+
+ ColumnPtr a_result;
+ ASSERT_TRUE(extract_object_key(root_variants, root_nullable,
StringRef("a"), &a_result).ok());
+ EXPECT_EQ(assert_cast<const
ColumnNullable&>(*a_result).get_null_map_data(),
+ (NullMap {0, 1, 0, 1, 1}));
+}
+
+TEST(VariantColumnReaderTest,
UnshreddedCachedMissingStillValidatesEveryObject) {
+ VariantBatchBuilder builder;
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("other"));
+ row.add_int(1);
+ root.finish();
+ row.finish();
+ VariantBatchBuilder batch = builder.finish_batch();
+ const VariantRef valid = batch.value_at(0);
+ std::string truncated(valid.value.data, valid.value.size - 1);
+
+ MutableColumns fields;
+ fields.push_back(nullable_strings({StringRef(valid.metadata.data,
valid.metadata.size),
+ StringRef(valid.metadata.data,
valid.metadata.size)},
+ {0, 0}));
+ fields.push_back(nullable_strings({valid.value, StringRef(truncated)}, {0,
0}));
+ auto physical = root_wrapper(std::move(fields), {0, 0});
+ auto output =
make_nullable(std::make_shared<DataTypeVariantV2>())->create_column();
+ ASSERT_TRUE(materialize_variant_rows(unshredded_schema(), *physical,
output).ok());
+ const auto& root_nullable = assert_cast<const ColumnNullable&>(*output);
+ const auto& root_variants =
+ assert_cast<const
ColumnVariantV2&>(root_nullable.get_nested_column());
+
+ ColumnPtr result;
+ const Status status =
+ extract_object_key(root_variants, root_nullable,
StringRef("missing"), &result);
+ EXPECT_FALSE(status.ok());
+ EXPECT_EQ(status.code(), ErrorCode::INVALID_ARGUMENT);
+ EXPECT_NE(status.to_string().find("Truncated Variant value"),
std::string::npos);
+ EXPECT_FALSE(result);
+}
+
+TEST(VariantColumnReaderTest,
UnshreddedIntegerLeafReportsCorruptPayloadOnFallback) {
+ VariantBatchBuilder builder;
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("time_us"));
+ row.add_int(7);
+ root.finish();
+ row.finish();
+ VariantBatchBuilder batch = builder.finish_batch();
+ const VariantRef encoded = batch.value_at(0);
+ VariantRef encoded_integer;
+ ASSERT_TRUE(encoded.object_find(StringRef("time_us"), &encoded_integer));
+
+ std::string invalid_value(encoded.value.data, encoded.value.size);
+ const size_t integer_offset = encoded_integer.value.data -
encoded.value.data;
+ ASSERT_LT(integer_offset, invalid_value.size());
+ invalid_value[integer_offset] = static_cast<char>(
+ static_cast<uint8_t>(VariantPrimitiveId::INT64) <<
VARIANT_VALUE_HEADER_SHIFT);
+
+ MutableColumns fields;
+ fields.push_back(
+ nullable_strings({StringRef(encoded.metadata.data,
encoded.metadata.size)}, {0}));
+ fields.push_back(nullable_strings({StringRef(invalid_value)}, {0}));
+ auto physical = root_wrapper(std::move(fields));
+ auto output =
make_nullable(std::make_shared<DataTypeVariantV2>())->create_column();
+ ASSERT_TRUE(materialize_variant_rows(unshredded_schema(), *physical,
output).ok());
+ const auto& root_nullable = assert_cast<const ColumnNullable&>(*output);
+ const auto& root_variants =
+ assert_cast<const
ColumnVariantV2&>(root_nullable.get_nested_column());
+ const std::array segments
{VariantElementV2PathSegment::object_key(StringRef("time_us"))};
+ std::unique_ptr<ResolvedVariantElementV2Path> path;
+ ASSERT_TRUE(resolve_variant_element_v2_path(segments, &path).ok());
+ ColumnPtr result;
+ const Status extract_status = extract_variant_element_v2(
+ root_variants, *path,
+ std::span<const uint8_t>(root_nullable.get_null_map_data().data(),
+ root_nullable.get_null_map_data().size()),
+ &result);
+ EXPECT_FALSE(extract_status.ok());
+ EXPECT_EQ(extract_status.code(), ErrorCode::INVALID_ARGUMENT);
+ EXPECT_NE(extract_status.to_string().find("Truncated Variant value"),
std::string::npos);
+ EXPECT_FALSE(result);
+}
+
+TEST(VariantColumnReaderTest, UnshreddedMixedScalarLeafFallsBackToSubtree) {
+ VariantBatchBuilder builder;
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("value"));
+ row.add_int(7);
+ root.finish();
+ row.finish();
+ }
+ {
+ auto row = builder.begin_row();
+ auto root = row.start_object();
+ root.add_key(StringRef("value"));
+ row.add_string(StringRef("seven"));
+ root.finish();
+ row.finish();
+ }
+ VariantBatchBuilder batch = builder.finish_batch();
+ std::vector<StringRef> metadata;
+ std::vector<StringRef> values;
+ for (size_t row = 0; row < batch.num_rows(); ++row) {
+ const VariantRef value = batch.value_at(row);
+ metadata.emplace_back(value.metadata.data, value.metadata.size);
+ values.push_back(value.value);
+ }
+
+ MutableColumns fields;
+ fields.push_back(nullable_strings(metadata, {0, 0}));
+ fields.push_back(nullable_strings(values, {0, 0}));
+ auto physical = root_wrapper(std::move(fields), {0, 0});
+ auto output =
make_nullable(std::make_shared<DataTypeVariantV2>())->create_column();
+ ASSERT_TRUE(materialize_variant_rows(unshredded_schema(), *physical,
output).ok());
+ const auto& root_nullable = assert_cast<const ColumnNullable&>(*output);
+ const auto& root_variants =
+ assert_cast<const
ColumnVariantV2&>(root_nullable.get_nested_column());
+ const std::array segments
{VariantElementV2PathSegment::object_key(StringRef("value"))};
+ std::unique_ptr<ResolvedVariantElementV2Path> path;
+ ASSERT_TRUE(resolve_variant_element_v2_path(segments, &path).ok());
+ ColumnPtr result;
+ ASSERT_TRUE(extract_variant_element_v2(
+ root_variants, *path,
+ std::span<const
uint8_t>(root_nullable.get_null_map_data().data(),
+
root_nullable.get_null_map_data().size()),
+ &result)
+ .ok());
+
+ const auto& result_variants = assert_cast<const ColumnVariantV2&>(
+ assert_cast<const ColumnNullable&>(*result).get_nested_column());
+ ASSERT_TRUE(result_variants.is_shredded());
+ EXPECT_EQ(result_variants.get_value_ref(0).get_int(), 7);
+ EXPECT_EQ(result_variants.get_value_ref(1).get_string(),
StringRef("seven"));
+}
+
TEST(VariantColumnReaderTest,
RequiredPhysicalGroupAppendsToNullableExternalSlot) {
const std::array<char, 2> int_seven {
static_cast<char>(static_cast<uint8_t>(VariantPrimitiveId::INT8)
diff --git a/be/test/util/variant/variant_value_test.cpp
b/be/test/util/variant/variant_value_test.cpp
index 0d8cc8e6cdb..4c4e44ca579 100644
--- a/be/test/util/variant/variant_value_test.cpp
+++ b/be/test/util/variant/variant_value_test.cpp
@@ -388,6 +388,13 @@ TEST(VariantValueTest,
ObjectLookupSortedAndUnsortedMetadata) {
EXPECT_FALSE(found.get_bool());
EXPECT_FALSE(sorted_ref.object_find(StringRef("missing", 7), &found));
+ std::string invalid_sorted_metadata = metadata({"a", "b", "c"}, true);
+ invalid_sorted_metadata[3] = 2;
+ invalid_sorted_metadata[4] = 1;
+ const std::string invalid_target = object_value({1}, {0}, {true_value});
+ EXPECT_THROW(value_ref(invalid_sorted_metadata,
invalid_target).object_find_by_id(1, &found),
+ Exception);
+
const std::string unsorted_metadata = metadata({"z", "a", "m"}, false);
const std::string unsorted_object =
object_value({1, 2, 0}, {0, 1, 2},
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]