Coverage Report

Created: 2026-07-21 06:57

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
60.2k
    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
40.9k
    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
240k
    virtual void set_batch_size(size_t batch_size) {
190
240k
        _batch_size = std::max<size_t>(1, batch_size);
191
240k
        if (_data_reader.reader != nullptr) {
192
119k
            _data_reader.reader->set_batch_size(_batch_size);
193
119k
        }
194
240k
    }
195
196
#ifdef BE_TEST
197
    size_t TEST_batch_size() const { return _batch_size; }
198
    bool TEST_current_data_file_is_immutable() const {
199
        DORIS_CHECK(_current_task != nullptr);
200
        DORIS_CHECK(_current_task->data_file != nullptr);
201
        DORIS_CHECK(_current_file_description.has_value());
202
        DORIS_CHECK(_current_task->data_file->is_immutable ==
203
                    _current_file_description->is_immutable);
204
        return _current_task->data_file->is_immutable;
205
    }
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
89.3k
    virtual bool current_split_pruned() const { return _current_split_pruned; }
214
303k
    virtual bool current_split_uses_metadata_count() const {
215
303k
        return _current_split_uses_metadata_count;
216
303k
    }
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
181k
    virtual Status get_block(Block* block, bool* eos) {
239
181k
        SCOPED_TIMER(_profile.exec_timer);
240
181k
        DORIS_CHECK(block->columns() == _projected_columns.size());
241
181k
        block->clear_column_data(_projected_columns.size());
242
243
249k
        while (true) {
244
249k
            if (*eos) {
245
0
                return Status::OK();
246
0
            }
247
249k
            if (_io_ctx != nullptr && _io_ctx->should_stop) {
248
10
                *eos = true;
249
10
                return Status::OK();
250
10
            }
251
249k
            if (!_data_reader.reader) {
252
118k
                if (_is_table_level_count_active()) {
253
302
                    RETURN_IF_ERROR(_read_table_level_count(block, eos));
254
302
                    return Status::OK();
255
302
                }
256
118k
                RETURN_IF_ERROR(create_next_reader(eos));
257
118k
                if (!_data_reader.reader) {
258
59.4k
                    DCHECK(*eos);
259
59.4k
                    return Status::OK();
260
59.4k
                }
261
118k
            }
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
190k
            if (!_aggregate_pushdown_tried) {
268
58.7k
                SCOPED_TIMER(_profile.pushdown_agg_timer);
269
58.7k
                bool pushed_down = false;
270
58.7k
                const auto status = _try_materialize_aggregate_pushdown_rows(block, &pushed_down);
271
58.7k
                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
58.7k
                if (pushed_down) {
280
1.27k
                    return Status::OK();
281
1.27k
                }
282
58.7k
            }
283
284
188k
            bool current_eof = false;
285
188k
            _data_reader.block_template.clear_column_data(
286
188k
                    cast_set<int64_t>(_data_reader.file_block_layout.size()));
287
188k
            size_t current_rows = 0;
288
188k
            RETURN_IF_ERROR(_data_reader.reader->get_block(&_data_reader.block_template,
289
188k
                                                           &current_rows, &current_eof));
290
188k
            const bool stopped_during_read = _io_ctx != nullptr && _io_ctx->should_stop;
291
188k
            if (current_rows == 0) {
292
68.2k
                if (current_eof) {
293
56.9k
                    _current_reader_reached_eof = !stopped_during_read;
294
56.9k
                    RETURN_IF_ERROR(close_current_reader());
295
56.9k
                }
296
68.2k
                continue;
297
68.2k
            }
298
188k
            DCHECK_EQ(_data_reader.block_template.columns(), _data_reader.file_block_layout.size())
299
0
                    << _data_reader.block_template.dump_structure();
300
120k
#ifndef NDEBUG
301
120k
            RETURN_IF_ERROR(_check_file_block_columns("after file reader get_block", current_rows));
302
120k
#endif
303
120k
            DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
304
120k
            RETURN_IF_ERROR(finalize_chunk(block, current_rows));
305
120k
#ifndef NDEBUG
306
120k
            RETURN_IF_ERROR(
307
120k
                    _check_table_block_columns("after finalize_chunk", block, current_rows));
308
120k
#endif
309
120k
            if (current_eof) {
310
18
                _current_reader_reached_eof = !stopped_during_read;
311
18
                RETURN_IF_ERROR(close_current_reader());
312
18
            }
313
120k
            return Status::OK();
314
120k
        }
315
181k
    }
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
40.8k
    virtual Status close() {
320
40.8k
        if (_data_reader.reader) {
321
634
            RETURN_IF_ERROR(close_current_reader());
322
634
        }
323
40.8k
        _current_task.reset();
324
40.8k
        _current_file_description.reset();
325
40.8k
        _remaining_table_level_count = -1;
326
40.8k
        _current_split_uses_metadata_count = false;
327
40.8k
        return Status::OK();
328
40.8k
    }
329
330
81.4k
    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
22.0k
    virtual Status validate_projected_columns(const ProjectedColumnBuildContext& context) const {
339
22.0k
        (void)context;
340
22.0k
        return Status::OK();
341
22.0k
    }
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
28.7k
    void mark_current_data_file_immutable() {
350
28.7k
        DORIS_CHECK(_current_task != nullptr);
351
28.7k
        DORIS_CHECK(_current_task->data_file != nullptr);
352
28.7k
        DORIS_CHECK(_current_file_description.has_value());
353
28.7k
        _current_task->data_file->is_immutable = true;
354
28.7k
        _current_file_description->is_immutable = true;
355
28.7k
    }
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
31.4k
                                               DeleteFileDesc* desc, bool* has_delete_file) {
363
31.4k
        *has_delete_file = false;
364
31.4k
        return Status::OK();
365
31.4k
    }
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
6.38k
    virtual TableColumnMappingMode mapping_mode() const { return TableColumnMappingMode::BY_NAME; }
372
34.8k
    virtual Status annotate_file_schema(std::vector<ColumnDefinition>* file_schema) {
373
34.8k
        DORIS_CHECK(file_schema != nullptr);
374
34.8k
        return Status::OK();
375
34.8k
    }
376
377
    // Open the concrete reader for the current split/task and build the file-local scan request.
