Coverage Report

Created: 2026-08-21 12:10

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
    bool metadata_pruning_safe = true;
86
};
87
88
struct ScanTask {
89
480
    virtual ~ScanTask() = default;
90
91
    std::unique_ptr<io::FileDescription> data_file;
92
};
93
94
struct ProjectedColumnBuildContext {
95
    const TFileScanRangeParams* scan_params = nullptr;
96
    const TFileRangeDesc* range = nullptr;
97
    RuntimeState* runtime_state = nullptr;
98
    const SlotDescriptor* slot_desc = nullptr;
99
    std::optional<ColumnDefinition> schema_column = std::nullopt;
100
    size_t next_file_column_idx = 0;
101
};
102
103
struct ReadProfile {
104
    RuntimeProfile::Counter* total_timer = nullptr;
105
    RuntimeProfile::Counter* init_timer = nullptr;
106
    RuntimeProfile::Counter* num_delete_files = nullptr;
107
    RuntimeProfile::Counter* num_delete_rows = nullptr;
108
    RuntimeProfile::Counter* parse_delete_file_time = nullptr;
109
    RuntimeProfile::Counter* decoded_dv_cache_hit_count = nullptr;
110
    RuntimeProfile::Counter* decoded_dv_cache_miss_count = nullptr;
111
    RuntimeProfile::Counter* dv_file_cache_hit_count = nullptr;
112
    RuntimeProfile::Counter* dv_file_cache_miss_count = nullptr;
113
    RuntimeProfile::Counter* dv_file_cache_peer_read_count = nullptr;
114
    RuntimeProfile::Counter* exec_timer = nullptr;
115
    RuntimeProfile::Counter* prepare_split_timer = nullptr;
116
    RuntimeProfile::Counter* finalize_timer = nullptr;
117
    RuntimeProfile::Counter* create_reader_timer = nullptr;
118
    RuntimeProfile::Counter* pushdown_agg_timer = nullptr;
119
    RuntimeProfile::Counter* open_reader_timer = nullptr;
120
    RuntimeProfile::Counter* refresh_conjuncts_timer = nullptr;
121
    RuntimeProfile::Counter* runtime_filter_partition_prune_timer = nullptr;
122
    RuntimeProfile::Counter* runtime_filter_partition_pruned_range_counter = nullptr;
123
    RuntimeProfile::Counter* close_timer = nullptr;
124
    RuntimeProfile::Counter* file_reader_total_timer = nullptr;
125
    RuntimeProfile::Counter* file_reader_init_timer = nullptr;
126
    RuntimeProfile::Counter* file_reader_schema_timer = nullptr;
127
    RuntimeProfile::Counter* file_reader_mapper_timer = nullptr;
128
    RuntimeProfile::Counter* file_reader_open_timer = nullptr;
129
    RuntimeProfile::Counter* file_reader_refresh_timer = nullptr;
130
    RuntimeProfile::Counter* file_reader_get_block_timer = nullptr;
131
    RuntimeProfile::Counter* file_reader_aggregate_timer = nullptr;
132
    RuntimeProfile::Counter* file_reader_close_timer = nullptr;
133
};
134
135
struct TableReadOptions {
136
    // Columns need to be read from file and output by table reader. They are all in table/global
137
    // schema semantics.
138
    const std::vector<ColumnDefinition> projected_columns;
139
    // All complex conjuncts from scan operator
140
    const VExprContextSPtrs conjuncts;
141
    // File format of the underlying data files, needed for reader initialization and reader-level
142
    // filter pushdown.
143
    const FileFormat format;
144
    TFileScanRangeParams* scan_params;
145
    std::shared_ptr<io::IOContext> io_ctx;
146
    RuntimeState* runtime_state;
147
    RuntimeProfile* scanner_profile;
148
    // File formats without complete self-describing metadata, such as CSV, Text, and JSON, need
149
    // the FE-planned physical file slots to build their file-local schema and deserialize values.
150
    const std::vector<SlotDescriptor*>* file_slot_descs = nullptr;
151
    // Push-down aggregate type.
152
    const TPushAggOp::type push_down_agg_type = TPushAggOp::type::NONE;
153
    // Table/global indices of explicit COUNT arguments. nullopt means an old FE did not send the
154
    // semantic argument field, while an explicit empty vector means COUNT(*)/COUNT(1). Keeping
155
    // those states separate prevents a rolling-upgrade plan from being reinterpreted by a new BE.
156
    const std::optional<std::vector<GlobalIndex>> push_down_count_columns = std::nullopt;
157
    // Initial digest of predicates available during scanner open. Scanner-driven splits override it
158
    // with SplitReadOptions::condition_cache_digest after collecting late-arrival runtime filters.
159
    // A zero digest disables condition cache.
160
    uint64_t condition_cache_digest = 0;
161
};
162
163
struct SplitReadOptions {
164
    // 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.
165
    std::map<std::string, Field> partition_values;
166
    // Latest scanner conjuncts rewritten to table/global column indices. Runtime filters may
167
    // arrive after TableReader::init(), so scanner-driven splits replace the initial snapshot.
168
    // nullopt preserves the initial snapshot for standalone TableReader callers.
169
    std::optional<VExprContextSPtrs> conjuncts = std::nullopt;
170
    // Independent clones used for partition pruning because evaluation prepares and opens them
171
    // against a synthetic partition block before the file reader opens its row-level conjuncts.
172
    VExprContextSPtrs partition_prune_conjuncts;
173
    // Table-level COUNT may emit one metadata-derived batch and resume on a later scheduler turn.
174
    // It is safe only after every runtime filter assigned to the scanner has arrived; otherwise a
175
    // filter could arrive after synthetic rows have already been returned and those rows cannot be
176
    // retracted. Standalone TableReader callers have no scanner runtime-filter lifecycle.
177
    bool all_runtime_filters_applied = true;
178
    // Digest for the exact scanner conjunct snapshot attached to this split. FileScannerV2 rebuilds
179
    // it after collecting late-arrival RFs, so different RF payloads cannot share a cache entry. A
180
    // zero value explicitly disables condition cache for this split.
181
    std::optional<uint64_t> condition_cache_digest;
182
    ShardedKVCache* cache = nullptr;
183
    TFileRangeDesc current_range;
184
    FileFormat current_split_format = FileFormat::PARQUET;
185
    std::optional<GlobalRowIdContext> global_rowid_context;
186
};
187
188
// Base class for table-level readers.
189
// This layer owns common table-level orchestration, such as split iteration, dynamic partition
190
// pruning, delete handling and conversion from file-local blocks to table-schema blocks. Concrete
191
// table-format readers only need to provide format-specific hooks for opening readers and parsing
192
// split metadata.
193
class TableReader {
194
public:
195
600
    virtual ~TableReader() = default;
196
197
    // Initialize common runtime options for the table reader. Subclasses may call this from their
198
    // own init(options); table-format schema and split metadata are provided later per split.
199
    virtual Status init(TableReadOptions&& options);
200
201
    // FileScannerV2 adjusts this before each get_block() using an adaptive bytes-per-row estimate.
202
    // Store it here as well as forwarding to the current reader so newly opened split readers start
203
    // with the latest predicted batch size.
204
1.04k
    virtual void set_batch_size(size_t batch_size) {
205
1.04k
        _batch_size = std::max<size_t>(1, batch_size);
206
1.04k
        if (_data_reader.reader != nullptr) {
207
469
            _data_reader.reader->set_batch_size(_batch_size);
208
469
        }
209
1.04k
    }
210
211
#ifdef BE_TEST
212
    size_t TEST_batch_size() const { return _batch_size; }
213
    void TEST_set_condition_cache_hit_count(int64_t hits) { _condition_cache_hit_count = hits; }
214
    bool TEST_current_data_file_is_immutable() const {
215
        DORIS_CHECK(_current_task != nullptr);
216
        DORIS_CHECK(_current_task->data_file != nullptr);
217
        DORIS_CHECK(_current_file_description.has_value());
218
        DORIS_CHECK(_current_task->data_file->is_immutable ==
219
                    _current_file_description->is_immutable);
220
        return _current_task->data_file->is_immutable;
221
    }
222
#endif
223
224
    // Prepare for reading a new split/task.
225
    // 1. Pass a new split/task to reader, which will be used in subsequent open_reader() to initialize the underlying file reader.
226
    // 2. Parse delete predicates from split/task information, which will be used for later dynamic filtering and delete handling.
227
    virtual Status prepare_split(const SplitReadOptions& options);
228
229
    // Refresh row-level predicates for an already prepared split. Physical readers that support
230
    // this operation decide the safe boundary at which the new immutable request becomes active.
231
    virtual Status refresh_conjuncts(VExprContextSPtrs conjuncts);
232
233
388
    virtual bool current_split_pruned() const { return _current_split_pruned; }
234
1.24k
    virtual bool current_split_uses_metadata_count() const {
235
1.24k
        return _current_split_uses_metadata_count;
236
1.24k
    }
237
238
    // Discard the active split after the caller decides an error is ignorable, for example a
239
    // stale external-table file listing that returns NOT_FOUND. The next prepare_split() must start
240
    // with no concrete reader or split-local state left from the failed split.
241
1
    virtual Status abort_split() {
242
        // Ignored open failures still spend time closing partially initialized readers. Include
243
        // that recovery path in the common lifecycle profile so NOT_FOUND cannot become invisible.
244
1
        SCOPED_TIMER(_profile.total_timer);
245
1
        SCOPED_TIMER(_profile.close_timer);
246
1
        if (_data_reader.reader != nullptr) {
247
1
            RETURN_IF_ERROR(close_current_reader());
248
1
        } else {
249
0
            _current_task.reset();
250
0
            _current_file_description.reset();
251
0
        }
252
1
        _delete_rows = nullptr;
253
1
        _remaining_table_level_count = -1;
254
1
        _remaining_file_level_count = -1;
255
1
        _current_split_uses_metadata_count = false;
256
1
        _current_split_pruned = false;
257
1
        return Status::OK();
258
1
    }
259
260
    // Public entry point for reading a table-schema block. The base class opens the current reader,
261
    // advances across EOF, and closes exhausted readers. Subclasses provide protected hooks for
262
    // table-format-specific behavior.
263
985
    virtual Status get_block(Block* block, bool* eos) {
264
985
        SCOPED_TIMER(_profile.total_timer);
265
985
        SCOPED_TIMER(_profile.exec_timer);
266
985
        DORIS_CHECK(block->columns() == _projected_columns.size());
267
985
        block->clear_column_data(_projected_columns.size());
268
269
1.33k
        while (true) {
270
1.33k
            if (*eos) {
271
0
                return Status::OK();
272
0
            }
273
1.33k
            if (_io_ctx != nullptr && _io_ctx->should_stop) {
274
0
                *eos = true;
275
0
                return Status::OK();
276
0
            }
277
1.33k
            if (!_data_reader.reader) {
278
803
                if (_is_table_level_count_active()) {
279
6
                    RETURN_IF_ERROR(_read_table_level_count(block, eos));
280
6
                    return Status::OK();
281
6
                }
282
797
                if (_is_file_level_count_active()) {
283
4
                    RETURN_IF_ERROR(_read_file_level_count(block, eos));
284
4
                    return Status::OK();
285
4
                }
286
793
                RETURN_IF_ERROR(create_next_reader(eos));
287
790
                if (!_data_reader.reader) {
288
347
                    DCHECK(*eos);
289
347
                    return Status::OK();
290
347
                }
291
790
            }
292
293
            // Materialize a reduced row set for upper aggregate operators when aggregate
294
            // pushdown can be applied. This is not the final aggregate result: COUNT emits
295
            // `count` default rows for the upper COUNT(*), and MIN/MAX emits two rows containing
296
            // file-level min/max values for the upper MIN/MAX.
297
972
            if (!_aggregate_pushdown_tried) {
298
443
                SCOPED_TIMER(_profile.pushdown_agg_timer);
299
443
                bool pushed_down = false;
300
443
                const auto status = _try_materialize_aggregate_pushdown_rows(block, &pushed_down);
301
443
                if (!status.ok()) {
302
1
                    if (_io_ctx != nullptr && _io_ctx->should_stop &&
303
1
                        status.is<ErrorCode::END_OF_FILE>()) {
304
1
                        *eos = true;
305
1
                        return Status::OK();
306
1
                    }
307
0
                    return status;
308
1
                }
309
442
                if (pushed_down) {
310
11
                    return Status::OK();
311
11
                }
312
442
            }
313
314
960
            bool current_eof = false;
315
960
            _data_reader.block_template.clear_column_data(
316
960
                    cast_set<int64_t>(_data_reader.file_block_layout.size()));
317
960
            size_t current_rows = 0;
318
960
            {
319
960
                SCOPED_TIMER(_profile.file_reader_total_timer);
320
960
                SCOPED_TIMER(_profile.file_reader_get_block_timer);
321
960
                RETURN_IF_ERROR(_data_reader.reader->get_block(&_data_reader.block_template,
322
960
                                                               &current_rows, &current_eof));
323
960
            }
324
960
            const bool stopped_during_read = _io_ctx != nullptr && _io_ctx->should_stop;
325
960
            if (current_rows == 0) {
326
347
                if (current_eof) {
327
345
                    _current_reader_reached_eof = !stopped_during_read;
328
345
                    RETURN_IF_ERROR(close_current_reader());
329
345
                }
330
347
                continue;
331
347
            }
332
960
            DCHECK_EQ(_data_reader.block_template.columns(), _data_reader.file_block_layout.size())
333
0
                    << _data_reader.block_template.dump_structure();
334
613
#ifndef NDEBUG
335
613
            RETURN_IF_ERROR(_check_file_block_columns("after file reader get_block", current_rows));
336
613
#endif
337
613
            DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
338
613
            RETURN_IF_ERROR(finalize_chunk(block, current_rows));
339
610
#ifndef NDEBUG
340
610
            RETURN_IF_ERROR(
341
610
                    _check_table_block_columns("after finalize_chunk", block, current_rows));
342
610
#endif
343
610
            if (current_eof) {
344
21
                _current_reader_reached_eof = !stopped_during_read;
345
21
                RETURN_IF_ERROR(close_current_reader());
346
21
            }
347
610
            return Status::OK();
348
610
        }
349
985
    }
350
351
    // Close the table reader and the currently active file reader. Subclasses that hold additional
352
    // table-format resources should override this and call TableReader::close() first.
353
447
    virtual Status close() {
354
447
        SCOPED_TIMER(_profile.total_timer);
355
447
        SCOPED_TIMER(_profile.close_timer);
356
447
        if (_data_reader.reader) {
357
67
            RETURN_IF_ERROR(close_current_reader());
358
67
        }
359
447
        _current_task.reset();
360
447
        _current_file_description.reset();
361
447
        _remaining_table_level_count = -1;
362
447
        _remaining_file_level_count = -1;
363
447
        _current_split_uses_metadata_count = false;
364
447
        return Status::OK();
365
447
    }
366
367
568
    virtual int64_t condition_cache_hit_count() const { return _condition_cache_hit_count; }
368
369
    virtual std::string debug_string() const;
370
371
    virtual Status annotate_projected_column(const TFileScanSlotInfo& slot_info,
372
                                             ProjectedColumnBuildContext* context,
373
                                             ColumnDefinition* column) const;
374
375
280
    virtual Status validate_projected_columns(const ProjectedColumnBuildContext& context) const {
376
280
        (void)context;
377
280
        return Status::OK();
378
280
    }
379
380
protected:
381
    // TableReader keeps the active file description both in the scan task and separately for
382
    // creating the physical reader. Table-format readers must update both copies when their
383
    // snapshot protocol guarantees that a file path is never overwritten with different bytes.
384
    // This guarantee lets readers safely build cache keys without mtime; it must not be used for
385
    // ordinary Hive/TVF files whose paths may be overwritten in place.
386
97
    void mark_current_data_file_immutable() {
387
97
        DORIS_CHECK(_current_task != nullptr);
388
97
        DORIS_CHECK(_current_task->data_file != nullptr);
389
97
        DORIS_CHECK(_current_file_description.has_value());
390
97
        _current_task->data_file->is_immutable = true;
391
97
        _current_file_description->is_immutable = true;
392
97
    }
393
394
    std::optional<ColumnDefinition> _find_table_column_by_field_id(
395
            int32_t field_id, DataTypePtr type, bool include_historical_schemas) const;
396
    std::optional<std::vector<ColumnDefinition>> _find_table_column_path_by_field_id(
397
            int32_t field_id, DataTypePtr leaf_type, bool include_historical_schemas) const;
398
    std::optional<std::vector<ColumnDefinition>> _find_table_column_identity_path_by_field_id(
399
            int32_t field_id, bool include_historical_schemas) const;
400
401
    // Parse deletion vector information from table format specific file description.
402
    virtual Status _parse_deletion_vector_file(const TTableFormatFileDesc& t_desc,
403
386
                                               DeleteFileDesc* desc, bool* has_delete_file) {
404
386
        *has_delete_file = false;
405
386
        return Status::OK();
406
386
    }
407
408
    // Advance to the next reader. This closes the current reader first and then opens the next
409
    // concrete reader. Subclasses should not duplicate this loop.
410
    Status create_next_reader(bool* eos);
411
    virtual Status create_file_reader(std::unique_ptr<FileReader>* reader);
412
365
    virtual TableColumnMappingMode mapping_mode() const { return TableColumnMappingMode::BY_NAME; }
413
369
    virtual void configure_mapper_options(TableColumnMapperOptions*) const {}
414
442
    virtual Status annotate_file_schema(std::vector<ColumnDefinition>* file_schema) {
415
442
        DORIS_CHECK(file_schema != nullptr);
416
442
        return Status::OK();
417
442
    }
418
371
    virtual Status validate_file_mapping(const TableColumnMapper&) const { return Status::OK(); }
419
420
    // Open the concrete reader for the current split/task and build the file-local scan request.
421
446
    virtual Status open_reader() {
422
446
        SCOPED_TIMER(_profile.open_reader_timer);
423
        // 1. Get file schema and create column mapping.
424
446
        std::vector<ColumnDefinition> file_schema;
425
446
        {
426
446
            SCOPED_TIMER(_profile.file_reader_total_timer);
427
446
            SCOPED_TIMER(_profile.file_reader_schema_timer);
428
446
            RETURN_IF_ERROR(_data_reader.reader->get_schema(&file_schema));
429
446
        }
430
        // For Paimon/Hudi, FE can provide field ids through `history_schema_info`. Annotate the
431
        // file schema before column mapping when the table format maps columns by field id.
432
446
        RETURN_IF_ERROR(annotate_file_schema(&file_schema));
433
446
        _data_reader.file_schema = file_schema;
434
446
        _mapper_options.mode = mapping_mode();
435
446
        configure_mapper_options(&_mapper_options);
436
437
446
        {
438
446
            SCOPED_TIMER(_profile.file_reader_total_timer);
439
446
            SCOPED_TIMER(_profile.file_reader_mapper_timer);
440
446
            _data_reader.column_mapper = _data_reader.reader->create_column_mapper(_mapper_options);
441
446
        }
442
446
        DORIS_CHECK(_data_reader.column_mapper != nullptr);
443
446
        RETURN_IF_ERROR(_data_reader.column_mapper->create_mapping(_projected_columns,
444
446
                                                                   _partition_values, file_schema));
445
446
        DORIS_CHECK(_data_reader.column_mapper->mappings().size() == _projected_columns.size());
446
447
        // 2. Build table filters based on conjuncts and column predicates.
448
446
        RETURN_IF_ERROR(_build_table_filters_from_conjuncts());
449
450
        // 3. Create file scan request based on column mapping and table filters, then open file
451
        // reader with the request. File scan request carries row-level expression filters and
452
        // file-level pruning hints. Only expression filters decide returned rows.
453
446
        auto file_request = std::make_shared<FileScanRequest>();
454
446
        RETURN_IF_ERROR(_data_reader.column_mapper->create_scan_request(
455
446
                _table_filters, _projected_columns, file_request.get(), _runtime_state));
456
446
        _constant_pruning_safe_filter_count =
457
446
                std::min(_constant_pruning_safe_filter_count,
458
446
                         file_request->constant_pruning_safe_table_filter_count);
459
446
        bool constant_filter_pruned_split = false;
460
446
        RETURN_IF_ERROR(_evaluate_constant_filters(&constant_filter_pruned_split));
461
446
        if (constant_filter_pruned_split) {
462
1
            RETURN_IF_ERROR(close_current_reader());
463
1
            return Status::OK();
464
1
        }
465
445
        RETURN_IF_ERROR(validate_file_mapping(*_data_reader.column_mapper));
466
        // COUNT(*) has no semantic column argument, but Nereids retains a minimum-width scan slot
467
        // so the scan node still has an output tuple. Record only the current non-predicate file
468
        // columns before table-format hooks add row-position or equality-delete dependencies. This
469
        // marker is independent of aggregate eligibility: with position deletes, for example,
470
        // metadata COUNT must fall back to reading rows, but an arbitrary unsupported TIME_MILLIS
471
        // placeholder still must not be validated or decoded merely to carry the surviving count.
472
        // Pending runtime filters may later target this retained slot, so placeholder values are
473
        // safe only after every filter for the split has arrived.
474
444
        if (_push_down_agg_type == TPushAggOp::type::COUNT &&
475
444
            _push_down_count_columns.has_value() && _push_down_count_columns->empty() &&
476
444
            _all_runtime_filters_applied_for_split) {
477
11
            file_request->count_star_placeholder_columns.reserve(
478
11
                    file_request->non_predicate_columns.size());
479
11
            for (const auto& column : file_request->non_predicate_columns) {
480
8
                file_request->count_star_placeholder_columns.push_back(column.column_id());
481
8
            }
482
11
        }
483
444
        RETURN_IF_ERROR(customize_file_scan_request(file_request.get()));
484
443
        RETURN_IF_ERROR(_open_local_filter_exprs(*file_request));
485
443
        _data_reader.file_block_layout.clear();
486
443
        _data_reader.block_template.clear();
487
443
        _file_scan_request.reset();
488
443
        _data_reader.file_block_layout.resize(file_request->block_column_count());
489
490
        // 4. Build file block layout from file schema and column mapping. The layout describes
491
        // the block returned by file reader before table-column materialization.
492
443
        auto add_file_block_column = [&](const LocalColumnIndex& projection,
493
2.37k
                                         LocalIndex block_position) -> Status {
494
2.37k
            DORIS_CHECK(block_position.value() < _data_reader.file_block_layout.size());
495
2.37k
            const auto file_column_id = projection.column_id();
496
2.37k
            const auto* field = _find_column_definition(_data_reader.file_schema, file_column_id);
497
2.37k
            DORIS_CHECK(field != nullptr);
498
499
2.37k
            ColumnDefinition projected_field;
500
2.37k
            RETURN_IF_ERROR(project_column_definition(*field, projection, &projected_field));
501
2.37k
            _data_reader.file_block_layout[block_position.value()] = {
502
2.37k
                    .file_column_id = file_column_id,
503
2.37k
                    .name = projected_field.name,
504
2.37k
                    .type = projected_field.type,
505
2.37k
            };
506
2.37k
            DORIS_CHECK(_data_reader.file_block_layout[block_position.value()].type != nullptr);
507
2.37k
            return Status::OK();
508
2.37k
        };
509
443
        for (const auto& projection : file_request->predicate_columns) {
510
93
            RETURN_IF_ERROR(add_file_block_column(
511
93
                    projection, file_request->local_positions.at(projection.column_id())));
512
93
        }
513
2.27k
        for (const auto& projection : file_request->non_predicate_columns) {
514
2.27k
            RETURN_IF_ERROR(add_file_block_column(
515
2.27k
                    projection, file_request->non_predicate_position(projection.column_id())));
516
2.27k
        }
517
518
        // 5. Prepare block template from file block layout. The block template stores the block
519
        // returned by file reader before table-column materialization.
520
443
        _data_reader.block_template.reserve(_data_reader.file_block_layout.size());
521
2.37k
        for (const auto& column : _data_reader.file_block_layout) {
522
2.37k
            _data_reader.block_template.insert(
523
2.37k
                    {column.type->create_column(), column.type, column.name});
524
2.37k
        }
525
443
        if (VLOG_DEBUG_IS_ON) {
526
0
            VLOG_DEBUG << "TableReader debug: " << debug_string();
527
0
        }
528
443
        RETURN_IF_ERROR(_open_mapping_exprs());
529
443
        {
530
443
            SCOPED_TIMER(_profile.file_reader_total_timer);
531
443
            SCOPED_TIMER(_profile.file_reader_open_timer);
532
443
            RETURN_IF_ERROR(_data_reader.reader->open(file_request));
533
443
        }
534
443
        _file_scan_request = std::move(file_request);
535
443
        RETURN_IF_ERROR(_init_reader_condition_cache(*_file_scan_request));
536
443
        return Status::OK();
537
443
    }
538
539
    Status _build_table_filters_from_conjuncts();
540
    Status _evaluate_partition_prune_conjuncts(const VExprContextSPtrs& conjuncts,
541
                                               bool* can_filter_all);
542
    static bool _is_safe_to_pre_execute(const VExprContextSPtr& conjunct);
543
    Status _build_partition_prune_block(Block* block) const;
544
    Status _open_local_filter_exprs(const FileScanRequest& file_request);
545
    Status _init_reader_condition_cache(const FileScanRequest& file_request);
546
    void _finalize_reader_condition_cache();
547
    bool _should_enable_condition_cache(const FileScanRequest& file_request) const;
548
549
446
    Status _evaluate_constant_filters(bool* can_filter_all) {
550
446
        DORIS_CHECK(can_filter_all != nullptr);
551
446
        DORIS_CHECK_LE(_constant_pruning_safe_filter_count, _table_filters.size());
552
446
        *can_filter_all = false;
553
        // The bound was derived from the original `_conjuncts` order, which includes slotless
554
        // expressions omitted from `_table_filters`. Iterating only this prefix therefore cannot
555
        // skip an unsafe row-level predicate and pre-execute a later constant predicate.
556
483
        for (size_t i = 0; i < _constant_pruning_safe_filter_count; ++i) {
557
38
            const auto& table_filter = _table_filters[i];
558
38
            if (table_filter.conjunct == nullptr) {
559
0
                continue;
560
0
            }
561
38
            DORIS_CHECK(_is_safe_to_pre_execute(table_filter.conjunct));
562
            // RuntimeFilterExpr does not implement execute_column_impl(); it is evaluated by the
563
            // row-level filter path through execute_filter(). Constant split pruning uses
564
            // VExprContext::execute() on a one-row synthetic block, so runtime filters must not be
565
            // pre-executed here even when their referenced slot maps to a constant value.
566
38
            if (table_filter.conjunct->root()->is_rf_wrapper() ||
567
38
                !_table_filter_has_only_constant_entries(table_filter)) {
568
36
                continue;
569
36
            }
570
2
            Block eval_block;
571
2
            RETURN_IF_ERROR(_build_constant_filter_block(table_filter, &eval_block));
572
2
            RowDescriptor row_desc;
573
2
            RETURN_IF_ERROR(table_filter.conjunct->prepare(_runtime_state, row_desc));
574
2
            RETURN_IF_ERROR(table_filter.conjunct->open(_runtime_state));
575
2
            int result_column_id = -1;
576
2
            RETURN_IF_ERROR(table_filter.conjunct->execute(&eval_block, &result_column_id));
577
2
            DORIS_CHECK(result_column_id >= 0);
578
2
            if (_filter_result_filters_all(eval_block.get_by_position(result_column_id).column)) {
579
1
                *can_filter_all = true;
580
1
                return Status::OK();
581
1
            }
582
2
        }
583
445
        return Status::OK();
584
446
    }
585
586
33
    bool _table_filter_has_only_constant_entries(const TableFilter& table_filter) const {
587
33
        const auto& filter_entries = _data_reader.column_mapper->filter_entries();
588
33
        for (const auto global_index : table_filter.global_indices) {
589
33
            const auto entry_it = filter_entries.find(global_index);
590
33
            if (entry_it == filter_entries.end() || !entry_it->second.is_constant()) {
591
31
                return false;
592
31
            }
593
33
        }
594
2
        return !table_filter.global_indices.empty();
595
33
    }
596
597
2
    Status _build_constant_filter_block(const TableFilter& table_filter, Block* eval_block) {
598
2
        DORIS_CHECK(eval_block != nullptr);
599
2
        eval_block->clear();
600
2
        const auto& mappings = _data_reader.column_mapper->mappings();
601
2
        const auto& filter_entries = _data_reader.column_mapper->filter_entries();
602
2
        DORIS_CHECK(mappings.size() == _projected_columns.size());
603
4
        for (size_t column_idx = 0; column_idx < mappings.size(); ++column_idx) {
604
2
            const auto global_index = GlobalIndex(column_idx);
605
2
            const auto& mapping = mappings[column_idx];
606
2
            const auto entry_it = filter_entries.find(global_index);
607
2
            const bool referenced_by_filter =
608
2
                    std::find(table_filter.global_indices.begin(),
609
2
                              table_filter.global_indices.end(),
610
2
                              global_index) != table_filter.global_indices.end();
611
2
            if (referenced_by_filter && entry_it != filter_entries.end() &&
612
2
                entry_it->second.is_constant()) {
613
2
                ColumnPtr constant_column;
614
2
                RETURN_IF_ERROR(_materialize_constant_filter_column(
615
2
                        entry_it->second.constant_index(), &constant_column));
616
2
                eval_block->insert({std::move(constant_column), mapping.table_type,
617
2
                                    mapping.table_column_name});
618
2
            } else {
619
0
                eval_block->insert({mapping.table_type->create_column_const_with_default_value(1),
620
0
                                    mapping.table_type, mapping.table_column_name});
621
0
            }
622
2
        }
623
2
        return Status::OK();
624
2
    }
625
626
2
    Status _materialize_constant_filter_column(ConstantIndex constant_index, ColumnPtr* column) {
627
2
        DORIS_CHECK(column != nullptr);
628
2
        const auto& constant_entry = _data_reader.column_mapper->constant_map().get(constant_index);
629
2
        DORIS_CHECK(constant_entry.expr != nullptr);
630
2
        DORIS_CHECK(constant_entry.type != nullptr);
631
2
        RowDescriptor row_desc;
632
2
        RETURN_IF_ERROR(constant_entry.expr->prepare(_runtime_state, row_desc));
633
2
        RETURN_IF_ERROR(constant_entry.expr->open(_runtime_state));
634
2
        Block eval_block;
635
2
        eval_block.insert({constant_entry.type->create_column_const_with_default_value(1),
636
2
                           constant_entry.type, "__table_reader_constant_filter"});
637
2
        int result_column_id = -1;
638
2
        RETURN_IF_ERROR(constant_entry.expr->execute(&eval_block, &result_column_id));
639
2
        DORIS_CHECK(result_column_id >= 0);
640
2
        *column = eval_block.get_by_position(result_column_id).column;
641
2
        DORIS_CHECK((*column)->size() == 1);
642
2
        return Status::OK();
643
2
    }
644
645
2
    static bool _filter_result_filters_all(const ColumnPtr& filter_column) {
646
2
        DORIS_CHECK(filter_column.get() != nullptr);
647
2
        DORIS_CHECK(filter_column->size() == 1);
648
2
        return !filter_column->get_bool(0);
649
2
    }
650
651
452
    virtual Status customize_file_scan_request(FileScanRequest* file_request) {
652
452
        return _append_delete_predicate(file_request);
653
452
    }
654
655
1.37k
    bool _is_table_level_count_active() const { return _remaining_table_level_count >= 0; }
656
657
797
    bool _is_file_level_count_active() const { return _remaining_file_level_count >= 0; }
658
659
14
    Status _materialize_count_rows(size_t rows, Block* block) const {
660
14
        DORIS_CHECK(block != nullptr);
661
14
        DORIS_CHECK(block->columns() > 0 || rows == 0);
662
28
        for (size_t column_idx = 0; column_idx < block->columns(); ++column_idx) {
663
14
            auto column = block->get_by_position(column_idx).type->create_column();
664
14
            if (auto* nullable = check_and_get_column<ColumnNullable>(*column)) {
665
                // Metadata COUNT emits synthetic input rows for the unchanged upper aggregate.
666
                // They must be non-NULL for COUNT(nullable_col), and constructing them explicitly
667
                // also keeps every nullable null map boolean-valid in debug/ASAN block checks.
668
14
                nullable->get_nested_column().insert_many_defaults(rows);
669
14
                nullable->get_null_map_data().resize_fill(rows, 0);
670
14
            } else {
671
0
                column->insert_many_defaults(rows);
672
0
            }
673
14
            block->replace_by_position(column_idx, std::move(column));
674
14
        }
675
14
        return Status::OK();
676
14
    }
677
678
14
    Status _materialize_next_count_batch(int64_t* remaining_rows, Block* block) const {
679
14
        DORIS_CHECK(remaining_rows != nullptr);
680
14
        DORIS_CHECK(*remaining_rows > 0);
681
14
        const int64_t batch_size = _runtime_state == nullptr
682
14
                                           ? *remaining_rows
683
14
                                           : static_cast<int64_t>(_runtime_state->batch_size());
684
14
        const auto rows = std::min(*remaining_rows, batch_size);
685
14
        RETURN_IF_ERROR(_materialize_count_rows(cast_set<size_t>(rows), block));
686
14
        *remaining_rows -= rows;
687
14
        return Status::OK();
688
14
    }
689
690
10
    Status _read_count_batch(int64_t* remaining_rows, Block* block, bool* eos) {
691
10
        DORIS_CHECK(block != nullptr);
692
10
        DORIS_CHECK(eos != nullptr);
693
10
        DORIS_CHECK(_push_down_agg_type == TPushAggOp::type::COUNT);
694
10
        DORIS_CHECK(remaining_rows != nullptr);
695
10
        DORIS_CHECK(*remaining_rows >= 0);
696
10
        if (*remaining_rows == 0) {
697
4
            *remaining_rows = -1;
698
4
            _current_task.reset();
699
4
            *eos = true;
700
4
            return Status::OK();
701
4
        }
702
6
        RETURN_IF_ERROR(_materialize_next_count_batch(remaining_rows, block));
703
6
        *eos = false;
704
6
        return Status::OK();
705
6
    }
706
707
6
    Status _read_table_level_count(Block* block, bool* eos) {
708
6
        return _read_count_batch(&_remaining_table_level_count, block, eos);
709
6
    }
710
711
4
    Status _read_file_level_count(Block* block, bool* eos) {
712
4
        return _read_count_batch(&_remaining_file_level_count, block, eos);
713
4
    }
714
715
    void _append_file_scan_column(FileScanRequest* request, LocalColumnId column_id,
716
70
                                  std::vector<LocalColumnIndex>* scan_columns) {
717
70
        DORIS_CHECK(request != nullptr);
718
70
        DORIS_CHECK(scan_columns != nullptr);
719
70
        FileScanRequestBuilder builder(request);
720
70
        Status status;
721
70
        if (scan_columns == &request->predicate_columns) {
722
59
            status = builder.add_predicate_column(column_id);
723
59
        } else {
724
11
            DORIS_CHECK(scan_columns == &request->non_predicate_columns);
725
11
            status = builder.add_non_predicate_column(column_id);
726
11
        }
727
70
        DORIS_CHECK(status.ok()) << status.to_string();
728
70
        if (column_id == LocalColumnId(ROW_POSITION_COLUMN_ID) &&
729
70
            _find_column_definition(_data_reader.file_schema, column_id) == nullptr) {
730
39
            _data_reader.file_schema.push_back(row_position_column_definition());
731
39
        }
732
70
    }
733
734
    // 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.
735
452
    Status _append_delete_predicate(FileScanRequest* request) {
736
452
        DORIS_CHECK(request != nullptr);
737
452
        if ((_delete_rows == nullptr || _delete_rows->empty()) &&
738
452
            (_deletion_vector == nullptr || _deletion_vector->isEmpty())) {
739
436
            return Status::OK();
740
436
        }
741
16
        const auto row_position_column_id = LocalColumnId(ROW_POSITION_COLUMN_ID);
742
16
        _append_file_scan_column(request, row_position_column_id, &request->predicate_columns);
743
744
16
        const auto block_position = request->local_positions.at(row_position_column_id);
745
16
        auto append_predicate = [&](auto& deleted_rows) {
746
16
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
747
16
            delete_predicate->add_child(VSlotRef::create_shared(
748
16
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
749
16
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
750
16
            request->delete_conjuncts.push_back(
751
16
                    VExprContext::create_shared(std::move(delete_predicate)));
752
16
        };
_ZZN5doris6format11TableReader24_append_delete_predicateEPNS0_15FileScanRequestEENKUlRT_E_clISt6vectorIlSaIlEEEEDaS5_
Line
Count
Source
745
12
        auto append_predicate = [&](auto& deleted_rows) {
746
12
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
747
12
            delete_predicate->add_child(VSlotRef::create_shared(
748
12
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
749
12
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
750
12
            request->delete_conjuncts.push_back(
751
12
                    VExprContext::create_shared(std::move(delete_predicate)));
752
12
        };
_ZZN5doris6format11TableReader24_append_delete_predicateEPNS0_15FileScanRequestEENKUlRT_E_clIN7roaring12Roaring64MapEEEDaS5_
Line
Count
Source
745
4
        auto append_predicate = [&](auto& deleted_rows) {
746
4
            auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows);
747
4
            delete_predicate->add_child(VSlotRef::create_shared(
748
4
                    cast_set<int>(block_position.value()), cast_set<int>(block_position.value()),
749
4
                    -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME));
750
4
            request->delete_conjuncts.push_back(
751
4
                    VExprContext::create_shared(std::move(delete_predicate)));
752
4
        };
753
16
        if (_delete_rows != nullptr && !_delete_rows->empty()) {
754
12
            append_predicate(*_delete_rows);
755
12
        }
756
16
        if (_deletion_vector != nullptr && !_deletion_vector->isEmpty()) {
757
4
            append_predicate(*_deletion_vector);
758
4
        }
759
16
        return Status::OK();
760
452
    }
761
762
    // Close the current concrete reader. This hook is called by both create_next_reader() and
763
    // close(), so it should remain idempotent.
764
446
    virtual Status close_current_reader() {
765
446
        _finalize_reader_condition_cache();
766
446
        {
767
446
            SCOPED_TIMER(_profile.file_reader_total_timer);
768
446
            SCOPED_TIMER(_profile.file_reader_close_timer);
769
446
            RETURN_IF_ERROR(_data_reader.reader->close());
770
446
        }
771
446
        _data_reader.reader.reset();
772
446
        if (_data_reader.column_mapper != nullptr) {
773
445
            _data_reader.column_mapper->clear();
774
445
            _data_reader.column_mapper.reset();
775
445
        }
776
446
        _table_filters.clear();
777
446
        _constant_pruning_safe_filter_count = 0;
778
446
        _data_reader.file_schema.clear();
779
446
        _data_reader.file_block_layout.clear();
780
446
        _data_reader.block_template.clear();
781
446
        _file_scan_request.reset();
782
446
        _current_task.reset();
783
446
        _current_file_description.reset();
784
446
        _current_reader_reached_eof = false;
785
446
        return Status::OK();
786
446
    }
787
788
2
    void _record_scan_rows(size_t rows) {
789
2
        if (_io_ctx != nullptr && _io_ctx->file_reader_stats != nullptr) {
790
2
            _io_ctx->file_reader_stats->read_rows += rows;
791
2
        }
792
2
    }
793
794
    // Finalize file-local block to table/global schema block.
795
613
    Status finalize_chunk(Block* block, const size_t rows) {
796
613
        SCOPED_TIMER(_profile.finalize_timer);
797
613
        size_t idx = 0;
798
613
        const auto& mappings = _data_reader.column_mapper->mappings();
799
5.28k
        for (const auto& mapping : mappings) {
800
5.28k
            ColumnPtr column;
801
5.28k
            RETURN_IF_ERROR(_materialize_mapping_column(mapping, &_data_reader.block_template, rows,
802
5.28k
                                                        &column, idx + 1 == mappings.size()));
803
5.28k
            block->replace_by_position(idx, IColumn::mutate(std::move(column)));
804
5.28k
            idx++;
805
5.28k
        }
806
610
        RETURN_IF_ERROR(materialize_virtual_columns(block));
807
        // Enforce CHAR/VARCHAR length declared by the table schema after all file-to-table
808
        // materialization has finished.
809
610
        RETURN_IF_ERROR(_truncate_char_or_varchar_columns(block));
810
610
        return Status::OK();
811
610
    }
812
813
    // Materialize virtual columns in the table block, such as Iceberg _row_id and
814
    // _last_updated_sequence_number. This runs after normal column materialization so finalize
815
    // expressions can reference those virtual columns.
816
549
    virtual Status materialize_virtual_columns(Block* table_block) { return Status::OK(); }
817
818
#ifndef NDEBUG
819
613
    Status _check_file_block_columns(std::string_view stage, size_t rows) {
820
613
        DORIS_CHECK(_data_reader.block_template.columns() == _data_reader.file_block_layout.size());
821
5.91k
        for (size_t idx = 0; idx < _data_reader.block_template.columns(); ++idx) {
822
5.30k
            const auto& file_block_column = _data_reader.file_block_layout[idx];
823
5.30k
            const auto& column_with_type = _data_reader.block_template.get_by_position(idx);
824
5.30k
            const auto* column = column_with_type.column.get();
825
5.30k
            try {
826
5.30k
                if (column == nullptr) {
827
0
                    auto st = Status::InternalError(
828
0
                            "Invalid file block column {} at {}: file_column_id={}, name='{}', "
829
0
                            "type={}, column=null, expected_rows={}, reader={}",
830
0
                            idx, stage, file_block_column.file_column_id.value(),
831
0
                            file_block_column.name,
832
0
                            file_block_column.type == nullptr ? "null"
833
0
                                                              : file_block_column.type->get_name(),
834
0
                            rows, debug_string());
835
0
                    LOG(WARNING) << st;
836
0
                    return st;
837
0
                }
838
5.30k
                column->sanity_check();
839
5.30k
                auto st = column_with_type.check_type_and_column_match();
840
5.30k
                if (!st.ok()) {
841
0
                    auto contextual_status = Status::InternalError(
842
0
                            "Invalid file block column {} at {}: file_column_id={}, name='{}', "
843
0
                            "type={}, column={}, column_size={}, expected_rows={}, error={}, "
844
0
                            "reader={}",
845
0
                            idx, stage, file_block_column.file_column_id.value(),
846
0
                            file_block_column.name,
847
0
                            file_block_column.type == nullptr ? "null"
848
0
                                                              : file_block_column.type->get_name(),
849
0
                            column->get_name(), column->size(), rows, st.to_string(),
850
0
                            debug_string());
851
0
                    LOG(WARNING) << contextual_status;
852
0
                    return contextual_status;
853
0
                }
854
5.30k
            } catch (const Exception& e) {
855
0
                auto st = Status::InternalError(
856
0
                        "Invalid file block column {} at {}: file_column_id={}, name='{}', "
857
0
                        "type={}, column={}, column_size={}, expected_rows={}, error={}, "
858
0
                        "reader={}",
859
0
                        idx, stage, file_block_column.file_column_id.value(),
860
0
                        file_block_column.name,
861
0
                        file_block_column.type == nullptr ? "null"
862
0
                                                          : file_block_column.type->get_name(),
863
0
                        column == nullptr ? "null" : column->get_name(),
864
0
                        column == nullptr ? 0 : column->size(), rows, e.to_string(),
865
0
                        debug_string());
866
0
                LOG(WARNING) << st;
867
0
                return st;
868
0
            } catch (const std::exception& e) {
869
0
                auto st = Status::InternalError(
870
0
                        "Invalid file block column {} at {}: file_column_id={}, name='{}', "
871
0
                        "type={}, column={}, column_size={}, expected_rows={}, error={}, "
872
0
                        "reader={}",
873
0
                        idx, stage, file_block_column.file_column_id.value(),
874
0
                        file_block_column.name,
875
0
                        file_block_column.type == nullptr ? "null"
876
0
                                                          : file_block_column.type->get_name(),
877
0
                        column == nullptr ? "null" : column->get_name(),
878
0
                        column == nullptr ? 0 : column->size(), rows, e.what(), debug_string());
879
0
                LOG(WARNING) << st;
880
0
                return st;
881
0
            }
882
5.30k
        }
883
613
        return Status::OK();
884
613
    }
885
886
610
    Status _check_table_block_columns(std::string_view stage, const Block* block, size_t rows) {
887
610
        DORIS_CHECK(block != nullptr);
888
610
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
889
5.89k
        for (size_t idx = 0; idx < block->columns(); ++idx) {
890
5.28k
            const auto& mapping = _data_reader.column_mapper->mappings()[idx];
891
5.28k
            const auto& column_with_type = block->get_by_position(idx);
892
5.28k
            const auto* column = column_with_type.column.get();
893
5.28k
            try {
894
5.28k
                if (column == nullptr) {
895
0
                    auto st = Status::InternalError(
896
0
                            "Invalid table block column {} at {}: table_column='{}', "
897
0
                            "global_index={}, type={}, column=null, expected_rows={}, mapping={}",
898
0
                            idx, stage, mapping.table_column_name, mapping.global_index.value(),
899
0
                            mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(),
900
0
                            rows, mapping.debug_string());
901
0
                    LOG(WARNING) << st;
902
0
                    return st;
903
0
                }
904
5.28k
                column->sanity_check();
905
5.28k
                auto st = column_with_type.check_type_and_column_match();
906
5.28k
                if (!st.ok()) {
907
0
                    auto contextual_status = Status::InternalError(
908
0
                            "Invalid table block column {} at {}: table_column='{}', "
909
0
                            "global_index={}, type={}, column={}, column_size={}, "
910
0
                            "expected_rows={}, error={}, mapping={}",
911
0
                            idx, stage, mapping.table_column_name, mapping.global_index.value(),
912
0
                            mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(),
913
0
                            column->get_name(), column->size(), rows, st.to_string(),
914
0
                            mapping.debug_string());
915
0
                    LOG(WARNING) << contextual_status;
916
0
                    return contextual_status;
917
0
                }
918
5.28k
            } catch (const Exception& e) {
919
0
                auto st = Status::InternalError(
920
0
                        "Invalid table block column {} at {}: table_column='{}', global_index={}, "
921
0
                        "type={}, column={}, column_size={}, expected_rows={}, error={}, "
922
0
                        "mapping={}",
923
0
                        idx, stage, mapping.table_column_name, mapping.global_index.value(),
924
0
                        mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(),
925
0
                        column == nullptr ? "null" : column->get_name(),
926
0
                        column == nullptr ? 0 : column->size(), rows, e.to_string(),
927
0
                        mapping.debug_string());
928
0
                LOG(WARNING) << st;
929
0
                return st;
930
0
            } catch (const std::exception& e) {
931
0
                auto st = Status::InternalError(
932
0
                        "Invalid table block column {} at {}: table_column='{}', global_index={}, "
933
0
                        "type={}, column={}, column_size={}, expected_rows={}, error={}, "
934
0
                        "mapping={}",
935
0
                        idx, stage, mapping.table_column_name, mapping.global_index.value(),
936
0
                        mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(),
937
0
                        column == nullptr ? "null" : column->get_name(),
938
0
                        column == nullptr ? 0 : column->size(), rows, e.what(),
939
0
                        mapping.debug_string());
940
0
                LOG(WARNING) << st;
941
0
                return st;
942
0
            }
943
5.28k
        }
944
610
        return Status::OK();
945
610
    }
946
#endif
947
948
610
    Status _truncate_char_or_varchar_columns(Block* block) {
949
610
        DORIS_CHECK(block != nullptr);
950
610
        if (_runtime_state == nullptr ||
951
610
            !_runtime_state->query_options().truncate_char_or_varchar_columns) {
952
610
            return Status::OK();
953
610
        }
954
0
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
955
0
        for (size_t idx = 0; idx < _data_reader.column_mapper->mappings().size(); ++idx) {
956
0
            const auto& mapping = _data_reader.column_mapper->mappings()[idx];
957
0
            if (!_should_truncate_char_or_varchar_column(mapping)) {
958
0
                continue;
959
0
            }
960
0
            const auto target_len =
961
0
                    assert_cast<const DataTypeString*>(remove_nullable(mapping.table_type).get())
962
0
                            ->len();
963
0
            _truncate_char_or_varchar_column(block, idx, target_len);
964
0
        }
965
0
        return Status::OK();
966
610
    }
967
968
    // Return true when the table schema has a bounded CHAR/VARCHAR length that is stricter than
969
    // the file-side type. Examples:
970
    // - table VARCHAR(10), file VARCHAR(20): truncate to 10;
971
    // - table VARCHAR(10), file STRING: truncate to 10 because STRING has no declared bound;
972
    // - table STRING, any file type: no truncation because the target has no bound.
973
5
    static bool _should_truncate_char_or_varchar_column(const ColumnMapping& mapping) {
974
5
        if (mapping.table_type == nullptr) {
975
0
            return false;
976
0
        }
977
5
        const auto table_type = remove_nullable(mapping.table_type);
978
5
        const auto primitive_type = table_type->get_primitive_type();
979
5
        if (primitive_type != TYPE_VARCHAR && primitive_type != TYPE_CHAR) {
980
1
            return false;
981
1
        }
982
4
        const auto target_len = assert_cast<const DataTypeString*>(table_type.get())->len();
983
4
        if (target_len <= 0) {
984
0
            return false;
985
0
        }
986
4
        if (mapping.file_type == nullptr) {
987
0
            return true;
988
0
        }
989
4
        const auto file_type = remove_nullable(mapping.file_type);
990
4
        DORIS_CHECK(file_type != nullptr);
991
4
        int file_len = -1;
992
4
        if (file_type->get_primitive_type() == TYPE_VARCHAR ||
993
4
            file_type->get_primitive_type() == TYPE_CHAR ||
994
4
            file_type->get_primitive_type() == TYPE_STRING) {
995
3
            file_len = assert_cast<const DataTypeString*>(file_type.get())->len();
996
3
        }
997
998
4
        return file_len < 0 || target_len < file_len;
999
4
    }
1000
1001
    // Truncate a materialized CHAR/VARCHAR column in place by reusing the vectorized substring
1002
    // implementation: substring(column, 1, len). Nullable columns are unwrapped before substring
1003
    // execution and wrapped back with the original null map afterward, because substring operates
1004
    // on the nested string payload only.
1005
1
    static void _truncate_char_or_varchar_column(Block* block, size_t idx, int len) {
1006
1
        DORIS_CHECK(block != nullptr);
1007
1
        auto int_type = std::make_shared<DataTypeInt32>();
1008
1
        const auto num_columns_without_result = cast_set<uint32_t>(block->columns());
1009
1
        auto& target = block->get_by_position(idx);
1010
1
        const bool is_nullable = target.type->is_nullable();
1011
1
        ColumnPtr input_column = target.column;
1012
1
        ColumnPtr null_map_column;
1013
1
        if (is_nullable) {
1014
1
            const auto* nullable_column = assert_cast<const ColumnNullable*>(target.column.get());
1015
1
            input_column = nullable_column->get_nested_column_ptr();
1016
1
            null_map_column = nullable_column->get_null_map_column_ptr();
1017
1
        }
1018
1
        block->replace_by_position(idx, std::move(input_column));
1019
1
        block->insert({int_type->create_column_const(block->rows(), to_field<TYPE_INT>(1)),
1020
1
                       int_type, "const 1"});
1021
1
        block->insert({int_type->create_column_const(block->rows(), to_field<TYPE_INT>(len)),
1022
1
                       int_type, "const len"});
1023
1
        block->insert({nullptr, std::make_shared<DataTypeString>(), "result"});
1024
1025
1
        ColumnNumbers temp_arguments(3);
1026
1
        temp_arguments[0] = cast_set<uint32_t>(idx);
1027
1
        temp_arguments[1] = num_columns_without_result;
1028
1
        temp_arguments[2] = num_columns_without_result + 1;
1029
1
        const uint32_t result_column_id = num_columns_without_result + 2;
1030
1
        SubstringUtil::substring_execute(*block, temp_arguments, result_column_id, block->rows());
1031
1032
1
        ColumnPtr result_column = block->get_by_position(result_column_id).column;
1033
1
        if (is_nullable) {
1034
1
            result_column = ColumnNullable::create(std::move(result_column), null_map_column);
1035
1
        }
1036
1
        block->replace_by_position(idx, std::move(result_column));
1037
1
        block->erase_tail(num_columns_without_result);
1038
1
    }
1039
1040
443
    Status _try_materialize_aggregate_pushdown_rows(Block* block, bool* pushed_down) {
1041
443
        DORIS_CHECK(block != nullptr);
1042
443
        DORIS_CHECK(pushed_down != nullptr);
1043
443
        *pushed_down = false;
1044
443
        block->clear_column_data(_projected_columns.size());
1045
443
        _aggregate_pushdown_tried = true;
1046
443
        if (!_supports_aggregate_pushdown(_push_down_agg_type)) {
1047
430
            return Status::OK();
1048
430
        }
1049
1050
13
        FileAggregateRequest file_request;
1051
13
        RETURN_IF_ERROR(_build_file_aggregate_request(_push_down_agg_type, &file_request));
1052
13
        FileAggregateResult file_result;
1053
13
        Status status;
1054
13
        {
1055
13
            SCOPED_TIMER(_profile.file_reader_total_timer);
1056
13
            SCOPED_TIMER(_profile.file_reader_aggregate_timer);
1057
13
            status = _data_reader.reader->get_aggregate_result(file_request, &file_result);
1058
13
        }
1059
13
        if (status.is<ErrorCode::NOT_IMPLEMENTED_ERROR>()) {
1060
1
            return Status::OK();
1061
1
        }
1062
12
        RETURN_IF_ERROR(status);
1063
11
        if (_push_down_agg_type == TPushAggOp::type::COUNT) {
1064
8
            DORIS_CHECK(file_result.count >= 0);
1065
            // The upper aggregate consumes synthetic input rows, but emitting the whole metadata
1066
            // count in one block bypasses the runtime batch contract and can allocate by file size.
1067
            // Keep the remaining cardinality as split state and expose at most one batch per call.
1068
8
            _remaining_file_level_count = file_result.count;
1069
8
            _current_split_uses_metadata_count = true;
1070
8
            if (_remaining_file_level_count > 0) {
1071
8
                RETURN_IF_ERROR(_materialize_next_count_batch(&_remaining_file_level_count, block));
1072
8
            }
1073
8
        } else {
1074
3
            RETURN_IF_ERROR(
1075
3
                    _materialize_aggregate_pushdown_rows(_push_down_agg_type, file_result, block));
1076
3
        }
1077
11
        *pushed_down = true;
1078
11
        RETURN_IF_ERROR(close_current_reader());
1079
11
        return Status::OK();
1080
11
    }
1081
1082
458
    virtual bool _supports_aggregate_pushdown(TPushAggOp::type agg_type) const {
1083
        // Only COUNT and MIN/MAX can be push down.
1084
458
        if (agg_type != TPushAggOp::type::COUNT && agg_type != TPushAggOp::type::MINMAX) {
1085
410
            return false;
1086
410
        }
1087
        // Aggregate pushdown returns reduced synthetic rows and may close the physical reader
1088
        // before the next scheduler turn. If a runtime filter is still pending, those rows could
1089
        // escape before the filter arrives and cannot later be reconstructed from real file rows.
1090
        // This is the same irreversibility constraint as table-level metadata COUNT, and applies
1091
        // to COUNT and MIN/MAX for Parquet/ORC as well as COUNT for text readers.
1092
48
        if (!_all_runtime_filters_applied_for_split) {
1093
3
            return false;
1094
3
        }
1095
        // Scanner owns the original conjunct list and evaluates it after TableReader finalizes
1096
        // rows. Even a slotless conjunct that cannot become a TableFilter must see every source
1097
        // row before an aggregate reduces the stream to synthetic COUNT/MINMAX rows.
1098
45
        if (!_conjuncts.empty()) {
1099
5
            return false;
1100
5
        }
1101
        // Only support aggregate pushdown when there is no delete or filter, so
1102
        // the reduced rows consumed by the upper aggregate remain semantically equivalent to a
1103
        // normal scan.
1104
40
        if ((_delete_rows != nullptr && !_delete_rows->empty()) ||
1105
40
            (_deletion_vector != nullptr && !_deletion_vector->isEmpty())) {
1106
5
            return false;
1107
5
        }
1108
35
        if (!_table_filters.empty()) {
1109
0
            return false;
1110
0
        }
1111
35
        if (agg_type == TPushAggOp::type::COUNT) {
1112
            // Old FEs do not serialize push_down_count_slot_ids. During the supported BE-first
1113
            // rolling upgrade, nullopt therefore means "COUNT semantics are unknown", not
1114
            // COUNT(*). Fall back to reading rows until the FE explicitly sends either an empty
1115
            // list for COUNT(*) or one slot for COUNT(col).
1116
24
            if (!_push_down_count_columns.has_value()) {
1117
3
                return false;
1118
3
            }
1119
            // COUNT(*) needs no column metadata. COUNT(col) currently supports one direct file
1120
            // column; multiple COUNT arguments fall back to the normal scan so every upper
1121
            // aggregate receives the original rows.
1122
21
            if (_push_down_count_columns->empty()) {
1123
10
                return true;
1124
10
            }
1125
11
            if (_push_down_count_columns->size() != 1) {
1126
1
                return false;
1127
1
            }
1128
10
            const auto& mapping = _push_down_count_mapping();
1129
            // Metadata COUNT skips TableReader's normal materialization path. Only a trivial
1130
            // mapping is safe: for example, a nullable Parquet INT mapped to a NOT NULL table
1131
            // BIGINT normally needs both an INT->BIGINT cast and nullability validation. Counting
1132
            // footer values directly would bypass both operations and could hide invalid data.
1133
10
            return mapping.file_local_id.has_value() && mapping.file_type != nullptr &&
1134
10
                   mapping.table_type != nullptr && mapping.is_trivial &&
1135
10
                   mapping.virtual_column_type == TableVirtualColumnType::INVALID &&
1136
10
                   mapping.default_expr == nullptr;
1137
11
        }
1138
        // For MIN/MAX, only support direct file-to-table column mappings. The two emitted rows
1139
        // must be enough for the upper MIN/MAX aggregate without evaluating default expressions or
1140
        // virtual columns.
1141
13
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
1142
13
            if (!mapping.file_local_id.has_value() ||
1143
13
                mapping.virtual_column_type != TableVirtualColumnType::INVALID ||
1144
13
                mapping.default_expr != nullptr || mapping.file_type == nullptr ||
1145
13
                mapping.table_type == nullptr) {
1146
1
                return false;
1147
1
            }
1148
12
            if (!_can_push_down_minmax_for_mapping(mapping)) {
1149
2
                return false;
1150
2
            }
1151
12
        }
1152
8
        return true;
1153
11
    }
1154
1155
5.30k
    static ColumnPtr _detach_column(ColumnPtr column) {
1156
5.30k
        DORIS_CHECK(column.get() != nullptr);
1157
5.30k
        return IColumn::mutate(std::move(column));
1158
5.30k
    }
1159
1160
595
    static ColumnPtr _take_and_detach_block_column(Block* block, int position) {
1161
595
        DORIS_CHECK(block != nullptr);
1162
595
        DORIS_CHECK(position >= 0 && position < static_cast<int>(block->columns()));
1163
595
        auto& source = block->get_by_position(position);
1164
595
        ColumnPtr column = source.column;
1165
        // The final mapping no longer needs the file block. Release its COW owner before mutate(),
1166
        // otherwise nested MAP/STRING columns are deep-copied and a multi-GB payload can OOM.
1167
595
        block->replace_by_position(position, source.type->create_column());
1168
595
        return _detach_column(std::move(column));
1169
595
    }
1170
1171
    static Status _align_column_nullability(ColumnPtr* column, const DataTypePtr& table_type,
1172
246
                                            const NullMap* nullable_parent_null_map = nullptr) {
1173
246
        DORIS_CHECK(column != nullptr);
1174
246
        DORIS_CHECK(column->get() != nullptr);
1175
246
        DORIS_CHECK(table_type != nullptr);
1176
        // Must return non-const column
1177
246
        *column = (*column)->convert_to_full_column_if_const();
1178
246
        if (table_type->is_nullable()) {
1179
92
            const auto& nested_type =
1180
92
                    assert_cast<const DataTypeNullable&>(*table_type).get_nested_type();
1181
92
            if (!(*column)->is_nullable()) {
1182
2
                RETURN_IF_ERROR(
1183
2
                        _align_column_nullability(column, nested_type, nullable_parent_null_map));
1184
2
                *column = make_nullable(*column);
1185
2
                return Status::OK();
1186
2
            }
1187
90
            const auto& nullable_column = assert_cast<const ColumnNullable&>(**column);
1188
90
            ColumnPtr nested_column = nullable_column.get_nested_column_ptr();
1189
90
            NullMap combined_null_map;
1190
90
            const NullMap* nested_parent_null_map = &nullable_column.get_null_map_data();
1191
90
            if (nullable_parent_null_map != nullptr) {
1192
52
                const auto& own_null_map = nullable_column.get_null_map_data();
1193
52
                DORIS_CHECK(nullable_parent_null_map->size() == own_null_map.size());
1194
                // Required descendants are hidden when either this nullable container or any
1195
                // inherited nullable ancestor masks the row, so preserve the union recursively.
1196
52
                combined_null_map.resize(own_null_map.size());
1197
178
                for (size_t i = 0; i < own_null_map.size(); ++i) {
1198
126
                    combined_null_map[i] = own_null_map[i] || (*nullable_parent_null_map)[i];
1199
126
                }
1200
52
                nested_parent_null_map = &combined_null_map;
1201
52
            }
1202
90
            RETURN_IF_ERROR(
1203
90
                    _align_column_nullability(&nested_column, nested_type, nested_parent_null_map));
1204
90
            *column = ColumnNullable::create(nested_column,
1205
90
                                             nullable_column.get_null_map_column_ptr());
1206
90
            return Status::OK();
1207
90
        }
1208
154
        if ((*column)->is_nullable()) {
1209
24
            const auto& nullable_column = assert_cast<const ColumnNullable&>(**column);
1210
24
            if (nullable_column.has_null()) {
1211
11
                const auto& null_map = nullable_column.get_null_map_data();
1212
11
                if (nullable_parent_null_map == nullptr ||
1213
11
                    nullable_parent_null_map->size() != null_map.size()) {
1214
1
                    return Status::InternalError(
1215
1
                            "Default expression produced NULL for non-nullable table column");
1216
1
                }
1217
23
                for (size_t i = 0; i < null_map.size(); ++i) {
1218
                    // A required child may contain a physical NULL placeholder only when its
1219
                    // nullable parent masks that row from the logical value.
1220
17
                    if (null_map[i] && !(*nullable_parent_null_map)[i]) {
1221
4
                        return Status::InternalError(
1222
4
                                "Default expression produced NULL for non-nullable table column");
1223
4
                    }
1224
17
                }
1225
10
            }
1226
19
            ColumnPtr nested_column = nullable_column.get_nested_column_ptr();
1227
19
            RETURN_IF_ERROR(_align_column_nullability(&nested_column, table_type,
1228
19
                                                      nullable_parent_null_map));
1229
19
            *column = nested_column;
1230
19
            return Status::OK();
1231
19
        }
1232
130
        if (const auto* array_type = typeid_cast<const DataTypeArray*>(table_type.get())) {
1233
2
            const auto& array_column = assert_cast<const ColumnArray&>(**column);
1234
2
            ColumnPtr nested_column = array_column.get_data_ptr();
1235
2
            NullMap descendant_parent_null_map;
1236
            // Collection entries use offset coordinates, so inherited row masks must be projected
1237
            // only when a required descendant can consume them. This avoids scratch proportional
1238
            // to all array entries for the common all-required schema.
1239
2
            const NullMap* descendant_parent_null_map_ptr = nullptr;
1240
2
            if (_requires_collection_parent_null_map(
1241
2
                        nullable_parent_null_map, nested_column, array_type->get_nested_type(),
1242
2
                        array_column.size(), array_column.get_offsets())) {
1243
1
                descendant_parent_null_map_ptr =
1244
1
                        _project_collection_parent_null_map_for_hidden_entries(
1245
1
                                nullptr, nullable_parent_null_map, array_column.size(),
1246
1
                                array_column.get_offsets(), nested_column->size(),
1247
1
                                &descendant_parent_null_map);
1248
1
            }
1249
2
            RETURN_IF_ERROR(_align_column_nullability(&nested_column, array_type->get_nested_type(),
1250
2
                                                      descendant_parent_null_map_ptr));
1251
2
            *column = ColumnArray::create(nested_column, array_column.get_offsets_ptr());
1252
2
            return Status::OK();
1253
2
        }
1254
128
        if (const auto* map_type = typeid_cast<const DataTypeMap*>(table_type.get())) {
1255
1
            const auto& map_column = assert_cast<const ColumnMap&>(**column);
1256
1
            ColumnPtr key_column = map_column.get_keys_ptr();
1257
1
            ColumnPtr value_column = map_column.get_values_ptr();
1258
1
            NullMap descendant_parent_null_map;
1259
1
            const NullMap* descendant_parent_null_map_ptr = nullptr;
1260
1
            if (_requires_collection_parent_null_map(nullable_parent_null_map, key_column,
1261
1
                                                     map_type->get_key_type(), map_column.size(),
1262
1
                                                     map_column.get_offsets()) ||
1263
1
                _requires_collection_parent_null_map(nullable_parent_null_map, value_column,
1264
1
                                                     map_type->get_value_type(), map_column.size(),
1265
1
                                                     map_column.get_offsets())) {
1266
                // Keys and values share offsets, so one projected mask safely covers both streams.
1267
0
                descendant_parent_null_map_ptr =
1268
0
                        _project_collection_parent_null_map_for_hidden_entries(
1269
0
                                nullptr, nullable_parent_null_map, map_column.size(),
1270
0
                                map_column.get_offsets(), key_column->size(),
1271
0
                                &descendant_parent_null_map);
1272
0
            }
1273
1
            RETURN_IF_ERROR(_align_column_nullability(&key_column, map_type->get_key_type(),
1274
1
                                                      descendant_parent_null_map_ptr));
1275
1
            RETURN_IF_ERROR(_align_column_nullability(&value_column, map_type->get_value_type(),
1276
1
                                                      descendant_parent_null_map_ptr));
1277
1
            *column = ColumnMap::create(key_column, value_column, map_column.get_offsets_ptr());
1278
1
            return Status::OK();
1279
1
        }
1280
127
        if (const auto* struct_type = typeid_cast<const DataTypeStruct*>(table_type.get())) {
1281
14
            const auto& struct_column = assert_cast<const ColumnStruct&>(**column);
1282
14
            Columns columns = struct_column.get_columns_copy();
1283
14
            DORIS_CHECK(columns.size() == struct_type->get_elements().size());
1284
40
            for (size_t i = 0; i < columns.size(); ++i) {
1285
26
                RETURN_IF_ERROR(_align_column_nullability(&columns[i], struct_type->get_element(i),
1286
26
                                                          nullable_parent_null_map));
1287
26
            }
1288
14
            *column = ColumnStruct::create(columns);
1289
14
            return Status::OK();
1290
14
        }
1291
113
        return Status::OK();
1292
127
    }
1293
1294
    static Status _execute_default_expr_without_root_type_check(
1295
            const VExprContextSPtr& default_expr, const Block* block,
1296
16
            ColumnWithTypeAndName* result_data) {
1297
16
        DORIS_CHECK(default_expr != nullptr);
1298
16
        DORIS_CHECK(block != nullptr);
1299
16
        DORIS_CHECK(result_data != nullptr);
1300
16
        ColumnPtr result_column;
1301
16
        Status st;
1302
16
        RETURN_IF_CATCH_EXCEPTION({
1303
16
            st = default_expr->root()->execute_column_impl(default_expr.get(), block, nullptr,
1304
16
                                                           block->rows(), result_column);
1305
16
        });
1306
16
        RETURN_IF_ERROR(st);
1307
16
        DORIS_CHECK(result_column.get() != nullptr);
1308
16
        if (result_column->size() != block->rows()) {
1309
0
            return Status::InternalError(
1310
0
                    "Default expr {} return column size {} not equal to expected size {}",
1311
0
                    default_expr->expr_name(), result_column->size(), block->rows());
1312
0
        }
1313
16
        result_data->column = result_column;
1314
16
        result_data->type = default_expr->execute_type(block);
1315
16
        result_data->name = default_expr->expr_name();
1316
16
        return Status::OK();
1317
16
    }
1318
1319
    Status _cast_column_to_type(ColumnPtr* column, const DataTypePtr& file_type,
1320
                                const DataTypePtr& table_type,
1321
26
                                const std::string& column_name) const {
1322
26
        DORIS_CHECK(column != nullptr);
1323
26
        DORIS_CHECK(column->get() != nullptr);
1324
26
        DORIS_CHECK(file_type != nullptr);
1325
26
        DORIS_CHECK(table_type != nullptr);
1326
26
        if (file_type->equals(*table_type) ||
1327
26
            remove_nullable(file_type)->equals(*remove_nullable(table_type))) {
1328
15
            return Status::OK();
1329
15
        }
1330
1331
11
        DataTypePtr input_type = file_type;
1332
        // Cast wrappers unwrap nullable inputs according to the declared input type, so keep the
1333
        // root nullability of the declared input aligned with the actual column shape. When the
1334
        // runtime column is nullable, also keep the cast target nullable; the caller applies the
1335
        // table's final nullability after value conversion. Casting a nullable runtime column
1336
        // directly to a non-nullable target would pass ColumnNullable to CastToImpl.
1337
11
        if ((*column)->is_nullable() && !input_type->is_nullable()) {
1338
0
            input_type = make_nullable(input_type);
1339
11
        } else if (!(*column)->is_nullable() && input_type->is_nullable()) {
1340
1
            input_type = remove_nullable(input_type);
1341
1
        }
1342
11
        DataTypePtr cast_type = table_type;
1343
11
        if ((*column)->is_nullable() && !cast_type->is_nullable()) {
1344
8
            cast_type = make_nullable(cast_type);
1345
8
        }
1346
11
        Block cast_block;
1347
11
        cast_block.insert({*column, input_type, column_name});
1348
11
        auto slot_ref = VSlotRef::create_shared(0, 0, -1, input_type, column_name);
1349
11
        auto cast_expr = Cast::create_shared(cast_type);
1350
11
        cast_expr->add_child(std::move(slot_ref));
1351
11
        auto cast_ctx = VExprContext::create_shared(std::move(cast_expr));
1352
11
        RowDescriptor row_desc;
1353
11
        RETURN_IF_ERROR(cast_ctx->prepare(_runtime_state, row_desc));
1354
11
        RETURN_IF_ERROR(cast_ctx->open(_runtime_state));
1355
11
        ColumnPtr cast_column;
1356
11
        RETURN_IF_ERROR(cast_ctx->execute(&cast_block, cast_column));
1357
11
        *column = std::move(cast_column);
1358
11
        return Status::OK();
1359
11
    }
1360
1361
    Status _try_materialize_scalar_cast_with_runtime_nullability(const ColumnMapping& mapping,
1362
                                                                 const Block* current_block,
1363
                                                                 ColumnPtr* column,
1364
5.27k
                                                                 bool* handled) const {
1365
5.27k
        DORIS_CHECK(column != nullptr);
1366
5.27k
        DORIS_CHECK(handled != nullptr);
1367
5.27k
        *handled = false;
1368
5.27k
        if (mapping.projection == nullptr || !mapping.file_local_id.has_value() ||
1369
5.27k
            !mapping.child_mappings.empty()) {
1370
690
            return Status::OK();
1371
690
        }
1372
1373
4.58k
        const auto& root = mapping.projection->root();
1374
4.58k
        if (root == nullptr || root->node_type() != TExprNodeType::CAST_EXPR) {
1375
4.56k
            return Status::OK();
1376
4.56k
        }
1377
20
        DORIS_CHECK(root->get_num_children() == 1);
1378
20
        const auto* slot = dynamic_cast<const VSlotRef*>(root->get_child(0).get());
1379
20
        DORIS_CHECK(slot != nullptr);
1380
20
        DORIS_CHECK(current_block != nullptr);
1381
20
        DORIS_CHECK(slot->column_id() >= 0);
1382
20
        DORIS_CHECK(cast_set<size_t>(slot->column_id()) < current_block->columns());
1383
20
        const auto& source = current_block->get_by_position(slot->column_id());
1384
20
        DORIS_CHECK(source.column.get() != nullptr);
1385
20
        DORIS_CHECK(slot->data_type() != nullptr);
1386
20
        DORIS_CHECK(mapping.table_type != nullptr);
1387
20
        const bool runtime_input_mismatch =
1388
20
                source.column->is_nullable() != slot->data_type()->is_nullable();
1389
20
        const bool nullable_input_to_required_table =
1390
20
                source.column->is_nullable() && !mapping.table_type->is_nullable();
1391
20
        if (!runtime_input_mismatch && !nullable_input_to_required_table) {
1392
8
            return Status::OK();
1393
8
        }
1394
1395
        // File readers can return a nullable runtime column even when the physical schema marks the
1396
        // leaf required. A pre-built Cast binds to the declared file type and can therefore pass a
1397
        // ColumnNullable to a non-nullable CastToImpl. Rebuild only when that runtime shape differs
1398
        // from the declared input, or when a declared nullable file field maps to a required table
1399
        // field. Keep the cast target nullable while converting values, then let
1400
        // _align_column_nullability() reject an actual NULL before removing the wrapper.
1401
12
        ColumnPtr result_column = source.column;
1402
12
        RETURN_IF_ERROR(_cast_column_to_type(&result_column, slot->data_type(), mapping.table_type,
1403
12
                                             mapping.file_column_name));
1404
12
        RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type));
1405
11
        *column = _detach_column(std::move(result_column));
1406
11
        *handled = true;
1407
11
        return Status::OK();
1408
12
    }
1409
1410
    Status _materialize_present_child_mapping_column(
1411
            const ColumnMapping& mapping, const ColumnPtr& file_column, const size_t rows,
1412
57
            ColumnPtr* column, const NullMap* nullable_parent_null_map = nullptr) {
1413
57
        DORIS_CHECK(column != nullptr);
1414
57
        DORIS_CHECK(mapping.file_type != nullptr);
1415
57
        DORIS_CHECK(mapping.table_type != nullptr);
1416
57
        *column = file_column;
1417
57
        if (!mapping.is_trivial) {
1418
28
            if (!mapping.child_mappings.empty()) {
1419
16
                RETURN_IF_ERROR(_materialize_complex_mapping_column(mapping, *column, rows, column,
1420
16
                                                                    nullable_parent_null_map));
1421
16
            } else {
1422
12
                RETURN_IF_ERROR(_cast_column_to_type(column, mapping.file_type, mapping.table_type,
1423
12
                                                     mapping.file_column_name));
1424
12
            }
1425
28
        }
1426
54
        RETURN_IF_ERROR(
1427
54
                _align_column_nullability(column, mapping.table_type, nullable_parent_null_map));
1428
50
        return Status::OK();
1429
54
    }
1430
1431
    Status _materialize_default_or_missing_column(
1432
            const ColumnMapping& mapping, const Block* current_block, const size_t rows,
1433
36
            ColumnPtr* column, const NullMap* nullable_parent_null_map = nullptr) {
1434
36
        DORIS_CHECK(mapping.table_type != nullptr);
1435
36
        DORIS_CHECK(column != nullptr);
1436
36
        if (mapping.default_expr != nullptr) {
1437
16
            Block synthetic_block;
1438
16
            const Block* eval_block = current_block;
1439
16
            if (eval_block == nullptr || eval_block->rows() != rows) {
1440
                // Nested ARRAY/MAP children use element/entry cardinality rather than the root
1441
                // block's row count. Iceberg initial defaults are typed literals, so a synthetic
1442
                // block with the desired row count is sufficient and avoids a top-level
1443
                // ConstantMap dependency for nested mappings.
1444
14
                synthetic_block.insert(
1445
14
                        {mapping.table_type->create_column_const_with_default_value(rows),
1446
14
                         mapping.table_type, "__table_reader_nested_default_rows"});
1447
14
                eval_block = &synthetic_block;
1448
14
            }
1449
16
            ColumnWithTypeAndName result;
1450
16
            RETURN_IF_ERROR(_execute_default_expr_without_root_type_check(mapping.default_expr,
1451
16
                                                                          eval_block, &result));
1452
16
            ColumnPtr result_column = result.column;
1453
16
            RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type,
1454
16
                                                      nullable_parent_null_map));
1455
16
            *column = _detach_column(std::move(result_column));
1456
16
            return Status::OK();
1457
16
        }
