Coverage Report

Created: 2026-07-20 20:03

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format_v2/table_reader.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 <bvar/status.h>
21
22
#include <algorithm>
23
#include <exception>
24
#include <map>
25
#include <memory>
26
#include <optional>
27
#include <string>
28
#include <string_view>
29
#include <utility>
30
#include <vector>
31
32
#include "common/cast_set.h"
33
#include "common/exception.h"
34
#include "common/logging.h"
35
#include "common/status.h"
36
#include "core/assert_cast.h"
37
#include "core/block/block.h"
38
#include "core/column/column_array.h"
39
#include "core/column/column_const.h"
40
#include "core/column/column_map.h"
41
#include "core/column/column_nullable.h"
42
#include "core/column/column_struct.h"
43
#include "core/column/column_vector.h"
44
#include "core/data_type/data_type.h"
45
#include "core/data_type/data_type_array.h"
46
#include "core/data_type/data_type_map.h"
47
#include "core/data_type/data_type_nullable.h"
48
#include "core/data_type/data_type_number.h"
49
#include "core/data_type/data_type_string.h"
50
#include "core/data_type/data_type_struct.h"
51
#include "core/field.h"
52
#include "exec/common/stringop_substring.h"
53
#include "exprs/vexpr.h"
54
#include "exprs/vexpr_context.h"
55
#include "exprs/vexpr_fwd.h"
56
#include "exprs/vslot_ref.h"
57
#include "format/table/deletion_vector.h"
58
#include "format_v2/column_data.h"
59
#include "format_v2/column_mapper.h"
60
#include "format_v2/expr/cast.h"
61
#include "format_v2/expr/delete_predicate.h"
62
#include "format_v2/file_reader.h"
63
#include "format_v2/parquet/reader/column_reader.h"
64
#include "format_v2/schema_projection.h"
65
#include "gen_cpp/PlanNodes_types.h"
66
#include "io/io_common.h"
67
#include "runtime/descriptors.h"
68
#include "storage/segment/condition_cache.h"
69
70
namespace doris {
71
class Block;
72
struct DeleteFileDesc;
73
class RuntimeState;
74
} // namespace doris
75
76
namespace doris::format {
77
78
using DeleteRows = std::vector<int64_t>;
79
80
// Row-level predicates on table/global schema. They are rewritten to file-local expressions when
81
// possible, and remain the source of row-level filtering after localization.
82
struct TableFilter {
83
    VExprContextSPtr conjunct;
84
    std::vector<GlobalIndex> global_indices;
85
};
86
87
struct ScanTask {
88
135
    virtual ~ScanTask() = default;
89
90
    std::unique_ptr<io::FileDescription> data_file;
91
};
92
93
struct ProjectedColumnBuildContext {
94
    const TFileScanRangeParams* scan_params = nullptr;
95
    const TFileRangeDesc* range = nullptr;
96
    RuntimeState* runtime_state = nullptr;
97
    std::optional<ColumnDefinition> schema_column = std::nullopt;
98
    size_t next_file_column_idx = 0;
99
};
100
101
struct ReadProfile {
102
    RuntimeProfile::Counter* num_delete_files = nullptr;
103
    RuntimeProfile::Counter* num_delete_rows = nullptr;
104
    RuntimeProfile::Counter* parse_delete_file_time = nullptr;
105
    RuntimeProfile::Counter* decoded_dv_cache_hit_count = nullptr;
106
    RuntimeProfile::Counter* decoded_dv_cache_miss_count = nullptr;
107
    RuntimeProfile::Counter* dv_file_cache_hit_count = nullptr;
108
    RuntimeProfile::Counter* dv_file_cache_miss_count = nullptr;
109
    RuntimeProfile::Counter* dv_file_cache_peer_read_count = nullptr;
110
    RuntimeProfile::Counter* exec_timer = nullptr;
111
    RuntimeProfile::Counter* prepare_split_timer = nullptr;
112
    RuntimeProfile::Counter* finalize_timer = nullptr;
113
    RuntimeProfile::Counter* create_reader_timer = nullptr;
114
    RuntimeProfile::Counter* pushdown_agg_timer = nullptr;
115
    RuntimeProfile::Counter* open_reader_timer = nullptr;
116
    RuntimeProfile::Counter* runtime_filter_partition_prune_timer = nullptr;
117
    RuntimeProfile::Counter* runtime_filter_partition_pruned_range_counter = nullptr;
118
};
119
120
struct TableReadOptions {
121
    // Columns need to be read from file and output by table reader. They are all in table/global
122
    // schema semantics.
123
    const std::vector<ColumnDefinition> projected_columns;
124
    // All complex conjuncts from scan operator
125
    const VExprContextSPtrs conjuncts;
126
    // File format of the underlying data files, needed for reader initialization and reader-level
127
    // filter pushdown.
128
    const FileFormat format;
129
    TFileScanRangeParams* scan_params;
130
    std::shared_ptr<io::IOContext> io_ctx;
131
    RuntimeState* runtime_state;
132
    RuntimeProfile* scanner_profile;
133
    // File formats without complete self-describing metadata, such as CSV, Text, and JSON, need
134
    // the FE-planned physical file slots to build their file-local schema and deserialize values.
135
    const std::vector<SlotDescriptor*>* file_slot_descs = nullptr;
136
    // Push-down aggregate type.
137
    const TPushAggOp::type push_down_agg_type = TPushAggOp::type::NONE;
138
    // Table/global indices of explicit COUNT arguments. nullopt means an old FE did not send the
139
    // semantic argument field, while an explicit empty vector means COUNT(*)/COUNT(1). Keeping
140
    // those states separate prevents a rolling-upgrade plan from being reinterpreted by a new BE.
141
    const std::optional<std::vector<GlobalIndex>> push_down_count_columns = std::nullopt;
142
    // Initial digest of predicates available during scanner open. Scanner-driven splits override it
143
    // with SplitReadOptions::condition_cache_digest after collecting late-arrival runtime filters.
144
    // A zero digest disables condition cache.
145
    uint64_t condition_cache_digest = 0;
146
};
147
148
struct SplitReadOptions {
149
    // Split-level information for reader initialization, which may include file path, partition values, delete file info, etc. The content is table format specific and opaque to table reader base class; it's the responsibility of the concrete table reader implementation to parse necessary information for reader initialization and filter pushdown.
150
    std::map<std::string, Field> partition_values;
151
    // Latest scanner conjuncts rewritten to table/global column indices. Runtime filters may
152
    // arrive after TableReader::init(), so scanner-driven splits replace the initial snapshot.
153
    // nullopt preserves the initial snapshot for standalone TableReader callers.
154
    std::optional<VExprContextSPtrs> conjuncts = std::nullopt;
155
    // Independent clones used for partition pruning because evaluation prepares and opens them
156
    // against a synthetic partition block before the file reader opens its row-level conjuncts.
157
    VExprContextSPtrs partition_prune_conjuncts;
158
    // Table-level COUNT may emit one metadata-derived batch and resume on a later scheduler turn.
159
    // It is safe only after every runtime filter assigned to the scanner has arrived; otherwise a
160
    // filter could arrive after synthetic rows have already been returned and those rows cannot be
161
    // retracted. Standalone TableReader callers have no scanner runtime-filter lifecycle.
162
    bool all_runtime_filters_applied = true;
163
    // Digest for the exact scanner conjunct snapshot attached to this split. FileScannerV2 rebuilds
164
    // it after collecting late-arrival RFs, so different RF payloads cannot share a cache entry. A
165
    // zero value explicitly disables condition cache for this split.
166
    std::optional<uint64_t> condition_cache_digest;
167
    ShardedKVCache* cache = nullptr;
168
    TFileRangeDesc current_range;
169
    FileFormat current_split_format = FileFormat::PARQUET;
170
    std::optional<GlobalRowIdContext> global_rowid_context;
171
};
172
173
// Base class for table-level readers.
174
// This layer owns common table-level orchestration, such as split iteration, dynamic partition
175
// pruning, delete handling and conversion from file-local blocks to table-schema blocks. Concrete
176
// table-format readers only need to provide format-specific hooks for opening readers and parsing
177
// split metadata.
178
class TableReader {
179
public:
180
185
    virtual ~TableReader() = default;
181
182
    // Initialize common runtime options for the table reader. Subclasses may call this from their
183
    // own init(options); table-format schema and split metadata are provided later per split.
184
    virtual Status init(TableReadOptions&& options);
185
186
    // FileScannerV2 adjusts this before each get_block() using an adaptive bytes-per-row estimate.
187
    // Store it here as well as forwarding to the current reader so newly opened split readers start
188
    // with the latest predicted batch size.
189
10
    virtual void set_batch_size(size_t batch_size) {
190
10
        _batch_size = std::max<size_t>(1, batch_size);
191
10
        if (_data_reader.reader != nullptr) {
192
0
            _data_reader.reader->set_batch_size(_batch_size);
193
0
        }
194
10
    }
195
196
#ifdef BE_TEST
197
9
    size_t TEST_batch_size() const { return _batch_size; }
198
4
    bool TEST_current_data_file_is_immutable() const {
199
4
        DORIS_CHECK(_current_task != nullptr);
200
4
        DORIS_CHECK(_current_task->data_file != nullptr);
201
4
        DORIS_CHECK(_current_file_description.has_value());
202
4
        DORIS_CHECK(_current_task->data_file->is_immutable ==
203
4
                    _current_file_description->is_immutable);
204
4
        return _current_task->data_file->is_immutable;
205
4
    }
206
#endif
207
208
    // Prepare for reading a new split/task.
209
    // 1. Pass a new split/task to reader, which will be used in subsequent open_reader() to initialize the underlying file reader.
210
    // 2. Parse delete predicates from split/task information, which will be used for later dynamic filtering and delete handling.
211
    virtual Status prepare_split(const SplitReadOptions& options);
212
213
65
    virtual bool current_split_pruned() const { return _current_split_pruned; }
214
5
    virtual bool current_split_uses_metadata_count() const {
215
5
        return _current_split_uses_metadata_count;
216
5
    }
217
218
    // Discard the active split after the caller decides an error is ignorable, for example a
219
    // stale external-table file listing that returns NOT_FOUND. The next prepare_split() must start
220
    // with no concrete reader or split-local state left from the failed split.
221
1
    virtual Status abort_split() {
222
1
        if (_data_reader.reader != nullptr) {
223
1
            RETURN_IF_ERROR(close_current_reader());
224
1
        } else {
225
0
            _current_task.reset();
226
0
            _current_file_description.reset();
227
0
        }
228
1
        _delete_rows = nullptr;
229
1
        _remaining_table_level_count = -1;
230
1
        _current_split_uses_metadata_count = false;
231
1
        _current_split_pruned = false;
232
1
        return Status::OK();
233
1
    }
234
235
    // Public entry point for reading a table-schema block. The base class opens the current reader,
236
    // advances across EOF, and closes exhausted readers. Subclasses provide protected hooks for
237
    // table-format-specific behavior.
238
145
    virtual Status get_block(Block* block, bool* eos) {
239
145
        SCOPED_TIMER(_profile.exec_timer);
240
145
        DORIS_CHECK(block->columns() == _projected_columns.size());
241
145
        block->clear_column_data(_projected_columns.size());
242
243
178
        while (true) {
244
178
            if (*eos) {
245
0
                return Status::OK();
246
0
            }
247
178
            if (_io_ctx != nullptr && _io_ctx->should_stop) {
248
0
                *eos = true;
249
0
                return Status::OK();
250
0
            }
251
178
            if (!_data_reader.reader) {
252
150
                if (_is_table_level_count_active()) {
253
6
                    RETURN_IF_ERROR(_read_table_level_count(block, eos));
254
6
                    return Status::OK();
255
6
                }
256
144
                RETURN_IF_ERROR(create_next_reader(eos));
257
143
                if (!_data_reader.reader) {
258
36
                    DCHECK(*eos);
259
36
                    return Status::OK();
260
36
                }
261
143
            }
262
263
            // Materialize a reduced row set for upper aggregate operators when aggregate
264
            // pushdown can be applied. This is not the final aggregate result: COUNT emits
265
            // `count` default rows for the upper COUNT(*), and MIN/MAX emits two rows containing
266
            // file-level min/max values for the upper MIN/MAX.
267
135
            if (!_aggregate_pushdown_tried) {
268
107
                SCOPED_TIMER(_profile.pushdown_agg_timer);
269
107
                bool pushed_down = false;
270
107
                const auto status = _try_materialize_aggregate_pushdown_rows(block, &pushed_down);
271
107
                if (!status.ok()) {
272
1
                    if (_io_ctx != nullptr && _io_ctx->should_stop &&
273
1
                        status.is<ErrorCode::END_OF_FILE>()) {
274
1
                        *eos = true;
275
1
                        return Status::OK();
276
1
                    }
277
0
                    return status;
278
1
                }
279
106
                if (pushed_down) {
280
10
                    return Status::OK();
281
10
                }
282
106
            }
283
284
124
            bool current_eof = false;
285
124
            _data_reader.block_template.clear_column_data(
286
124
                    cast_set<int64_t>(_data_reader.file_block_layout.size()));
287
124
            size_t current_rows = 0;
288
124
            RETURN_IF_ERROR(_data_reader.reader->get_block(&_data_reader.block_template,
289
124
                                                           &current_rows, &current_eof));
290
124
            const bool stopped_during_read = _io_ctx != nullptr && _io_ctx->should_stop;
291
124
            if (current_rows == 0) {
292
33
                if (current_eof) {
293
33
                    _current_reader_reached_eof = !stopped_during_read;
294
33
                    RETURN_IF_ERROR(close_current_reader());
295
33
                }
296
33
                continue;
297
33
            }
298
124
            DCHECK_EQ(_data_reader.block_template.columns(), _data_reader.file_block_layout.size())
299
0
                    << _data_reader.block_template.dump_structure();
300
91
#ifndef NDEBUG
301
91
            RETURN_IF_ERROR(_check_file_block_columns("after file reader get_block", current_rows));
302
91
#endif
303
91
            DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
304
91
            RETURN_IF_ERROR(finalize_chunk(block, current_rows));
305
90
#ifndef NDEBUG
306
90
            RETURN_IF_ERROR(
307
90
                    _check_table_block_columns("after finalize_chunk", block, current_rows));
308
90
#endif
309
90
            if (current_eof) {
310
18
                _current_reader_reached_eof = !stopped_during_read;
311
18
                RETURN_IF_ERROR(close_current_reader());
312
18
            }
313
90
            return Status::OK();
314
90
        }
315
145
    }
316
317
    // Close the table reader and the currently active file reader. Subclasses that hold additional
318
    // table-format resources should override this and call TableReader::close() first.
319
107
    virtual Status close() {
320
107
        if (_data_reader.reader) {
321
44
            RETURN_IF_ERROR(close_current_reader());
322
44
        }
323
107
        _current_task.reset();
324
107
        _current_file_description.reset();
325
107
        _remaining_table_level_count = -1;
326
107
        _current_split_uses_metadata_count = false;
327
107
        return Status::OK();
328
107
    }
329
330
4
    int64_t condition_cache_hit_count() const { return _condition_cache_hit_count; }
331
332
    virtual std::string debug_string() const;
333
334
    virtual Status annotate_projected_column(const TFileScanSlotInfo& slot_info,
335
                                             ProjectedColumnBuildContext* context,
336
                                             ColumnDefinition* column) const;
337
338
0
    virtual Status validate_projected_columns(const ProjectedColumnBuildContext& context) const {
339
0
        (void)context;
340
0
        return Status::OK();
341
0
    }
342
343
protected:
344
    // TableReader keeps the active file description both in the scan task and separately for
345
    // creating the physical reader. Table-format readers must update both copies when their
346
    // snapshot protocol guarantees that a file path is never overwritten with different bytes.
347
    // This guarantee lets readers safely build cache keys without mtime; it must not be used for
348
    // ordinary Hive/TVF files whose paths may be overwritten in place.
349
61
    void mark_current_data_file_immutable() {
350
61
        DORIS_CHECK(_current_task != nullptr);
351
61
        DORIS_CHECK(_current_task->data_file != nullptr);
352
61
        DORIS_CHECK(_current_file_description.has_value());
353
61
        _current_task->data_file->is_immutable = true;
354
61
        _current_file_description->is_immutable = true;
355
61
    }
356
357
    std::optional<ColumnDefinition> _find_current_table_column_by_field_id(int32_t field_id,
358
                                                                           DataTypePtr type) const;
359
360
    // Parse deletion vector information from table format specific file description.
361
    virtual Status _parse_deletion_vector_file(const TTableFormatFileDesc& t_desc,
362
77
                                               DeleteFileDesc* desc, bool* has_delete_file) {
363
77
        *has_delete_file = false;
364
77
        return Status::OK();
365
77
    }
366
367
    // Advance to the next reader. This closes the current reader first and then opens the next
368
    // concrete reader. Subclasses should not duplicate this loop.
369
    Status create_next_reader(bool* eos);
370
    virtual Status create_file_reader(std::unique_ptr<FileReader>* reader);
371
65
    virtual TableColumnMappingMode mapping_mode() const { return TableColumnMappingMode::BY_NAME; }
372
105
    virtual Status annotate_file_schema(std::vector<ColumnDefinition>* file_schema) {
373
105
        DORIS_CHECK(file_schema != nullptr);
374
105
        return Status::OK();
375
105
    }
376
377
    // Open the concrete reader for the current split/task and build the file-local scan request.
378
108
    virtual Status open_reader() {
379
108
        SCOPED_TIMER(_profile.open_reader_timer);
380
        // 1. Get file schema and create column mapping.
381
108
        std::vector<ColumnDefinition> file_schema;
382
108
        RETURN_IF_ERROR(_data_reader.reader->get_schema(&file_schema));
383
        // For Paimon/Hudi, FE can provide field ids through `history_schema_info`. Annotate the
384
        // file schema before column mapping when the table format maps columns by field id.
385
108
        RETURN_IF_ERROR(annotate_file_schema(&file_schema));
386
108
        _data_reader.file_schema = file_schema;
387
108
        _mapper_options.mode = mapping_mode();
388
389
108
        _data_reader.column_mapper = _data_reader.reader->create_column_mapper(_mapper_options);
390
108
        DORIS_CHECK(_data_reader.column_mapper != nullptr);
391
108
        RETURN_IF_ERROR(_data_reader.column_mapper->create_mapping(_projected_columns,
392
108
                                                                   _partition_values, file_schema));
393
108
        DORIS_CHECK(_data_reader.column_mapper->mappings().size() == _projected_columns.size());
394
395
        // 2. Build table filters based on conjuncts and column predicates.
396
108
        RETURN_IF_ERROR(_build_table_filters_from_conjuncts());
397
398
        // 3. Create file scan request based on column mapping and table filters, then open file
399
        // reader with the request. File scan request carries row-level expression filters and
400
        // file-level pruning hints. Only expression filters decide returned rows.
401
108
        auto file_request = std::make_shared<FileScanRequest>();
402
108
        RETURN_IF_ERROR(_data_reader.column_mapper->create_scan_request(
403
108
                _table_filters, _projected_columns, file_request.get(), _runtime_state));
404
108
        bool constant_filter_pruned_split = false;
405
108
        RETURN_IF_ERROR(_evaluate_constant_filters(&constant_filter_pruned_split));
406
108
        if (constant_filter_pruned_split) {
407
1
            RETURN_IF_ERROR(close_current_reader());
408
1
            return Status::OK();
409
1
        }
410
        // COUNT(*) has no semantic column argument, but Nereids retains a minimum-width scan slot
411
        // so the scan node still has an output tuple. Record only the current non-predicate file
412
        // columns before table-format hooks add row-position or equality-delete dependencies. This
413
        // marker is independent of aggregate eligibility: with position deletes, for example,
414
        // metadata COUNT must fall back to reading rows, but an arbitrary unsupported TIME_MILLIS
415
        // placeholder still must not be validated or decoded merely to carry the surviving count.
416
107
        if (_push_down_agg_type == TPushAggOp::type::COUNT &&
417
107
            _push_down_count_columns.has_value() && _push_down_count_columns->empty()) {
418
11
            file_request->count_star_placeholder_columns.reserve(
419
11
                    file_request->non_predicate_columns.size());
420
11
            for (const auto& column : file_request->non_predicate_columns) {
421
8
                file_request->count_star_placeholder_columns.push_back(column.column_id());
422
8
            }
423
11
        }
424
107
        RETURN_IF_ERROR(customize_file_scan_request(file_request.get()));
425
107
        RETURN_IF_ERROR(_open_local_filter_exprs(*file_request));
426
107
        _data_reader.file_block_layout.clear();
427
107
        _data_reader.block_template.clear();
428
107
        _data_reader.file_block_layout.resize(file_request->local_positions.size());
429
430
        // 4. Build file block layout from file schema and column mapping. The layout describes
431
        // the block returned by file reader before table-column materialization.
432
151
        for (const auto& [file_column_id, block_position] : file_request->local_positions) {
433
151
            DORIS_CHECK(block_position.value() < _data_reader.file_block_layout.size());
434
151
            const auto* field = _find_column_definition(_data_reader.file_schema, file_column_id);
435
151
            DORIS_CHECK(field != nullptr);
436
437
151
            ColumnDefinition projected_field;
438
151
            {
439
151
                auto it = std::find_if(
440
151
                        file_request->non_predicate_columns.begin(),
441
151
                        file_request->non_predicate_columns.end(),
442
158
                        [&](const LocalColumnIndex& p) { return p.column_id() == file_column_id; });
443
151
                if (it != file_request->non_predicate_columns.end()) {
444
94
                    RETURN_IF_ERROR(project_column_definition(*field, *it, &projected_field));
445
94
                }
446
151
            }
447
151
            {
448
151
                auto it = std::find_if(
449
151
                        file_request->predicate_columns.begin(),
450
151
                        file_request->predicate_columns.end(),
451
151
                        [&](const LocalColumnIndex& p) { return p.column_id() == file_column_id; });
452
151
                if (it != file_request->predicate_columns.end()) {
453
57
                    RETURN_IF_ERROR(project_column_definition(*field, *it, &projected_field));
454
57
                }
455
151
            }
456
151
            _data_reader.file_block_layout[block_position.value()] = {
457
151
                    .file_column_id = file_column_id,
458
151
                    .name = projected_field.name,
459
151
                    .type = projected_field.type,
460
151
            };
461
151
            DORIS_CHECK(_data_reader.file_block_layout[block_position.value()].type != nullptr);
462
151
        }
463
464
        // 5. Prepare block template from file block layout. The block template stores the block
465
        // returned by file reader before table-column materialization.
466
107
        _data_reader.block_template.reserve(_data_reader.file_block_layout.size());
467
151
        for (const auto& column : _data_reader.file_block_layout) {
468
151
            _data_reader.block_template.insert(
469
151
                    {column.type->create_column(), column.type, column.name});
470
151
        }
471
107
        if (VLOG_DEBUG_IS_ON) {
472
0
            VLOG_DEBUG << "TableReader debug: " << debug_string();
473
0
        }
474
107
        RETURN_IF_ERROR(_open_mapping_exprs());
475
107
        RETURN_IF_ERROR(_data_reader.reader->open(file_request));
476
107
        RETURN_IF_ERROR(_init_reader_condition_cache(*file_request));
477
107
        return Status::OK();
478
107
    }
479
480
    Status _build_table_filters_from_conjuncts();
481
    Status _evaluate_partition_prune_conjuncts(const VExprContextSPtrs& conjuncts,
482
                                               bool* can_filter_all);
483
    static bool _is_safe_to_pre_execute(const VExprContextSPtr& conjunct);
484
    Status _build_partition_prune_block(Block* block) const;
485
    Status _open_local_filter_exprs(const FileScanRequest& file_request);
486
    Status _init_reader_condition_cache(const FileScanRequest& file_request);
487
    void _finalize_reader_condition_cache();
488
    bool _should_enable_condition_cache(const FileScanRequest& file_request) const;
489
490
108
    Status _evaluate_constant_filters(bool* can_filter_all) {
491
108
        DORIS_CHECK(can_filter_all != nullptr);
492
108
        DORIS_CHECK_LE(_constant_pruning_safe_filter_count, _table_filters.size());
493
108
        *can_filter_all = false;
494
        // The bound was derived from the original `_conjuncts` order, which includes slotless
495
        // expressions omitted from `_table_filters`. Iterating only this prefix therefore cannot
496
        // skip an unsafe row-level predicate and pre-execute a later constant predicate.
497
136
        for (size_t i = 0; i < _constant_pruning_safe_filter_count; ++i) {
498
29
            const auto& table_filter = _table_filters[i];
499
29
            if (table_filter.conjunct == nullptr) {
500
0
                continue;
501
0
            }
502
29
            DORIS_CHECK(_is_safe_to_pre_execute(table_filter.conjunct));
503
            // RuntimeFilterExpr does not implement execute_column_impl(); it is evaluated by the
504
            // row-level filter path through execute_filter(). Constant split pruning uses
505
            // VExprContext::execute() on a one-row synthetic block, so runtime filters must not be
506
            // pre-executed here even when their referenced slot maps to a constant value.
507
29
            if (table_filter.conjunct->root()->is_rf_wrapper() ||
508
29
                !_table_filter_has_only_constant_entries(table_filter)) {
509
27
                continue;
510
27
            }
511
2
            Block eval_block;
512
2
            RETURN_IF_ERROR(_build_constant_filter_block(table_filter, &eval_block));
513
2
            RowDescriptor row_desc;
514
2
            RETURN_IF_ERROR(table_filter.conjunct->prepare(_runtime_state, row_desc));
515
2
            RETURN_IF_ERROR(table_filter.conjunct->open(_runtime_state));
516
2
            int result_column_id = -1;
517
2
            RETURN_IF_ERROR(table_filter.conjunct->execute(&eval_block, &result_column_id));
518
2
            DORIS_CHECK(result_column_id >= 0);
519
2
            if (_filter_result_filters_all(eval_block.get_by_position(result_column_id).column)) {
520
1
                *can_filter_all = true;
521
1
                return Status::OK();
522
1
            }
523
2
        }
524
107
        return Status::OK();
525
108
    }
526
527
24
    bool _table_filter_has_only_constant_entries(const TableFilter& table_filter) const {
528
24
        const auto& filter_entries = _data_reader.column_mapper->filter_entries();
529
24
        for (const auto global_index : table_filter.global_indices) {
530
24
            const auto entry_it = filter_entries.find(global_index);
531
24
            if (entry_it == filter_entries.end() || !entry_it->second.is_constant()) {
532
22
                return false;
533
22
            }
534
24
        }
535
2
        return !table_filter.global_indices.empty();
536
24
    }
537
538
2
    Status _build_constant_filter_block(const TableFilter& table_filter, Block* eval_block) {
539
2
        DORIS_CHECK(eval_block != nullptr);
540
2
        eval_block->clear();
541
2
        const auto& mappings = _data_reader.column_mapper->mappings();
542
2
        const auto& filter_entries = _data_reader.column_mapper->filter_entries();
543
2
        DORIS_CHECK(mappings.size() == _projected_columns.size());
544
4
        for (size_t column_idx = 0; column_idx < mappings.size(); ++column_idx) {
545
2
            const auto global_index = GlobalIndex(column_idx);
546
2
            const auto& mapping = mappings[column_idx];
547
2
            const auto entry_it = filter_entries.find(global_index);
548
2
            const bool referenced_by_filter =
549
2
                    std::find(table_filter.global_indices.begin(),
550
2
                              table_filter.global_indices.end(),
551
2
                              global_index) != table_filter.global_indices.end();
552
2
            if (referenced_by_filter && entry_it != filter_entries.end() &&
553
2
                entry_it->second.is_constant()) {
554
2
                ColumnPtr constant_column;
555
2
                RETURN_IF_ERROR(_materialize_constant_filter_column(
556
2
                        entry_it->second.constant_index(), &constant_column));
557
2
                eval_block->insert({std::move(constant_column), mapping.table_type,
558
2
                                    mapping.table_column_name});
559
2
            } else {
560
0
                eval_block->insert({mapping.table_type->create_column_const_with_default_value(1),
561
0
                                    mapping.table_type, mapping.table_column_name});
562
0
            }
563
2
        }
564
2
        return Status::OK();
565
2
    }
566
567
2
    Status _materialize_constant_filter_column(ConstantIndex constant_index, ColumnPtr* column) {
568
2
        DORIS_CHECK(column != nullptr);
569
2
        const auto& constant_entry = _data_reader.column_mapper->constant_map().get(constant_index);
570
2
        DORIS_CHECK(constant_entry.expr != nullptr);
571
2
        DORIS_CHECK(constant_entry.type != nullptr);
572
2
        RowDescriptor row_desc;
573
2
        RETURN_IF_ERROR(constant_entry.expr->prepare(_runtime_state, row_desc));
574
2
        RETURN_IF_ERROR(constant_entry.expr->open(_runtime_state));
575
2
        Block eval_block;
576
2
        eval_block.insert({constant_entry.type->create_column_const_with_default_value(1),
577
2
                           constant_entry.type, "__table_reader_constant_filter"});
578
2
        int result_column_id = -1;
579
2
        RETURN_IF_ERROR(constant_entry.expr->execute(&eval_block, &result_column_id));
580
2
        DORIS_CHECK(result_column_id >= 0);
581
2
        *column = eval_block.get_by_position(result_column_id).column;
582
2
        DORIS_CHECK((*column)->size() == 1);
583
2
        return Status::OK();
584
2
    }
585
586
2
    static bool _filter_result_filters_all(const ColumnPtr& filter_column) {
587
2
        DORIS_CHECK(filter_column.get() != nullptr);
588
2
        DORIS_CHECK(filter_column->size() == 1);
589
2
        return !filter_column->get_bool(0);
590
2
    }
591
592
108
    virtual Status customize_file_scan_request(FileScanRequest* file_request) {
593
108
        return _append_delete_predicate(file_request);
594
108
    }
595
596
337
    bool _is_table_level_count_active() const { return _remaining_table_level_count >= 0; }
597
598
11
    Status _materialize_count_rows(size_t rows, Block* block) const {
599
11
        DORIS_CHECK(block != nullptr);
600
11
        DORIS_CHECK(block->columns() > 0 || rows == 0);
601
22
        for (size_t column_idx = 0; column_idx < block->columns(); ++column_idx) {
602
11
            auto column = block->get_by_position(column_idx).type->create_column();
603
11
            column->resize(rows);
604
11
            block->replace_by_position(column_idx, std::move(column));
605
11
        }
606
11
        return Status::OK();
607
11
    }
608
609
6
    Status _read_table_level_count(Block* block, bool* eos) {
610
6
        DORIS_CHECK(block != nullptr);
611
6
        DORIS_CHECK(eos != nullptr);
612
6
        DORIS_CHECK(_push_down_agg_type == TPushAggOp::type::COUNT);
613
6
        DORIS_CHECK(_remaining_table_level_count >= 0);
614
6
        if (_remaining_table_level_count == 0) {
615
2
            _remaining_table_level_count = -1;
616
2
            _current_task.reset();
617
2
            *eos = true;
618
2
            return Status::OK();
619
2
        }
620
621
4
        const int64_t batch_size = _runtime_state == nullptr
622
4
                                           ? _remaining_table_level_count
623
4
                                           : static_cast<int64_t>(_runtime_state->batch_size());
624
4
        const auto rows = std::min(_remaining_table_level_count, batch_size);
625
4
        RETURN_IF_ERROR(_materialize_count_rows(cast_set<size_t>(rows), block));
626
4
        _remaining_table_level_count -= rows;
627
4
        *eos = false;
628
4
        return Status::OK();
629
4
    }
630
631
    void _append_file_scan_column(FileScanRequest* request, LocalColumnId column_id,
632
45
                                  std::vector<LocalColumnIndex>* scan_columns) {
633
45
        DORIS_CHECK(request != nullptr);
634
45
        DORIS_CHECK(scan_columns != nullptr);
635
45
        FileScanRequestBuilder builder(request);
636
45
        Status status;
637
45
        if (scan_columns == &request->predicate_columns) {
638
34
            status = builder.add_predicate_column(column_id);
639
34
        } else {
640
11
            DORIS_CHECK(scan_columns == &request->non_predicate_columns);
641
11
            status = builder.add_non_predicate_column(column_id);
642
11
        }
643
45
        DORIS_CHECK(status.ok()) << status.to_string();
644
45
        if (column_id == LocalColumnId(ROW_POSITION_COLUMN_ID) &&
645
45
            _find_column_definition(_data_reader.file_schema, column_id) == nullptr) {
646
30
            _data_reader.file_schema.push_back(row_position_column_definition());
647
30
        }
648
45
    }
649
650
    // Append DeletePredicate to file scan request if there are deletes. The predicate will be evaluated in file reader level and filter out deleted rows before returning data to table reader.
651
108
    Status _append_delete_predicate(FileScanRequest* request) {
652
108
        DORIS_CHECK(request != nullptr);
653
108
        if ((_delete_rows == nullptr || _delete_rows->empty()) &&
654
108
            (_deletion_vector == nullptr || _deletion_vector->isEmpty())) {
655
92
            return Status::OK();
656
92
        }
657
16
        const auto row_position_column_id = LocalColumnId(ROW_POSITION_COLUMN_ID);
658
16
        _append_file_scan_column(request, row_position_column_id, &request->predicate_columns);
659
660
16
        const auto block_position = request->local_positions.at(row_position_column_id);
661
16
        auto append_predicate = [&](auto& deleted_rows) {
662
16
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
663
16
            delete_predicate->add_child(VSlotRef::create_shared(
664
16
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
665
16
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
666
16
            request->delete_conjuncts.push_back(
667
16
                    VExprContext::create_shared(std::move(delete_predicate)));
668
16
        };
_ZZN5doris6format11TableReader24_append_delete_predicateEPNS0_15FileScanRequestEENKUlRT_E_clISt6vectorIlSaIlEEEEDaS5_
Line
Count
Source
661
12
        auto append_predicate = [&](auto& deleted_rows) {
662
12
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
663
12
            delete_predicate->add_child(VSlotRef::create_shared(
664
12
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
665
12
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
666
12
            request->delete_conjuncts.push_back(
667
12
                    VExprContext::create_shared(std::move(delete_predicate)));
668
12
        };
_ZZN5doris6format11TableReader24_append_delete_predicateEPNS0_15FileScanRequestEENKUlRT_E_clIN7roaring12Roaring64MapEEEDaS5_
Line
Count
Source
661
4
        auto append_predicate = [&](auto& deleted_rows) {
662
4
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
663
4
            delete_predicate->add_child(VSlotRef::create_shared(
664
4
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
665
4
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
666
4
            request->delete_conjuncts.push_back(
667
4
                    VExprContext::create_shared(std::move(delete_predicate)));
668
4
        };
669
16
        if (_delete_rows != nullptr && !_delete_rows->empty()) {
670
12
            append_predicate(*_delete_rows);
671
12
        }
672
16
        if (_deletion_vector != nullptr && !_deletion_vector->isEmpty()) {
673
4
            append_predicate(*_deletion_vector);
674
4
        }
675
16
        return Status::OK();
676
108
    }
677
678
    // Close the current concrete reader. This hook is called by both create_next_reader() and
679
    // close(), so it should remain idempotent.
680
107
    virtual Status close_current_reader() {
681
107
        _finalize_reader_condition_cache();
682
107
        RETURN_IF_ERROR(_data_reader.reader->close());
683
107
        _data_reader.reader.reset();
684
107
        if (_data_reader.column_mapper != nullptr) {
685
106
            _data_reader.column_mapper->clear();
686
106
            _data_reader.column_mapper.reset();
687
106
        }
688
107
        _table_filters.clear();
689
107
        _constant_pruning_safe_filter_count = 0;
690
107
        _data_reader.file_schema.clear();
691
107
        _data_reader.file_block_layout.clear();
692
107
        _data_reader.block_template.clear();
693
107
        _current_task.reset();
694
107
        _current_file_description.reset();
695
107
        _current_reader_reached_eof = false;
696
107
        return Status::OK();
697
107
    }
698
699
2
    void _record_scan_rows(size_t rows) {
700
2
        if (_io_ctx != nullptr && _io_ctx->file_reader_stats != nullptr) {
701
2
            _io_ctx->file_reader_stats->read_rows += rows;
702
2
        }
703
2
    }
704
705
    // Finalize file-local block to table/global schema block.
706
91
    Status finalize_chunk(Block* block, const size_t rows) {
707
91
        SCOPED_TIMER(_profile.finalize_timer);
708
91
        size_t idx = 0;
709
123
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
710
123
            ColumnPtr column;
711
123
            RETURN_IF_ERROR(_materialize_mapping_column(mapping, &_data_reader.block_template, rows,
712
123
                                                        &column));
713
122
            block->replace_by_position(idx, IColumn::mutate(std::move(column)));
714
122
            idx++;
715
122
        }
716
90
        RETURN_IF_ERROR(materialize_virtual_columns(block));
717
        // Enforce CHAR/VARCHAR length declared by the table schema after all file-to-table
718
        // materialization has finished.
719
90
        RETURN_IF_ERROR(_truncate_char_or_varchar_columns(block));
720
90
        return Status::OK();
721
90
    }
722
723
    // Materialize virtual columns in the table block, such as Iceberg _row_id and
724
    // _last_updated_sequence_number. This runs after normal column materialization so finalize
725
    // expressions can reference those virtual columns.
726
56
    virtual Status materialize_virtual_columns(Block* table_block) { return Status::OK(); }
727
728
#ifndef NDEBUG
729
91
    Status _check_file_block_columns(std::string_view stage, size_t rows) {
730
91
        DORIS_CHECK(_data_reader.block_template.columns() == _data_reader.file_block_layout.size());
731
219
        for (size_t idx = 0; idx < _data_reader.block_template.columns(); ++idx) {
732
128
            const auto& file_block_column = _data_reader.file_block_layout[idx];
733
128
            const auto& column_with_type = _data_reader.block_template.get_by_position(idx);
734
128
            const auto* column = column_with_type.column.get();
735
128
            try {
736
128
                if (column == nullptr) {
737
0
                    auto st = Status::InternalError(
738
0
                            "Invalid file block column {} at {}: file_column_id={}, name='{}', "
739
0
                            "type={}, column=null, expected_rows={}, reader={}",
740
0
                            idx, stage, file_block_column.file_column_id.value(),
741
0
                            file_block_column.name,
742
0
                            file_block_column.type == nullptr ? "null"
743
0
                                                              : file_block_column.type->get_name(),
744
0
                            rows, debug_string());
745
0
                    LOG(WARNING) << st;
746
0
                    return st;
747
0
                }
748
128
                column->sanity_check();
749
128
                auto st = column_with_type.check_type_and_column_match();
750
128
                if (!st.ok()) {
751
0
                    auto contextual_status = Status::InternalError(
752
0
                            "Invalid file block column {} at {}: file_column_id={}, name='{}', "
753
0
                            "type={}, column={}, column_size={}, expected_rows={}, error={}, "
754
0
                            "reader={}",
755
0
                            idx, stage, file_block_column.file_column_id.value(),
756
0
                            file_block_column.name,
757
0
                            file_block_column.type == nullptr ? "null"
758
0
                                                              : file_block_column.type->get_name(),
759
0
                            column->get_name(), column->size(), rows, st.to_string(),
760
0
                            debug_string());
761
0
                    LOG(WARNING) << contextual_status;
762
0
                    return contextual_status;
763
0
                }
764
128
            } catch (const Exception& e) {
765
0
                auto st = Status::InternalError(
766
0
                        "Invalid file block column {} at {}: file_column_id={}, name='{}', "
767
0
                        "type={}, column={}, column_size={}, expected_rows={}, error={}, "
768
0
                        "reader={}",
769
0
                        idx, stage, file_block_column.file_column_id.value(),
770
0
                        file_block_column.name,
771
0
                        file_block_column.type == nullptr ? "null"
772
0
                                                          : file_block_column.type->get_name(),
773
0
                        column == nullptr ? "null" : column->get_name(),
774
0
                        column == nullptr ? 0 : column->size(), rows, e.to_string(),
775
0
                        debug_string());
776
0
                LOG(WARNING) << st;
777
0
                return st;
778
0
            } catch (const std::exception& e) {
779
0
                auto st = Status::InternalError(
780
0
                        "Invalid file block column {} at {}: file_column_id={}, name='{}', "
781
0
                        "type={}, column={}, column_size={}, expected_rows={}, error={}, "
782
0
                        "reader={}",
783
0
                        idx, stage, file_block_column.file_column_id.value(),
784
0
                        file_block_column.name,
785
0
                        file_block_column.type == nullptr ? "null"
786
0
                                                          : file_block_column.type->get_name(),
787
0
                        column == nullptr ? "null" : column->get_name(),
788
0
                        column == nullptr ? 0 : column->size(), rows, e.what(), debug_string());
789
0
                LOG(WARNING) << st;
790
0
                return st;
791
0
            }
792
128
        }
793
91
        return Status::OK();
794
91
    }
795
796
90
    Status _check_table_block_columns(std::string_view stage, const Block* block, size_t rows) {
797
90
        DORIS_CHECK(block != nullptr);
798
90
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
799
212
        for (size_t idx = 0; idx < block->columns(); ++idx) {
800
122
            const auto& mapping = _data_reader.column_mapper->mappings()[idx];
801
122
            const auto& column_with_type = block->get_by_position(idx);
802
122
            const auto* column = column_with_type.column.get();
803
122
            try {
804
122
                if (column == nullptr) {
805
0
                    auto st = Status::InternalError(
806
0
                            "Invalid table block column {} at {}: table_column='{}', "
807
0
                            "global_index={}, type={}, column=null, expected_rows={}, mapping={}",
808
0
                            idx, stage, mapping.table_column_name, mapping.global_index.value(),
809
0
                            mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(),
810
0
                            rows, mapping.debug_string());
811
0
                    LOG(WARNING) << st;
812
0
                    return st;
813
0
                }
814
122
                column->sanity_check();
815
122
                auto st = column_with_type.check_type_and_column_match();
816
122
                if (!st.ok()) {
817
0
                    auto contextual_status = Status::InternalError(
818
0
                            "Invalid table block column {} at {}: table_column='{}', "
819
0
                            "global_index={}, type={}, column={}, column_size={}, "
820
0
                            "expected_rows={}, error={}, mapping={}",
821
0
                            idx, stage, mapping.table_column_name, mapping.global_index.value(),
822
0
                            mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(),
823
0
                            column->get_name(), column->size(), rows, st.to_string(),
824
0
                            mapping.debug_string());
825
0
                    LOG(WARNING) << contextual_status;
826
0
                    return contextual_status;
827
0
                }
828
122
            } catch (const Exception& e) {
829
0
                auto st = Status::InternalError(
830
0
                        "Invalid table block column {} at {}: table_column='{}', global_index={}, "
831
0
                        "type={}, column={}, column_size={}, expected_rows={}, error={}, "
832
0
                        "mapping={}",
833
0
                        idx, stage, mapping.table_column_name, mapping.global_index.value(),
834
0
                        mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(),
835
0
                        column == nullptr ? "null" : column->get_name(),
836
0
                        column == nullptr ? 0 : column->size(), rows, e.to_string(),
837
0
                        mapping.debug_string());
838
0
                LOG(WARNING) << st;
839
0
                return st;
840
0
            } catch (const std::exception& e) {
841
0
                auto st = Status::InternalError(
842
0
                        "Invalid table block column {} at {}: table_column='{}', global_index={}, "
843
0
                        "type={}, column={}, column_size={}, expected_rows={}, error={}, "
844
0
                        "mapping={}",
845
0
                        idx, stage, mapping.table_column_name, mapping.global_index.value(),
846
0
                        mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(),
847
0
                        column == nullptr ? "null" : column->get_name(),
848
0
                        column == nullptr ? 0 : column->size(), rows, e.what(),
849
0
                        mapping.debug_string());
850
0
                LOG(WARNING) << st;
851
0
                return st;
852
0
            }
853
122
        }
854
90
        return Status::OK();
855
90
    }
856
#endif
857
858
90
    Status _truncate_char_or_varchar_columns(Block* block) {
859
90
        DORIS_CHECK(block != nullptr);
860
90
        if (_runtime_state == nullptr ||
861
90
            !_runtime_state->query_options().truncate_char_or_varchar_columns) {
862
90
            return Status::OK();
863
90
        }
864
0
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
865
0
        for (size_t idx = 0; idx < _data_reader.column_mapper->mappings().size(); ++idx) {
866
0
            const auto& mapping = _data_reader.column_mapper->mappings()[idx];
867
0
            if (!_should_truncate_char_or_varchar_column(mapping)) {
868
0
                continue;
869
0
            }
870
0
            const auto target_len =
871
0
                    assert_cast<const DataTypeString*>(remove_nullable(mapping.table_type).get())
872
0
                            ->len();
873
0
            _truncate_char_or_varchar_column(block, idx, target_len);
874
0
        }
875
0
        return Status::OK();
876
90
    }
877
878
    // Return true when the table schema has a bounded CHAR/VARCHAR length that is stricter than
879
    // the file-side type. Examples:
880
    // - table VARCHAR(10), file VARCHAR(20): truncate to 10;
881
    // - table VARCHAR(10), file STRING: truncate to 10 because STRING has no declared bound;
882
    // - table STRING, any file type: no truncation because the target has no bound.
883
5
    static bool _should_truncate_char_or_varchar_column(const ColumnMapping& mapping) {
884
5
        if (mapping.table_type == nullptr) {
885
0
            return false;
886
0
        }
887
5
        const auto table_type = remove_nullable(mapping.table_type);
888
5
        const auto primitive_type = table_type->get_primitive_type();
889
5
        if (primitive_type != TYPE_VARCHAR && primitive_type != TYPE_CHAR) {
890
1
            return false;
891
1
        }
892
4
        const auto target_len = assert_cast<const DataTypeString*>(table_type.get())->len();
893
4
        if (target_len <= 0) {
894
0
            return false;
895
0
        }
896
4
        if (mapping.file_type == nullptr) {
897
0
            return true;
898
0
        }
899
4
        const auto file_type = remove_nullable(mapping.file_type);
900
4
        DORIS_CHECK(file_type != nullptr);
901
4
        int file_len = -1;
902
4
        if (file_type->get_primitive_type() == TYPE_VARCHAR ||
903
4
            file_type->get_primitive_type() == TYPE_CHAR ||
904
4
            file_type->get_primitive_type() == TYPE_STRING) {
905
3
            file_len = assert_cast<const DataTypeString*>(file_type.get())->len();
906
3
        }
907
908
4
        return file_len < 0 || target_len < file_len;
909
4
    }
910
911
    // Truncate a materialized CHAR/VARCHAR column in place by reusing the vectorized substring
912
    // implementation: substring(column, 1, len). Nullable columns are unwrapped before substring
913
    // execution and wrapped back with the original null map afterward, because substring operates
914
    // on the nested string payload only.
915
1
    static void _truncate_char_or_varchar_column(Block* block, size_t idx, int len) {
916
1
        DORIS_CHECK(block != nullptr);
917
1
        auto int_type = std::make_shared<DataTypeInt32>();
918
1
        const auto num_columns_without_result = cast_set<uint32_t>(block->columns());
919
1
        auto& target = block->get_by_position(idx);
920
1
        const bool is_nullable = target.type->is_nullable();
921
1
        ColumnPtr input_column = target.column;
922
1
        ColumnPtr null_map_column;
923
1
        if (is_nullable) {
924
1
            const auto* nullable_column = assert_cast<const ColumnNullable*>(target.column.get());
925
1
            input_column = nullable_column->get_nested_column_ptr();
926
1
            null_map_column = nullable_column->get_null_map_column_ptr();
927
1
        }
928
1
        block->replace_by_position(idx, std::move(input_column));
929
1
        block->insert({int_type->create_column_const(block->rows(), to_field<TYPE_INT>(1)),
930
1
                       int_type, "const 1"});
931
1
        block->insert({int_type->create_column_const(block->rows(), to_field<TYPE_INT>(len)),
932
1
                       int_type, "const len"});
933
1
        block->insert({nullptr, std::make_shared<DataTypeString>(), "result"});
934
935
1
        ColumnNumbers temp_arguments(3);
936
1
        temp_arguments[0] = cast_set<uint32_t>(idx);
937
1
        temp_arguments[1] = num_columns_without_result;
938
1
        temp_arguments[2] = num_columns_without_result + 1;
939
1
        const uint32_t result_column_id = num_columns_without_result + 2;
940
1
        SubstringUtil::substring_execute(*block, temp_arguments, result_column_id, block->rows());
941
942
1
        ColumnPtr result_column = block->get_by_position(result_column_id).column;
943
1
        if (is_nullable) {
944
1
            result_column = ColumnNullable::create(std::move(result_column), null_map_column);
945
1
        }
946
1
        block->replace_by_position(idx, std::move(result_column));
947
1
        block->erase_tail(num_columns_without_result);
948
1
    }
949
950
107
    Status _try_materialize_aggregate_pushdown_rows(Block* block, bool* pushed_down) {
951
107
        DORIS_CHECK(block != nullptr);
952
107
        DORIS_CHECK(pushed_down != nullptr);
953
107
        *pushed_down = false;
954
107
        block->clear_column_data(_projected_columns.size());
955
107
        _aggregate_pushdown_tried = true;
956
107
        if (!_supports_aggregate_pushdown(_push_down_agg_type)) {
957
95
            return Status::OK();
958
95
        }
959
960
12
        FileAggregateRequest file_request;
961
12
        RETURN_IF_ERROR(_build_file_aggregate_request(_push_down_agg_type, &file_request));
962
12
        FileAggregateResult file_result;
963
12
        const auto status = _data_reader.reader->get_aggregate_result(file_request, &file_result);
964
12
        if (status.is<ErrorCode::NOT_IMPLEMENTED_ERROR>()) {
965
1
            return Status::OK();
966
1
        }
967
11
        RETURN_IF_ERROR(status);
968
10
        RETURN_IF_ERROR(
969
10
                _materialize_aggregate_pushdown_rows(_push_down_agg_type, file_result, block));
970
10
        if (_push_down_agg_type == TPushAggOp::type::COUNT) {
971
7
            _current_split_uses_metadata_count = true;
972
7
        }
973
10
        *pushed_down = true;
974
10
        RETURN_IF_ERROR(close_current_reader());
975
10
        return Status::OK();
976
10
    }
977
978
119
    virtual bool _supports_aggregate_pushdown(TPushAggOp::type agg_type) const {
979
        // Only COUNT and MIN/MAX can be push down.
980
119
        if (agg_type != TPushAggOp::type::COUNT && agg_type != TPushAggOp::type::MINMAX) {
981
76
            return false;
982
76
        }
983
        // Aggregate pushdown returns reduced synthetic rows and may close the physical reader
984
        // before the next scheduler turn. If a runtime filter is still pending, those rows could
985
        // escape before the filter arrives and cannot later be reconstructed from real file rows.
986
        // This is the same irreversibility constraint as table-level metadata COUNT, and applies
987
        // to COUNT and MIN/MAX for Parquet/ORC as well as COUNT for text readers.
988
43
        if (!_all_runtime_filters_applied_for_split) {
989
2
            return false;
990
2
        }
991
        // Scanner owns the original conjunct list and evaluates it after TableReader finalizes
992
        // rows. Even a slotless conjunct that cannot become a TableFilter must see every source
993
        // row before an aggregate reduces the stream to synthetic COUNT/MINMAX rows.
994
41
        if (!_conjuncts.empty()) {
995
5
            return false;
996
5
        }
997
        // Only support aggregate pushdown when there is no delete or filter, so
998
        // the reduced rows consumed by the upper aggregate remain semantically equivalent to a
999
        // normal scan.
1000
36
        if ((_delete_rows != nullptr && !_delete_rows->empty()) ||
1001
36
            (_deletion_vector != nullptr && !_deletion_vector->isEmpty())) {
1002
3
            return false;
1003
3
        }
1004
33
        if (!_table_filters.empty()) {
1005
0
            return false;
1006
0
        }
1007
33
        if (agg_type == TPushAggOp::type::COUNT) {
1008
            // Old FEs do not serialize push_down_count_slot_ids. During the supported BE-first
1009
            // rolling upgrade, nullopt therefore means "COUNT semantics are unknown", not
1010
            // COUNT(*). Fall back to reading rows until the FE explicitly sends either an empty
1011
            // list for COUNT(*) or one slot for COUNT(col).
1012
22
            if (!_push_down_count_columns.has_value()) {
1013
3
                return false;
1014
3
            }
1015
            // COUNT(*) needs no column metadata. COUNT(col) currently supports one direct file
1016
            // column; multiple COUNT arguments fall back to the normal scan so every upper
1017
            // aggregate receives the original rows.
1018
19
            if (_push_down_count_columns->empty()) {
1019
10
                return true;
1020
10
            }
1021
9
            if (_push_down_count_columns->size() != 1) {
1022
1
                return false;
1023
1
            }
1024
8
            const auto& mapping = _push_down_count_mapping();
1025
            // Metadata COUNT skips TableReader's normal materialization path. Only a trivial
1026
            // mapping is safe: for example, a nullable Parquet INT mapped to a NOT NULL table
1027
            // BIGINT normally needs both an INT->BIGINT cast and nullability validation. Counting
1028
            // footer values directly would bypass both operations and could hide invalid data.
1029
8
            return mapping.file_local_id.has_value() && mapping.file_type != nullptr &&
1030
8
                   mapping.table_type != nullptr && mapping.is_trivial &&
1031
8
                   mapping.virtual_column_type == TableVirtualColumnType::INVALID &&
1032
8
                   mapping.default_expr == nullptr;
1033
9
        }
1034
        // For MIN/MAX, only support direct file-to-table column mappings. The two emitted rows
1035
        // must be enough for the upper MIN/MAX aggregate without evaluating default expressions or
1036
        // virtual columns.
1037
13
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
1038
13
            if (!mapping.file_local_id.has_value() ||
1039
13
                mapping.virtual_column_type != TableVirtualColumnType::INVALID ||
1040
13
                mapping.default_expr != nullptr || mapping.file_type == nullptr ||
1041
13
                mapping.table_type == nullptr) {
1042
1
                return false;
1043
1
            }
1044
12
            if (!_can_push_down_minmax_for_mapping(mapping)) {
1045
2
                return false;
1046
2
            }
1047
12
        }
1048
8
        return true;
1049
11
    }
1050
1051
121
    static ColumnPtr _detach_column(ColumnPtr column) {
1052
121
        DORIS_CHECK(column.get() != nullptr);
1053
121
        return IColumn::mutate(std::move(column));
1054
121
    }
1055
1056
62
    static Status _align_column_nullability(ColumnPtr* column, const DataTypePtr& table_type) {
1057
62
        DORIS_CHECK(column != nullptr);
1058
62
        DORIS_CHECK(column->get() != nullptr);
1059
62
        DORIS_CHECK(table_type != nullptr);
1060
        // Must return non-const column
1061
62
        *column = (*column)->convert_to_full_column_if_const();
1062
62
        if (table_type->is_nullable()) {
1063
25
            const auto& nested_type =
1064
25
                    assert_cast<const DataTypeNullable&>(*table_type).get_nested_type();
1065
25
            if (!(*column)->is_nullable()) {
1066
2
                RETURN_IF_ERROR(_align_column_nullability(column, nested_type));
1067
2
                *column = make_nullable(*column);
1068
2
                return Status::OK();
1069
2
            }
1070
23
            const auto& nullable_column = assert_cast<const ColumnNullable&>(**column);
1071
23
            ColumnPtr nested_column = nullable_column.get_nested_column_ptr();
1072
23
            RETURN_IF_ERROR(_align_column_nullability(&nested_column, nested_type));
1073
23
            *column = ColumnNullable::create(nested_column,
1074
23
                                             nullable_column.get_null_map_column_ptr());
1075
23
            return Status::OK();
1076
23
        }
1077
37
        if ((*column)->is_nullable()) {
1078
0
            const auto& nullable_column = assert_cast<const ColumnNullable&>(**column);
1079
0
            if (nullable_column.has_null()) {
1080
0
                return Status::InternalError(
1081
0
                        "Default expression produced NULL for non-nullable table column");
1082
0
            }
1083
0
            ColumnPtr nested_column = nullable_column.get_nested_column_ptr();
1084
0
            RETURN_IF_ERROR(_align_column_nullability(&nested_column, table_type));
1085
0
            *column = nested_column;
1086
0
            return Status::OK();
1087
0
        }
1088
37
        if (const auto* array_type = typeid_cast<const DataTypeArray*>(table_type.get())) {
1089
1
            const auto& array_column = assert_cast<const ColumnArray&>(**column);
1090
1
            ColumnPtr nested_column = array_column.get_data_ptr();
1091
1
            RETURN_IF_ERROR(
1092
1
                    _align_column_nullability(&nested_column, array_type->get_nested_type()));
1093
1
            *column = ColumnArray::create(nested_column, array_column.get_offsets_ptr());
1094
1
            return Status::OK();
1095
1
        }
1096
36
        if (const auto* map_type = typeid_cast<const DataTypeMap*>(table_type.get())) {
1097
0
            const auto& map_column = assert_cast<const ColumnMap&>(**column);
1098
0
            ColumnPtr key_column = map_column.get_keys_ptr();
1099
0
            ColumnPtr value_column = map_column.get_values_ptr();
1100
0
            RETURN_IF_ERROR(_align_column_nullability(&key_column, map_type->get_key_type()));
1101
0
            RETURN_IF_ERROR(_align_column_nullability(&value_column, map_type->get_value_type()));
1102
0
            *column = ColumnMap::create(key_column, value_column, map_column.get_offsets_ptr());
1103
0
            return Status::OK();
1104
0
        }
1105
36
        if (const auto* struct_type = typeid_cast<const DataTypeStruct*>(table_type.get())) {
1106
5
            const auto& struct_column = assert_cast<const ColumnStruct&>(**column);
1107
5
            Columns columns = struct_column.get_columns_copy();
1108
5
            DORIS_CHECK(columns.size() == struct_type->get_elements().size());
1109
16
            for (size_t i = 0; i < columns.size(); ++i) {
1110
11
                RETURN_IF_ERROR(
1111
11
                        _align_column_nullability(&columns[i], struct_type->get_element(i)));
1112
11
            }
1113
5
            *column = ColumnStruct::create(columns);
1114
5
            return Status::OK();
1115
5
        }
1116
31
        return Status::OK();
1117
36
    }
1118
1119
    static Status _execute_default_expr_without_root_type_check(
1120
            const VExprContextSPtr& default_expr, const Block* block,
1121
7
            ColumnWithTypeAndName* result_data) {
1122
7
        DORIS_CHECK(default_expr != nullptr);
1123
7
        DORIS_CHECK(block != nullptr);
1124
7
        DORIS_CHECK(result_data != nullptr);
1125
7
        ColumnPtr result_column;
1126
7
        Status st;
1127
7
        RETURN_IF_CATCH_EXCEPTION({
1128
7
            st = default_expr->root()->execute_column_impl(default_expr.get(), block, nullptr,
1129
7
                                                           block->rows(), result_column);
1130
7
        });
1131
7
        RETURN_IF_ERROR(st);
1132
7
        DORIS_CHECK(result_column.get() != nullptr);
1133
7
        if (result_column->size() != block->rows()) {
1134
0
            return Status::InternalError(
1135
0
                    "Default expr {} return column size {} not equal to expected size {}",
1136
0
                    default_expr->expr_name(), result_column->size(), block->rows());
1137
0
        }
1138
7
        result_data->column = result_column;
1139
7
        result_data->type = default_expr->execute_type(block);
1140
7
        result_data->name = default_expr->expr_name();
1141
7
        return Status::OK();
1142
7
    }
1143
1144
    Status _cast_column_to_type(ColumnPtr* column, const DataTypePtr& file_type,
1145
                                const DataTypePtr& table_type,
1146
1
                                const std::string& column_name) const {
1147
1
        DORIS_CHECK(column != nullptr);
1148
1
        DORIS_CHECK(column->get() != nullptr);
1149
1
        DORIS_CHECK(file_type != nullptr);
1150
1
        DORIS_CHECK(table_type != nullptr);
1151
1
        if (file_type->equals(*table_type)) {
1152
0
            return Status::OK();
1153
0
        }
1154
1155
1
        DataTypePtr input_type = file_type;
1156
        // Cast wrappers unwrap nullable inputs according to the declared input type, so keep the
1157
        // root nullability of the declared type aligned with the actual column shape.
1158
1
        if ((*column)->is_nullable() && !input_type->is_nullable()) {
1159
0
            input_type = make_nullable(input_type);
1160
1
        } else if (!(*column)->is_nullable() && input_type->is_nullable()) {
1161
1
            input_type = remove_nullable(input_type);
1162
1
        }
1163
1
        Block cast_block;
1164
1
        cast_block.insert({*column, input_type, column_name});
1165
1
        auto slot_ref = VSlotRef::create_shared(0, 0, -1, input_type, column_name);
1166
1
        auto cast_expr = Cast::create_shared(table_type);
1167
1
        cast_expr->add_child(std::move(slot_ref));
1168
1
        auto cast_ctx = VExprContext::create_shared(std::move(cast_expr));
1169
1
        RowDescriptor row_desc;
1170
1
        RETURN_IF_ERROR(cast_ctx->prepare(_runtime_state, row_desc));
1171
1
        RETURN_IF_ERROR(cast_ctx->open(_runtime_state));
1172
1
        ColumnPtr cast_column;
1173
1
        RETURN_IF_ERROR(cast_ctx->execute(&cast_block, cast_column));
1174
1
        *column = std::move(cast_column);
1175
1
        return Status::OK();
1176
1
    }
1177
1178
    Status _materialize_present_child_mapping_column(const ColumnMapping& mapping,
1179
                                                     const ColumnPtr& file_column,
1180
18
                                                     const size_t rows, ColumnPtr* column) {
1181
18
        DORIS_CHECK(column != nullptr);
1182
18
        DORIS_CHECK(mapping.file_type != nullptr);
1183
18
        DORIS_CHECK(mapping.table_type != nullptr);
1184
18
        *column = file_column;
1185
18
        if (!mapping.is_trivial) {
1186
5
            if (!mapping.child_mappings.empty()) {
1187
5
                RETURN_IF_ERROR(
1188
5
                        _materialize_complex_mapping_column(mapping, *column, rows, column));
1189
5
            } else {
1190
0
                RETURN_IF_ERROR(_cast_column_to_type(column, mapping.file_type, mapping.table_type,
1191
0
                                                     mapping.file_column_name));
1192
0
            }
1193
5
        }
1194
18
        RETURN_IF_ERROR(_align_column_nullability(column, mapping.table_type));
1195
18
        return Status::OK();
1196
18
    }
1197
1198
    Status _materialize_mapping_column(const ColumnMapping& mapping, Block* current_block,
1199
127
                                       const size_t rows, ColumnPtr* column) {
1200
127
        if (!mapping.is_trivial && mapping.file_local_id.has_value() &&
1201
127
            !mapping.child_mappings.empty()) {
1202
5
            DCHECK(mapping.projection != nullptr);
1203
5
            int res_id;
1204
5
            auto st = mapping.projection->execute(current_block, &res_id);
1205
5
            if (!st.ok()) {
1206
0
                return Status::InternalError(
1207
0
                        "Failed to execute complex mapping projection for table column '{}' "
1208
0
                        "(global_index={}, file_local_id={}, rows={}): {}, mapping={}",
1209
0
                        mapping.table_column_name, mapping.global_index.value(),
1210
0
                        *mapping.file_local_id, rows, st.to_string(), mapping.debug_string());
1211
0
            }
1212
5
            ColumnPtr result_column = current_block->get_by_position(res_id).column;
1213
5
            RETURN_IF_ERROR(
1214
5
                    _materialize_complex_mapping_column(mapping, result_column, rows, column));
1215
5
            return Status::OK();
1216
5
        }
1217
122
        if (mapping.projection != nullptr) {
1218
102
            int res_id;
1219
102
            auto st = mapping.projection->execute(current_block, &res_id);
1220
102
            if (!st.ok()) {
1221
1
                std::string file_local_id = "null";
1222
1
                if (mapping.file_local_id.has_value()) {
1223
1
                    file_local_id = std::to_string(*mapping.file_local_id);
1224
1
                }
1225
1
                return Status::InternalError(
1226
1
                        "Failed to execute mapping projection for table column '{}' "
1227
1
                        "(global_index={}, file_local_id={}, rows={}): {}, mapping={}",
1228
1
                        mapping.table_column_name, mapping.global_index.value(), file_local_id,
1229
1
                        rows, st.to_string(), mapping.debug_string());
1230
1
            }
1231
101
            ColumnPtr result_column = current_block->get_by_position(res_id).column;
1232
101
            *column = _detach_column(std::move(result_column));
1233
101
            return Status::OK();
1234
102
        }
1235
20
        if (mapping.default_expr != nullptr) {
1236
7
            if (current_block->rows() == rows) {
1237
0
                ColumnWithTypeAndName result;
1238
0
                RETURN_IF_ERROR(_execute_default_expr_without_root_type_check(
1239
0
                        mapping.default_expr, current_block, &result));
1240
0
                ColumnPtr result_column = result.column;
1241
0
                RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type));
1242
0
                *column = _detach_column(std::move(result_column));
1243
7
            } else {
1244
7
                DORIS_CHECK(mapping.constant_index.has_value());
1245
7
                Block eval_block;
1246
7
                eval_block.insert({mapping.table_type->create_column_const_with_default_value(rows),
1247
7
                                   mapping.table_type, "__table_reader_const_rows"});
1248
7
                ColumnWithTypeAndName result;
1249
7
                RETURN_IF_ERROR(_execute_default_expr_without_root_type_check(
1250
7
                        mapping.default_expr, &eval_block, &result));
1251
7
                ColumnPtr result_column = result.column;
1252
7
                RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type));
1253
7
                *column = _detach_column(std::move(result_column));
1254
7
            }
1255
7
            return Status::OK();
1256
7
        }
1257
13
        ColumnPtr result_column = mapping.table_type->create_column_const_with_default_value(rows);
1258
13
        *column = _detach_column(std::move(result_column));
1259
13
        return Status::OK();
1260
20
    }
1261
1262
    Status _materialize_complex_mapping_column(const ColumnMapping& mapping,
1263
                                               const ColumnPtr& file_column, const size_t rows,
1264
10
                                               ColumnPtr* column) {
1265
10
        DORIS_CHECK(mapping.table_type != nullptr);
1266
10
        DORIS_CHECK(file_column.get() != nullptr);
1267
10
        const auto table_type = remove_nullable(mapping.table_type);
1268
10
        switch (table_type->get_primitive_type()) {
1269
7
        case TYPE_STRUCT:
1270
7
            RETURN_IF_ERROR(_materialize_struct_mapping_column(mapping, file_column, rows, column));
1271
7
            break;
1272
7
        case TYPE_ARRAY:
1273
2
            RETURN_IF_ERROR(_materialize_array_mapping_column(mapping, file_column, rows, column));
1274
2
            break;
1275
2
        case TYPE_MAP:
1276
1
            RETURN_IF_ERROR(_materialize_map_mapping_column(mapping, file_column, rows, column));
1277
1
            break;
1278
1
        default:
1279
0
            *column = _detach_column(file_column);
1280
0
            break;
1281
10
        }
1282
10
        return Status::OK();
1283
10
    }
1284
1285
    static std::vector<const ColumnMapping*> _present_child_mappings_in_file_order(
1286
7
            const std::vector<ColumnMapping>& child_mappings) {
1287
7
        std::vector<const ColumnMapping*> result;
1288
7
        result.reserve(child_mappings.size());
1289
16
        for (const auto& child_mapping : child_mappings) {
1290
16
            if (child_mapping.file_local_id.has_value()) {
1291
11
                result.push_back(&child_mapping);
1292
11
            }
1293
16
        }
1294
7
        std::ranges::sort(result, [](const ColumnMapping* lhs, const ColumnMapping* rhs) {
1295
6
            DORIS_CHECK(lhs->file_local_id.has_value());
1296
6
            DORIS_CHECK(rhs->file_local_id.has_value());
1297
6
            return *lhs->file_local_id < *rhs->file_local_id;
1298
6
        });
1299
7
        return result;
1300
7
    }
1301
1302
    static size_t _file_child_ordinal_for_mapping(
1303
            const ColumnMapping& mapping, const ColumnMapping& child_mapping,
1304
11
            const std::vector<const ColumnMapping*>& file_ordered_children) {
1305
11
        DORIS_CHECK(child_mapping.file_local_id.has_value());
1306
11
        if (!mapping.projected_file_children.empty()) {
1307
7
            const auto child_it = std::ranges::find_if(
1308
10
                    mapping.projected_file_children, [&](const ColumnDefinition& file_child) {
1309
10
                        return file_child.file_local_id() == *child_mapping.file_local_id;
1310
10
                    });
1311
7
            DORIS_CHECK(child_it != mapping.projected_file_children.end());
1312
7
            return static_cast<size_t>(
1313
7
                    std::distance(mapping.projected_file_children.begin(), child_it));
1314
7
        }
1315
4
        const auto child_it = std::ranges::find(file_ordered_children, &child_mapping);
1316
4
        DORIS_CHECK(child_it != file_ordered_children.end());
1317
4
        return static_cast<size_t>(std::distance(file_ordered_children.begin(), child_it));
1318
11
    }
1319
1320
    static std::vector<const ColumnMapping*> _child_mappings_in_table_type_order(
1321
7
            const ColumnMapping& mapping, const DataTypeStruct& table_type) {
1322
7
        std::vector<const ColumnMapping*> result;
1323
7
        result.reserve(mapping.child_mappings.size());
1324
23
        for (size_t child_idx = 0; child_idx < table_type.get_elements().size(); ++child_idx) {
1325
16
            const auto& child_name = table_type.get_element_name(child_idx);
1326
16
            const auto child_it = std::ranges::find_if(
1327
28
                    mapping.child_mappings, [&](const ColumnMapping& child_mapping) {
1328
28
                        return child_mapping.table_column_name == child_name;
1329
28
                    });
1330
16
            DORIS_CHECK(child_it != mapping.child_mappings.end())
1331
0
                    << mapping.debug_string() << ", table_child_name=" << child_name;
1332
16
            result.push_back(&*child_it);
1333
16
        }
1334
7
        return result;
1335
7
    }
1336
1337
    static const IColumn* _nested_column_if_nullable(const ColumnPtr& column,
1338
12
                                                     const NullMap** null_map) {
1339
12
        DORIS_CHECK(column.get() != nullptr);
1340
12
        if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1341
8
            if (null_map != nullptr) {
1342
8
                *null_map = &nullable_column->get_null_map_data();
1343
8
            }
1344
8
            return &nullable_column->get_nested_column();
1345
8
        }
1346
4
        return column.get();
1347
12
    }
1348
1349
    Status _materialize_struct_mapping_column(const ColumnMapping& mapping,
1350
                                              const ColumnPtr& file_column, const size_t rows,
1351
7
                                              ColumnPtr* column) {
1352
7
        DORIS_CHECK(mapping.table_type != nullptr);
1353
7
        const auto* table_type =
1354
7
                assert_cast<const DataTypeStruct*>(remove_nullable(mapping.table_type).get());
1355
7
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1356
7
        const NullMap* parent_null_map = nullptr;
1357
7
        const auto* nested_file_column =
1358
7
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1359
7
        const auto* file_struct = assert_cast<const ColumnStruct*>(nested_file_column);
1360
7
        DORIS_CHECK(table_type->get_elements().size() == mapping.child_mappings.size());
1361
1362
7
        Columns child_columns;
1363
7
        child_columns.reserve(mapping.child_mappings.size());
1364
7
        const auto file_ordered_children =
1365
7
                _present_child_mappings_in_file_order(mapping.child_mappings);
1366
7
        const auto table_ordered_children =
1367
7
                _child_mappings_in_table_type_order(mapping, *table_type);
1368
16
        for (const auto* child_mapping : table_ordered_children) {
1369
16
            DORIS_CHECK(child_mapping != nullptr);
1370
16
            if (!child_mapping->file_local_id.has_value()) {
1371
5
                child_columns.push_back(
1372
5
                        child_mapping->table_type->create_column_const_with_default_value(rows)
1373
5
                                ->convert_to_full_column_if_const());
1374
5
                continue;
1375
5
            }
1376
11
            const auto file_child_idx =
1377
11
                    _file_child_ordinal_for_mapping(mapping, *child_mapping, file_ordered_children);
1378
11
            DORIS_CHECK(file_child_idx < file_struct->get_columns().size());
1379
11
            ColumnPtr child_column = file_struct->get_column_ptr(file_child_idx);
1380
11
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(*child_mapping, child_column,
1381
11
                                                                      rows, &child_column));
1382
11
            child_columns.push_back(std::move(child_column));
1383
11
        }
1384
7
        MutableColumns mutable_child_columns;
1385
7
        mutable_child_columns.reserve(child_columns.size());
1386
16
        for (auto& child_column : child_columns) {
1387
16
            mutable_child_columns.push_back(IColumn::mutate(std::move(child_column)));
1388
16
        }
1389
7
        auto result = ColumnStruct::create(std::move(mutable_child_columns));
1390
7
        if (mapping.table_type->is_nullable()) {
1391
5
            auto null_map = ColumnUInt8::create();
1392
5
            auto& null_map_data = null_map->get_data();
1393
5
            null_map_data.resize(rows);
1394
5
            if (parent_null_map != nullptr) {
1395
5
                DORIS_CHECK(parent_null_map->size() == rows);
1396
5
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
1397
5
            } else {
1398
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
1399
0
            }
1400
5
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
1401
5
        } else {
1402
2
            *column = std::move(result);
1403
2
        }
1404
7
        return Status::OK();
1405
7
    }