378
60.0k
    virtual Status open_reader() {
379
60.0k
        SCOPED_TIMER(_profile.open_reader_timer);
380
        // 1. Get file schema and create column mapping.
381
60.0k
        std::vector<ColumnDefinition> file_schema;
382
60.0k
        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
60.0k
        RETURN_IF_ERROR(annotate_file_schema(&file_schema));
386
60.0k
        _data_reader.file_schema = file_schema;
387
60.0k
        _mapper_options.mode = mapping_mode();
388
389
60.0k
        _data_reader.column_mapper = _data_reader.reader->create_column_mapper(_mapper_options);
390
60.0k
        DORIS_CHECK(_data_reader.column_mapper != nullptr);
391
60.0k
        RETURN_IF_ERROR(_data_reader.column_mapper->create_mapping(_projected_columns,
392
60.0k
                                                                   _partition_values, file_schema));
393
60.0k
        DORIS_CHECK(_data_reader.column_mapper->mappings().size() == _projected_columns.size());
394
395
        // 2. Build table filters based on conjuncts and column predicates.
396
60.0k
        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
60.0k
        auto file_request = std::make_shared<FileScanRequest>();
402
60.0k
        RETURN_IF_ERROR(_data_reader.column_mapper->create_scan_request(
403
60.0k
                _table_filters, _projected_columns, file_request.get(), _runtime_state));
404
60.0k
        bool constant_filter_pruned_split = false;
405
60.0k
        RETURN_IF_ERROR(_evaluate_constant_filters(&constant_filter_pruned_split));
406
60.0k
        if (constant_filter_pruned_split) {
407
1.26k
            RETURN_IF_ERROR(close_current_reader());
408
1.26k
            return Status::OK();
409
1.26k
        }
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
58.7k
        if (_push_down_agg_type == TPushAggOp::type::COUNT &&
417
58.7k
            _push_down_count_columns.has_value() && _push_down_count_columns->empty()) {
418
1.95k
            file_request->count_star_placeholder_columns.reserve(
419
1.95k
                    file_request->non_predicate_columns.size());
420
1.95k
            for (const auto& column : file_request->non_predicate_columns) {
421
1.94k
                file_request->count_star_placeholder_columns.push_back(column.column_id());
422
1.94k
            }
423
1.95k
        }
424
58.7k
        RETURN_IF_ERROR(customize_file_scan_request(file_request.get()));
425
58.7k
        RETURN_IF_ERROR(_open_local_filter_exprs(*file_request));
426
58.7k
        _data_reader.file_block_layout.clear();
427
58.7k
        _data_reader.block_template.clear();
428
58.7k
        _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
381k
        for (const auto& [file_column_id, block_position] : file_request->local_positions) {
433
381k
            DORIS_CHECK(block_position.value() < _data_reader.file_block_layout.size());
434
381k
            const auto* field = _find_column_definition(_data_reader.file_schema, file_column_id);
435
381k
            DORIS_CHECK(field != nullptr);
436
437
381k
            ColumnDefinition projected_field;
438
381k
            {
439
381k
                auto it = std::find_if(
440
381k
                        file_request->non_predicate_columns.begin(),
441
381k
                        file_request->non_predicate_columns.end(),
442
7.90M
                        [&](const LocalColumnIndex& p) { return p.column_id() == file_column_id; });
443
381k
                if (it != file_request->non_predicate_columns.end()) {
444
348k
                    RETURN_IF_ERROR(project_column_definition(*field, *it, &projected_field));
445
348k
                }
446
381k
            }
447
381k
            {
448
381k
                auto it = std::find_if(
449
381k
                        file_request->predicate_columns.begin(),
450
381k
                        file_request->predicate_columns.end(),
451
381k
                        [&](const LocalColumnIndex& p) { return p.column_id() == file_column_id; });
452
381k
                if (it != file_request->predicate_columns.end()) {
453
32.5k
                    RETURN_IF_ERROR(project_column_definition(*field, *it, &projected_field));
454
32.5k
                }
455
381k
            }
456
381k
            _data_reader.file_block_layout[block_position.value()] = {
457
381k
                    .file_column_id = file_column_id,
458
381k
                    .name = projected_field.name,
459
381k
                    .type = projected_field.type,
460
381k
            };
461
381k
            DORIS_CHECK(_data_reader.file_block_layout[block_position.value()].type != nullptr);
462
381k
        }
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
58.7k
        _data_reader.block_template.reserve(_data_reader.file_block_layout.size());
467
381k
        for (const auto& column : _data_reader.file_block_layout) {
468
381k
            _data_reader.block_template.insert(
469
381k
                    {column.type->create_column(), column.type, column.name});
470
381k
        }
471
58.7k
        if (VLOG_DEBUG_IS_ON) {
472
0
            VLOG_DEBUG << "TableReader debug: " << debug_string();
473
0
        }
474
58.7k
        RETURN_IF_ERROR(_open_mapping_exprs());
475
58.7k
        RETURN_IF_ERROR(_data_reader.reader->open(file_request));
476
58.7k
        RETURN_IF_ERROR(_init_reader_condition_cache(*file_request));
477
58.7k
        return Status::OK();
478
58.7k
    }
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
59.9k
    Status _evaluate_constant_filters(bool* can_filter_all) {
491
59.9k
        DORIS_CHECK(can_filter_all != nullptr);
492
59.9k
        DORIS_CHECK_LE(_constant_pruning_safe_filter_count, _table_filters.size());
493
59.9k
        *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
96.6k
        for (size_t i = 0; i < _constant_pruning_safe_filter_count; ++i) {
498
37.9k
            const auto& table_filter = _table_filters[i];
499
37.9k
            if (table_filter.conjunct == nullptr) {
500
0
                continue;
501
0
            }
502
37.9k
            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
37.9k
            if (table_filter.conjunct->root()->is_rf_wrapper() ||
508
37.9k
                !_table_filter_has_only_constant_entries(table_filter)) {
509
34.5k
                continue;
510
34.5k
            }
511
3.39k
            Block eval_block;
512
3.39k
            RETURN_IF_ERROR(_build_constant_filter_block(table_filter, &eval_block));
513
3.39k
            RowDescriptor row_desc;
514
3.39k
            RETURN_IF_ERROR(table_filter.conjunct->prepare(_runtime_state, row_desc));
515
3.39k
            RETURN_IF_ERROR(table_filter.conjunct->open(_runtime_state));
516
3.39k
            int result_column_id = -1;
517
3.39k
            RETURN_IF_ERROR(table_filter.conjunct->execute(&eval_block, &result_column_id));
518
3.39k
            DORIS_CHECK(result_column_id >= 0);
519
3.39k
            if (_filter_result_filters_all(eval_block.get_by_position(result_column_id).column)) {
520
1.26k
                *can_filter_all = true;
521
1.26k
                return Status::OK();
522
1.26k
            }
523
3.39k
        }
524
58.6k
        return Status::OK();
525
59.9k
    }
526
527
32.7k
    bool _table_filter_has_only_constant_entries(const TableFilter& table_filter) const {
528
32.7k
        const auto& filter_entries = _data_reader.column_mapper->filter_entries();
529
32.9k
        for (const auto global_index : table_filter.global_indices) {
530
32.9k
            const auto entry_it = filter_entries.find(global_index);
531
32.9k
            if (entry_it == filter_entries.end() || !entry_it->second.is_constant()) {
532
29.2k
                return false;
533
29.2k
            }
534
32.9k
        }
535
3.43k
        return !table_filter.global_indices.empty();
536
32.7k
    }
537
538
3.43k
    Status _build_constant_filter_block(const TableFilter& table_filter, Block* eval_block) {
539
3.43k
        DORIS_CHECK(eval_block != nullptr);
540
3.43k
        eval_block->clear();
541
3.43k
        const auto& mappings = _data_reader.column_mapper->mappings();
542
3.43k
        const auto& filter_entries = _data_reader.column_mapper->filter_entries();
543
3.43k
        DORIS_CHECK(mappings.size() == _projected_columns.size());
544
14.7k
        for (size_t column_idx = 0; column_idx < mappings.size(); ++column_idx) {
545
11.3k
            const auto global_index = GlobalIndex(column_idx);
546
11.3k
            const auto& mapping = mappings[column_idx];
547
11.3k
            const auto entry_it = filter_entries.find(global_index);
548
11.3k
            const bool referenced_by_filter =
549
11.3k
                    std::find(table_filter.global_indices.begin(),
550
11.3k
                              table_filter.global_indices.end(),
551
11.3k
                              global_index) != table_filter.global_indices.end();
552
11.3k
            if (referenced_by_filter && entry_it != filter_entries.end() &&
553
11.3k
                entry_it->second.is_constant()) {
554
3.67k
                ColumnPtr constant_column;
555
3.67k
                RETURN_IF_ERROR(_materialize_constant_filter_column(
556
3.67k
                        entry_it->second.constant_index(), &constant_column));
557
3.67k
                eval_block->insert({std::move(constant_column), mapping.table_type,
558
3.67k
                                    mapping.table_column_name});
559
7.62k
            } else {
560
7.62k
                eval_block->insert({mapping.table_type->create_column_const_with_default_value(1),
561
7.62k
                                    mapping.table_type, mapping.table_column_name});
562
7.62k
            }
563
11.3k
        }
564
3.43k
        return Status::OK();
565
3.43k
    }