1458
20
        ColumnPtr result_column = mapping.table_type->create_column_const_with_default_value(rows);
1459
20
        RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type,
1460
20
                                                  nullable_parent_null_map));
1461
20
        *column = _detach_column(std::move(result_column));
1462
20
        return Status::OK();
1463
20
    }
1464
1465
    Status _materialize_mapping_column(const ColumnMapping& mapping, Block* current_block,
1466
                                       const size_t rows, ColumnPtr* column,
1467
5.29k
                                       bool take_projection_result = false) {
1468
5.29k
        if (!mapping.is_trivial && mapping.file_local_id.has_value() &&
1469
5.29k
            !mapping.child_mappings.empty()) {
1470
18
            DCHECK(mapping.projection != nullptr);
1471
18
            int res_id;
1472
18
            auto st = mapping.projection->execute(current_block, &res_id);
1473
18
            if (!st.ok()) {
1474
0
                return Status::InternalError(
1475
0
                        "Failed to execute complex mapping projection for table column '{}' "
1476
0
                        "(global_index={}, file_local_id={}, rows={}): {}, mapping={}",
1477
0
                        mapping.table_column_name, mapping.global_index.value(),
1478
0
                        *mapping.file_local_id, rows, st.to_string(), mapping.debug_string());
1479
0
            }
1480
18
            ColumnPtr result_column = take_projection_result
1481
18
                                              ? _take_and_detach_block_column(current_block, res_id)
1482
18
                                              : current_block->get_by_position(res_id).column;
1483
18
            RETURN_IF_ERROR(
1484
18
                    _materialize_complex_mapping_column(mapping, result_column, rows, column));
1485
15
            return Status::OK();
1486
18
        }
1487
5.27k
        bool runtime_nullability_cast_handled = false;
1488
5.27k
        RETURN_IF_ERROR(_try_materialize_scalar_cast_with_runtime_nullability(
1489
5.27k
                mapping, current_block, column, &runtime_nullability_cast_handled));
1490
5.27k
        if (runtime_nullability_cast_handled) {
1491
11
            return Status::OK();
1492
11
        }
1493
5.26k
        if (mapping.projection != nullptr) {
1494
5.24k
            int res_id;
1495
5.24k
            auto st = mapping.projection->execute(current_block, &res_id);
1496
5.24k
            if (!st.ok()) {
1497
0
                std::string file_local_id = "null";
1498
0
                if (mapping.file_local_id.has_value()) {
1499
0
                    file_local_id = std::to_string(*mapping.file_local_id);
1500
0
                }
1501
0
                return Status::InternalError(
1502
0
                        "Failed to execute mapping projection for table column '{}' "
1503
0
                        "(global_index={}, file_local_id={}, rows={}): {}, mapping={}",
1504
0
                        mapping.table_column_name, mapping.global_index.value(), file_local_id,
1505
0
                        rows, st.to_string(), mapping.debug_string());
1506
0
            }
1507
5.24k
            if (take_projection_result) {
1508
580
                *column = _take_and_detach_block_column(current_block, res_id);
1509
4.66k
            } else {
1510
4.66k
                ColumnPtr result_column = current_block->get_by_position(res_id).column;
1511
4.66k
                *column = _detach_column(std::move(result_column));
1512
4.66k
            }
1513
5.24k
            return Status::OK();
1514
5.24k
        }
1515
23
        return _materialize_default_or_missing_column(mapping, current_block, rows, column);
1516
5.26k
    }
