hubgeter commented on code in PR #66758:
URL: https://github.com/apache/doris/pull/66758#discussion_r3782336359
##########
be/src/format_v2/parquet/reader/variant_column_reader.cpp:
##########
@@ -514,6 +516,200 @@ void import_unshredded_variant_range(const
ParquetColumnSchema& schema, const Co
appender.append(std::span<const VariantRef>(encoded_rows));
}
+struct DirectResidualSeekResult {
+ ColumnPtr column;
+ int64_t selected_value_bytes = 0;
+};
+
+struct UnshreddedMetadataCache {
+ static constexpr uint32_t NO_METADATA =
std::numeric_limits<uint32_t>::max();
+
+ DorisVector<VariantMetadataRef> metadatas;
+ DorisVector<uint32_t> row_metadata_ids;
+
+ size_t allocated_bytes() const {
+ return metadatas.capacity() * sizeof(VariantMetadataRef) +
+ row_metadata_ids.capacity() * sizeof(uint32_t);
+ }
+};
+
+std::shared_ptr<const UnshreddedMetadataCache> build_unshredded_metadata_cache(
+ const ParquetColumnSchema& schema, const IColumn& physical, size_t
metadata_index) {
+ 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);
+ if (structure.tuple_size() != schema.children.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} physical
field count mismatch",
+ schema.name);
+ }
+
+ using MetadataIndex =
+ 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>>>;
+ auto cache = std::make_shared<UnshreddedMetadataCache>();
+ cache->row_metadata_ids.resize(physical.size(),
UnshreddedMetadataCache::NO_METADATA);
+ MetadataIndex metadata_index_by_value;
+ std::optional<uint32_t> previous_metadata_id;
+
+ auto same_metadata = [&](uint32_t id, StringRef bytes) {
+ const VariantMetadataRef cached = cache->metadatas[id];
+ return StringRef(cached.data, cached.size) == bytes;
+ };
+ 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 =
cell_at(structure.get_column(metadata_index), row);
+ if (metadata_cell.is_null) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} has
null metadata at row {}",
+ schema.name, row);
+ }
+ const StringRef metadata_bytes =
metadata_cell.column->get_data_at(row);
+ if (previous_metadata_id.has_value() &&
+ same_metadata(*previous_metadata_id, metadata_bytes)) {
+ cache->row_metadata_ids[row] = *previous_metadata_id;
+ continue;
+ }
+
+ const std::string_view key(metadata_bytes.data == nullptr ? "" :
metadata_bytes.data,
+ metadata_bytes.size);
+ if (!metadata_index_by_value.empty()) {
+ if (const auto found = metadata_index_by_value.find(key);
+ found != metadata_index_by_value.end()) {
+ cache->row_metadata_ids[row] = found->second;
+ previous_metadata_id = found->second;
+ continue;
+ }
+ }
+
+ VariantMetadataRef metadata {.data = metadata_bytes.data, .size =
metadata_bytes.size};
+ validate_variant_metadata(metadata);
+ if (cache->metadatas.size() == std::numeric_limits<uint32_t>::max()) {
+ throw Exception(ErrorCode::INVALID_ARGUMENT,
+ "Parquet Variant metadata dictionary exceeds the
uint32 id limit");
+ }
+ if (cache->metadatas.size() == 1 && metadata_index_by_value.empty()) {
+ const VariantMetadataRef first = cache->metadatas.front();
+ metadata_index_by_value.emplace(
+ std::string_view(first.data == nullptr ? "" : first.data,
first.size), 0);
+ }
+ const auto id = static_cast<uint32_t>(cache->metadatas.size());
+ cache->metadatas.push_back(metadata);
+ if (!metadata_index_by_value.empty()) {
+ metadata_index_by_value.emplace(key, id);
+ }
+ cache->row_metadata_ids[row] = id;
+ previous_metadata_id = id;
+ }
+ return cache;
+}
+
+DirectResidualSeekResult seek_unshredded_variant_path(
+ const ParquetColumnSchema& schema, const IColumn& physical, size_t
value_index,
+ const UnshreddedMetadataCache& metadata_cache,
+ std::span<const VariantShreddedPathSegment> path) {
+ 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);
+ if (structure.tuple_size() != schema.children.size()) {
+ throw Exception(ErrorCode::CORRUPTION, "Parquet Variant {} physical
field count mismatch",
+ schema.name);
+ }
+
+ DORIS_CHECK_EQ(metadata_cache.row_metadata_ids.size(), physical.size());
+ if (!metadata_cache.metadatas.empty() &&
+ path.size() > std::numeric_limits<size_t>::max() /
metadata_cache.metadatas.size()) {
+ throw Exception(ErrorCode::INVALID_ARGUMENT,
+ "Parquet Variant direct-seek path cache size overflows
size_t");
+ }
+ DorisVector<int64_t> object_ids(metadata_cache.metadatas.size() *
path.size(), -1);
+ for (size_t metadata_id = 0; metadata_id <
metadata_cache.metadatas.size(); ++metadata_id) {
+ for (size_t position = 0; position < path.size(); ++position) {
+ if (path[position].kind ==
VariantShreddedPathSegment::Kind::OBJECT_KEY) {
+ object_ids[metadata_id * path.size() + position] =
+
metadata_cache.metadatas[metadata_id].find_key(path[position].key);
+ }
+ }
+ }
+ DorisVector<VariantRef> selected_rows;
+ selected_rows.reserve(physical.size());
+ auto nulls = ColumnUInt8::create();
+ nulls->reserve(physical.size());
+ std::vector<uint32_t> object_offset_scratch;
+ int64_t selected_value_bytes = 0;
+
+ auto add_selected_bytes = [&](size_t bytes) {
+ DORIS_CHECK_LE(bytes,
static_cast<size_t>(std::numeric_limits<int64_t>::max() -
+ selected_value_bytes));
+ selected_value_bytes += static_cast<int64_t>(bytes);
+ };
+ auto append_missing = [&](VariantMetadataRef metadata) {
+ selected_rows.push_back({.metadata = metadata,
+ .value = {VARIANT_NULL_VALUE.data(),
VARIANT_NULL_VALUE.size()}});
+ nulls->insert_value(1);
+ add_selected_bytes(VARIANT_NULL_VALUE.size());
+ };
+ for (size_t row = 0; row < physical.size(); ++row) {
+ if (outer_nullable != nullptr &&
outer_nullable->get_null_map_data()[row] != 0) {
+ append_missing(
+ {.data = VARIANT_EMPTY_METADATA.data(), .size =
VARIANT_EMPTY_METADATA.size()});
+ continue;
+ }
+
+ const uint32_t metadata_id = metadata_cache.row_metadata_ids[row];
+ DORIS_CHECK_NE(metadata_id, UnshreddedMetadataCache::NO_METADATA);
+ DORIS_CHECK_LT(metadata_id, metadata_cache.metadatas.size());
+ const VariantMetadataRef metadata =
metadata_cache.metadatas[metadata_id];
+ const Cell value_cell = cell_at(structure.get_column(value_index),
row);
+ if (value_cell.is_null) {
+ append_missing(metadata);
+ continue;
+ }
+
+ VariantRef current {.metadata = metadata, .value =
value_cell.column->get_data_at(row)};
Review Comment:
Fixed in 715e82f7eea. Direct traversal now validates the exact root envelope
and scalar semantics before any miss, validates every accessed container table,
and recursively validates the selected subtree while still skipping unrelated
sibling payloads. Focused malformed-root and container tests were added.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]