1406
1407
    Status _materialize_array_mapping_column(const ColumnMapping& mapping,
1408
                                             const ColumnPtr& file_column, const size_t rows,
1409
2
                                             ColumnPtr* column) {
1410
2
        DORIS_CHECK(mapping.child_mappings.size() == 1);
1411
2
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1412
2
        const NullMap* parent_null_map = nullptr;
1413
2
        const auto* nested_file_column =
1414
2
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1415
2
        const auto* file_array = assert_cast<const ColumnArray*>(nested_file_column);
1416
2
        ColumnPtr nested_column = file_array->get_data_ptr();
1417
2
        const auto& element_mapping = mapping.child_mappings[0];
1418
2
        RETURN_IF_ERROR(_materialize_present_child_mapping_column(
1419
2
                element_mapping, nested_column, nested_column->size(), &nested_column));
1420
2
        auto offsets_column = file_array->get_offsets_ptr()->convert_to_full_column_if_const();
1421
2
        auto result = ColumnArray::create(IColumn::mutate(std::move(nested_column)),
1422
2
                                          IColumn::mutate(std::move(offsets_column)));
1423
2
        if (mapping.table_type->is_nullable()) {
1424
2
            auto null_map = ColumnUInt8::create();
1425
2
            auto& null_map_data = null_map->get_data();
1426
2
            null_map_data.resize(rows);
1427
2
            if (parent_null_map != nullptr) {
1428
2
                DORIS_CHECK(parent_null_map->size() == rows);
1429
2
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
1430
2
            } else {
1431
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
1432
0
            }
1433
2
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
1434
2
        } else {
1435
0
            *column = std::move(result);
1436
0
        }
1437
2
        return Status::OK();
1438
2
    }