566
567
3.68k
    Status _materialize_constant_filter_column(ConstantIndex constant_index, ColumnPtr* column) {
568
3.68k
        DORIS_CHECK(column != nullptr);
569
3.68k
        const auto& constant_entry = _data_reader.column_mapper->constant_map().get(constant_index);
570
3.68k
        DORIS_CHECK(constant_entry.expr != nullptr);
571
3.68k
        DORIS_CHECK(constant_entry.type != nullptr);
572
3.68k
        RowDescriptor row_desc;
573
3.68k
        RETURN_IF_ERROR(constant_entry.expr->prepare(_runtime_state, row_desc));
574
3.68k
        RETURN_IF_ERROR(constant_entry.expr->open(_runtime_state));
575
3.68k
        Block eval_block;
576
3.68k
        eval_block.insert({constant_entry.type->create_column_const_with_default_value(1),
577
3.68k
                           constant_entry.type, "__table_reader_constant_filter"});
578
3.68k
        int result_column_id = -1;
579
3.68k
        RETURN_IF_ERROR(constant_entry.expr->execute(&eval_block, &result_column_id));
580
3.68k
        DORIS_CHECK(result_column_id >= 0);
581
3.68k
        *column = eval_block.get_by_position(result_column_id).column;
582
3.68k
        DORIS_CHECK((*column)->size() == 1);
583
3.68k
        return Status::OK();
584
3.68k
    }
585
586
3.44k
    static bool _filter_result_filters_all(const ColumnPtr& filter_column) {
587
3.44k
        DORIS_CHECK(filter_column.get() != nullptr);
588
3.44k
        DORIS_CHECK(filter_column->size() == 1);
589
3.44k
        return !filter_column->get_bool(0);
590
3.44k
    }
591
592
58.6k
    virtual Status customize_file_scan_request(FileScanRequest* file_request) {
593
58.6k
        return _append_delete_predicate(file_request);
594
58.6k
    }
595
596
209k
    bool _is_table_level_count_active() const { return _remaining_table_level_count >= 0; }
597
598
1.40k
    Status _materialize_count_rows(size_t rows, Block* block) const {
599
1.40k
        DORIS_CHECK(block != nullptr);
600
1.40k
        DORIS_CHECK(block->columns() > 0 || rows == 0);
601
2.81k
        for (size_t column_idx = 0; column_idx < block->columns(); ++column_idx) {
602
1.40k
            auto column = block->get_by_position(column_idx).type->create_column();
603
1.40k
            column->resize(rows);
604
1.40k
            block->replace_by_position(column_idx, std::move(column));
605
1.40k
        }
606
1.40k
        return Status::OK();
607
1.40k
    }
608
609
302
    Status _read_table_level_count(Block* block, bool* eos) {
610
302
        DORIS_CHECK(block != nullptr);
611
302
        DORIS_CHECK(eos != nullptr);
612
302
        DORIS_CHECK(_push_down_agg_type == TPushAggOp::type::COUNT);
613
302
        DORIS_CHECK(_remaining_table_level_count >= 0);
614
302
        if (_remaining_table_level_count == 0) {
615
150
            _remaining_table_level_count = -1;
616
150
            _current_task.reset();
617
150
            *eos = true;
618
150
            return Status::OK();
619
150
        }
620
621
152
        const int64_t batch_size = _runtime_state == nullptr
622
152
                                           ? _remaining_table_level_count
623
152
                                           : static_cast<int64_t>(_runtime_state->batch_size());
624
152
        const auto rows = std::min(_remaining_table_level_count, batch_size);
625
152
        RETURN_IF_ERROR(_materialize_count_rows(cast_set<size_t>(rows), block));
626
152
        _remaining_table_level_count -= rows;
627
152
        *eos = false;
628
152
        return Status::OK();
629
152
    }