1517
1518
    Status _materialize_complex_mapping_column(const ColumnMapping& mapping,
1519
                                               const ColumnPtr& file_column, const size_t rows,
1520
                                               ColumnPtr* column,
1521
34
                                               const NullMap* nullable_parent_null_map = nullptr) {
1522
34
        DORIS_CHECK(mapping.table_type != nullptr);
1523
34
        DORIS_CHECK(file_column.get() != nullptr);
1524
34
        const auto table_type = remove_nullable(mapping.table_type);
1525
34
        switch (table_type->get_primitive_type()) {
1526
22
        case TYPE_STRUCT:
1527
22
            RETURN_IF_ERROR(_materialize_struct_mapping_column(mapping, file_column, rows, column,
1528
22
                                                               nullable_parent_null_map));
1529
19
            break;
1530
19
        case TYPE_ARRAY:
1531
7
            RETURN_IF_ERROR(_materialize_array_mapping_column(mapping, file_column, rows, column,
1532
7
                                                              nullable_parent_null_map));
1533
4
            break;
1534
5
        case TYPE_MAP:
1535
5
            RETURN_IF_ERROR(_materialize_map_mapping_column(mapping, file_column, rows, column,
1536
5
                                                            nullable_parent_null_map));
1537
5
            break;
1538
5
        default:
1539
0
            *column = _detach_column(file_column);
1540
0
            break;
1541
34
        }
1542
28
        return Status::OK();
1543
34
    }