1439
1440
    Status _materialize_map_mapping_column(const ColumnMapping& mapping,
1441
                                           const ColumnPtr& file_column, const size_t rows,
1442
3
                                           ColumnPtr* column) {
1443
3
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1444
3
        const NullMap* parent_null_map = nullptr;
1445
3
        const auto* nested_file_column =
1446
3
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1447
3
        const auto* file_map = assert_cast<const ColumnMap*>(nested_file_column);
1448
3
        ColumnPtr key_column = file_map->get_keys_ptr();
1449
3
        ColumnPtr value_column = file_map->get_values_ptr();
1450
1451
3
        const ColumnMapping* key_mapping = nullptr;
1452
3
        const ColumnMapping* value_mapping = nullptr;
1453
5
        for (const auto& child_mapping : mapping.child_mappings) {
1454
5
            if (!child_mapping.file_local_id.has_value()) {
1455
0
                continue;
1456
0
            }
1457
5
            if (*child_mapping.file_local_id == 0) {
1458
2
                key_mapping = &child_mapping;
1459
3
            } else if (*child_mapping.file_local_id == 1) {
1460
3
                value_mapping = &child_mapping;
1461
3
            }
1462
5
        }
1463
1464
3
        if (key_mapping != nullptr) {
1465
2
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(
1466
2
                    *key_mapping, key_column, key_column->size(), &key_column));
1467
2
        }