630
631
    void _append_file_scan_column(FileScanRequest* request, LocalColumnId column_id,
632
12.8k
                                  std::vector<LocalColumnIndex>* scan_columns) {
633
12.8k
        DORIS_CHECK(request != nullptr);
634
12.8k
        DORIS_CHECK(scan_columns != nullptr);
635
12.8k
        FileScanRequestBuilder builder(request);
636
12.8k
        Status status;
637
12.8k
        if (scan_columns == &request->predicate_columns) {
638
10.6k
            status = builder.add_predicate_column(column_id);
639
10.6k
        } else {
640
2.11k
            DORIS_CHECK(scan_columns == &request->non_predicate_columns);
641
2.11k
            status = builder.add_non_predicate_column(column_id);
642
2.11k
        }
643
12.8k
        DORIS_CHECK(status.ok()) << status.to_string();
644
12.8k
        if (column_id == LocalColumnId(ROW_POSITION_COLUMN_ID) &&
645
12.8k
            _find_column_definition(_data_reader.file_schema, column_id) == nullptr) {
646
7.46k
            _data_reader.file_schema.push_back(row_position_column_definition());
647
7.46k
        }
648
12.8k
    }
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
58.5k
    Status _append_delete_predicate(FileScanRequest* request) {
652
58.5k
        DORIS_CHECK(request != nullptr);
653
58.5k
        if ((_delete_rows == nullptr || _delete_rows->empty()) &&
654
58.5k
            (_deletion_vector == nullptr || _deletion_vector->isEmpty())) {
655
52.7k
            return Status::OK();
656
52.7k
        }
657
5.78k
        const auto row_position_column_id = LocalColumnId(ROW_POSITION_COLUMN_ID);
658
5.78k
        _append_file_scan_column(request, row_position_column_id, &request->predicate_columns);
659
660
5.78k
        const auto block_position = request->local_positions.at(row_position_column_id);
661
5.82k
        auto append_predicate = [&](auto& deleted_rows) {
662
5.82k
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
663
5.82k
            delete_predicate->add_child(VSlotRef::create_shared(
664
5.82k
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
665
5.82k
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
666
5.82k
            request->delete_conjuncts.push_back(
667
5.82k
                    VExprContext::create_shared(std::move(delete_predicate)));
668
5.82k
        };
_ZZN5doris6format11TableReader24_append_delete_predicateEPNS0_15FileScanRequestEENKUlRT_E_clISt6vectorIlSaIlEEEEDaS5_
Line
Count
Source
661
1.65k
        auto append_predicate = [&](auto& deleted_rows) {
662
1.65k
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
663
1.65k
            delete_predicate->add_child(VSlotRef::create_shared(
664
1.65k
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
665
1.65k
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
666
1.65k
            request->delete_conjuncts.push_back(
667
1.65k
                    VExprContext::create_shared(std::move(delete_predicate)));
668
1.65k
        };
_ZZN5doris6format11TableReader24_append_delete_predicateEPNS0_15FileScanRequestEENKUlRT_E_clIN7roaring12Roaring64MapEEEDaS5_
Line
Count
Source
661
4.17k
        auto append_predicate = [&](auto& deleted_rows) {
662
4.17k
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
663
4.17k
            delete_predicate->add_child(VSlotRef::create_shared(
664
4.17k
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
665
4.17k
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
666
4.17k
            request->delete_conjuncts.push_back(
667
4.17k
                    VExprContext::create_shared(std::move(delete_predicate)));
668
4.17k
        };
669
5.78k
        if (_delete_rows != nullptr && !_delete_rows->empty()) {
670
1.65k
            append_predicate(*_delete_rows);
671
1.65k
        }
672
5.78k
        if (_deletion_vector != nullptr && !_deletion_vector->isEmpty()) {
673
4.18k
            append_predicate(*_deletion_vector);
674
4.18k
        }
675
5.78k
        return Status::OK();
676
58.5k
    }
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
60.0k
    virtual Status close_current_reader() {
681
60.0k
        _finalize_reader_condition_cache();
682
60.0k
        RETURN_IF_ERROR(_data_reader.reader->close());
683
60.0k
        _data_reader.reader.reset();
684
60.0k
        if (_data_reader.column_mapper != nullptr) {
685
60.0k
            _data_reader.column_mapper->clear();
686
60.0k
            _data_reader.column_mapper.reset();
687
60.0k
        }
688
60.0k
        _table_filters.clear();
689
60.0k
        _constant_pruning_safe_filter_count = 0;
690
60.0k
        _data_reader.file_schema.clear();
691
60.0k
        _data_reader.file_block_layout.clear();
692
60.0k
        _data_reader.block_template.clear();
693
60.0k
        _current_task.reset();
694
60.0k
        _current_file_description.reset();
695
60.0k
        _current_reader_reached_eof = false;
696
60.0k
        return Status::OK();
697
60.0k
    }
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
120k
    Status finalize_chunk(Block* block, const size_t rows) {
707
120k
        SCOPED_TIMER(_profile.finalize_timer);
708
120k
        size_t idx = 0;
709
663k
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
710
663k
            ColumnPtr column;
711
663k
            RETURN_IF_ERROR(_materialize_mapping_column(mapping, &_data_reader.block_template, rows,
712
663k
                                                        &column));
713
663k
            block->replace_by_position(idx, IColumn::mutate(std::move(column)));
714
663k
            idx++;
715
663k
        }
716
120k
        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
120k
        RETURN_IF_ERROR(_truncate_char_or_varchar_columns(block));
720
120k
        return Status::OK();
721
120k
    }
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
93.5k
    virtual Status materialize_virtual_columns(Block* table_block) { return Status::OK(); }
727
728
#ifndef NDEBUG
729
120k
    Status _check_file_block_columns(std::string_view stage, size_t rows) {
730
120k
        DORIS_CHECK(_data_reader.block_template.columns() == _data_reader.file_block_layout.size());
731
772k
        for (size_t idx = 0; idx < _data_reader.block_template.columns(); ++idx) {
732
652k
            const auto& file_block_column = _data_reader.file_block_layout[idx];
733
652k
            const auto& column_with_type = _data_reader.block_template.get_by_position(idx);
734
652k
            const auto* column = column_with_type.column.get();
735
652k
            try {
736
652k
                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
652k
                column->sanity_check();
749
652k
                auto st = column_with_type.check_type_and_column_match();
750
652k
                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
652k
            } 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
652k
        }
793
120k
        return Status::OK();
794
120k
    }
795
796
120k
    Status _check_table_block_columns(std::string_view stage, const Block* block, size_t rows) {
797
120k
        DORIS_CHECK(block != nullptr);
798
120k
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
799
784k
        for (size_t idx = 0; idx < block->columns(); ++idx) {
800
663k
            const auto& mapping = _data_reader.column_mapper->mappings()[idx];
801
663k
            const auto& column_with_type = block->get_by_position(idx);
802
663k
            const auto* column = column_with_type.column.get();
803
663k
            try {
804
663k
                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
663k
                column->sanity_check();
815
663k
                auto st = column_with_type.check_type_and_column_match();
816
663k
                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
663k
            } 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
663k
        }
854
120k
        return Status::OK();
855
120k
    }
856
#endif
857
858
120k
    Status _truncate_char_or_varchar_columns(Block* block) {
859
120k
        DORIS_CHECK(block != nullptr);
860
120k
        if (_runtime_state == nullptr ||
861
120k
            !_runtime_state->query_options().truncate_char_or_varchar_columns) {
862
120k
            return Status::OK();
863
120k
        }
864
14
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
865
50
        for (size_t idx = 0; idx < _data_reader.column_mapper->mappings().size(); ++idx) {
866
36
            const auto& mapping = _data_reader.column_mapper->mappings()[idx];
867
36
            if (!_should_truncate_char_or_varchar_column(mapping)) {
868
12
                continue;
869
12
            }
870
24
            const auto target_len =
871
24
                    assert_cast<const DataTypeString*>(remove_nullable(mapping.table_type).get())
872
24
                            ->len();
873
24
            _truncate_char_or_varchar_column(block, idx, target_len);
874
24
        }
875
14
        return Status::OK();
876
120k
    }
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
41
    static bool _should_truncate_char_or_varchar_column(const ColumnMapping& mapping) {
884
41
        if (mapping.table_type == nullptr) {
885
0
            return false;
886
0
        }
887
41
        const auto table_type = remove_nullable(mapping.table_type);
888
41
        const auto primitive_type = table_type->get_primitive_type();
889
41
        if (primitive_type != TYPE_VARCHAR && primitive_type != TYPE_CHAR) {
890
13
            return false;
891
13
        }
892
28
        const auto target_len = assert_cast<const DataTypeString*>(table_type.get())->len();
893
28
        if (target_len <= 0) {
894
0
            return false;
895
0
        }
896
28
        if (mapping.file_type == nullptr) {
897
0
            return true;
898
0
        }
899
28
        const auto file_type = remove_nullable(mapping.file_type);
900
28
        DORIS_CHECK(file_type != nullptr);
901
28
        int file_len = -1;
902
28
        if (file_type->get_primitive_type() == TYPE_VARCHAR ||
903
28
            file_type->get_primitive_type() == TYPE_CHAR ||
904
28
            file_type->get_primitive_type() == TYPE_STRING) {
905
27
            file_len = assert_cast<const DataTypeString*>(file_type.get())->len();
906
27
        }
907
908
28
        return file_len < 0 || target_len < file_len;
909
28
    }
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
25
    static void _truncate_char_or_varchar_column(Block* block, size_t idx, int len) {
916
25
        DORIS_CHECK(block != nullptr);
917
25
        auto int_type = std::make_shared<DataTypeInt32>();
918
25
        const auto num_columns_without_result = cast_set<uint32_t>(block->columns());
919
25
        auto& target = block->get_by_position(idx);
920
25
        const bool is_nullable = target.type->is_nullable();
921
25
        ColumnPtr input_column = target.column;
922
25
        ColumnPtr null_map_column;
923
25
        if (is_nullable) {
924
25
            const auto* nullable_column = assert_cast<const ColumnNullable*>(target.column.get());
925
25
            input_column = nullable_column->get_nested_column_ptr();
926
25
            null_map_column = nullable_column->get_null_map_column_ptr();
927
25
        }
928
25
        block->replace_by_position(idx, std::move(input_column));
929
25
        block->insert({int_type->create_column_const(block->rows(), to_field<TYPE_INT>(1)),
930
25
                       int_type, "const 1"});
931
25
        block->insert({int_type->create_column_const(block->rows(), to_field<TYPE_INT>(len)),
932
25
                       int_type, "const len"});
933
25
        block->insert({nullptr, std::make_shared<DataTypeString>(), "result"});
934
935
25
        ColumnNumbers temp_arguments(3);
936
25
        temp_arguments[0] = cast_set<uint32_t>(idx);
937
25
        temp_arguments[1] = num_columns_without_result;
938
25
        temp_arguments[2] = num_columns_without_result + 1;
939
25
        const uint32_t result_column_id = num_columns_without_result + 2;
940
25
        SubstringUtil::substring_execute(*block, temp_arguments, result_column_id, block->rows());
941
942
25
        ColumnPtr result_column = block->get_by_position(result_column_id).column;
943
25
        if (is_nullable) {
944
25
            result_column = ColumnNullable::create(std::move(result_column), null_map_column);
945
25
        }
946
25
        block->replace_by_position(idx, std::move(result_column));
947
25
        block->erase_tail(num_columns_without_result);
948
25
    }
949
950
58.6k
    Status _try_materialize_aggregate_pushdown_rows(Block* block, bool* pushed_down) {
951
58.6k
        DORIS_CHECK(block != nullptr);
952
58.6k
        DORIS_CHECK(pushed_down != nullptr);
953
58.6k
        *pushed_down = false;
954
58.6k
        block->clear_column_data(_projected_columns.size());
955
58.6k
        _aggregate_pushdown_tried = true;
956
58.6k
        if (!_supports_aggregate_pushdown(_push_down_agg_type)) {
957
57.4k
            return Status::OK();
958
57.4k
        }
959
960
1.22k
        FileAggregateRequest file_request;
961
1.22k
        RETURN_IF_ERROR(_build_file_aggregate_request(_push_down_agg_type, &file_request));
962
1.22k
        FileAggregateResult file_result;
963
1.22k
        const auto status = _data_reader.reader->get_aggregate_result(file_request, &file_result);
964
1.22k
        if (status.is<ErrorCode::NOT_IMPLEMENTED_ERROR>()) {
965
5
            return Status::OK();
966
5
        }
967
1.21k
        RETURN_IF_ERROR(status);
968
1.21k
        RETURN_IF_ERROR(
969
1.21k
                _materialize_aggregate_pushdown_rows(_push_down_agg_type, file_result, block));
970
1.25k
        if (_push_down_agg_type == TPushAggOp::type::COUNT) {
971
1.25k
            _current_split_uses_metadata_count = true;
972
1.25k
        }
973
1.21k
        *pushed_down = true;
974
1.21k
        RETURN_IF_ERROR(close_current_reader());
975
1.21k
        return Status::OK();
976
1.21k
    }
977
978
60.0k
    virtual bool _supports_aggregate_pushdown(TPushAggOp::type agg_type) const {
979
        // Only COUNT and MIN/MAX can be push down.
980
60.0k
        if (agg_type != TPushAggOp::type::COUNT && agg_type != TPushAggOp::type::MINMAX) {
981
56.6k
            return false;
982
56.6k
        }
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
3.32k
        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
3.31k
        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
3.31k
        if ((_delete_rows != nullptr && !_delete_rows->empty()) ||
1001
3.31k
            (_deletion_vector != nullptr && !_deletion_vector->isEmpty())) {
1002
597
            return false;
1003
597
        }
1004
2.71k
        if (!_table_filters.empty()) {
1005
0
            return false;
1006
0
        }
1007
2.71k
        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
2.64k
            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
2.64k
            if (_push_down_count_columns->empty()) {
1019
2.63k
                return true;
1020
2.63k
            }
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
149
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
1038
149
            if (!mapping.file_local_id.has_value() ||
1039
149
                mapping.virtual_column_type != TableVirtualColumnType::INVALID ||
1040
149
                mapping.default_expr != nullptr || mapping.file_type == nullptr ||
1041
149
                mapping.table_type == nullptr) {
1042
9
                return false;
1043
9
            }
1044
140
            if (!_can_push_down_minmax_for_mapping(mapping)) {
1045
46
                return false;
1046
46
            }
1047
140
        }
1048
16
        return true;
1049
71
    }
1050
1051
654k
    static ColumnPtr _detach_column(ColumnPtr column) {
1052
654k
        DORIS_CHECK(column.get() != nullptr);
1053
654k
        return IColumn::mutate(std::move(column));
1054
654k
    }
1055
1056
94.5k
    static Status _align_column_nullability(ColumnPtr* column, const DataTypePtr& table_type) {
1057
94.5k
        DORIS_CHECK(column != nullptr);
1058
94.5k
        DORIS_CHECK(column->get() != nullptr);
1059
94.5k
        DORIS_CHECK(table_type != nullptr);
1060
        // Must return non-const column
1061
94.5k
        *column = (*column)->convert_to_full_column_if_const();
1062
94.5k
        if (table_type->is_nullable()) {
1063
47.2k
            const auto& nested_type =
1064
47.2k
                    assert_cast<const DataTypeNullable&>(*table_type).get_nested_type();
1065
47.2k
            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
47.2k
            const auto& nullable_column = assert_cast<const ColumnNullable&>(**column);
1071
47.2k
            ColumnPtr nested_column = nullable_column.get_nested_column_ptr();
1072
47.2k
            RETURN_IF_ERROR(_align_column_nullability(&nested_column, nested_type));
1073
47.2k
            *column = ColumnNullable::create(nested_column,
1074
47.2k
                                             nullable_column.get_null_map_column_ptr());
1075
47.2k
            return Status::OK();
1076
47.2k
        }
1077
47.3k
        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
47.3k
        if (const auto* array_type = typeid_cast<const DataTypeArray*>(table_type.get())) {
1089
129
            const auto& array_column = assert_cast<const ColumnArray&>(**column);
1090
129
            ColumnPtr nested_column = array_column.get_data_ptr();
1091
129
            RETURN_IF_ERROR(
1092
129
                    _align_column_nullability(&nested_column, array_type->get_nested_type()));
1093
129
            *column = ColumnArray::create(nested_column, array_column.get_offsets_ptr());
1094
129
            return Status::OK();
1095
129
        }
1096
47.1k
        if (const auto* map_type = typeid_cast<const DataTypeMap*>(table_type.get())) {
1097
8
            const auto& map_column = assert_cast<const ColumnMap&>(**column);
1098
8
            ColumnPtr key_column = map_column.get_keys_ptr();
1099
8
            ColumnPtr value_column = map_column.get_values_ptr();
1100
8
            RETURN_IF_ERROR(_align_column_nullability(&key_column, map_type->get_key_type()));
1101
8
            RETURN_IF_ERROR(_align_column_nullability(&value_column, map_type->get_value_type()));
1102
8
            *column = ColumnMap::create(key_column, value_column, map_column.get_offsets_ptr());
1103
8
            return Status::OK();
1104
8
        }
1105
47.1k
        if (const auto* struct_type = typeid_cast<const DataTypeStruct*>(table_type.get())) {
1106
5.62k
            const auto& struct_column = assert_cast<const ColumnStruct&>(**column);
1107
5.62k
            Columns columns = struct_column.get_columns_copy();
1108
5.62k
            DORIS_CHECK(columns.size() == struct_type->get_elements().size());
1109
19.0k
            for (size_t i = 0; i < columns.size(); ++i) {
1110
13.3k
                RETURN_IF_ERROR(
1111
13.3k
                        _align_column_nullability(&columns[i], struct_type->get_element(i)));
1112
13.3k
            }
1113
5.62k
            *column = ColumnStruct::create(columns);
1114
5.62k
            return Status::OK();
1115
5.62k
        }
1116
41.5k
        return Status::OK();
1117
47.1k
    }
1118
1119
    static Status _execute_default_expr_without_root_type_check(
1120
            const VExprContextSPtr& default_expr, const Block* block,
1121
16.7k
            ColumnWithTypeAndName* result_data) {
1122
16.7k
        DORIS_CHECK(default_expr != nullptr);
1123
16.7k
        DORIS_CHECK(block != nullptr);
1124
16.7k
        DORIS_CHECK(result_data != nullptr);
1125
16.7k
        ColumnPtr result_column;
1126
16.7k
        Status st;
1127
16.7k
        RETURN_IF_CATCH_EXCEPTION({
1128
16.7k
            st = default_expr->root()->execute_column_impl(default_expr.get(), block, nullptr,
1129
16.7k
                                                           block->rows(), result_column);
1130
16.7k
        });
1131
16.7k
        RETURN_IF_ERROR(st);
1132
16.7k
        DORIS_CHECK(result_column.get() != nullptr);
1133
16.7k
        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
16.7k
        result_data->column = result_column;
1139
16.7k
        result_data->type = default_expr->execute_type(block);
1140
16.7k
        result_data->name = default_expr->expr_name();
1141
16.7k
        return Status::OK();
1142
16.7k
    }
1143
1144
    Status _cast_column_to_type(ColumnPtr* column, const DataTypePtr& file_type,
1145
                                const DataTypePtr& table_type,
1146
8.77k
                                const std::string& column_name) const {
1147
8.77k
        DORIS_CHECK(column != nullptr);
1148
8.77k
        DORIS_CHECK(column->get() != nullptr);
1149
8.77k
        DORIS_CHECK(file_type != nullptr);
1150
8.77k
        DORIS_CHECK(table_type != nullptr);
1151
8.77k
        if (file_type->equals(*table_type)) {
1152
0
            return Status::OK();
1153
0
        }
1154
1155
8.77k
        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
8.77k
        if ((*column)->is_nullable() && !input_type->is_nullable()) {
1159
0
            input_type = make_nullable(input_type);
1160
8.77k
        } else if (!(*column)->is_nullable() && input_type->is_nullable()) {
1161
1
            input_type = remove_nullable(input_type);
1162
1
        }
1163
8.77k
        Block cast_block;
1164
8.77k
        cast_block.insert({*column, input_type, column_name});
1165
8.77k
        auto slot_ref = VSlotRef::create_shared(0, 0, -1, input_type, column_name);
1166
8.77k
        auto cast_expr = Cast::create_shared(table_type);
1167
8.77k
        cast_expr->add_child(std::move(slot_ref));
1168
8.77k
        auto cast_ctx = VExprContext::create_shared(std::move(cast_expr));
1169
8.77k
        RowDescriptor row_desc;
1170
8.77k
        RETURN_IF_ERROR(cast_ctx->prepare(_runtime_state, row_desc));
1171
8.77k
        RETURN_IF_ERROR(cast_ctx->open(_runtime_state));
1172
8.77k
        ColumnPtr cast_column;
1173
8.77k
        RETURN_IF_ERROR(cast_ctx->execute(&cast_block, cast_column));
1174
8.77k
        *column = std::move(cast_column);
1175
8.77k
        return Status::OK();
1176
8.77k
    }
1177
1178
    Status _materialize_present_child_mapping_column(const ColumnMapping& mapping,
1179
                                                     const ColumnPtr& file_column,
1180
17.0k
                                                     const size_t rows, ColumnPtr* column) {
1181
17.0k
        DORIS_CHECK(column != nullptr);
1182
17.0k
        DORIS_CHECK(mapping.file_type != nullptr);
1183
17.0k
        DORIS_CHECK(mapping.table_type != nullptr);
1184
17.0k
        *column = file_column;
1185
17.0k
        if (!mapping.is_trivial) {
1186
10.7k
            if (!mapping.child_mappings.empty()) {
1187
1.99k
                RETURN_IF_ERROR(
1188
1.99k
                        _materialize_complex_mapping_column(mapping, *column, rows, column));
1189
8.77k
            } else {
1190
8.77k
                RETURN_IF_ERROR(_cast_column_to_type(column, mapping.file_type, mapping.table_type,
1191
8.77k
                                                     mapping.file_column_name));
1192
8.77k
            }
1193
10.7k
        }
1194
17.0k
        RETURN_IF_ERROR(_align_column_nullability(column, mapping.table_type));
1195
17.0k
        return Status::OK();
1196
17.0k
    }
1197
1198
    Status _materialize_mapping_column(const ColumnMapping& mapping, Block* current_block,
1199
663k
                                       const size_t rows, ColumnPtr* column) {
1200
663k
        if (!mapping.is_trivial && mapping.file_local_id.has_value() &&
1201
663k
            !mapping.child_mappings.empty()) {
1202
9.59k
            DCHECK(mapping.projection != nullptr);
1203
9.59k
            int res_id;
1204
9.59k
            auto st = mapping.projection->execute(current_block, &res_id);
1205
9.59k
            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
9.59k
            ColumnPtr result_column = current_block->get_by_position(res_id).column;
1213
9.59k
            RETURN_IF_ERROR(
1214
9.59k
                    _materialize_complex_mapping_column(mapping, result_column, rows, column));
1215
9.59k
            return Status::OK();
1216
9.59k
        }
1217
654k
        if (mapping.projection != nullptr) {
1218
634k
            int res_id;
1219
634k
            auto st = mapping.projection->execute(current_block, &res_id);
1220
634k
            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
634k
            ColumnPtr result_column = current_block->get_by_position(res_id).column;
1232
634k
            *column = _detach_column(std::move(result_column));
1233
634k
            return Status::OK();
1234
634k
        }
1235
19.5k
        if (mapping.default_expr != nullptr) {
1236
16.7k
            if (current_block->rows() == rows) {
1237
14.7k
                ColumnWithTypeAndName result;
1238
14.7k
                RETURN_IF_ERROR(_execute_default_expr_without_root_type_check(
1239
14.7k
                        mapping.default_expr, current_block, &result));
1240
14.7k
                ColumnPtr result_column = result.column;
1241
14.7k
                RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type));
1242
14.7k
                *column = _detach_column(std::move(result_column));
1243
14.7k
            } else {
1244
1.99k
                DORIS_CHECK(mapping.constant_index.has_value());
1245
1.99k
                Block eval_block;
1246
1.99k
                eval_block.insert({mapping.table_type->create_column_const_with_default_value(rows),
1247
1.99k
                                   mapping.table_type, "__table_reader_const_rows"});
1248
1.99k
                ColumnWithTypeAndName result;
1249
1.99k
                RETURN_IF_ERROR(_execute_default_expr_without_root_type_check(
1250
1.99k
                        mapping.default_expr, &eval_block, &result));
1251
1.99k
                ColumnPtr result_column = result.column;
1252
1.99k
                RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type));
1253
1.99k
                *column = _detach_column(std::move(result_column));
1254
1.99k
            }
1255
16.7k
            return Status::OK();
1256
16.7k
        }
1257
2.86k
        ColumnPtr result_column = mapping.table_type->create_column_const_with_default_value(rows);
1258
2.86k
        *column = _detach_column(std::move(result_column));
1259
2.86k
        return Status::OK();
1260
19.5k
    }