1544
1545
    static std::vector<const ColumnMapping*> _present_child_mappings_in_file_order(
1546
27
            const std::vector<ColumnMapping>& child_mappings) {
1547
27
        std::vector<const ColumnMapping*> result;
1548
27
        result.reserve(child_mappings.size());
1549
46
        for (const auto& child_mapping : child_mappings) {
1550
46
            if (child_mapping.file_local_id.has_value()) {
1551
33
                result.push_back(&child_mapping);
1552
33
            }
1553
46
        }
1554
27
        std::ranges::sort(result, [](const ColumnMapping* lhs, const ColumnMapping* rhs) {
1555
14
            DORIS_CHECK(lhs->file_local_id.has_value());
1556
14
            DORIS_CHECK(rhs->file_local_id.has_value());
1557
14
            return *lhs->file_local_id < *rhs->file_local_id;
1558
14
        });
1559
27
        return result;
1560
27
    }
1561
1562
    static size_t _file_child_ordinal_for_mapping(
1563
            const ColumnMapping& mapping, const ColumnMapping& child_mapping,
1564
32
            const std::vector<const ColumnMapping*>& file_ordered_children) {
1565
32
        DORIS_CHECK(child_mapping.file_local_id.has_value());
1566
32
        if (!mapping.projected_file_children.empty()) {
1567
21
            const auto child_it = std::ranges::find_if(
1568
28
                    mapping.projected_file_children, [&](const ColumnDefinition& file_child) {
1569
28
                        return file_child.file_local_id() == *child_mapping.file_local_id;
1570
28
                    });
1571
21
            DORIS_CHECK(child_it != mapping.projected_file_children.end());
1572
21
            return static_cast<size_t>(
1573
21
                    std::distance(mapping.projected_file_children.begin(), child_it));
1574
21
        }
1575
11
        const auto child_it = std::ranges::find(file_ordered_children, &child_mapping);
1576
11
        DORIS_CHECK(child_it != file_ordered_children.end());
1577
11
        return static_cast<size_t>(std::distance(file_ordered_children.begin(), child_it));
1578
32
    }
