Coverage Report

Created: 2026-09-30 16:28

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format_v2/parquet/parquet_scan.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
//   http://www.apache.org/licenses/LICENSE-2.0
9
// Unless required by applicable law or agreed to in writing,
10
// software distributed under the License is distributed on an
11
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
12
// KIND, either express or implied.  See the License for the
13
// specific language governing permissions and limitations
14
// under the License.
15
16
#pragma once
17
18
#include <gen_cpp/parquet_types.h>
19
20
#include <cstddef>
21
#include <cstdint>
22
#include <map>
23
#include <memory>
24
#include <optional>
25
#include <unordered_map>
26
#include <utility>
27
#include <vector>
28
29
#include "common/status.h"
30
#include "core/column/column.h"
31
#include "format_v2/file_reader.h"
32
#include "format_v2/parquet/parquet_profile.h"
33
#include "format_v2/parquet/parquet_statistics.h"
34
#include "format_v2/parquet/reader/column_reader.h"
35
#include "format_v2/parquet/selection_vector.h"
36
#include "runtime/runtime_profile.h"
37
#include "storage/segment/condition_cache.h"
38
39
namespace cctz {
40
class time_zone;
41
} // namespace cctz
42
43
namespace doris {
44
class Block;
45
class RuntimeState;
46
47
namespace format {
48
struct FileScanRequest;
49
} // namespace format
50
} // namespace doris
51
52
namespace doris::format::parquet {
53
54
struct ParquetFileContext;
55
struct ParquetColumnSchema;
56
struct ParquetPageCacheRange;
57
struct ParquetScanRange;
58
class NativeParquetMetadata;
59
60
namespace detail {
61
struct PredicateConjunctStage {
62
    VExprContextSPtr owner_context;
63
    VExprSPtr expression;
64
    std::vector<size_t> required_positions;
65
};
66
67
struct PredicateConjunctSchedule {
68
    std::map<size_t, VExprContextSPtrs> single_column_conjuncts;
69
    VExprContextSPtrs remaining_conjuncts;
70
    std::vector<PredicateConjunctStage> remaining_stages;
71
    bool supports_lazy_materialization = true;
72
};
73
74
struct AdaptivePredicateStats {
75
    double cost_per_input_row_ns = 0;
76
    double survival_ratio = 1;
77
    size_t samples = 0;
78
};
79
80
std::vector<size_t> order_adaptive_predicates(
81
        const std::vector<size_t>& positions,
82
        const std::unordered_map<size_t, AdaptivePredicateStats>& stats);
83
std::vector<size_t> adaptive_prefetch_prefix(
84
        const std::vector<size_t>& ordered_positions,
85
        const std::unordered_map<size_t, AdaptivePredicateStats>& stats,
86
        double minimum_reach_probability);
87
bool should_sample_adaptive_predicate(size_t samples, size_t batch_sequence);
88
Status validate_ephemeral_expr_result_column(size_t original_columns, int result_column_id,
89
                                             size_t current_columns);
90
Status build_native_prefetch_ranges(
91
        const tparquet::FileMetaData& metadata,
92
        const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
93
        const std::vector<format::LocalColumnIndex>& scan_columns, int row_group_idx,
94
        size_t file_size, bool parquet_816_padding, std::vector<ParquetPageCacheRange>* ranges);
95
Status select_native_row_groups_by_scan_range(const tparquet::FileMetaData& metadata,
96
                                              const ParquetScanRange& scan_range,
97
                                              std::vector<int64_t>* row_group_first_rows,
98
                                              std::vector<int>* selected_row_groups);
99
} // namespace detail
100
101
// ============================================================================
102
// ============================================================================
103
104
struct ParquetScanRange {
105
    int64_t start_offset = 0;
106
    int64_t size = -1;      // -1 means read the whole file
107
    int64_t file_size = -1; // -1 means unknown
108
};
109
110
struct RowGroupReadPlan {
111
    int row_group_id = -1;                 // row group id
112
    int64_t first_file_row = 0;            // first file row for this row group (0-based)
113
    int64_t row_group_rows = 0;            // row count of this row group
114
    std::vector<RowRange> selected_ranges; // row ranges to read after page-index pruning
115
    std::map<int, ParquetPageSkipPlan>
116
            page_skip_plans; // leaf_column_id -> data pages that can be skipped completely
117
    // Deferred planning transfers parsed indexes to execution so narrowed scans never issue the
118
    // same remote index reads a second time while opening the row group.
119
    std::unordered_map<int, tparquet::OffsetIndex> offset_indexes;
120
    // Footer statistics are cheap and eager. Remote dictionary/Bloom/page-index probes fill the
121
    // remaining fields only when this row group reaches the scheduler.
122
    bool expensive_pruning_pending = false;
123
};
124
125
struct RowGroupScanPlan {
126
    std::vector<RowGroupReadPlan> row_groups; // row groups selected after pruning
127
    ParquetPruningStats pruning_stats;        // pruning statistics
128
    bool enable_bloom_filter = false;
129
};
130
131
// ============================================================================
132
// ============================================================================
133
134
Status plan_parquet_row_groups(const NativeParquetMetadata& metadata,
135
                               const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
136
                               const format::FileScanRequest& request,
137
                               const ParquetScanRange& scan_range, bool enable_bloom_filter,
138
                               RowGroupScanPlan* plan, const cctz::time_zone* timezone = nullptr,
139
                               const RuntimeState* runtime_state = nullptr,
140
                               ParquetFileContext* file_context = nullptr,
141
                               const ParquetColumnReaderProfile& column_reader_profile = {});
142
143
Status finalize_parquet_row_group_plans(
144
        const NativeParquetMetadata& metadata,
145
        const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
146
        const format::FileScanRequest& request, bool enable_bloom_filter, RowGroupScanPlan* plan,
147
        const cctz::time_zone* timezone, const RuntimeState* runtime_state,
148
        ParquetFileContext* file_context, const ParquetColumnReaderProfile& column_reader_profile,
149
        const ParquetProfile* parquet_profile = nullptr);
150
151
IColumn::Filter selection_to_filter(const SelectionVector& selection, uint16_t selected_rows,
152
                                    int64_t batch_rows);
153
154
uint16_t apply_compact_filter_to_selection(const IColumn::Filter& filter,
155
                                           SelectionVector* selection, uint16_t selected_rows);
156
157
Status execute_batch_filters(const format::FileScanRequest& request, int64_t batch_rows,
158
                             Block* file_block, SelectionVector* selection, uint16_t* selected_rows,
159
                             int64_t* conjunct_filtered_rows = nullptr);
160
161
// ============================================================================
162
// ============================================================================
163
//   while true:
164
//     3. read_current_row_group_batch(batch_rows)
165
// ============================================================================
166
class ParquetScanScheduler {
167
public:
168
    static constexpr int64_t DEFAULT_READ_BATCH_SIZE = 4096;
169
170
    void set_plan(RowGroupScanPlan plan);
171
459
    void set_page_skip_profile(ParquetPageSkipProfile page_skip_profile) {
172
459
        _page_skip_profile = page_skip_profile;
173
459
    }
174
459
    void set_scan_profile(ParquetScanProfile scan_profile) { _scan_profile = scan_profile; }
175
332
    void set_pruning_profile(const ParquetProfile* parquet_profile) {
176
332
        _parquet_profile = parquet_profile;
177
332
    }
178
471
    void set_merge_read_options(RuntimeProfile* profile, int64_t merge_read_slice_size) {
179
471
        _profile = profile;
180
471
        _merge_read_slice_size = merge_read_slice_size;
181
471
    }
182
459
    void set_global_rowid_context(std::optional<format::GlobalRowIdContext> context) {
183
459
        _global_rowid_context = context;
184
459
    }
185
    void set_condition_cache_context(std::shared_ptr<ConditionCacheContext> ctx);
186
472
    void set_timezone(const cctz::time_zone* timezone) { _timezone = timezone; }
187
83
    void set_int96_timezone(std::optional<const cctz::time_zone*> timezone) {
188
83
        _int96_timezone = timezone;
189
83
    }
190
471
    void set_enable_strict_mode(bool enable_strict_mode) {
191
471
        _enable_strict_mode = enable_strict_mode;
192
471
    }
193
472
    void set_runtime_state(RuntimeState* runtime_state) { _runtime_state = runtime_state; }
194
    void set_scan_request(std::shared_ptr<format::FileScanRequest> request);
195
    void queue_scan_request(std::shared_ptr<format::FileScanRequest> request);
196
    // Release row-group readers before the owning RuntimeProfile is reported. Native readers
197
    // publish their accumulated page/decode statistics from their destructor.
198
236
    void close() { reset_current_row_group(); }
199
    // Upper scanner owns adaptive memory feedback; scheduler only applies the current row cap when
200
    // splitting selected row ranges into physical read batches.
201
477
    void set_batch_size(size_t batch_size) {
202
477
        _batch_size = batch_size == 0 ? 1 : static_cast<int64_t>(batch_size);
203
477
    }
204
    void reset();
205
460
    bool empty() const { return _row_group_plans.empty(); }
206
2
    int64_t condition_cache_filtered_rows() const { return _condition_cache_filtered_rows; }
207
1.28k
    int64_t predicate_filtered_rows() const { return _predicate_filtered_rows; }
208
1.96k
    int64_t raw_rows_read() const { return _raw_rows_read; }
209
210
    Status read_next_batch(ParquetFileContext& file_context,
211
                           const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
212
                           Block* file_block, size_t* rows, bool* eof);
213
214
private:
215
    static constexpr size_t PROFILE_FLUSH_BATCH_INTERVAL = 16;
216
217
    void reset_current_row_group();
218
    void activate_pending_scan_request_at_row_group_boundary();
219
    void flush_current_reader_profiles();
220
    bool finish_current_reader_batch_profiles();
221
    const detail::PredicateConjunctSchedule& predicate_conjunct_schedule(
222
            const format::FileScanRequest& request);
223
    std::vector<format::LocalColumnIndex> adaptive_predicate_prefetch_columns(
224
            const format::FileScanRequest& request);
225
226
    Status open_next_row_group(ParquetFileContext& file_context,
227
                               const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
228
                               const format::FileScanRequest& request, bool* has_row_group);
229
230
    Status skip_current_row_group_rows(int64_t rows);
231
    Status flush_pending_non_predicate_skip_rows();
232
233
    Status read_filter_columns(int64_t batch_rows, const format::FileScanRequest& request,
234
                               Block* file_block, SelectionVector* selection,
235
                               uint16_t* selected_rows, int64_t* conjunct_filtered_rows,
236
                               bool* predicate_columns_filtered);
237
238
    Status prepare_current_dictionary_filters(
239
            ParquetFileContext& file_context,
240
            const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
241
            const format::FileScanRequest& request, int row_group_idx,
242
            const tparquet::RowGroup& row_group_metadata);
243
244
    Status prefetch_current_row_group_columns(
245
            ParquetFileContext& file_context,
246
            const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
247
            const format::FileScanRequest& request,
248
            const std::vector<format::LocalColumnIndex>& scan_columns, bool* prefetched);
249
250
    Status read_current_row_group_batch(
251
            ParquetFileContext& file_context,
252
            const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
253
            int64_t batch_rows, const format::FileScanRequest& request,
254
            int64_t batch_first_file_row, Block* file_block, size_t* rows);
255
256
    Status materialize_pending_predicate_batch(const format::FileScanRequest& request,
257
                                               Block* file_block, size_t* rows);
258
259
    void mark_condition_cache_granules(const SelectionVector& selection, uint16_t selected_rows,
260
                                       int64_t batch_first_file_row);
261
262
    std::vector<RowGroupReadPlan> _row_group_plans; // row group queue to scan
263
    size_t _next_row_group_plan_idx = 0;            // index of the next row group to process
264
265
    bool _has_current_row_group = false;
266
    // Readers retain pointers into this immutable row-group map, so it must outlive both maps below.
267
    std::unordered_map<int, tparquet::OffsetIndex> _current_offset_indexes;
268
    // File-local ids are signed because virtual columns use reserved negative values. Keeping the
269
    // typed id as the map key prevents GLOBAL_ROWID_COLUMN_ID from wrapping to a storage ColumnId.
270
    std::map<format::LocalColumnId, std::unique_ptr<ParquetColumnReader>>
271
            _current_predicate_columns; // predicate ColumnReaders
272
    std::map<format::LocalColumnId, std::unique_ptr<ParquetColumnReader>>
273
            _current_non_predicate_columns; // non-predicate ColumnReaders
274
    std::map<format::LocalColumnId, IColumn::Filter>
275
            _current_dictionary_filters; // local id -> dict entry bitmap
276
    std::map<format::LocalColumnId, std::vector<std::pair<VExprContextSPtr, VExprSPtr>>>
277
            _current_dictionary_residual_conjuncts; // local id -> row-level residual conjuncts
278
    int64_t _current_row_group_rows = 0;            // current row group row count
279
    int _current_row_group_id = -1;                 // current row group id in parquet metadata
280
    int64_t _current_row_group_rows_read = 0;       // rows read in the current row group (cursor)
281
    int64_t _current_row_group_first_row = 0;       // first file row of the current row group
282
    std::vector<RowRange>
283
            _current_selected_ranges; // selected ranges for the current row group after page-index pruning
284
    size_t _current_range_idx = 0;        // current selected_range index
285
    int64_t _current_range_rows_read = 0; // rows read in the current range
286
    // Predicate readers move immediately because they decide which rows survive. Non-predicate
287
    // readers can lag behind across fully filtered batches and range gaps; the lag is flushed once
288
    // before the next surviving batch is materialized, or discarded with the row group.
289
    int64_t _pending_non_predicate_skip_rows = 0;
290
    // Empty predicate batches may widen their physical probe. If the first non-empty probe finds
291
    // more rows than the caller's cap, keep its narrow predicate result here and materialize lazy
292
    // columns in capped physical slices on subsequent calls.
293
    int64_t _pending_predicate_batch_rows = 0;
294
    int64_t _pending_predicate_batch_rows_consumed = 0;
295
    size_t _pending_predicate_selected_offset = 0;
296
    std::vector<SelectionVector::Index> _pending_predicate_selection;
297
    std::map<size_t, ColumnPtr> _pending_predicate_columns;
298
    SelectionVector _pending_output_selection;
299
300
    bool _current_predicate_prefetched = false;
301
    bool _current_non_predicate_prefetched = false;
302
    bool _current_merge_range_active = false;
303
    ParquetPageSkipProfile _page_skip_profile;
304
    ParquetScanProfile _scan_profile;
305
    const ParquetProfile* _parquet_profile = nullptr;
306
    RuntimeProfile* _profile = nullptr;
307
    int64_t _merge_read_slice_size = -1;
308
    std::optional<format::GlobalRowIdContext> _global_rowid_context;
309
    const cctz::time_zone* _timezone = nullptr;
310
    std::optional<const cctz::time_zone*> _int96_timezone;
311
    bool _enable_strict_mode = false;
312
    bool _enable_bloom_filter = false;
313
    RuntimeState* _runtime_state = nullptr;
314
    int64_t _batch_size = DEFAULT_READ_BATCH_SIZE;
315
    // Batch control scratch is scheduler-owned so adaptive row caps change logical sizes without
316
    // reallocating selection indices, dense filter bytes, or compacted-column positions.
317
    SelectionVector _selection;
318
    std::vector<uint32_t> _read_column_positions_scratch;
319
    const format::FileScanRequest* _predicate_schedule_request = nullptr;
320
    std::shared_ptr<format::FileScanRequest> _active_request;
321
    std::shared_ptr<format::FileScanRequest> _pending_request;
322
    bool _remaining_plans_need_replanning = false;
323
    detail::PredicateConjunctSchedule _predicate_schedule;
324
    std::vector<size_t> _predicate_positions_scratch;
325
    std::unordered_map<size_t, size_t> _predicate_indices_by_position_scratch;
326
    std::unordered_set<size_t> _materialized_predicate_positions_scratch;
327
    std::vector<size_t> _ordered_predicate_positions_scratch;
328
    std::unordered_map<uint32_t, std::vector<SelectionVector::Index>>
329
            _predicate_column_selection_scratch;
330
    IColumn::Filter _predicate_compaction_filter_scratch;
331
    size_t _predicate_batch_sequence = 0;
332
    size_t _batches_since_profile_flush = 0;
333
    std::unordered_map<size_t, detail::AdaptivePredicateStats> _predicate_runtime_stats;
334
    double _predicate_survival_ratio = -1;
335
    std::shared_ptr<ConditionCacheContext> _condition_cache_ctx;
336
    int64_t _condition_cache_filtered_rows = 0;
337
    int64_t _predicate_filtered_rows = 0;
338
    int64_t _raw_rows_read = 0;
339
};
340
341
} // namespace doris::format::parquet