Coverage Report

Created: 2026-08-06 07:53

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/parquet/parquet_predicate.h
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
#pragma once
19
20
#include <gen_cpp/parquet_types.h>
21
22
#include <cmath>
23
#include <cstring>
24
#include <vector>
25
26
#include "cctz/time_zone.h"
27
#include "core/data_type/data_type_decimal.h"
28
#include "core/data_type/primitive_type.h"
29
#include "exec/common/endian.h"
30
#include "format/format_common.h"
31
#include "format/parquet/parquet_block_split_bloom_filter.h"
32
#include "format/parquet/parquet_column_convert.h"
33
#include "format/parquet/parquet_common.h"
34
#include "format/parquet/schema_desc.h"
35
#include "io/fs/file_reader.h"
36
#include "storage/olap_scan_common.h"
37
#include "storage/segment/row_ranges.h"
38
#include "util/timezone_utils.h"
39
40
namespace doris {
41
class ParquetPredicate {
42
private:
43
66
    static inline bool _is_ascii(uint8_t byte) { return byte < 128; }
44
45
17
    static int _common_prefix(const std::string& encoding_min, const std::string& encoding_max) {
46
17
        size_t min_length = std::min(encoding_min.size(), encoding_max.size());
47
17
        int common_length = 0;
48
52
        while (common_length < min_length &&
49
52
               encoding_min[common_length] == encoding_max[common_length]) {
50
35
            common_length++;
51
35
        }
52
17
        return common_length;
53
17
    }
54
55
118
    static bool _try_read_old_utf8_stats(std::string& encoding_min, std::string& encoding_max) {
56
118
        if (encoding_min == encoding_max) {
57
            // If min = max, then there is a single value only
58
            // No need to modify, just use min
59
101
            encoding_max = encoding_min;
60
101
            return true;
61
101
        } else {
62
17
            int common_prefix_length = _common_prefix(encoding_min, encoding_max);
63
64
            // For min we can retain all-ASCII, because this produces a strictly lower value.
65
17
            int min_good_length = common_prefix_length;
66
37
            while (min_good_length < encoding_min.size() &&
67
37
                   _is_ascii(static_cast<uint8_t>(encoding_min[min_good_length]))) {
68
20
                min_good_length++;
69
20
            }
70
71
            // For max we can be sure only of the part matching the min. When they differ, we can consider only one next, and only if both are ASCII
72
17
            int max_good_length = common_prefix_length;
73
17
            if (max_good_length < encoding_max.size() && max_good_length < encoding_min.size() &&
74
17
                _is_ascii(static_cast<uint8_t>(encoding_min[max_good_length])) &&
75
17
                _is_ascii(static_cast<uint8_t>(encoding_max[max_good_length]))) {
76
10
                max_good_length++;
77
10
            }
78
            // Incrementing 127 would overflow. Incrementing within non-ASCII can have side-effects.
79
22
            while (max_good_length > 0 &&
80
22
                   (static_cast<uint8_t>(encoding_max[max_good_length - 1]) == 127 ||
81
18
                    !_is_ascii(static_cast<uint8_t>(encoding_max[max_good_length - 1])))) {
82
5
                max_good_length--;
83
5
            }
84
17
            if (max_good_length == 0) {
85
                // We can return just min bound, but code downstream likely expects both are present or both are absent.
86
4
                return false;
87
4
            }
88
89
13
            encoding_min.resize(min_good_length);
90
13
            encoding_max.resize(max_good_length);
91
13
            if (max_good_length > 0) {
92
13
                encoding_max[max_good_length - 1]++;
93
13
            }
94
13
            return true;
95
17
        }
96
118
    }
97
98
0
    static SortOrder _determine_sort_order(const tparquet::SchemaElement& parquet_schema) {
99
0
        tparquet::Type::type physical_type = parquet_schema.type;
100
0
        const tparquet::LogicalType& logical_type = parquet_schema.logicalType;
101
102
        // Assume string type is SortOrder::SIGNED, use ParquetPredicate::_try_read_old_utf8_stats() to handle it.
103
0
        if (logical_type.__isset.STRING &&
104
0
            (physical_type == tparquet::Type::BYTE_ARRAY ||
105
0
             physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY)) {
106
0
            return SortOrder::SIGNED;
107
0
        }
108
109
0
        if (logical_type.__isset.INTEGER) {
110
0
            if (logical_type.INTEGER.isSigned) {
111
0
                return SortOrder::SIGNED;
112
0
            } else {
113
0
                return SortOrder::UNSIGNED;
114
0
            }
115
0
        } else if (logical_type.__isset.DATE) {
116
0
            return SortOrder::SIGNED;
117
0
        } else if (logical_type.__isset.ENUM) {
118
0
            return SortOrder::UNSIGNED;
119
0
        } else if (logical_type.__isset.BSON) {
120
0
            return SortOrder::UNSIGNED;
121
0
        } else if (logical_type.__isset.JSON) {
122
0
            return SortOrder::UNSIGNED;
123
0
        } else if (logical_type.__isset.STRING) {
124
0
            return SortOrder::UNSIGNED;
125
0
        } else if (logical_type.__isset.DECIMAL) {
126
0
            return SortOrder::UNKNOWN;
127
0
        } else if (logical_type.__isset.MAP) {
128
0
            return SortOrder::UNKNOWN;
129
0
        } else if (logical_type.__isset.LIST) {
130
0
            return SortOrder::UNKNOWN;
131
0
        } else if (logical_type.__isset.TIME) {
132
0
            return SortOrder::SIGNED;
133
0
        } else if (logical_type.__isset.TIMESTAMP) {
134
0
            return SortOrder::SIGNED;
135
0
        } else if (logical_type.__isset.UNKNOWN) {
136
0
            return SortOrder::UNKNOWN;
137
0
        } else {
138
0
            switch (physical_type) {
139
0
            case tparquet::Type::BOOLEAN:
140
0
            case tparquet::Type::INT32:
141
0
            case tparquet::Type::INT64:
142
0
            case tparquet::Type::FLOAT:
143
0
            case tparquet::Type::DOUBLE:
144
0
                return SortOrder::SIGNED;
145
0
            case tparquet::Type::BYTE_ARRAY:
146
0
            case tparquet::Type::FIXED_LEN_BYTE_ARRAY:
147
0
                return SortOrder::UNSIGNED;
148
0
            case tparquet::Type::INT96:
149
0
                return SortOrder::UNKNOWN;
150
0
            default:
151
0
                return SortOrder::UNKNOWN;
152
0
            }
153
0
        }
154
0
    }
155
156
public:
157
    static constexpr int BLOOM_FILTER_MAX_HEADER_LENGTH = 64;
158
    struct ColumnStat {
159
        std::string encoded_min_value;
160
        std::string encoded_max_value;
161
        bool has_null;
162
        bool is_all_null;
163
        const FieldSchema* col_schema;
164
        const cctz::time_zone* ctz;
165
        std::unique_ptr<ParquetBlockSplitBloomFilter> bloom_filter;
166
        std::function<bool(ParquetPredicate::ColumnStat*, const int)>* get_stat_func = nullptr;
167
        std::function<bool(ParquetPredicate::ColumnStat*, const int)>* get_bloom_filter_func =
168
                nullptr;
169
    };
170
171
8
    static bool bloom_filter_supported(PrimitiveType type) {
172
        // Only support types where physical type == logical type (no conversion needed)
173
        // For types like DATEV2, DATETIMEV2, DECIMAL, Parquet stores them in physical format
174
        // (INT32, INT64, etc.) but Doris uses different internal representations.
175
        // Bloom filter works with physical bytes, but we only have logical type values,
176
        // and there's no reverse conversion (logical -> physical) available.
177
        // TINYINT/SMALLINT also need conversion via LittleIntPhysicalConverter.
178
8
        switch (type) {
179
0
        case TYPE_BOOLEAN:
180
8
        case TYPE_INT:
181
8
        case TYPE_BIGINT:
182
8
        case TYPE_FLOAT:
183
8
        case TYPE_DOUBLE:
184
8
        case TYPE_CHAR:
185
8
        case TYPE_VARCHAR:
186
8
        case TYPE_STRING:
187
8
            return true;
188
0
        default:
189
0
            return false;
190
8
        }
191
8
    }
192
193
    struct PageIndexStat {
194
        // Indicates whether the page index information in this column can be used.
195
        bool available = false;
196
        int64_t num_of_pages;
197
        std::vector<std::string> encoded_min_value;
198
        std::vector<std::string> encoded_max_value;
199
        std::vector<bool> has_null;
200
        std::vector<bool> is_all_null;
201
        const FieldSchema* col_schema;
202
203
        // Record the row range corresponding to each page.
204
        std::vector<segment_v2::RowRange> ranges;
205
    };
206
207
    struct CachedPageIndexStat {
208
        const cctz::time_zone* ctz;
209
        std::map<int, PageIndexStat> stats;
210
        std::function<bool(PageIndexStat**, int)> get_stat_func;
211
        RowRange row_group_range;
212
    };
213
214
    // The encoded Parquet min-max value is parsed into `fields`;
215
    // Can be used in row groups and page index statistics.
216
    static Status parse_min_max_value(const FieldSchema* col_schema, const std::string& encoded_min,
217
                                      const std::string& encoded_max, const cctz::time_zone& ctz,
218
6.41k
                                      Field* min_field, Field* max_field) {
219
6.41k
        auto logical_data_type = remove_nullable(col_schema->data_type);
220
6.41k
        auto converter = parquet::PhysicalToLogicalConverter::get_converter(
221
6.41k
                col_schema, logical_data_type, logical_data_type, &ctz);
222
6.41k
        ColumnPtr physical_column;
223
6.41k
        switch (col_schema->parquet_schema.type) {
224
2
        case tparquet::Type::type::BOOLEAN: {
225
2
            auto physical_col = ColumnUInt8::create();
226
2
            physical_col->get_data().data();
227
2
            physical_col->resize(2);
228
2
            physical_col->get_data()[0] = *reinterpret_cast<const bool*>(encoded_min.data());
229
2
            physical_col->get_data()[1] = *reinterpret_cast<const bool*>(encoded_max.data());
230
2
            physical_column = std::move(physical_col);
231
2
            break;
232
0
        }
233
6.01k
        case tparquet::Type::type::INT32: {
234
6.01k
            auto physical_col = ColumnInt32::create();
235
6.01k
            physical_col->resize(2);
236
237
6.01k
            physical_col->get_data()[0] = *reinterpret_cast<const int32_t*>(encoded_min.data());
238
6.01k
            physical_col->get_data()[1] = *reinterpret_cast<const int32_t*>(encoded_max.data());
239
240
6.01k
            physical_column = std::move(physical_col);
241
6.01k
            break;
242
0
        }
243
205
        case tparquet::Type::type::INT64: {
244
205
            auto physical_col = ColumnInt64::create();
245
205
            physical_col->resize(2);
246
205
            physical_col->get_data()[0] = *reinterpret_cast<const int64_t*>(encoded_min.data());
247
205
            physical_col->get_data()[1] = *reinterpret_cast<const int64_t*>(encoded_max.data());
248
205
            physical_column = std::move(physical_col);
249
205
            break;
250
0
        }
251
14
        case tparquet::Type::type::FLOAT: {
252
14
            auto physical_col = ColumnFloat32::create();
253
14
            physical_col->resize(2);
254
14
            physical_col->get_data()[0] = *reinterpret_cast<const float*>(encoded_min.data());
255
14
            physical_col->get_data()[1] = *reinterpret_cast<const float*>(encoded_max.data());
256
14
            physical_column = std::move(physical_col);
257
14
            break;
258
0
        }
259
0
        case tparquet::Type::type::DOUBLE: {
260
0
            auto physical_col = ColumnFloat64 ::create();
261
0
            physical_col->resize(2);
262
0
            physical_col->get_data()[0] = *reinterpret_cast<const double*>(encoded_min.data());
263
0
            physical_col->get_data()[1] = *reinterpret_cast<const double*>(encoded_max.data());
264
0
            physical_column = std::move(physical_col);
265
0
            break;
266
0
        }
267
177
        case tparquet::Type::type::BYTE_ARRAY: {
268
177
            auto physical_col = ColumnString::create();
269
177
            physical_col->insert_data(encoded_min.data(), encoded_min.size());
270
177
            physical_col->insert_data(encoded_max.data(), encoded_max.size());
271
177
            physical_column = std::move(physical_col);
272
177
            break;
273
0
        }
274
2
        case tparquet::Type::type::FIXED_LEN_BYTE_ARRAY: {
275
2
            auto physical_col = ColumnUInt8::create();
276
2
            physical_col->resize(2 * col_schema->parquet_schema.type_length);
277
2
            DCHECK(col_schema->parquet_schema.type_length == encoded_min.length());
278
2
            DCHECK(col_schema->parquet_schema.type_length == encoded_max.length());
279
280
2
            auto ptr = physical_col->get_data().data();
281
2
            memcpy(ptr, encoded_min.data(), encoded_min.length());
282
2
            memcpy(ptr + encoded_min.length(), encoded_max.data(), encoded_max.length());
283
2
            physical_column = std::move(physical_col);
284
2
            break;
285
0
        }
286
0
        case tparquet::Type::type::INT96: {
287
0
            auto physical_col = ColumnInt8::create();
288
0
            physical_col->resize(2 * sizeof(ParquetInt96));
289
0
            DCHECK(sizeof(ParquetInt96) == encoded_min.length());
290
0
            DCHECK(sizeof(ParquetInt96) == encoded_max.length());
291
292
0
            auto ptr = physical_col->get_data().data();
293
0
            memcpy(ptr, encoded_min.data(), encoded_min.length());
294
0
            memcpy(ptr + encoded_min.length(), encoded_max.data(), encoded_max.length());
295
0
            physical_column = std::move(physical_col);
296
0
            break;
297
0
        }
298
6.41k
        }
299
300
6.42k
        ColumnPtr logical_column;
301
6.42k
        if (converter->is_consistent()) {
302
6.27k
            logical_column = physical_column;
303
6.27k
        } else {
304
148
            logical_column = logical_data_type->create_column();
305
148
            RETURN_IF_ERROR(converter->physical_convert(physical_column, logical_column));
306
148
        }
307
308
6.42k
        DCHECK(logical_column->size() == 2);
309
6.42k
        *min_field = logical_column->operator[](0);
310
6.42k
        *max_field = logical_column->operator[](1);
311
312
6.42k
        auto logical_prim_type = logical_data_type->get_primitive_type();
313
314
6.42k
        if (logical_prim_type == TYPE_FLOAT) {
315
11
            auto& min_value = min_field->get<TYPE_FLOAT>();
316
11
            auto& max_value = max_field->get<TYPE_FLOAT>();
317
318
11
            if (std::isnan(min_value) || std::isnan(max_value)) {
319
1
                return Status::DataQualityError("Can not use this parquet min/max value.");
320
1
            }
321
            // Updating min to -0.0 and max to +0.0 to ensure that no 0.0 values would be skipped
322
10
            if (std::signbit(min_value) == 0 && min_value == 0.0F) {
323
0
                min_value = -0.0F;
324
0
            }
325
10
            if (std::signbit(max_value) != 0 && max_value == -0.0F) {
326
0
                max_value = 0.0F;
327
0
            }
328
6.40k
        } else if (logical_prim_type == TYPE_DOUBLE) {
329
0
            auto& min_value = min_field->get<TYPE_DOUBLE>();
330
0
            auto& max_value = max_field->get<TYPE_DOUBLE>();
331
332
0
            if (std::isnan(min_value) || std::isnan(max_value)) {
333
0
                return Status::DataQualityError("Can not use this parquet min/max value.");
334
0
            }
335
            // Updating min to -0.0 and max to +0.0 to ensure that no 0.0 values would be skipped
336
0
            if (std::signbit(min_value) == 0 && min_value == 0.0F) {
337
0
                min_value = -0.0F;
338
0
            }
339
0
            if (std::signbit(max_value) != 0 && max_value == -0.0F) {
340
0
                max_value = 0.0F;
341
0
            }
342
6.40k
        } else if (col_schema->parquet_schema.type == tparquet::Type::type::INT96 ||
343
6.40k
                   logical_prim_type == TYPE_DATETIMEV2) {
344
145
            auto min_value = min_field->get<TYPE_DATETIMEV2>();
345
145
            auto max_value = min_field->get<TYPE_DATETIMEV2>();
346
347
            // From Trino: Parquet INT96 timestamp values were compared incorrectly
348
            // for the purposes of producing statistics by older parquet writers,
349
            // so PARQUET-1065 deprecated them. The result is that any writer that produced stats
350
            // was producing unusable incorrect values, except the special case where min == max
351
            // and an incorrect ordering would not be material to the result.
352
            // PARQUET-1026 made binary stats available and valid in that special case.
353
145
            if (min_value != max_value) {
354
0
                return Status::DataQualityError("invalid min/max value");
355
0
            }
356
145
        }
357
358
6.41k
        return Status::OK();
359
6.42k
    }
360
361
    static Status read_column_stats(const FieldSchema* col_schema,
362
                                    const tparquet::ColumnMetaData& column_meta_data,
363
                                    std::unordered_map<tparquet::Type::type, bool>* ignored_stats,
364
972
                                    const std::string& file_created_by, ColumnStat* ans_stat) {
365
972
        auto& statistic = column_meta_data.statistics;
366
367
972
        if (!statistic.__isset.null_count) [[unlikely]] {
368
2
            return Status::DataQualityError("This parquet Column meta no set null_count.");
369
2
        }
370
970
        ans_stat->has_null = statistic.null_count > 0;
371
970
        ans_stat->is_all_null = statistic.null_count == column_meta_data.num_values;
372
970
        if (ans_stat->is_all_null) {
373
8
            return Status::OK();
374
8
        }
375
962
        auto prim_type = remove_nullable(col_schema->data_type)->get_primitive_type();
376
377
        // Min-max of statistic is plain-encoded value
378
964
        if (statistic.__isset.min_value && statistic.__isset.max_value) {
379
964
            ColumnOrderName column_order =
380
964
                    col_schema->physical_type == tparquet::Type::INT96 ||
381
964
                                    col_schema->parquet_schema.logicalType.__isset.UNKNOWN
382
964
                            ? ColumnOrderName::UNDEFINED
383
964
                            : ColumnOrderName::TYPE_DEFINED_ORDER;
384
964
            if ((statistic.min_value != statistic.max_value) &&
385
964
                (column_order != ColumnOrderName::TYPE_DEFINED_ORDER)) {
386
0
                return Status::DataQualityError("Can not use this parquet min/max value.");
387
0
            }
388
964
            ans_stat->encoded_min_value = statistic.min_value;
389
964
            ans_stat->encoded_max_value = statistic.max_value;
390
391
964
            if (prim_type == TYPE_VARCHAR || prim_type == TYPE_CHAR || prim_type == TYPE_STRING) {
392
105
                auto encoded_min_copy = ans_stat->encoded_min_value;
393
105
                auto encoded_max_copy = ans_stat->encoded_max_value;
394
105
                if (!_try_read_old_utf8_stats(encoded_min_copy, encoded_max_copy)) {
395
0
                    return Status::DataQualityError("Can not use this parquet min/max value.");
396
0
                }
397
105
                ans_stat->encoded_min_value = encoded_min_copy;
398
105
                ans_stat->encoded_max_value = encoded_max_copy;
399
105
            }
400
401
18.4E
        } else if (statistic.__isset.min && statistic.__isset.max) {
402
0
            bool max_equals_min = statistic.min == statistic.max;
403
404
0
            SortOrder sort_order = _determine_sort_order(col_schema->parquet_schema);
405
0
            bool sort_orders_match = SortOrder::SIGNED == sort_order;
406
0
            if (!sort_orders_match && !max_equals_min) {
407
0
                return Status::NotSupported("Can not use this parquet min/max value.");
408
0
            }
409
410
0
            bool should_ignore_corrupted_stats = false;
411
0
            if (ignored_stats != nullptr) {
412
0
                if (ignored_stats->count(col_schema->physical_type) == 0) {
413
0
                    if (CorruptStatistics::should_ignore_statistics(file_created_by,
414
0
                                                                    col_schema->physical_type)) {
415
0
                        ignored_stats->emplace(col_schema->physical_type, true);
416
0
                        should_ignore_corrupted_stats = true;
417
0
                    } else {
418
0
                        ignored_stats->emplace(col_schema->physical_type, false);
419
0
                    }
420
0
                } else if (ignored_stats->at(col_schema->physical_type)) {
421
0
                    should_ignore_corrupted_stats = true;
422
0
                }
423
0
            } else if (CorruptStatistics::should_ignore_statistics(file_created_by,
424
0
                                                                   col_schema->physical_type)) {
425
0
                should_ignore_corrupted_stats = true;
426
0
            }
427
428
0
            if (should_ignore_corrupted_stats) {
429
0
                return Status::DataQualityError("Error statistics, should ignore.");
430
0
            }
431
432
0
            ans_stat->encoded_min_value = statistic.min;
433
0
            ans_stat->encoded_max_value = statistic.max;
434
18.4E
        } else {
435
18.4E
            return Status::DataQualityError("This parquet file not set min/max value");
436
18.4E
        }
437
438
964
        return Status::OK();
439
962
    }
440
441
    static Status read_bloom_filter(const tparquet::ColumnMetaData& column_meta_data,
442
                                    io::FileReaderSPtr file_reader, io::IOContext* io_ctx,
443
12
                                    ColumnStat* ans_stat) {
444
12
        size_t size;
445
12
        if (!column_meta_data.__isset.bloom_filter_offset) {
446
0
            return Status::NotSupported("Can not use this parquet bloom filter.");
447
0
        }
448
449
12
        if (column_meta_data.__isset.bloom_filter_length &&
450
12
            column_meta_data.bloom_filter_length > 0) {
451
12
            size = column_meta_data.bloom_filter_length;
452
12
        } else {
453
0
            size = BLOOM_FILTER_MAX_HEADER_LENGTH;
454
0
        }
455
12
        size_t bytes_read = 0;
456
12
        std::vector<uint8_t> header_buffer(size);
457
12
        RETURN_IF_ERROR(file_reader->read_at(column_meta_data.bloom_filter_offset,
458
12
                                             Slice(header_buffer.data(), size), &bytes_read,
459
12
                                             io_ctx));
460
461
12
        tparquet::BloomFilterHeader t_bloom_filter_header;
462
12
        uint32_t t_bloom_filter_header_size = static_cast<uint32_t>(bytes_read);
463
12
        RETURN_IF_ERROR(deserialize_thrift_msg(header_buffer.data(), &t_bloom_filter_header_size,
464
12
                                               true, &t_bloom_filter_header));
465
466
        // TODO the bloom filter could be encrypted, too, so need to double check that this is NOT the case
467
12
        if (!t_bloom_filter_header.algorithm.__isset.BLOCK ||
468
12
            !t_bloom_filter_header.compression.__isset.UNCOMPRESSED ||
469
12
            !t_bloom_filter_header.hash.__isset.XXHASH) {
470
0
            return Status::NotSupported("Can not use this parquet bloom filter.");
471
0
        }
472
473
12
        ans_stat->bloom_filter = std::make_unique<ParquetBlockSplitBloomFilter>();
474
475
12
        std::vector<uint8_t> data_buffer(t_bloom_filter_header.numBytes);
476
12
        RETURN_IF_ERROR(file_reader->read_at(
477
12
                column_meta_data.bloom_filter_offset + t_bloom_filter_header_size,
478
12
                Slice(data_buffer.data(), t_bloom_filter_header.numBytes), &bytes_read, io_ctx));
479
480
12
        RETURN_IF_ERROR(ans_stat->bloom_filter->init(
481
12
                reinterpret_cast<const char*>(data_buffer.data()), t_bloom_filter_header.numBytes,
482
12
                segment_v2::HashStrategyPB::XX_HASH_64));
483
484
12
        return Status::OK();
485
12
    }
486
};
487
488
} // namespace doris