1261
1262
    Status _materialize_complex_mapping_column(const ColumnMapping& mapping,
1263
                                               const ColumnPtr& file_column, const size_t rows,
1264
11.5k
                                               ColumnPtr* column) {
1265
11.5k
        DORIS_CHECK(mapping.table_type != nullptr);
1266
11.5k
        DORIS_CHECK(file_column.get() != nullptr);
1267
11.5k
        const auto table_type = remove_nullable(mapping.table_type);
1268
11.5k
        switch (table_type->get_primitive_type()) {
1269
3.59k
        case TYPE_STRUCT:
1270
3.59k
            RETURN_IF_ERROR(_materialize_struct_mapping_column(mapping, file_column, rows, column));
1271
3.59k
            break;
1272
4.43k
        case TYPE_ARRAY:
1273
4.43k
            RETURN_IF_ERROR(_materialize_array_mapping_column(mapping, file_column, rows, column));
1274
4.43k
            break;
1275
4.43k
        case TYPE_MAP:
1276
3.55k
            RETURN_IF_ERROR(_materialize_map_mapping_column(mapping, file_column, rows, column));
1277
3.55k
            break;
1278
3.55k
        default:
1279
0
            *column = _detach_column(file_column);
1280
0
            break;
1281
11.5k
        }
1282
11.5k
        return Status::OK();
1283
11.5k
    }