1468
3
        if (value_mapping != nullptr) {
1469
3
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(
1470
3
                    *value_mapping, value_column, value_column->size(), &value_column));
1471
3
        }
1472
3
        auto offsets_column = file_map->get_offsets_ptr()->convert_to_full_column_if_const();
1473
3
        auto result = ColumnMap::create(IColumn::mutate(std::move(key_column)),
1474
3
                                        IColumn::mutate(std::move(value_column)),
1475
3
                                        IColumn::mutate(std::move(offsets_column)));
1476
3
        if (mapping.table_type->is_nullable()) {
1477
1
            auto null_map = ColumnUInt8::create();
1478
1
            auto& null_map_data = null_map->get_data();
1479
1
            null_map_data.resize(rows);
1480
1
            if (parent_null_map != nullptr) {
1481
1
                DORIS_CHECK(parent_null_map->size() == rows);
1482
1
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
1483
1
            } else {
1484
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
1485
0
            }
1486
1
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
1487
2
        } else {
1488
2
            *column = std::move(result);
1489
2
        }
1490
3
        return Status::OK();
1491
3
    }
1492
1493
107
    Status _open_mapping_exprs() {
1494
107
        RowDescriptor row_desc;
1495
140
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
1496
140
            if (mapping.projection != nullptr) {
1497
120
                RETURN_IF_ERROR(mapping.projection->prepare(_runtime_state, row_desc));
1498
120
                RETURN_IF_ERROR(mapping.projection->open(_runtime_state));
1499
120
            }
1500
140
            if (mapping.default_expr != nullptr) {
1501
7
                RETURN_IF_ERROR(mapping.default_expr->prepare(_runtime_state, row_desc));
1502
7
                RETURN_IF_ERROR(mapping.default_expr->open(_runtime_state));
1503
7
            }
1504
140
        }