1579
1580
    static std::vector<const ColumnMapping*> _child_mappings_in_table_type_order(
1581
27
            const ColumnMapping& mapping, const DataTypeStruct& table_type) {
1582
27
        std::vector<const ColumnMapping*> result;
1583
27
        result.reserve(mapping.child_mappings.size());
1584
73
        for (size_t child_idx = 0; child_idx < table_type.get_elements().size(); ++child_idx) {
1585
46
            const auto& child_name = table_type.get_element_name(child_idx);
1586
46
            const auto child_it = std::ranges::find_if(
1587
68
                    mapping.child_mappings, [&](const ColumnMapping& child_mapping) {
1588
68
                        return child_mapping.table_column_name == child_name;
1589
68
                    });
1590
46
            DORIS_CHECK(child_it != mapping.child_mappings.end())
1591
0
                    << mapping.debug_string() << ", table_child_name=" << child_name;
1592
46
            result.push_back(&*child_it);
1593
46
        }
1594
27
        return result;
1595
27
    }
1596
1597
    static const IColumn* _nested_column_if_nullable(const ColumnPtr& column,
1598
47
                                                     const NullMap** null_map) {
1599
47
        DORIS_CHECK(column.get() != nullptr);
1600
47
        if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1601
43
            if (null_map != nullptr) {
1602
43
                *null_map = &nullable_column->get_null_map_data();
1603
43
            }
1604
43
            return &nullable_column->get_nested_column();
1605
43
        }