1284
1285
    static std::vector<const ColumnMapping*> _present_child_mappings_in_file_order(
1286
3.59k
            const std::vector<ColumnMapping>& child_mappings) {
1287
3.59k
        std::vector<const ColumnMapping*> result;
1288
3.59k
        result.reserve(child_mappings.size());
1289
9.13k
        for (const auto& child_mapping : child_mappings) {
1290
9.13k
            if (child_mapping.file_local_id.has_value()) {
1291
5.52k
                result.push_back(&child_mapping);
1292
5.52k
            }
1293
9.13k
        }
1294
4.02k
        std::ranges::sort(result, [](const ColumnMapping* lhs, const ColumnMapping* rhs) {
1295
4.02k
            DORIS_CHECK(lhs->file_local_id.has_value());
1296
4.02k
            DORIS_CHECK(rhs->file_local_id.has_value());
1297
4.02k
            return *lhs->file_local_id < *rhs->file_local_id;
1298
4.02k
        });
1299
3.59k
        return result;
1300
3.59k
    }
1301
1302
    static size_t _file_child_ordinal_for_mapping(
1303
            const ColumnMapping& mapping, const ColumnMapping& child_mapping,
1304
5.52k
            const std::vector<const ColumnMapping*>& file_ordered_children) {
1305
5.52k
        DORIS_CHECK(child_mapping.file_local_id.has_value());
1306
5.52k
        if (!mapping.projected_file_children.empty()) {
1307
5.51k
            const auto child_it = std::ranges::find_if(
1308
8.48k
                    mapping.projected_file_children, [&](const ColumnDefinition& file_child) {
1309
8.48k
                        return file_child.file_local_id() == *child_mapping.file_local_id;
1310
8.48k
                    });
1311
5.51k
            DORIS_CHECK(child_it != mapping.projected_file_children.end());
1312
5.51k
            return static_cast<size_t>(
1313
5.51k
                    std::distance(mapping.projected_file_children.begin(), child_it));
1314
5.51k
        }
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
5.52k
    }
