Coverage Report

Created: 2026-07-29 15:18

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
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), &current)) {
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, &current)) {
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