1606
4
        return column.get();
1607
47
    }
1608
1609
    static bool _requires_parent_null_map_for_alignment(const ColumnPtr& column,
1610
0
                                                        const DataTypePtr& table_type) {
1611
0
        DORIS_CHECK(column.get() != nullptr);
1612
0
        DORIS_CHECK(table_type != nullptr);
1613
0
        if (table_type->is_nullable()) {
1614
0
            const auto& nested_type =
1615
0
                    assert_cast<const DataTypeNullable&>(*table_type).get_nested_type();
1616
0
            if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1617
0
                return _requires_parent_null_map_for_alignment(
1618
0
                        nullable_column->get_nested_column_ptr(), nested_type);
1619
0
            }
1620
0
            return _requires_parent_null_map_for_alignment(column, nested_type);
1621
0
        }
1622
0
        if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1623
0
            if (nullable_column->has_null()) {
1624
0
                return true;
1625
0
            }
1626
0
            return _requires_parent_null_map_for_alignment(nullable_column->get_nested_column_ptr(),
1627
0
                                                           table_type);
1628
0
        }
1629
0
        if (const auto* array_type = typeid_cast<const DataTypeArray*>(table_type.get())) {
1630
0
            const auto& array_column = assert_cast<const ColumnArray&>(*column);
1631
0
            return _requires_parent_null_map_for_alignment(array_column.get_data_ptr(),
1632
0
                                                           array_type->get_nested_type());
1633
0
        }
1634
0
        if (const auto* map_type = typeid_cast<const DataTypeMap*>(table_type.get())) {
1635
0
            const auto& map_column = assert_cast<const ColumnMap&>(*column);
1636
0
            return _requires_parent_null_map_for_alignment(map_column.get_keys_ptr(),
1637
0
                                                           map_type->get_key_type()) ||
1638
0
                   _requires_parent_null_map_for_alignment(map_column.get_values_ptr(),
1639
0
                                                           map_type->get_value_type());
1640
0
        }
1641
0
        if (const auto* struct_type = typeid_cast<const DataTypeStruct*>(table_type.get())) {
1642
0
            const auto& struct_column = assert_cast<const ColumnStruct&>(*column);
1643
0
            DORIS_CHECK(struct_column.tuple_size() == struct_type->get_elements().size());
1644
0
            for (size_t i = 0; i < struct_column.tuple_size(); ++i) {
1645
0
                if (_requires_parent_null_map_for_alignment(struct_column.get_column_ptr(i),
1646
0
                                                            struct_type->get_element(i))) {
1647
0
                    return true;
1648
0
                }
1649
0
            }
1650
0
        }
1651
0
        return false;
1652
0
    }
1653
1654
    static bool _requires_parent_null_map_for_alignment_at(const ColumnPtr& column,
1655
                                                           const DataTypePtr& table_type,
1656
17
                                                           const size_t row) {
1657
17
        DORIS_CHECK(column.get() != nullptr);
1658
17
        DORIS_CHECK(table_type != nullptr);
1659
17
        DORIS_CHECK(row < column->size());
1660
17
        if (table_type->is_nullable()) {
1661
5
            const auto& nested_type =
1662
5
                    assert_cast<const DataTypeNullable&>(*table_type).get_nested_type();
1663
5
            if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1664
                // A nearer nullable wrapper already protects its descendants at this entry, so an
1665
                // inherited collection mask cannot be needed there.
1666
5
                if (nullable_column->is_null_at(row)) {
1667
0
                    return false;
1668
0
                }
1669
5
                return _requires_parent_null_map_for_alignment_at(
1670
5
                        nullable_column->get_nested_column_ptr(), nested_type, row);
1671
5
            }
1672
0
            return _requires_parent_null_map_for_alignment_at(column, nested_type, row);
1673
5
        }
1674
12
        if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
1675
4
            if (nullable_column->is_null_at(row)) {
1676
2
                return true;
1677
2
            }
1678
2
            return _requires_parent_null_map_for_alignment_at(
1679
2
                    nullable_column->get_nested_column_ptr(), table_type, row);
1680
4
        }
1681
8
        if (const auto* array_type = typeid_cast<const DataTypeArray*>(table_type.get())) {
1682
0
            const auto& array_column = assert_cast<const ColumnArray&>(*column);
1683
0
            const auto& offsets = array_column.get_offsets();
1684
0
            const size_t begin = row == 0 ? 0 : offsets[row - 1];
1685
0
            const size_t end = offsets[row];
1686
0
            for (size_t child_row = begin; child_row < end; ++child_row) {
1687
0
                if (_requires_parent_null_map_for_alignment_at(array_column.get_data_ptr(),
1688
0
                                                               array_type->get_nested_type(),
1689
0
                                                               child_row)) {
1690
0
                    return true;
1691
0
                }
1692
0
            }
1693
0
            return false;
1694
0
        }
1695
8
        if (const auto* map_type = typeid_cast<const DataTypeMap*>(table_type.get())) {
1696
0
            const auto& map_column = assert_cast<const ColumnMap&>(*column);
1697
0
            const auto& offsets = map_column.get_offsets();
1698
0
            const size_t begin = row == 0 ? 0 : offsets[row - 1];
1699
0
            const size_t end = offsets[row];
1700
0
            for (size_t child_row = begin; child_row < end; ++child_row) {
1701
0
                if (_requires_parent_null_map_for_alignment_at(
1702
0
                            map_column.get_keys_ptr(), map_type->get_key_type(), child_row) ||
1703
0
                    _requires_parent_null_map_for_alignment_at(
1704
0
                            map_column.get_values_ptr(), map_type->get_value_type(), child_row)) {
1705
0
                    return true;
1706
0
                }
1707
0
            }
1708
0
            return false;
1709
0
        }
1710
8
        if (const auto* struct_type = typeid_cast<const DataTypeStruct*>(table_type.get())) {
1711
4
            const auto& struct_column = assert_cast<const ColumnStruct&>(*column);
1712
4
            DORIS_CHECK(struct_column.tuple_size() == struct_type->get_elements().size());
1713
6
            for (size_t i = 0; i < struct_column.tuple_size(); ++i) {
1714
4
                if (_requires_parent_null_map_for_alignment_at(struct_column.get_column_ptr(i),
1715
4
                                                               struct_type->get_element(i), row)) {
1716
2
                    return true;
1717
2
                }
1718
4
            }
1719
4
        }
1720
6
        return false;
1721
8
    }
1722
1723
    static bool _requires_collection_parent_null_map(const NullMap* parent_null_map,
1724
                                                     const ColumnPtr& column,
1725
                                                     const DataTypePtr& table_type) {
1726
        // Descendant null maps can be entry-sized. Scan them only when an inherited mask can
1727
        // actually hide a row; absent/all-clear masks cannot authorize any physical child NULL.
1728
        if (parent_null_map == nullptr ||
1729
            std::ranges::none_of(*parent_null_map, [](const auto value) { return value != 0; })) {
1730
            return false;
1731
        }
1732
        return _requires_parent_null_map_for_alignment(column, table_type);
1733
    }
1734
1735
    template <typename Offsets>
1736
    static bool _parent_null_map_hides_collection_entries(const NullMap* container_null_map,
1737
                                                          const NullMap* ancestor_null_map,
1738
                                                          const size_t rows,
1739
22
                                                          const Offsets& offsets) {
1740
22
        if (container_null_map == nullptr && ancestor_null_map == nullptr) {
1741
2
            return false;
1742
2
        }
1743
20
        DORIS_CHECK(container_null_map == nullptr || container_null_map->size() == rows);
1744
20
        DORIS_CHECK(ancestor_null_map == nullptr || ancestor_null_map->size() == rows);
1745
20
        DORIS_CHECK(offsets.size() == rows);
1746
20
        size_t begin = 0;
1747
69
        for (size_t row = 0; row < rows; ++row) {
1748
54
            const size_t end = offsets[row];
1749
54
            const bool hidden = (container_null_map != nullptr && (*container_null_map)[row]) ||
1750
54
                                (ancestor_null_map != nullptr && (*ancestor_null_map)[row]);
1751
            // A hidden collection row protects descendants only when its offset span is nonempty.
1752
54
            if (hidden && end > begin) {
1753
5
                return true;
1754
5
            }
1755
49
            begin = end;
1756
49
        }
1757
15
        return false;
1758
20
    }
1759
1760
    template <typename Offsets>
1761
    static bool _requires_collection_parent_null_map(const NullMap* parent_null_map,
1762
                                                     const ColumnPtr& column,
1763
                                                     const DataTypePtr& table_type,
1764
9
                                                     const size_t rows, const Offsets& offsets) {
1765
9
        if (parent_null_map == nullptr) {
1766
0
            return false;
1767
0
        }
1768
9
        DORIS_CHECK(parent_null_map->size() == rows);
1769
9
        DORIS_CHECK(offsets.size() == rows);
1770
9
        DORIS_CHECK(offsets.empty() || offsets.back() == column->size());
1771
9
        size_t begin = 0;
1772
22
        for (size_t row = 0; row < rows; ++row) {
1773
15
            const size_t end = offsets[row];
1774
15
            if ((*parent_null_map)[row]) {
1775
                // Only ancestor-hidden entries can consume this projection. Restricting the probe
1776
                // to their spans avoids scanning visible payload covered by nearer nullable masks.
1777
12
                for (size_t child_row = begin; child_row < end; ++child_row) {
1778
6
                    if (_requires_parent_null_map_for_alignment_at(column, table_type, child_row)) {
1779
2
                        return true;
1780
2
                    }
1781
6
                }
1782
8
            }
1783
13
            begin = end;
1784
13
        }
1785
7
        return false;
1786
9
    }
