be/src/exprs/function/function_variant_element_v2.cpp
Line | Count | Source |
1 | | // Licensed to the Apache Software Foundation (ASF) under one |
2 | | // or more contributor license agreements. See the NOTICE file |
3 | | // distributed with this work for additional information |
4 | | // regarding copyright ownership. The ASF licenses this file |
5 | | // to you under the Apache License, Version 2.0 (the |
6 | | // "License"); you may not use this file except in compliance |
7 | | // with the License. You may obtain a copy of the License at |
8 | | // |
9 | | // http://www.apache.org/licenses/LICENSE-2.0 |
10 | | // |
11 | | // Unless required by applicable law or agreed to in writing, |
12 | | // software distributed under the License is distributed on an |
13 | | // "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
14 | | // KIND, either express or implied. See the License for the |
15 | | // specific language governing permissions and limitations |
16 | | // under the License. |
17 | | |
18 | | #include "exprs/function/function_variant_element_v2.h" |
19 | | |
20 | | #include <limits> |
21 | | #include <utility> |
22 | | |
23 | | #include "common/check.h" |
24 | | #include "common/exception.h" |
25 | | #include "core/column/column_nullable.h" |
26 | | #include "core/column/column_vector.h" |
27 | | #include "core/column/variant_v2/column_variant_v2.h" |
28 | | #include "core/custom_allocator.h" |
29 | | #include "core/pod_array.h" |
30 | | #include "core/value/variant/variant_batch_builder.h" |
31 | | #include "core/value/variant/variant_parquet_encoding.h" |
32 | | |
33 | | namespace doris { |
34 | | |
35 | | namespace { |
36 | | |
37 | | struct OwnedPathSegment { |
38 | | VariantElementV2PathSegment::Kind kind; |
39 | | PaddedPODArray<char> key; |
40 | | int64_t index = 0; |
41 | | }; |
42 | | |
43 | | } // namespace |
44 | | |
45 | | struct ResolvedVariantElementV2Path::Impl { |
46 | | DorisVector<OwnedPathSegment> segments; |
47 | | }; |
48 | | |
49 | | namespace { |
50 | | |
51 | 20 | bool is_outer_null(std::span<const uint8_t> outer_nulls, size_t row) noexcept { |
52 | 20 | return !outer_nulls.empty() && outer_nulls[row] != 0; |
53 | 20 | } |
54 | | |
55 | | // Resolves object keys once per dense metadata id and traverses arrays without metadata lookup. |
56 | | class VariantPathV2BatchReader { |
57 | | public: |
58 | | VariantPathV2BatchReader(const ColumnVariantV2& source, |
59 | | const ResolvedVariantElementV2Path& path) |
60 | 13 | : _source(source.read_view()), |
61 | 13 | _path(path), |
62 | 13 | _resolved_metadata(_source.metadata_count(), 0) { |
63 | 13 | DORIS_CHECK_GT(_path.size(), 0); |
64 | 13 | if (_source.metadata_count() != 0 && |
65 | 13 | _path.size() > std::numeric_limits<size_t>::max() / _source.metadata_count()) { |
66 | 0 | throw Exception(ErrorCode::INVALID_ARGUMENT, |
67 | 0 | "Variant path cache size overflows size_t"); |
68 | 0 | } |
69 | 13 | _object_ids.assign(_source.metadata_count() * _path.size(), -1); |
70 | 13 | } |
71 | | |
72 | 19 | bool find_at(size_t row, VariantRef* const output) { |
73 | 19 | DORIS_CHECK(output != nullptr); |
74 | 19 | const uint32_t metadata_id = _source.metadata_id_at(row); |
75 | 19 | _resolve_metadata(metadata_id); |
76 | 19 | VariantRef current = _source.value_at(row); |
77 | 19 | const size_t cache_base = static_cast<size_t>(metadata_id) * _path.size(); |
78 | 43 | for (size_t position = 0; position < _path.size(); ++position) { |
79 | 30 | if (_path.kind_at(position) == VariantElementV2PathSegment::Kind::OBJECT_KEY) { |
80 | 25 | const int64_t field_id = _object_ids[cache_base + position]; |
81 | 25 | if (current.basic_type() != VariantBasicType::OBJECT || field_id < 0 || |
82 | 25 | !current.object_find_by_id(static_cast<uint32_t>(field_id), ¤t)) { |
83 | 4 | return false; |
84 | 4 | } |
85 | 25 | } else { |
86 | 5 | if (current.basic_type() != VariantBasicType::ARRAY) { |
87 | 0 | return false; |
88 | 0 | } |
89 | 5 | const int64_t requested_index = _path.array_index_at(position); |
90 | 5 | const int64_t element_count = current.num_elements(); |
91 | 5 | const int64_t resolved_index = |
92 | 5 | requested_index < 0 ? element_count + requested_index : requested_index; |
93 | 5 | if (resolved_index < 0 || resolved_index >= element_count) { |
94 | 2 | return false; |
95 | 2 | } |
96 | 3 | current = current.array_at(static_cast<uint32_t>(resolved_index)); |
97 | 3 | } |
98 | 30 | } |
99 | 13 | *output = current; |
100 | 13 | return true; |
101 | 19 | } |
102 | | |
103 | 13 | size_t metadata_count() const noexcept { return _source.metadata_count(); } |
104 | | |
105 | 19 | uint32_t metadata_id_at(size_t row) const { return _source.metadata_id_at(row); } |
106 | | |
107 | 18 | VariantMetadataRef metadata_at(uint32_t id) { return _source.metadata_at(id); } |
108 | | |
109 | | private: |
110 | 19 | void _resolve_metadata(uint32_t metadata_id) { |
111 | 19 | if (_resolved_metadata[metadata_id] != 0) { |
112 | 1 | return; |
113 | 1 | } |
114 | 18 | const VariantMetadataRef metadata = _source.metadata_at(metadata_id); |
115 | 18 | const size_t base = static_cast<size_t>(metadata_id) * _path.size(); |
116 | 47 | for (size_t position = 0; position < _path.size(); ++position) { |
117 | 29 | if (_path.kind_at(position) == VariantElementV2PathSegment::Kind::OBJECT_KEY) { |
118 | 24 | _object_ids[base + position] = metadata.find_key(_path.object_key_at(position)); |
119 | 24 | } |
120 | 29 | } |
121 | 18 | _resolved_metadata[metadata_id] = 1; |
122 | 18 | } |
123 | | |
124 | | ColumnVariantV2::ReadView _source; |
125 | | const ResolvedVariantElementV2Path& _path; |
126 | | DorisVector<int64_t> _object_ids; |
127 | | DorisVector<uint8_t> _resolved_metadata; |
128 | | }; |
129 | | |
130 | | Status extract_encoded_variant_element(const ColumnVariantV2& source, |
131 | | const ResolvedVariantElementV2Path& path, |
132 | | std::span<const uint8_t> outer_nulls, ColumnPtr* output); |
133 | | |
134 | | Status make_all_null_variant_element_result(size_t rows, ColumnPtr* output); |
135 | | |
136 | | } // namespace |
137 | | |
138 | | ResolvedVariantElementV2Path::ResolvedVariantElementV2Path(std::unique_ptr<Impl> impl) |
139 | 13 | : _impl(std::move(impl)) {} |
140 | | |
141 | 13 | ResolvedVariantElementV2Path::~ResolvedVariantElementV2Path() = default; |
142 | | ResolvedVariantElementV2Path::ResolvedVariantElementV2Path( |
143 | 0 | ResolvedVariantElementV2Path&&) noexcept = default; |
144 | | ResolvedVariantElementV2Path& ResolvedVariantElementV2Path::operator=( |
145 | 0 | ResolvedVariantElementV2Path&&) noexcept = default; |
146 | | |
147 | 268 | size_t ResolvedVariantElementV2Path::size() const noexcept { |
148 | 268 | return _impl->segments.size(); |
149 | 268 | } |
150 | | |
151 | 88 | VariantElementV2PathSegment::Kind ResolvedVariantElementV2Path::kind_at(size_t position) const { |
152 | 88 | DORIS_CHECK_LT(position, size()) << "Variant element path position is out of range"; |
153 | 88 | return _impl->segments[position].kind; |
154 | 88 | } |
155 | | |
156 | 24 | StringRef ResolvedVariantElementV2Path::object_key_at(size_t position) const { |
157 | 24 | DORIS_CHECK(kind_at(position) == VariantElementV2PathSegment::Kind::OBJECT_KEY) |
158 | 0 | << "Variant element path segment is not an object key"; |
159 | 24 | const auto& key = _impl->segments[position].key; |
160 | 24 | return {key.data(), key.size()}; |
161 | 24 | } |
162 | | |
163 | 5 | int64_t ResolvedVariantElementV2Path::array_index_at(size_t position) const { |
164 | 5 | DORIS_CHECK(kind_at(position) == VariantElementV2PathSegment::Kind::ARRAY_INDEX) |
165 | 0 | << "Variant element path segment is not an array index"; |
166 | 5 | return _impl->segments[position].index; |
167 | 5 | } |
168 | | |
169 | | Status resolve_variant_element_v2_path( |
170 | | std::span<const VariantElementV2PathSegment> segments, |
171 | | // Mutable smart-pointer output is published only after full path validation. |
172 | | // NOLINTNEXTLINE(readability-non-const-parameter) |
173 | 15 | std::unique_ptr<ResolvedVariantElementV2Path>* output) { |
174 | 15 | if (output == nullptr) { |
175 | 0 | return Status::InvalidArgument("Variant V2 resolved path output is null"); |
176 | 0 | } |
177 | 15 | if (segments.empty()) { |
178 | 1 | return Status::InvalidArgument("Variant V2 element path must not be empty"); |
179 | 1 | } |
180 | | |
181 | 14 | auto impl = std::make_unique<ResolvedVariantElementV2Path::Impl>(); |
182 | 14 | impl->segments.reserve(segments.size()); |
183 | 25 | for (const VariantElementV2PathSegment& segment : segments) { |
184 | 25 | OwnedPathSegment owned {.kind = segment.kind(), .key = {}, .index = segment.index()}; |
185 | 25 | if (segment.kind() == VariantElementV2PathSegment::Kind::OBJECT_KEY) { |
186 | 20 | if (segment.key().size != 0 && segment.key().data == nullptr) { |
187 | 1 | return Status::InvalidArgument("Variant V2 object path key has a null pointer"); |
188 | 1 | } |
189 | 19 | if (segment.key().size != 0) { |
190 | 19 | owned.key.assign(segment.key().data, segment.key().data + segment.key().size); |
191 | 19 | } |
192 | 19 | } |
193 | 24 | impl->segments.push_back(std::move(owned)); |
194 | 24 | } |
195 | | |
196 | 13 | auto candidate = std::unique_ptr<ResolvedVariantElementV2Path>( |
197 | 13 | new ResolvedVariantElementV2Path(std::move(impl))); |
198 | 13 | output->swap(candidate); |
199 | 13 | return Status::OK(); |
200 | 14 | } |
201 | | |
202 | | Status extract_variant_element_v2(const ColumnVariantV2& source, |
203 | | const ResolvedVariantElementV2Path& path, |
204 | | // Mutable smart-pointer output is published only on success. |
205 | | // NOLINTNEXTLINE(readability-non-const-parameter) |
206 | 14 | std::span<const uint8_t> outer_nulls, ColumnPtr* output) { |
207 | 14 | if (output == nullptr) { |
208 | 0 | return Status::InvalidArgument("Variant V2 element output is null"); |
209 | 0 | } |
210 | 14 | if (path.size() == 0) { |
211 | 0 | return Status::InvalidArgument("Variant V2 element path must not be empty"); |
212 | 0 | } |
213 | 14 | if (!outer_nulls.empty() && outer_nulls.size() != source.size()) { |
214 | 0 | return Status::InvalidArgument("Variant V2 outer null map has {} rows, expected {}", |
215 | 0 | outer_nulls.size(), source.size()); |
216 | 0 | } |
217 | | |
218 | 14 | ColumnPtr candidate; |
219 | 14 | try { |
220 | 14 | if (!source.is_typed()) { |
221 | 13 | RETURN_IF_ERROR(extract_encoded_variant_element(source, path, outer_nulls, &candidate)); |
222 | 13 | } else { |
223 | | // A typed Variant is one scalar root value per row. String payloads are strings, not |
224 | | // JSON documents, so every non-empty object/array path is absent for all typed roots. |
225 | 1 | RETURN_IF_ERROR(make_all_null_variant_element_result(source.size(), &candidate)); |
226 | 1 | } |
227 | 14 | } catch (const Exception& exception) { |
228 | 0 | if (exception.code() == ErrorCode::CORRUPTION) { |
229 | 0 | return Status::InvalidArgument("Invalid Variant V2 input: {}", exception.message()); |
230 | 0 | } |
231 | 0 | return exception.to_status(); |
232 | 0 | } |
233 | 14 | if (!candidate || candidate->size() != source.size()) { |
234 | 0 | return Status::InternalError("Variant V2 element kernel produced {} rows, expected {}", |
235 | 0 | candidate ? candidate->size() : 0, source.size()); |
236 | 0 | } |
237 | 14 | output->swap(candidate); |
238 | 14 | return Status::OK(); |
239 | 14 | } |
240 | | |
241 | | namespace { |
242 | | |
243 | | constexpr uint32_t UNMAPPED_METADATA = std::numeric_limits<uint32_t>::max(); |
244 | | |
245 | | struct ExtractedRows { |
246 | | PaddedPODArray<char> metadata_bytes; |
247 | | DorisVector<uint32_t> metadata_offsets {0}; |
248 | | DorisVector<uint32_t> metadata_ids; |
249 | | PaddedPODArray<char> value_bytes; |
250 | | DorisVector<uint32_t> value_offsets {0}; |
251 | | |
252 | 13 | ColumnVariantV2::EncodedDataView view() const { |
253 | 13 | return {.metadata_bytes = {metadata_bytes.data(), metadata_bytes.size()}, |
254 | 13 | .metadata_offsets = metadata_offsets, |
255 | 13 | .meta_ids = metadata_ids, |
256 | 13 | .value_bytes = {value_bytes.data(), value_bytes.size()}, |
257 | 13 | .value_offsets = value_offsets}; |
258 | 13 | } |
259 | | }; |
260 | | |
261 | | uint32_t append_bytes(PaddedPODArray<char>& destination, StringRef source, |
262 | 39 | std::string_view description) { |
263 | 39 | if (source.size > std::numeric_limits<uint32_t>::max() - destination.size()) { |
264 | 0 | throw Exception(ErrorCode::INVALID_ARGUMENT, |
265 | 0 | "Variant element {} exceeds the ColumnString uint32 byte limit", |
266 | 0 | description); |
267 | 0 | } |
268 | 39 | if (source.size != 0) { |
269 | 39 | destination.insert(source.data, source.data + source.size); |
270 | 39 | } |
271 | 39 | return static_cast<uint32_t>(destination.size()); |
272 | 39 | } |
273 | | |
274 | 19 | uint32_t append_metadata(ExtractedRows& rows, VariantMetadataRef metadata) { |
275 | 19 | rows.metadata_offsets.push_back( |
276 | 19 | append_bytes(rows.metadata_bytes, {metadata.data, metadata.size}, "metadata")); |
277 | 19 | return static_cast<uint32_t>(rows.metadata_offsets.size() - 2); |
278 | 19 | } |
279 | | |
280 | 20 | void append_value(ExtractedRows& rows, VariantRef value, uint32_t metadata_id) { |
281 | 20 | rows.value_offsets.push_back(append_bytes(rows.value_bytes, value.value, "value")); |
282 | 20 | rows.metadata_ids.push_back(metadata_id); |
283 | 20 | } |
284 | | |
285 | 13 | ColumnPtr wrap_result(ExtractedRows rows, MutableColumnPtr nulls) { |
286 | 13 | auto values = ColumnVariantV2::create(); |
287 | 13 | values->insert_encoded_rows(rows.view()); |
288 | 13 | return ColumnNullable::create(std::move(values), std::move(nulls)); |
289 | 13 | } |
290 | | |
291 | | Status extract_encoded_variant_element(const ColumnVariantV2& source, |
292 | | const ResolvedVariantElementV2Path& path, |
293 | | std::span<const uint8_t> outer_nulls, |
294 | 13 | ColumnPtr* const output) { |
295 | 13 | VariantPathV2BatchReader reader(source, path); |
296 | 13 | const size_t metadata_count = reader.metadata_count(); |
297 | 13 | DorisVector<uint32_t> output_metadata_ids(metadata_count, UNMAPPED_METADATA); |
298 | | |
299 | 13 | VariantBatchBuilder null_builder(VariantBatchBuilder::ReserveHint {.rows = 1}); |
300 | 13 | auto null_row = null_builder.begin_row(); |
301 | 13 | null_row.add_null(); |
302 | 13 | null_row.finish(); |
303 | 13 | VariantBatchBuilder null_block = null_builder.finish_batch(); |
304 | 13 | const VariantRef null_value = null_block.value_at(0); |
305 | | |
306 | 13 | ExtractedRows rows; |
307 | 13 | rows.metadata_ids.reserve(source.size()); |
308 | 13 | rows.value_offsets.reserve(source.size() + 1); |
309 | 13 | auto nulls = ColumnUInt8::create(); |
310 | 13 | nulls->reserve(source.size()); |
311 | 13 | uint32_t null_metadata_id = UNMAPPED_METADATA; |
312 | | |
313 | 13 | auto ensure_null_metadata = [&]() { |
314 | 1 | if (null_metadata_id == UNMAPPED_METADATA) { |
315 | 1 | null_metadata_id = append_metadata(rows, null_value.metadata); |
316 | 1 | } |
317 | 1 | return null_metadata_id; |
318 | 1 | }; |
319 | | |
320 | 33 | for (size_t row = 0; row < source.size(); ++row) { |
321 | 20 | if (is_outer_null(outer_nulls, row)) { |
322 | 1 | append_value(rows, null_value, ensure_null_metadata()); |
323 | 1 | nulls->insert_value(1); |
324 | 1 | continue; |
325 | 1 | } |
326 | | |
327 | 19 | const uint32_t source_metadata_id = reader.metadata_id_at(row); |
328 | 19 | if (output_metadata_ids[source_metadata_id] == UNMAPPED_METADATA) { |
329 | 18 | output_metadata_ids[source_metadata_id] = |
330 | 18 | append_metadata(rows, reader.metadata_at(source_metadata_id)); |
331 | 18 | } |
332 | | |
333 | 19 | VariantRef current; |
334 | 19 | if (reader.find_at(row, ¤t)) { |
335 | 13 | append_value(rows, current, output_metadata_ids[source_metadata_id]); |
336 | 13 | nulls->insert_value(0); |
337 | 13 | } else { |
338 | 6 | append_value(rows, null_value, output_metadata_ids[source_metadata_id]); |
339 | 6 | nulls->insert_value(1); |
340 | 6 | } |
341 | 19 | } |
342 | | |
343 | 13 | ColumnPtr candidate = wrap_result(std::move(rows), std::move(nulls)); |
344 | 13 | output->swap(candidate); |
345 | 13 | return Status::OK(); |
346 | 13 | } |
347 | | |
348 | 1 | Status make_all_null_variant_element_result(size_t rows, ColumnPtr* const output) { |
349 | 1 | VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows = rows}); |
350 | 2 | for (size_t row_index = 0; row_index < rows; ++row_index) { |
351 | 1 | auto row = builder.begin_row(); |
352 | 1 | row.add_null(); |
353 | 1 | row.finish(); |
354 | 1 | } |
355 | 1 | VariantBatchBuilder block = builder.finish_batch(); |
356 | 1 | auto values = ColumnVariantV2::create(); |
357 | 1 | values->insert_encoded_batch(block); |
358 | 1 | ColumnPtr candidate = ColumnNullable::create(std::move(values), ColumnUInt8::create(rows, 1)); |
359 | 1 | output->swap(candidate); |
360 | 1 | return Status::OK(); |
361 | 1 | } |
362 | | |
363 | | } // namespace |
364 | | |
365 | | } // namespace doris |