Coverage Report

Created: 2026-08-01 00:02

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