1787
1788
    template <typename Offsets>
1789
    static const NullMap* _project_collection_parent_null_map_for_hidden_entries(
1790
            const NullMap* container_null_map, const NullMap* ancestor_null_map, const size_t rows,
1791
22
            const Offsets& offsets, const size_t child_rows, NullMap* const projected_null_map) {
1792
22
        if (!_parent_null_map_hides_collection_entries(container_null_map, ancestor_null_map, rows,
1793
22
                                                       offsets)) {
1794
            // Nullable collection wrappers expose a null-map even when every row is present; avoid
1795
            // allocating entry-coordinate scratch unless a hidden row owns physical entries.
1796
17
            return nullptr;
1797
17
        }
1798
5
        projected_null_map->resize(child_rows);
1799
5
        std::fill(projected_null_map->begin(), projected_null_map->end(), 0);
1800
5
        size_t begin = 0;
1801
15
        for (size_t row = 0; row < rows; ++row) {
1802
10
            const size_t end = offsets[row];
1803
10
            const bool hidden = (container_null_map != nullptr && (*container_null_map)[row]) ||
1804
10
                                (ancestor_null_map != nullptr && (*ancestor_null_map)[row]);
1805
10
            if (hidden) {
1806
                // Collection masks use row coordinates; descendants need the same invariant
1807
                // projected through offsets so hidden physical payload cannot fail validation.
1808
5
                std::fill(projected_null_map->begin() + begin, projected_null_map->begin() + end,
1809
5
                          1);
1810
5
            }
1811
10
            begin = end;
1812
10
        }
1813
5
        DORIS_CHECK(begin == child_rows);
1814
5
        return projected_null_map;
1815
22
    }
1816
1817
    Status _materialize_struct_mapping_column(const ColumnMapping& mapping,
1818
                                              const ColumnPtr& file_column, const size_t rows,
1819
                                              ColumnPtr* column,
1820
28
                                              const NullMap* nullable_parent_null_map = nullptr) {
1821
28
        DORIS_CHECK(mapping.table_type != nullptr);
1822
28
        const auto* table_type =
1823
28
                assert_cast<const DataTypeStruct*>(remove_nullable(mapping.table_type).get());
1824
28
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1825
28
        const NullMap* parent_null_map = nullptr;
1826
28
        const auto* nested_file_column =
1827
28
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1828
28
        const auto* file_struct = assert_cast<const ColumnStruct*>(nested_file_column);
1829
28
        DORIS_CHECK(table_type->get_elements().size() == mapping.child_mappings.size());
1830
1831
28
        NullMap combined_parent_null_map;
1832
28
        const NullMap* descendant_parent_null_map = nullable_parent_null_map;
1833
28
        if (parent_null_map != nullptr) {
1834
26
            DORIS_CHECK(parent_null_map->size() == rows);
1835
26
            if (nullable_parent_null_map != nullptr) {
1836
3
                DORIS_CHECK(nullable_parent_null_map->size() == rows);
1837
3
            }
1838
26
            if (!mapping.table_type->is_nullable()) {
1839
7
                for (size_t i = 0; i < rows; ++i) {
1840
                    // A required nested container may drop its own NULL only when an ancestor
1841
                    // already hides that row; otherwise physical defaults become visible values.
1842
5
                    if ((*parent_null_map)[i] &&
1843
5
                        (nullable_parent_null_map == nullptr || !(*nullable_parent_null_map)[i])) {
1844
1
                        return Status::InternalError(
1845
1
                                "Source struct contains NULL for non-nullable table column");
1846
1
                    }
1847
5
                }
1848
3
            }
1849
25
            combined_parent_null_map.resize(rows);
1850
79
            for (size_t i = 0; i < rows; ++i) {
1851
54
                combined_parent_null_map[i] =
1852
54
                        (*parent_null_map)[i] ||
1853
54
                        (nullable_parent_null_map != nullptr && (*nullable_parent_null_map)[i]);
1854
54
            }
1855
25
            descendant_parent_null_map = &combined_parent_null_map;
1856
25
        }
1857
1858
27
        Columns child_columns;
1859
27
        child_columns.reserve(mapping.child_mappings.size());
1860
27
        const auto file_ordered_children =
1861
27
                _present_child_mappings_in_file_order(mapping.child_mappings);
1862
27
        const auto table_ordered_children =
1863
27
                _child_mappings_in_table_type_order(mapping, *table_type);
1864
45
        for (const auto* child_mapping : table_ordered_children) {
1865
45
            DORIS_CHECK(child_mapping != nullptr);
1866
45
            if (!child_mapping->file_local_id.has_value()) {
1867
13
                ColumnPtr child_column;
1868
13
                RETURN_IF_ERROR(_materialize_default_or_missing_column(
1869
13
                        *child_mapping, nullptr, rows, &child_column, descendant_parent_null_map));
1870
13
                child_column = child_column->convert_to_full_column_if_const();
1871
13
                child_columns.push_back(std::move(child_column));
1872
13
                continue;
1873
13
            }
1874
32
            const auto file_child_idx =
1875
32
                    _file_child_ordinal_for_mapping(mapping, *child_mapping, file_ordered_children);
1876
32
            DORIS_CHECK(file_child_idx < file_struct->get_columns().size());
1877
32
            ColumnPtr child_column = file_struct->get_column_ptr(file_child_idx);
1878
32
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(
1879
32
                    *child_mapping, child_column, rows, &child_column, descendant_parent_null_map));
1880
28
            child_columns.push_back(std::move(child_column));
1881
28
        }
1882
23
        MutableColumns mutable_child_columns;
1883
23
        mutable_child_columns.reserve(child_columns.size());
1884
39
        for (auto& child_column : child_columns) {
1885
39
            mutable_child_columns.push_back(IColumn::mutate(std::move(child_column)));
1886
39
        }
1887
23
        auto result = ColumnStruct::create(std::move(mutable_child_columns));
1888
23
        if (mapping.table_type->is_nullable()) {
1889
19
            auto null_map = ColumnUInt8::create();
1890
19
            auto& null_map_data = null_map->get_data();
1891
19
            null_map_data.resize(rows);
1892
19
            if (parent_null_map != nullptr) {
1893
19
                DORIS_CHECK(parent_null_map->size() == rows);
1894
19
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
1895
19
            } else {
1896
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
1897
0
            }
1898
19
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
1899
19
        } else {
1900
4
            *column = std::move(result);
1901
4
        }
1902
23
        return Status::OK();
1903
27
    }
1904
1905
    Status _materialize_array_mapping_column(const ColumnMapping& mapping,
1906
                                             const ColumnPtr& file_column, const size_t rows,
1907
                                             ColumnPtr* column,
1908
11
                                             const NullMap* nullable_parent_null_map = nullptr) {
1909
11
        DORIS_CHECK(mapping.child_mappings.size() == 1);
1910
11
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1911
11
        const NullMap* parent_null_map = nullptr;
1912
11
        const auto* nested_file_column =
1913
11
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1914
11
        if (parent_null_map != nullptr && !mapping.table_type->is_nullable()) {
1915
2
            DORIS_CHECK(parent_null_map->size() == rows);
1916
2
            if (nullable_parent_null_map != nullptr) {
1917
1
                DORIS_CHECK(nullable_parent_null_map->size() == rows);
1918
1
            }
1919
4
            for (size_t i = 0; i < rows; ++i) {
1920
                // ARRAY row masks cannot be forwarded to elements because they use different
1921
                // coordinates, so validate the container before dropping its nullable wrapper.
1922
3
                if ((*parent_null_map)[i] &&
1923
3
                    (nullable_parent_null_map == nullptr || !(*nullable_parent_null_map)[i])) {
1924
1
                    return Status::InternalError(
1925
1
                            "Source array contains NULL for non-nullable table column");
1926
1
                }
1927
3
            }
1928
2
        }
1929
10
        const auto* file_array = assert_cast<const ColumnArray*>(nested_file_column);
1930
10
        ColumnPtr nested_column = file_array->get_data_ptr();
1931
10
        auto element_mapping = mapping.child_mappings[0];
1932
        // Keep the descriptor type for schema matching. ARRAY's nullable element wrapper is a
1933
        // storage invariant, so add it only at the materialization boundary.
1934
10
        element_mapping.table_type = make_nullable(element_mapping.table_type);
1935
10
        NullMap descendant_parent_null_map;
1936
10
        const NullMap* descendant_parent_null_map_ptr =
1937
10
                _project_collection_parent_null_map_for_hidden_entries(
1938
10
                        parent_null_map, nullable_parent_null_map, rows, file_array->get_offsets(),
1939
10
                        nested_column->size(), &descendant_parent_null_map);
1940
10
        RETURN_IF_ERROR(_materialize_present_child_mapping_column(
1941
10
                element_mapping, nested_column, nested_column->size(), &nested_column,
1942
10
                descendant_parent_null_map_ptr));
1943
7
        auto offsets_column = file_array->get_offsets_ptr()->convert_to_full_column_if_const();
1944
7
        auto result = ColumnArray::create(IColumn::mutate(std::move(nested_column)),
1945
7
                                          IColumn::mutate(std::move(offsets_column)));
1946
7
        if (mapping.table_type->is_nullable()) {
1947
6
            auto null_map = ColumnUInt8::create();
1948
6
            auto& null_map_data = null_map->get_data();
1949
6
            null_map_data.resize(rows);
1950
6
            if (parent_null_map != nullptr) {
1951
6
                DORIS_CHECK(parent_null_map->size() == rows);
1952
6
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
1953
6
            } else {
1954
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
1955
0
            }
1956
6
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
1957
6
        } else {
1958
1
            *column = std::move(result);
1959
1
        }
1960
7
        return Status::OK();
1961
10
    }
1962
1963
    Status _materialize_map_mapping_column(const ColumnMapping& mapping,
1964
                                           const ColumnPtr& file_column, const size_t rows,
1965
                                           ColumnPtr* column,
1966
8
                                           const NullMap* nullable_parent_null_map = nullptr) {
1967
8
        const auto full_file_column = file_column->convert_to_full_column_if_const();
1968
8
        const NullMap* parent_null_map = nullptr;
1969
8
        const auto* nested_file_column =
1970
8
                _nested_column_if_nullable(full_file_column, &parent_null_map);
1971
8
        if (parent_null_map != nullptr && !mapping.table_type->is_nullable()) {
1972
0
            DORIS_CHECK(parent_null_map->size() == rows);
1973
0
            if (nullable_parent_null_map != nullptr) {
1974
0
                DORIS_CHECK(nullable_parent_null_map->size() == rows);
1975
0
            }
1976
0
            for (size_t i = 0; i < rows; ++i) {
1977
                // MAP row masks cannot be forwarded to entries because they use different
1978
                // coordinates, so validate the container before dropping its nullable wrapper.
1979
0
                if ((*parent_null_map)[i] &&
1980
0
                    (nullable_parent_null_map == nullptr || !(*nullable_parent_null_map)[i])) {
1981
0
                    return Status::InternalError(
1982
0
                            "Source map contains NULL for non-nullable table column");
1983
0
                }
1984
0
            }
1985
0
        }
1986
8
        const auto* file_map = assert_cast<const ColumnMap*>(nested_file_column);
1987
8
        ColumnPtr key_column = file_map->get_keys_ptr();
1988
8
        ColumnPtr value_column = file_map->get_values_ptr();
1989
8
        DORIS_CHECK(key_column->size() == value_column->size());
1990
8
        NullMap descendant_parent_null_map;
1991
8
        const NullMap* descendant_parent_null_map_ptr =
1992
8
                _project_collection_parent_null_map_for_hidden_entries(
1993
8
                        parent_null_map, nullable_parent_null_map, rows, file_map->get_offsets(),
1994
8
                        key_column->size(), &descendant_parent_null_map);
1995
1996
8
        const ColumnMapping* key_mapping = nullptr;
1997
8
        const ColumnMapping* value_mapping = nullptr;
1998
13
        for (const auto& child_mapping : mapping.child_mappings) {
1999
13
            if (!child_mapping.file_local_id.has_value()) {
2000
0
                continue;
2001
0
            }
2002
13
            if (*child_mapping.file_local_id == 0) {
2003
5
                key_mapping = &child_mapping;
2004
8
            } else if (*child_mapping.file_local_id == 1) {
2005
8
                value_mapping = &child_mapping;
2006
8
            }
2007
13
        }
2008
2009
8
        if (key_mapping != nullptr) {
2010
5
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(
2011
5
                    *key_mapping, key_column, key_column->size(), &key_column,
2012
5
                    descendant_parent_null_map_ptr));
2013
5
        } else {
2014
3
            const auto* table_map =
2015
3
                    assert_cast<const DataTypeMap*>(remove_nullable(mapping.table_type).get());
2016
            // Value-only projection retains the physical key stream to preserve entry offsets;
2017
            // align it under the entry mask so NULL placeholders from hidden Map rows stay hidden.
2018
3
            RETURN_IF_ERROR(_align_column_nullability(&key_column, table_map->get_key_type(),
2019
3
                                                      descendant_parent_null_map_ptr));
2020
3
        }
2021
8
        if (value_mapping != nullptr) {
2022
8
            RETURN_IF_ERROR(_materialize_present_child_mapping_column(
2023
8
                    *value_mapping, value_column, value_column->size(), &value_column,
2024
8
                    descendant_parent_null_map_ptr));
2025
8
        } else {
2026
0
            const auto* table_map =
2027
0
                    assert_cast<const DataTypeMap*>(remove_nullable(mapping.table_type).get());
2028
            // A retained structural value stream follows the same hidden-entry invariant as keys.
2029
0
            RETURN_IF_ERROR(_align_column_nullability(&value_column, table_map->get_value_type(),
2030
0
                                                      descendant_parent_null_map_ptr));
2031
0
        }
2032
8
        auto offsets_column = file_map->get_offsets_ptr()->convert_to_full_column_if_const();
2033
8
        auto result = ColumnMap::create(IColumn::mutate(std::move(key_column)),
2034
8
                                        IColumn::mutate(std::move(value_column)),
2035
8
                                        IColumn::mutate(std::move(offsets_column)));
2036
8
        if (mapping.table_type->is_nullable()) {
2037
6
            auto null_map = ColumnUInt8::create();
2038
6
            auto& null_map_data = null_map->get_data();
2039
6
            null_map_data.resize(rows);
2040
6
            if (parent_null_map != nullptr) {
2041
6
                DORIS_CHECK(parent_null_map->size() == rows);
2042
6
                null_map_data.assign(parent_null_map->begin(), parent_null_map->end());
2043
6
            } else {
2044
0
                std::fill(null_map_data.begin(), null_map_data.end(), 0);
2045
0
            }
2046
6
            *column = ColumnNullable::create(std::move(result), std::move(null_map));
2047
6
        } else {
2048
2
            *column = std::move(result);
2049
2
        }
2050
8
        return Status::OK();
2051
8
    }