1505
107
        return Status::OK();
1506
107
    }
1507
1508
    Status _build_file_aggregate_request(TPushAggOp::type agg_type,
1509
12
                                         FileAggregateRequest* request) const {
1510
12
        DORIS_CHECK(request != nullptr);
1511
12
        DORIS_CHECK(_supports_aggregate_pushdown(agg_type));
1512
12
        request->agg_type = agg_type;
1513
12
        request->columns.clear();
1514
12
        if (agg_type == TPushAggOp::type::COUNT) {
1515
8
            DORIS_CHECK(_push_down_count_columns.has_value());
1516
            // An empty explicit list is the semantic signal for COUNT(*). Do not inspect the
1517
            // mapping count: `SELECT COUNT(*) FROM t` may still project one nullable column because
1518
            // the planner keeps a placeholder slot. In a 10,000-row file where that arbitrary slot
1519
            // has 9,015 non-null values, passing the slot would ask Parquet/ORC metadata for
1520
            // COUNT(slot)=9,015 instead of the required row count 10,000.
1521
8
            if (!_push_down_count_columns->empty()) {
1522
3
                const auto& mapping = _push_down_count_mapping();
1523
3
                DORIS_CHECK(mapping.file_local_id.has_value());
1524
3
                FileAggregateRequest::Column column;
1525
3
                column.projection =
1526
3
                        LocalColumnIndex::top_level(LocalColumnId(*mapping.file_local_id));
1527
3
                request->columns.push_back(std::move(column));
1528
3
            }
1529
8
            return Status::OK();
1530
8
        }
1531
4
        request->columns.reserve(_data_reader.column_mapper->mappings().size());
1532
5
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
1533
5
            DORIS_CHECK(mapping.file_local_id.has_value());
1534
5
            FileAggregateRequest::Column column;
1535
5
            column.projection = LocalColumnIndex::top_level(LocalColumnId(*mapping.file_local_id));
1536
5
            if (!mapping.child_mappings.empty()) {
1537
1
                RETURN_IF_ERROR(build_aggregate_projection(mapping, &column.projection));
1538
1
            }
1539
5
            request->columns.push_back(std::move(column));
1540
5
        }