1319
1320
    static std::vector<const ColumnMapping*> _child_mappings_in_table_type_order(
1321
3.59k
            const ColumnMapping& mapping, const DataTypeStruct& table_type) {
1322
3.59k
        std::vector<const ColumnMapping*> result;
1323
3.59k
        result.reserve(mapping.child_mappings.size());
1324
12.7k
        for (size_t child_idx = 0; child_idx < table_type.get_elements().size(); ++child_idx) {
1325
9.13k
            const auto& child_name = table_type.get_element_name(child_idx);
1326
9.13k
            const auto child_it = std::ranges::find_if(
1327
17.5k
                    mapping.child_mappings, [&](const ColumnMapping& child_mapping) {
1328
17.5k
                        return child_mapping.table_column_name == child_name;
1329
17.5k
                    });
1330
9.13k
            DORIS_CHECK(child_it != mapping.child_mappings.end())
1331
0
                    << mapping.debug_string() << ", table_child_name=" << child_name;
1332
9.13k
            result.push_back(&*child_it);
1333
9.13k
        }
1334
3.59k
        return result;
1335
3.59k
    }
1336
1337
    static const IColumn* _nested_column_if_nullable(const ColumnPtr& column,
1338
11.5k
                                                     const NullMap** null_map) {
1339
11.5k
        DORIS_CHECK(column.get() != nullptr);
1340
11.5k
        if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1341
11.5k
            if (null_map != nullptr) {
1342
11.5k
                *null_map = &nullable_column->get_null_map_data();
1343
11.5k
            }
1344
11.5k
            return &nullable_column->get_nested_column();
1345
11.5k
        }
1346
4
        return column.get();
1347
11.5k
    }
1348
1349
    Status _materialize_struct_mapping_column(const ColumnMapping& mapping,
1350
                                              const ColumnPtr& file_column, const size_t rows,
1351
3.59k
                                              ColumnPtr* column) {
1352
3.59k
        DORIS_CHECK(mapping.table_type != nullptr);
1353
3.59k
        const auto* table_type =
1354
3.59k
                assert_cast<const DataTypeStruct*>(remove_nullable(mapping.table_type).get());
1355
3.59k
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1356
3.59k
        const NullMap* parent_null_map = nullptr;
1357
3.59k
        const auto* nested_file_column =
1358
3.59k
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1359
3.59k
        const auto* file_struct = assert_cast<const ColumnStruct*>(nested_file_column);
1360
3.59k
        DORIS_CHECK(table_type->get_elements().size() == mapping.child_mappings.size());
1361
1362
3.59k
        Columns child_columns;
1363
3.59k
        child_columns.reserve(mapping.child_mappings.size());
1364
3.59k
        const auto file_ordered_children =
1365
3.59k
                _present_child_mappings_in_file_order(mapping.child_mappings);
1366
3.59k
        const auto table_ordered_children =
1367
3.59k
                _child_mappings_in_table_type_order(mapping, *table_type);
1368
9.13k
        for (const auto* child_mapping : table_ordered_children) {
1369
9.13k
            DORIS_CHECK(child_mapping != nullptr);
1370
9.13k
            if (!child_mapping->file_local_id.has_value()) {
1371
3.61k
                child_columns.push_back(
1372
3.61k
                        child_mapping->table_type->create_column_const_with_default_value(rows)
1373
3.61k
                                ->convert_to_full_column_if_const());
1374
3.61k
                continue;
1375
3.61k
            }
1376
5.52k
            const auto file_child_idx =
1377
5.52k
                    _file_child_ordinal_for_mapping(mapping, *child_mapping, file_ordered_children);
1378
5.52k
            DORIS_CHECK(file_child_idx < file_struct->get_columns().size());
1379
5.52k
            ColumnPtr child_column = file_struct->get_column_ptr(file_child_idx);
1380
5.52k
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(*child_mapping, child_column,
1381
5.52k
                                                                      rows, &child_column));
1382
5.52k
            child_columns.push_back(std::move(child_column));
1383
5.52k
        }
1384
3.59k
        MutableColumns mutable_child_columns;
1385
3.59k
        mutable_child_columns.reserve(child_columns.size());
1386
9.13k
        for (auto& child_column : child_columns) {
1387
9.13k
            mutable_child_columns.push_back(IColumn::mutate(std::move(child_column)));
1388
9.13k
        }
1389
3.59k
        auto result = ColumnStruct::create(std::move(mutable_child_columns));
1390
3.59k
        if (mapping.table_type->is_nullable()) {
1391
3.59k
            auto null_map = ColumnUInt8::create();
1392
3.59k
            auto& null_map_data = null_map->get_data();
1393
3.59k
            null_map_data.resize(rows);
1394
3.59k
            if (parent_null_map != nullptr) {
1395
3.59k
                DORIS_CHECK(parent_null_map->size() == rows);
1396
3.59k
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
1397
3.59k
            } else {
1398
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
1399
0
            }
1400
3.59k
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
1401
3.59k
        } else {
1402
2
            *column = std::move(result);
1403
2
        }
1404
3.59k
        return Status::OK();
1405
3.59k
    }
1406
1407
    Status _materialize_array_mapping_column(const ColumnMapping& mapping,
1408
                                             const ColumnPtr& file_column, const size_t rows,
1409
4.43k
                                             ColumnPtr* column) {
1410
4.43k
        DORIS_CHECK(mapping.child_mappings.size() == 1);
1411
4.43k
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1412
4.43k
        const NullMap* parent_null_map = nullptr;
1413
4.43k
        const auto* nested_file_column =
1414
4.43k
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1415
4.43k
        const auto* file_array = assert_cast<const ColumnArray*>(nested_file_column);
1416
4.43k
        ColumnPtr nested_column = file_array->get_data_ptr();
1417
4.43k
        const auto& element_mapping = mapping.child_mappings[0];
1418
4.43k
        RETURN_IF_ERROR(_materialize_present_child_mapping_column(
1419
4.43k
                element_mapping, nested_column, nested_column->size(), &nested_column));
1420
4.43k
        auto offsets_column = file_array->get_offsets_ptr()->convert_to_full_column_if_const();
1421
4.43k
        auto result = ColumnArray::create(IColumn::mutate(std::move(nested_column)),
1422
4.43k
                                          IColumn::mutate(std::move(offsets_column)));
1423
4.44k
        if (mapping.table_type->is_nullable()) {
1424
4.44k
            auto null_map = ColumnUInt8::create();
1425
4.44k
            auto& null_map_data = null_map->get_data();
1426
4.44k
            null_map_data.resize(rows);
1427
4.44k
            if (parent_null_map != nullptr) {
1428
4.44k
                DORIS_CHECK(parent_null_map->size() == rows);
1429
4.44k
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
1430
4.44k
            } else {
1431
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
1432
0
            }
1433
4.44k
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
1434
18.4E
        } else {
1435
18.4E
            *column = std::move(result);
1436
18.4E
        }
1437
4.43k
        return Status::OK();
1438
4.43k
    }