2052
2053
3.53k
    Status _open_mapping_expr_tree(const ColumnMapping& mapping, const RowDescriptor& row_desc) {
2054
3.53k
        if (mapping.projection != nullptr) {
2055
2.31k
            RETURN_IF_ERROR(mapping.projection->prepare(_runtime_state, row_desc));
2056
2.31k
            RETURN_IF_ERROR(mapping.projection->open(_runtime_state));
2057
2.31k
        }
2058
3.53k
        if (mapping.default_expr != nullptr) {
2059
17
            RETURN_IF_ERROR(mapping.default_expr->prepare(_runtime_state, row_desc));
2060
17
            RETURN_IF_ERROR(mapping.default_expr->open(_runtime_state));
2061
17
        }
2062
3.53k
        for (const auto& child_mapping : mapping.child_mappings) {
2063
1.19k
            RETURN_IF_ERROR(_open_mapping_expr_tree(child_mapping, row_desc));
2064
1.19k
        }
2065
3.53k
        return Status::OK();
2066
3.53k
    }
2067
2068
443
    Status _open_mapping_exprs() {
2069
443
        RowDescriptor row_desc;
2070
2.33k
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
2071
2.33k
            RETURN_IF_ERROR(_open_mapping_expr_tree(mapping, row_desc));
2072
2.33k
        }
2073
443
        return Status::OK();
2074
443
    }
2075
2076
    Status _build_file_aggregate_request(TPushAggOp::type agg_type,
2077
13
                                         FileAggregateRequest* request) const {
2078
13
        DORIS_CHECK(request != nullptr);
2079
13
        DORIS_CHECK(_supports_aggregate_pushdown(agg_type));
2080
13
        request->agg_type = agg_type;
2081
13
        request->columns.clear();
2082
13
        if (agg_type == TPushAggOp::type::COUNT) {
2083
9
            DORIS_CHECK(_push_down_count_columns.has_value());
2084
            // An empty explicit list is the semantic signal for COUNT(*). Do not inspect the
2085
            // mapping count: `SELECT COUNT(*) FROM t` may still project one nullable column because
2086
            // the planner keeps a placeholder slot. In a 10,000-row file where that arbitrary slot
2087
            // has 9,015 non-null values, passing the slot would ask Parquet/ORC metadata for
2088
            // COUNT(slot)=9,015 instead of the required row count 10,000.
2089
9
            if (!_push_down_count_columns->empty()) {
2090
4
                const auto& mapping = _push_down_count_mapping();
2091
4
                DORIS_CHECK(mapping.file_local_id.has_value());
2092
4
                FileAggregateRequest::Column column;
2093
4
                column.projection =
2094
4
                        LocalColumnIndex::top_level(LocalColumnId(*mapping.file_local_id));
2095
4
                request->columns.push_back(std::move(column));
2096
4
            }
2097
9
            return Status::OK();
2098
9
        }
2099
4
        request->columns.reserve(_data_reader.column_mapper->mappings().size());
2100
5
        for (const auto& mapping : _data_reader.column_mapper->mappings()) {
2101
5
            DORIS_CHECK(mapping.file_local_id.has_value());
2102
5
            FileAggregateRequest::Column column;
2103
5
            column.projection = LocalColumnIndex::top_level(LocalColumnId(*mapping.file_local_id));
2104
5
            if (!mapping.child_mappings.empty()) {
2105
1
                RETURN_IF_ERROR(build_aggregate_projection(mapping, &column.projection));
2106
1
            }
2107
5
            request->columns.push_back(std::move(column));
2108
5
        }
2109
4
        return Status::OK();
2110
4
    }
2111
2112
14
    const ColumnMapping& _push_down_count_mapping() const {
2113
14
        DORIS_CHECK(_push_down_count_columns.has_value());
2114
14
        DORIS_CHECK(_push_down_count_columns->size() == 1);
2115
14
        const auto mapping_it =
2116
14
                std::ranges::find(_data_reader.column_mapper->mappings(),
2117
14
                                  _push_down_count_columns->front(), &ColumnMapping::global_index);
2118
        // FileScannerV2 translates FE SlotIds through the same projected-column list used to build
2119
        // the mapper, so a missing mapping is an FE/BE contract violation rather than a fallback.
2120
14
        DORIS_CHECK(mapping_it != _data_reader.column_mapper->mappings().end());
2121
14
        return *mapping_it;
2122
14
    }
2123
2124
    Status _materialize_aggregate_pushdown_rows(TPushAggOp::type agg_type,
2125
                                                const FileAggregateResult& file_result,
2126
3
                                                Block* block) {
2127
3
        DORIS_CHECK(agg_type == TPushAggOp::type::MINMAX);
2128
        // MIN/MAX pushdown emits two rows, min first and max second, for each projected column.
2129
        // The upper MIN/MAX aggregate consumes those two rows to produce the final aggregate value.
2130
3
        DORIS_CHECK(file_result.columns.size() == _data_reader.column_mapper->mappings().size());
2131
3
        DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size());
2132
3
        Block file_block;
2133
3
        file_block.reserve(_data_reader.file_block_layout.size());
2134
4
        for (const auto& column : _data_reader.file_block_layout) {
2135
4
            file_block.insert({column.type->create_column(), column.type, column.name});
2136
4
        }
2137
7
        for (size_t column_idx = 0; column_idx < file_result.columns.size(); ++column_idx) {
2138
4
            const auto& result_column = file_result.columns[column_idx];
2139
4
            if (!result_column.has_min || !result_column.has_max) {
2140
0
                return Status::NotSupported("Missing min/max aggregate result for column {}",
2141
0
                                            _projected_columns[column_idx].name);
2142
0
            }
2143
4
            bool found_file_column = false;
2144
5
            for (size_t block_position = 0; block_position < _data_reader.file_block_layout.size();
2145
5
                 ++block_position) {
2146
5
                if (_data_reader.file_block_layout[block_position].file_column_id ==
2147
5
                    file_result.columns[column_idx].projection.column_id()) {
2148
4
                    found_file_column = true;
2149
4
                    auto column = file_block.get_by_position(block_position)
2150
4
                                          .type->create_column()
2151
4
                                          ->assert_mutable();
2152
4
                    RETURN_IF_ERROR(_insert_aggregate_projection_value(
2153
4
                            file_result.columns[column_idx].projection, result_column.min_value,
2154
4
                            column.get()));
2155
4
                    RETURN_IF_ERROR(_insert_aggregate_projection_value(
2156
4
                            file_result.columns[column_idx].projection, result_column.max_value,
2157
4
                            column.get()));
2158
4
                    file_block.replace_by_position(block_position, std::move(column));
2159
4
                    break;
2160
4
                }
2161
5
            }
2162
4
            DORIS_CHECK(found_file_column);
2163
4
        }
2164
7
        for (size_t column_idx = 0; column_idx < _data_reader.column_mapper->mappings().size();
2165
4
             ++column_idx) {
2166
4
            ColumnPtr table_column;
2167
4
            RETURN_IF_ERROR(_materialize_mapping_column(
2168
4
                    _data_reader.column_mapper->mappings()[column_idx], &file_block, 2,
2169
4
                    &table_column,
2170
4
                    column_idx + 1 == _data_reader.column_mapper->mappings().size()));
2171
4
            block->replace_by_position(column_idx, std::move(table_column));
2172
4
        }
2173
3
        return Status::OK();
2174
3
    }
2175
2176
    struct FileBlockColumn {
2177
        LocalColumnId file_column_id = LocalColumnId::invalid();
2178
        std::string name;
2179
        DataTypePtr type;
2180
    };
2181
2182
    struct DataReader {
2183
        std::unique_ptr<FileReader> reader;
2184
        std::unique_ptr<TableColumnMapper> column_mapper;
2185
        // Schema of the data file, also including virtual column (row position).
2186
        std::vector<ColumnDefinition> file_schema;
2187
        // Layout of the block returned by file reader, determined by column mapping and file
2188
        // schema. It is used for file reader to materialize columns into correct type and position.
2189
        std::vector<FileBlockColumn> file_block_layout;
2190
        Block block_template;
2191
    };
2192
    DataReader _data_reader;
2193
    // Latest immutable request queued to the physical reader. The file-block layout remains fixed
2194
    // for the split even while predicates are refreshed at a reader-defined granule boundary.
2195
    std::shared_ptr<FileScanRequest> _file_scan_request;
2196
    std::vector<ColumnDefinition> _projected_columns;
2197
    std::unique_ptr<ScanTask> _current_task;
2198
    std::optional<io::FileDescription> _current_file_description;
2199
    // Range-level compression has higher priority than scan-param compression. TVF/load can keep
2200
    // the logical format as CSV/TEXT while carrying the concrete compression such as GZ or LZO on
2201
    // each TFileRangeDesc, matching the old FileScanner reader contract.
2202
    TFileCompressType::type _current_range_compress_type = TFileCompressType::UNKNOWN;
2203
    std::optional<TUniqueId> _current_range_load_id;
2204
    TFileRangeDesc _current_file_range_desc;
2205
    std::shared_ptr<io::FileSystemProperties> _system_properties;
2206
    // partition key -> value
2207
    std::map<std::string, Field> _partition_values;
2208
    // Predicates built from scan conjuncts before file-level localization.
2209
    std::vector<TableFilter> _table_filters;
2210
    // Number of localized filters before the first unsafe conjunct in the original row-level
2211
    // order. This differs from scanning `_table_filters` for safety because slotless predicates are
2212
    // intentionally absent from that vector but must still act as ordering barriers.
2213
    size_t _constant_pruning_safe_filter_count = 0;
2214
    VExprContextSPtrs _conjuncts;
2215
    ReadProfile _profile;
2216
    // Parsed from row-position based delete files, including position delete and deletion vector.
2217
    DeleteRows* _delete_rows = nullptr;
2218
    DeletionVector* _deletion_vector = nullptr;
2219
    TFileScanRangeParams* _scan_params;
2220
    std::shared_ptr<io::IOContext> _io_ctx;
2221
    RuntimeState* _runtime_state;
2222
    RuntimeProfile* _scanner_profile;
2223
    const std::vector<SlotDescriptor*>* _file_slot_descs = nullptr;
2224
    FileFormat _format;
2225
    TPushAggOp::type _push_down_agg_type = TPushAggOp::type::NONE;
2226
    std::optional<std::vector<GlobalIndex>> _push_down_count_columns;
2227
    size_t _batch_size = 0;
2228
    uint64_t _initial_condition_cache_digest = 0;
2229
    uint64_t _condition_cache_digest = 0;
2230
    // True only when prepare_split() received a digest for the exact conjunct snapshot used by
2231
    // this split. Standalone callers that only supplied TableReadOptions::condition_cache_digest
2232
    // keep the conservative runtime-filter guard.
2233
    bool _condition_cache_digest_covers_current_split = false;
2234
    segment_v2::ConditionCache::ExternalCacheKey _condition_cache_key;
2235
    std::shared_ptr<std::vector<bool>> _condition_cache;
2236
    std::shared_ptr<ConditionCacheContext> _condition_cache_ctx;
2237
    int64_t _condition_cache_hit_count = 0;
2238
    bool _current_reader_reached_eof = false;
2239
    int64_t _remaining_table_level_count = -1;
2240
    int64_t _remaining_file_level_count = -1;
2241
    // True only after the active split selects a table-level row-count shortcut or successfully
2242
    // materializes COUNT rows from file metadata. FileScannerV2 uses this result, rather than the
2243
    // raw aggregate opcode, to keep adaptive batching enabled for normal row-scan fallbacks.
2244
    bool _current_split_uses_metadata_count = false;
2245
    // Snapshot supplied by FileScannerV2 for the active split. It gates every shortcut that emits
2246
    // irreversible aggregate rows, not only the table-level row-count shortcut in prepare_split().
2247
    bool _all_runtime_filters_applied_for_split = true;
2248
    std::optional<GlobalRowIdContext> _global_rowid_context;
2249
    bool _aggregate_pushdown_tried = false;
2250
    bool _current_split_pruned = false;
2251
    TableColumnMapperOptions _mapper_options;
2252
2253
private:
2254
    static const ColumnDefinition* _find_column_definition(
2255
2.41k
            const std::vector<ColumnDefinition>& schema, LocalColumnId column_id) {
2256
17.9k
        for (const auto& field : schema) {
2257
17.9k
            if (field.file_local_id() == column_id.value()) {
2258
2.37k
                return &field;
2259
2.37k
            }
2260
17.9k
        }
2261
39
        return nullptr;
2262
2.41k
    }
2263
2264
14
    static bool _can_push_down_minmax_for_mapping(const ColumnMapping& mapping) {
2265
14
        if (mapping.child_mappings.empty()) {
2266
            // Direct mappings use a slot-ref projection to materialize the file column. The
2267
            // projection does not transform ordering; casts and other conversions are already
2268
            // represented by a non-trivial mapping and must fall back to row scanning.
2269
11
            return mapping.is_trivial;
2270
11
        }
2271
3
        const auto primitive_type = remove_nullable(mapping.file_type)->get_primitive_type();
2272
3
        if (primitive_type != TYPE_STRUCT) {
2273
1
            return false;
2274
1
        }
2275
2
        size_t mapped_children = 0;
2276
2
        const ColumnMapping* mapped_child = nullptr;
2277
2
        for (const auto& child_mapping : mapping.child_mappings) {
2278
2
            if (!child_mapping.file_local_id.has_value()) {
2279
0
                continue;
2280
0
            }
2281
2
            ++mapped_children;
2282
2
            mapped_child = &child_mapping;
2283
2
        }
2284
2
        return mapped_children == 1 && mapped_child != nullptr &&
2285
2
               _can_push_down_minmax_for_mapping(*mapped_child);
2286
3
    }
2287
2288
    static Status build_aggregate_projection(const ColumnMapping& mapping,
2289
2
                                             LocalColumnIndex* projection) {
2290
2
        DORIS_CHECK(projection != nullptr);
2291
2
        DORIS_CHECK(mapping.file_local_id.has_value());
2292
2
        *projection = LocalColumnIndex::local(*mapping.file_local_id);
2293
2
        projection->children.clear();
2294
2
        projection->project_all_children = true;
2295
2
        if (mapping.child_mappings.empty()) {
2296
1
            return Status::OK();
2297
1
        }
2298
1
        projection->project_all_children = false;
2299
1
        for (const auto& child_mapping : mapping.child_mappings) {
2300
1
            if (!child_mapping.file_local_id.has_value()) {
2301
0
                continue;
2302
0
            }
2303
1
            LocalColumnIndex child_projection;
2304
1
            RETURN_IF_ERROR(build_aggregate_projection(child_mapping, &child_projection));
2305
1
            projection->children.push_back(std::move(child_projection));
2306
1
        }
2307
1
        DORIS_CHECK(projection->children.size() == 1);
2308
1
        return Status::OK();
2309
1
    }
2310
2311
    static Status _insert_aggregate_projection_value(const LocalColumnIndex& projection,
2312
20
                                                     const Field& value, IColumn* column) {
2313
20
        DORIS_CHECK(column != nullptr);
2314
20
        if (auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) {
2315
10
            RETURN_IF_ERROR(_insert_aggregate_projection_value(
2316
10
                    projection, value, &nullable_column->get_nested_column()));
2317
10
            nullable_column->get_null_map_data().push_back(0);
2318
10
            return Status::OK();
2319
10
        }
2320
10
        if (projection.project_all_children || projection.children.empty()) {
2321
8
            column->insert(value);
2322
8
            return Status::OK();
2323
8
        }
2324
2
        auto* struct_column = assert_cast<ColumnStruct*>(column);
2325
2
        DORIS_CHECK(projection.children.size() == 1);
2326
2
        const auto& child_projection = projection.children[0];
2327
2
        DORIS_CHECK(struct_column->get_columns().size() == 1);
2328
2
        RETURN_IF_ERROR(_insert_aggregate_projection_value(child_projection, value,
2329
2
                                                           &struct_column->get_column(0)));
2330
2
        return Status::OK();
2331
2
    }
2332
2333
    // Parse a DV into its compressed bitmap. Position delete files continue to use _delete_rows.
2334
    Status _parse_delete_predicates(const SplitReadOptions& options);
2335
};
2336
2337
} // namespace doris::format