1541
4
        return Status::OK();
1542
4
    }
1543
1544
11
    const ColumnMapping& _push_down_count_mapping() const {
1545
11
        DORIS_CHECK(_push_down_count_columns.has_value());
1546
11
        DORIS_CHECK(_push_down_count_columns->size() == 1);
1547
11
        const auto mapping_it =
1548
11
                std::ranges::find(_data_reader.column_mapper->mappings(),
1549
11
                                  _push_down_count_columns->front(), &ColumnMapping::global_index);
1550
        // FileScannerV2 translates FE SlotIds through the same projected-column list used to build
1551
        // the mapper, so a missing mapping is an FE/BE contract violation rather than a fallback.
1552
11
        DORIS_CHECK(mapping_it != _data_reader.column_mapper->mappings().end());
1553
11
        return *mapping_it;
1554
11
    }
1555
1556
    Status _materialize_aggregate_pushdown_rows(TPushAggOp::type agg_type,
1557
                                                const FileAggregateResult& file_result,
1558
10
                                                Block* block) {
1559
10
        if (agg_type == TPushAggOp::type::COUNT) {
1560
            // COUNT pushdown is not a final count value. It emits `count` default rows so the
1561
            // upper COUNT(*) aggregate can count them and produce the final result, including
1562
            // zero rows when count is 0.
1563
7
            DORIS_CHECK(file_result.count >= 0);
1564
7
            return _materialize_count_rows(cast_set<size_t>(file_result.count), block);
1565
7
        }
1566
        // MIN/MAX pushdown emits two rows, min first and max second, for each projected column.
1567
        // The upper MIN/MAX aggregate consumes those two rows to produce the final aggregate value.
1568
3
        DORIS_CHECK(file_result.columns.size() == _data_reader.column_mapper->mappings().size());
1569
3
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
1570
3
        Block file_block;
1571
3
        file_block.reserve(_data_reader.file_block_layout.size());
1572
4
        for (const auto& column : _data_reader.file_block_layout) {
1573
4
            file_block.insert({column.type->create_column(), column.type, column.name});
1574
4
        }
1575
7
        for (size_t column_idx = 0; column_idx < file_result.columns.size(); ++column_idx) {
1576
4
            const auto& result_column = file_result.columns[column_idx];
1577
4
            if (!result_column.has_min || !result_column.has_max) {
1578
0
                return Status::NotSupported("Missing min/max aggregate result for column {}",
1579
0
                                            _projected_columns[column_idx].name);
1580
0
            }
1581
4
            bool found_file_column = false;
1582
5
            for (size_t block_position = 0; block_position < _data_reader.file_block_layout.size();
1583
5
                 ++block_position) {
1584
5
                if (_data_reader.file_block_layout[block_position].file_column_id ==
1585
5
                    file_result.columns[column_idx].projection.column_id()) {
1586
4
                    found_file_column = true;
1587
4
                    auto column = file_block.get_by_position(block_position)
1588
4
                                          .type->create_column()
1589
4
                                          ->assert_mutable();
1590
4
                    RETURN_IF_ERROR(_insert_aggregate_projection_value(
1591
4
                            file_result.columns[column_idx].projection, result_column.min_value,
1592
4
                            column.get()));
1593
4
                    RETURN_IF_ERROR(_insert_aggregate_projection_value(
1594
4
                            file_result.columns[column_idx].projection, result_column.max_value,
1595
4
                            column.get()));
1596
4
                    file_block.replace_by_position(block_position, std::move(column));
1597
4
                    break;
1598
4
                }
1599
5
            }
1600
4
            DORIS_CHECK(found_file_column);
1601
4
        }
1602
7
        for (size_t column_idx = 0; column_idx < _data_reader.column_mapper->mappings().size();
1603
4
             ++column_idx) {
1604
4
            ColumnPtr table_column;
1605
4
            RETURN_IF_ERROR(
1606
4
                    _materialize_mapping_column(_data_reader.column_mapper->mappings()[column_idx],
1607
4
                                                &file_block, 2, &table_column));
1608
4
            block->replace_by_position(column_idx, std::move(table_column));
1609
4
        }
1610
3
        return Status::OK();
1611
3
    }