1439
1440
    Status _materialize_map_mapping_column(const ColumnMapping& mapping,
1441
                                           const ColumnPtr& file_column, const size_t rows,
1442
3.55k
                                           ColumnPtr* column) {
1443
3.55k
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1444
3.55k
        const NullMap* parent_null_map = nullptr;
1445
3.55k
        const auto* nested_file_column =
1446
3.55k
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1447
3.55k
        const auto* file_map = assert_cast<const ColumnMap*>(nested_file_column);
1448
3.55k
        ColumnPtr key_column = file_map->get_keys_ptr();
1449
3.55k
        ColumnPtr value_column = file_map->get_values_ptr();
1450
1451
3.55k
        const ColumnMapping* key_mapping = nullptr;
1452
3.55k
        const ColumnMapping* value_mapping = nullptr;
1453
7.11k
        for (const auto& child_mapping : mapping.child_mappings) {
1454
7.11k
            if (!child_mapping.file_local_id.has_value()) {
1455
0
                continue;
1456
0
            }
1457
7.11k
            if (*child_mapping.file_local_id == 0) {
1458
3.55k
                key_mapping = &child_mapping;
1459
3.55k
            } else if (*child_mapping.file_local_id == 1) {
1460
3.55k
                value_mapping = &child_mapping;
1461
3.55k
            }
1462
7.11k
        }
1463
1464
3.55k
        if (key_mapping != nullptr) {
1465
3.55k
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(
1466
3.55k
                    *key_mapping, key_column, key_column->size(), &key_column));
1467
3.55k
        }
1468
3.55k
        if (value_mapping != nullptr) {
1469
3.55k
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(
1470
3.55k
                    *value_mapping, value_column, value_column->size(), &value_column));
1471
3.55k
        }
1472
3.55k
        auto offsets_column = file_map->get_offsets_ptr()->convert_to_full_column_if_const();
1473
3.55k
        auto result = ColumnMap::create(IColumn::mutate(std::move(key_column)),
1474
3.55k
                                        IColumn::mutate(std::move(value_column)),
1475
3.55k
                                        IColumn::mutate(std::move(offsets_column)));
1476
3.55k
        if (mapping.table_type->is_nullable()) {
1477
3.55k
            auto null_map = ColumnUInt8::create();
1478
3.55k
            auto& null_map_data = null_map->get_data();
1479
3.55k
            null_map_data.resize(rows);
1480
3.55k
            if (parent_null_map != nullptr) {
1481
3.55k
                DORIS_CHECK(parent_null_map->size() == rows);
1482
3.55k
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
1483
3.55k
            } else {
1484
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
1485
0
            }
1486
3.55k
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
1487
3.55k
        } else {
1488
0
            *column = std::move(result);
1489
0
        }
1490
3.55k
        return Status::OK();
1491
3.55k
    }
1492
1493
58.5k
    Status _open_mapping_exprs() {
1494
58.5k
        RowDescriptor row_desc;
1495
393k
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
1496
393k
            if (mapping.projection != nullptr) {
1497
373k
                RETURN_IF_ERROR(mapping.projection->prepare(_runtime_state, row_desc));
1498
373k
                RETURN_IF_ERROR(mapping.projection->open(_runtime_state));
1499
373k
            }
1500
393k
            if (mapping.default_expr != nullptr) {
1501
17.2k
                RETURN_IF_ERROR(mapping.default_expr->prepare(_runtime_state, row_desc));
1502
17.2k
                RETURN_IF_ERROR(mapping.default_expr->open(_runtime_state));
1503
17.2k
            }
1504
393k
        }
1505
58.5k
        return Status::OK();
1506
58.5k
    }
1507
1508
    Status _build_file_aggregate_request(TPushAggOp::type agg_type,
1509
1.28k
                                         FileAggregateRequest* request) const {
1510
1.28k
        DORIS_CHECK(request != nullptr);
1511
1.28k
        DORIS_CHECK(_supports_aggregate_pushdown(agg_type));
1512
1.28k
        request->agg_type = agg_type;
1513
1.28k
        request->columns.clear();
1514
1.28k
        if (agg_type == TPushAggOp::type::COUNT) {
1515
1.25k
            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
1.25k
            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
1.25k
            return Status::OK();
1530
1.25k
        }
1531
26
        request->columns.reserve(_data_reader.column_mapper->mappings().size());
1532
47
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
1533
47
            DORIS_CHECK(mapping.file_local_id.has_value());
1534
47
            FileAggregateRequest::Column column;
1535
47
            column.projection = LocalColumnIndex::top_level(LocalColumnId(*mapping.file_local_id));
1536
47
            if (!mapping.child_mappings.empty()) {
1537
1
                RETURN_IF_ERROR(build_aggregate_projection(mapping, &column.projection));
1538
1
            }
1539
47
            request->columns.push_back(std::move(column));
1540
47
        }
1541
26
        return Status::OK();
1542
26
    }
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
1.27k
                                                Block* block) {
1559
1.27k
        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
1.25k
            DORIS_CHECK(file_result.count >= 0);
1564
1.25k
            return _materialize_count_rows(cast_set<size_t>(file_result.count), block);
1565
1.25k
        }
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
21
        DORIS_CHECK(file_result.columns.size() == _data_reader.column_mapper->mappings().size());
1569
21
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
1570
21
        Block file_block;
1571
21
        file_block.reserve(_data_reader.file_block_layout.size());
1572
26
        for (const auto& column : _data_reader.file_block_layout) {
1573
26
            file_block.insert({column.type->create_column(), column.type, column.name});
1574
26
        }
1575
47
        for (size_t column_idx = 0; column_idx < file_result.columns.size(); ++column_idx) {
1576
26
            const auto& result_column = file_result.columns[column_idx];
1577
26
            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
26
            bool found_file_column = false;
1582
33
            for (size_t block_position = 0; block_position < _data_reader.file_block_layout.size();
1583
33
                 ++block_position) {
1584
33
                if (_data_reader.file_block_layout[block_position].file_column_id ==
1585
33
                    file_result.columns[column_idx].projection.column_id()) {
1586
26
                    found_file_column = true;
1587
26
                    auto column = file_block.get_by_position(block_position)
1588
26
                                          .type->create_column()
1589
26
                                          ->assert_mutable();
1590
26
                    RETURN_IF_ERROR(_insert_aggregate_projection_value(
1591
26
                            file_result.columns[column_idx].projection, result_column.min_value,
1592
26
                            column.get()));
1593
26
                    RETURN_IF_ERROR(_insert_aggregate_projection_value(
1594
26
                            file_result.columns[column_idx].projection, result_column.max_value,
1595
26
                            column.get()));
1596
26
                    file_block.replace_by_position(block_position, std::move(column));
1597
26
                    break;
1598
26
                }
1599
33
            }
1600
26
            DORIS_CHECK(found_file_column);
1601
26
        }
1602
47
        for (size_t column_idx = 0; column_idx < _data_reader.column_mapper->mappings().size();
1603
26
             ++column_idx) {
1604
26
            ColumnPtr table_column;
1605
26
            RETURN_IF_ERROR(
1606
26
                    _materialize_mapping_column(_data_reader.column_mapper->mappings()[column_idx],
1607
26
                                                &file_block, 2, &table_column));
1608
26
            block->replace_by_position(column_idx, std::move(table_column));
1609
26
        }
1610
21
        return Status::OK();
1611
21
    }
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
389k
            const std::vector<ColumnDefinition>& schema, LocalColumnId column_id) {
1689
8.42M
        for (const auto& field : schema) {
1690
8.42M
            if (field.file_local_id() == column_id.value()) {
1691
381k
                return &field;
1692
381k
            }
1693
8.42M
        }
1694
7.68k
        return nullptr;
1695
389k
    }
1696
1697
142
    static bool _can_push_down_minmax_for_mapping(const ColumnMapping& mapping) {
1698
142
        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
139
            return mapping.is_trivial;
1703
139
        }
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
108
                                                     const Field& value, IColumn* column) {
1746
108
        DORIS_CHECK(column != nullptr);
1747
108
        if (auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1748
54
            RETURN_IF_ERROR(_insert_aggregate_projection_value(
1749
54
                    projection, value, &nullable_column->get_nested_column()));
1750
54
            nullable_column->get_null_map_data().push_back(0);
1751
54
            return Status::OK();
1752
54
        }
1753
54
        if (projection.project_all_children || projection.children.empty()) {
1754
52
            column->insert(value);
1755
52
            return Status::OK();
1756
52
        }
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