1612
1613
    struct FileBlockColumn {
1614
        LocalColumnId file_column_id = LocalColumnId::invalid();
1615
        std::string name;
1616
        DataTypePtr type;
1617
    };
1618
1619
    struct DataReader {
1620
        std::unique_ptr<FileReader> reader;
1621
        std::unique_ptr<TableColumnMapper> column_mapper;
1622
        // Schema of the data file, also including virtual column (row position).
1623
        std::vector<ColumnDefinition> file_schema;
1624
        // Layout of the block returned by file reader, determined by column mapping and file
1625
        // schema. It is used for file reader to materialize columns into correct type and position.
1626
        std::vector<FileBlockColumn> file_block_layout;
1627
        Block block_template;
1628
    };
1629
    DataReader _data_reader;
1630
    std::vector<ColumnDefinition> _projected_columns;
1631
    std::unique_ptr<ScanTask> _current_task;
1632
    std::optional<io::FileDescription> _current_file_description;
1633
    // Range-level compression has higher priority than scan-param compression. TVF/load can keep
1634
    // the logical format as CSV/TEXT while carrying the concrete compression such as GZ or LZO on
1635
    // each TFileRangeDesc, matching the old FileScanner reader contract.
1636
    TFileCompressType::type _current_range_compress_type = TFileCompressType::UNKNOWN;
1637
    std::optional<TUniqueId> _current_range_load_id;
1638
    TFileRangeDesc _current_file_range_desc;
1639
    std::shared_ptr<io::FileSystemProperties> _system_properties;
1640
    // partition key -> value
1641
    std::map<std::string, Field> _partition_values;
1642
    // Predicates built from scan conjuncts before file-level localization.
1643
    std::vector<TableFilter> _table_filters;
1644
    // Number of localized filters before the first unsafe conjunct in the original row-level
1645
    // order. This differs from scanning `_table_filters` for safety because slotless predicates are
1646
    // intentionally absent from that vector but must still act as ordering barriers.
1647
    size_t _constant_pruning_safe_filter_count = 0;
1648
    VExprContextSPtrs _conjuncts;
1649
    ReadProfile _profile;
1650
    // Parsed from row-position based delete files, including position delete and deletion vector.
1651
    DeleteRows* _delete_rows = nullptr;
1652
    DeletionVector* _deletion_vector = nullptr;
1653
    TFileScanRangeParams* _scan_params;
1654
    std::shared_ptr<io::IOContext> _io_ctx;
1655
    RuntimeState* _runtime_state;
1656
    RuntimeProfile* _scanner_profile;
1657
    const std::vector<SlotDescriptor*>* _file_slot_descs = nullptr;
1658
    FileFormat _format;
1659
    TPushAggOp::type _push_down_agg_type = TPushAggOp::type::NONE;
1660
    std::optional<std::vector<GlobalIndex>> _push_down_count_columns;
1661
    size_t _batch_size = 0;
1662
    uint64_t _initial_condition_cache_digest = 0;
1663
    uint64_t _condition_cache_digest = 0;
1664
    // True only when prepare_split() received a digest for the exact conjunct snapshot used by
1665
    // this split. Standalone callers that only supplied TableReadOptions::condition_cache_digest
1666
    // keep the conservative runtime-filter guard.
1667
    bool _condition_cache_digest_covers_current_split = false;
1668
    segment_v2::ConditionCache::ExternalCacheKey _condition_cache_key;
1669
    std::shared_ptr<std::vector<bool>> _condition_cache;
1670
    std::shared_ptr<ConditionCacheContext> _condition_cache_ctx;
1671
    int64_t _condition_cache_hit_count = 0;
1672
    bool _current_reader_reached_eof = false;
1673
    int64_t _remaining_table_level_count = -1;
1674
    // True only after the active split selects a table-level row-count shortcut or successfully
1675
    // materializes COUNT rows from file metadata. FileScannerV2 uses this result, rather than the
1676
    // raw aggregate opcode, to keep adaptive batching enabled for normal row-scan fallbacks.
1677
    bool _current_split_uses_metadata_count = false;
1678
    // Snapshot supplied by FileScannerV2 for the active split. It gates every shortcut that emits
1679
    // irreversible aggregate rows, not only the table-level row-count shortcut in prepare_split().
1680
    bool _all_runtime_filters_applied_for_split = true;
1681
    std::optional<GlobalRowIdContext> _global_rowid_context;
1682
    bool _aggregate_pushdown_tried = false;
1683
    bool _current_split_pruned = false;
1684
    TableColumnMapperOptions _mapper_options;
1685
1686
private:
1687
    static const ColumnDefinition* _find_column_definition(
1688
183
            const std::vector<ColumnDefinition>& schema, LocalColumnId column_id) {
1689
321
        for (const auto& field : schema) {
1690
321
            if (field.file_local_id() == column_id.value()) {
1691
153
                return &field;
1692
153
            }
1693
321
        }
1694
30
        return nullptr;
1695
183
    }
1696
1697
14
    static bool _can_push_down_minmax_for_mapping(const ColumnMapping& mapping) {
1698
14
        if (mapping.child_mappings.empty()) {
1699
            // Direct mappings use a slot-ref projection to materialize the file column. The
1700
            // projection does not transform ordering; casts and other conversions are already
1701
            // represented by a non-trivial mapping and must fall back to row scanning.
1702
11
            return mapping.is_trivial;
1703
11
        }
1704
3
        const auto primitive_type = remove_nullable(mapping.file_type)->get_primitive_type();
1705
3
        if (primitive_type != TYPE_STRUCT) {
1706
1
            return false;
1707
1
        }
1708
2
        size_t mapped_children = 0;
1709
2
        const ColumnMapping* mapped_child = nullptr;
1710
2
        for (const auto& child_mapping : mapping.child_mappings) {
1711
2
            if (!child_mapping.file_local_id.has_value()) {
1712
0
                continue;
1713
0
            }
1714
2
            ++mapped_children;
1715
2
            mapped_child = &child_mapping;
1716
2
        }
1717
2
        return mapped_children == 1 && mapped_child != nullptr &&
1718
2
               _can_push_down_minmax_for_mapping(*mapped_child);
1719
3
    }
1720
1721
    static Status build_aggregate_projection(const ColumnMapping& mapping,
1722
2
                                             LocalColumnIndex* projection) {
1723
2
        DORIS_CHECK(projection != nullptr);
1724
2
        DORIS_CHECK(mapping.file_local_id.has_value());
1725
2
        *projection = LocalColumnIndex::local(*mapping.file_local_id);
1726
2
        projection->children.clear();
1727
2
        projection->project_all_children = true;
1728
2
        if (mapping.child_mappings.empty()) {
1729
1
            return Status::OK();
1730
1
        }
1731
1
        projection->project_all_children = false;
1732
1
        for (const auto& child_mapping : mapping.child_mappings) {
1733
1
            if (!child_mapping.file_local_id.has_value()) {
1734
0
                continue;
1735
0
            }
1736
1
            LocalColumnIndex child_projection;
1737
1
            RETURN_IF_ERROR(build_aggregate_projection(child_mapping, &child_projection));
1738
1
            projection->children.push_back(std::move(child_projection));
1739
1
        }
1740
1
        DORIS_CHECK(projection->children.size() == 1);
1741
1
        return Status::OK();
1742
1
    }
1743
1744
    static Status _insert_aggregate_projection_value(const LocalColumnIndex& projection,
1745
20
                                                     const Field& value, IColumn* column) {
1746
20
        DORIS_CHECK(column != nullptr);
1747
20
        if (auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1748
10
            RETURN_IF_ERROR(_insert_aggregate_projection_value(
1749
10
                    projection, value, &nullable_column->get_nested_column()));
1750
10
            nullable_column->get_null_map_data().push_back(0);
1751
10
            return Status::OK();
1752
10
        }
1753
10
        if (projection.project_all_children || projection.children.empty()) {
1754
8
            column->insert(value);
1755
8
            return Status::OK();
1756
8
        }
1757
2
        auto* struct_column = assert_cast<ColumnStruct*>(column);
1758
2
        DORIS_CHECK(projection.children.size() == 1);
1759
2
        const auto& child_projection = projection.children[0];
1760
2
        DORIS_CHECK(struct_column->get_columns().size() == 1);
1761
2
        RETURN_IF_ERROR(_insert_aggregate_projection_value(child_projection, value,
1762
2
                                                           &struct_column->get_column(0)));
1763
2
        return Status::OK();
1764
2
    }
1765
1766
    // Parse a DV into its compressed bitmap. Position delete files continue to use _delete_rows.
1767
    Status _parse_delete_predicates(const SplitReadOptions& options);
1768
};
1769
1770
} // namespace doris::format