Coverage Report

Created: 2026-09-28 17:01

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/table/iceberg_reader.cpp
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
#include "format/table/iceberg_reader.h"
19
20
#include <gen_cpp/Descriptors_types.h>
21
#include <gen_cpp/Metrics_types.h>
22
#include <gen_cpp/PlanNodes_types.h>
23
#include <gen_cpp/parquet_types.h>
24
#include <glog/logging.h>
25
#include <parallel_hashmap/phmap.h>
26
#include <rapidjson/document.h>
27
28
#include <algorithm>
29
#include <cstring>
30
#include <functional>
31
#include <memory>
32
33
#include "common/compiler_util.h" // IWYU pragma: keep
34
#include "common/consts.h"
35
#include "common/status.h"
36
#include "core/assert_cast.h"
37
#include "core/block/block.h"
38
#include "core/block/column_with_type_and_name.h"
39
#include "core/column/column.h"
40
#include "core/column/column_nullable.h"
41
#include "core/column/column_string.h"
42
#include "core/column/column_vector.h"
43
#include "core/data_type/data_type_factory.hpp"
44
#include "core/data_type/data_type_nullable.h"
45
#include "core/data_type/define_primitive_type.h"
46
#include "core/data_type/primitive_type.h"
47
#include "core/string_ref.h"
48
#include "exprs/aggregate/aggregate_function.h"
49
#include "format/format_common.h"
50
#include "format/generic_reader.h"
51
#include "format/orc/vorc_reader.h"
52
#include "format/parquet/schema_desc.h"
53
#include "format/parquet/vparquet_column_chunk_reader.h"
54
#include "format/table/deletion_vector_reader.h"
55
#include "format/table/iceberg/iceberg_orc_nested_column_utils.h"
56
#include "format/table/iceberg/iceberg_parquet_nested_column_utils.h"
57
#include "format/table/iceberg_scan_semantics.h"
58
#include "format/table/nested_column_access_helper.h"
59
#include "format/table/table_schema_change_helper.h"
60
#include "runtime/runtime_state.h"
61
#include "util/coding.h"
62
#include "util/string_util.h"
63
64
namespace cctz {
65
class time_zone;
66
} // namespace cctz
67
namespace doris {
68
class RowDescriptor;
69
class SlotDescriptor;
70
class TupleDescriptor;
71
72
namespace io {
73
struct IOContext;
74
} // namespace io
75
class VExprContext;
76
} // namespace doris
77
78
namespace doris {
79
namespace {
80
81
constexpr auto kIcebergOrcAttribute = "iceberg.id";
82
83
84
bool orc_subtree_has_iceberg_id(const orc::Type* type, const std::string& attribute) {
84
84
    if (type->hasAttributeKey(attribute)) {
85
18
        return true;
86
18
    }
87
102
    for (uint64_t idx = 0; idx < type->getSubtypeCount(); ++idx) {
88
54
        if (orc_subtree_has_iceberg_id(type->getSubtype(idx), attribute)) {
89
18
            return true;
90
18
        }
91
54
    }
92
48
    return false;
93
66
}
94
95
98
bool parquet_subtree_has_iceberg_id(const FieldSchema& field) {
96
98
    if (field.field_id >= 0) {
97
52
        return true;
98
52
    }
99
46
    return std::ranges::any_of(field.children, parquet_subtree_has_iceberg_id);
100
98
}
101
102
struct ParquetEqualityFieldPath {
103
    std::vector<const FieldSchema*> fields;
104
    std::vector<size_t> child_indexes;
105
};
106
107
bool find_parquet_equality_field_path_by_id(const FieldDescriptor* descriptor, int32_t field_id,
108
16
                                            ParquetEqualityFieldPath* result) {
109
16
    DORIS_CHECK(descriptor != nullptr);
110
16
    DORIS_CHECK(result != nullptr);
111
16
    const auto find = [field_id](const auto& self, const FieldSchema* field,
112
48
                                 ParquetEqualityFieldPath* path) -> bool {
113
48
        DORIS_CHECK(field != nullptr);
114
48
        path->fields.push_back(field);
115
48
        if (field->field_id == field_id) {
116
8
            return true;
117
8
        }
118
42
        for (size_t index = 0; index < field->children.size(); ++index) {
119
8
            path->child_indexes.push_back(index);
120
8
            if (self(self, &field->children[index], path)) {
121
6
                return true;
122
6
            }
123
2
            path->child_indexes.pop_back();
124
2
        }
125
34
        path->fields.pop_back();
126
34
        return false;
127
40
    };
128
48
    for (int index = 0; index < descriptor->size(); ++index) {
129
40
        if (find(find, descriptor->get_column(index), result)) {
130
8
            return true;
131
8
        }
132
40
    }
133
8
    return false;
134
16
}
135
136
bool find_parquet_equality_field_prefix_by_id_path(
137
        const FieldDescriptor* descriptor,
138
        const std::vector<const schema::external::TField*>& table_path,
139
8
        ParquetEqualityFieldPath* result) {
140
8
    DORIS_CHECK(descriptor != nullptr);
141
8
    DORIS_CHECK(result != nullptr);
142
8
    DORIS_CHECK(!table_path.empty());
143
8
    const std::vector<FieldSchema>* candidates = nullptr;
144
10
    for (size_t path_index = 0; path_index < table_path.size(); ++path_index) {
145
10
        const auto* table_field = table_path[path_index];
146
10
        DORIS_CHECK(table_field != nullptr);
147
10
        DORIS_CHECK(table_field->__isset.id);
148
10
        const FieldSchema* match = nullptr;
149
10
        size_t match_index = 0;
150
10
        const size_t candidate_count =
151
10
                candidates == nullptr ? cast_set<size_t>(descriptor->size()) : candidates->size();
152
34
        for (size_t candidate_index = 0; candidate_index < candidate_count; ++candidate_index) {
153
24
            const auto* candidate = candidates == nullptr
154
24
                                            ? descriptor->get_column(cast_set<int>(candidate_index))
155
24
                                            : &(*candidates)[candidate_index];
156
24
            if (candidate != nullptr && candidate->field_id == table_field->id) {
157
0
                match = candidate;
158
0
                match_index = candidate_index;
159
0
                break;
160
0
            }
161
24
        }
162
10
        if (match == nullptr) {
163
10
            const auto wrapper =
164
10
                    candidates == nullptr
165
10
                            ? TableSchemaChangeHelper::BuildTableInfoUtil::
166
8
                                      find_unique_idless_parquet_wrapper_index(
167
8
                                              *table_field, descriptor->get_fields_schema())
168
10
                            : TableSchemaChangeHelper::BuildTableInfoUtil::
169
2
                                      find_unique_idless_parquet_wrapper_index(*table_field,
170
2
                                                                               *candidates);
171
10
            if (wrapper.has_value()) {
172
2
                match_index = *wrapper;
173
2
                match = candidates == nullptr ? descriptor->get_column(cast_set<int>(match_index))
174
2
                                              : &(*candidates)[match_index];
175
2
            }
176
10
        }
177
10
        if (match == nullptr) {
178
8
            return false;
179
8
        }
180
2
        if (!result->fields.empty()) {
181
0
            result->child_indexes.push_back(match_index);
182
0
        }
183
2
        result->fields.push_back(match);
184
2
        candidates = &match->children;
185
2
    }
186
0
    return true;
187
8
}
188
189
std::vector<std::string> equality_field_name_candidates(const schema::external::TField& table_field,
190
42
                                                        const std::string* leaf_fallback) {
191
42
    std::vector<std::string> candidates;
192
42
    if (table_field.__isset.name_mapping) {
193
22
        candidates.insert(candidates.end(), table_field.name_mapping.begin(),
194
22
                          table_field.name_mapping.end());
195
22
        if (table_field.__isset.name_mapping_is_authoritative &&
196
22
            table_field.name_mapping_is_authoritative) {
197
20
            return candidates;
198
20
        }
199
22
    }
200
22
    if (table_field.__isset.name) {
201
22
        candidates.push_back(table_field.name);
202
22
    }
203
22
    if (leaf_fallback != nullptr) {
204
22
        candidates.push_back(*leaf_fallback);
205
22
    }
206
22
    return candidates;
207
42
}
208
209
bool find_parquet_equality_field_prefix_by_name_path(
210
        const FieldDescriptor* descriptor,
211
        const std::vector<const schema::external::TField*>& table_path,
212
14
        const std::string& leaf_fallback, ParquetEqualityFieldPath* result) {
213
14
    DORIS_CHECK(descriptor != nullptr);
214
14
    DORIS_CHECK(result != nullptr);
215
14
    DORIS_CHECK(!table_path.empty());
216
14
    const std::vector<FieldSchema>* children = nullptr;
217
34
    for (size_t path_index = 0; path_index < table_path.size(); ++path_index) {
218
22
        const auto* table_field = table_path[path_index];
219
22
        DORIS_CHECK(table_field != nullptr);
220
22
        const auto names = equality_field_name_candidates(
221
22
                *table_field, path_index + 1 == table_path.size() ? &leaf_fallback : nullptr);
222
22
        const FieldSchema* match = nullptr;
223
22
        size_t match_index = 0;
224
22
        const size_t child_count =
225
22
                children == nullptr ? cast_set<size_t>(descriptor->size()) : children->size();
226
30
        for (const auto& name : names) {
227
58
            for (size_t child_index = 0; child_index < child_count; ++child_index) {
228
48
                const auto* child = children == nullptr
229
48
                                            ? descriptor->get_column(cast_set<int>(child_index))
230
48
                                            : &(*children)[child_index];
231
48
                if (child != nullptr && iequal(child->name, name)) {
232
20
                    match = child;
233
20
                    match_index = child_index;
234
20
                    break;
235
20
                }
236
48
            }
237
30
            if (match != nullptr) {
238
20
                break;
239
20
            }
240
30
        }
241
22
        if (match == nullptr) {
242
2
            return false;
243
2
        }
244
20
        if (!result->fields.empty()) {
245
8
            result->child_indexes.push_back(match_index);
246
8
        }
247
20
        result->fields.push_back(match);
248
20
        children = &match->children;
249
20
    }
250
12
    return true;
251
14
}
252
253
struct OrcEqualityFieldPath {
254
    std::vector<const orc::Type*> fields;
255
    std::vector<std::string> names;
256
    std::vector<size_t> child_indexes;
257
};
258
259
bool find_orc_equality_field_path_by_id(const orc::Type* root, int32_t field_id,
260
16
                                        OrcEqualityFieldPath* result) {
261
16
    DORIS_CHECK(root != nullptr);
262
16
    DORIS_CHECK(result != nullptr);
263
16
    const auto find = [field_id](const auto& self, const orc::Type* field,
264
16
                                 const std::string& field_name,
265
44
                                 OrcEqualityFieldPath* path) -> bool {
266
44
        DORIS_CHECK(field != nullptr);
267
44
        path->fields.push_back(field);
268
44
        path->names.push_back(field_name);
269
44
        if (field->hasAttributeKey(kIcebergOrcAttribute) &&
270
44
            std::stoi(field->getAttributeValue(kIcebergOrcAttribute)) == field_id) {
271
6
            return true;
272
6
        }
273
40
        for (size_t index = 0; index < field->getSubtypeCount(); ++index) {
274
6
            path->child_indexes.push_back(index);
275
6
            if (self(self, field->getSubtype(index), field->getFieldName(index), path)) {
276
4
                return true;
277
4
            }
278
2
            path->child_indexes.pop_back();
279
2
        }
280
34
        path->fields.pop_back();
281
34
        path->names.pop_back();
282
34
        return false;
283
38
    };
284
48
    for (size_t index = 0; index < root->getSubtypeCount(); ++index) {
285
38
        if (find(find, root->getSubtype(index), root->getFieldName(index), result)) {
286
6
            return true;
287
6
        }
288
38
    }
289
10
    return false;
290
16
}
291
292
bool find_orc_equality_field_prefix_by_id_path(
293
        const orc::Type* root, const std::vector<const schema::external::TField*>& table_path,
294
8
        OrcEqualityFieldPath* result) {
295
8
    DORIS_CHECK(root != nullptr);
296
8
    DORIS_CHECK(result != nullptr);
297
8
    DORIS_CHECK(!table_path.empty());
298
8
    const orc::Type* parent = root;
299
10
    for (const auto* table_field : table_path) {
300
10
        DORIS_CHECK(table_field != nullptr);
301
10
        DORIS_CHECK(table_field->__isset.id);
302
10
        const orc::Type* match = nullptr;
303
10
        size_t match_index = 0;
304
34
        for (size_t candidate_index = 0; candidate_index < parent->getSubtypeCount();
305
24
             ++candidate_index) {
306
24
            const auto* candidate = parent->getSubtype(candidate_index);
307
24
            if (candidate->hasAttributeKey(kIcebergOrcAttribute) &&
308
24
                std::stoi(candidate->getAttributeValue(kIcebergOrcAttribute)) == table_field->id) {
309
0
                match = candidate;
310
0
                match_index = candidate_index;
311
0
                break;
312
0
            }
313
24
        }
314
10
        if (match == nullptr) {
315
10
            const auto wrapper = TableSchemaChangeHelper::BuildTableInfoUtil::
316
10
                    find_unique_idless_orc_wrapper_index(*table_field, parent,
317
10
                                                         kIcebergOrcAttribute);
318
10
            if (wrapper.has_value()) {
319
2
                match_index = *wrapper;
320
2
                match = parent->getSubtype(match_index);
321
2
            }
322
10
        }
323
10
        if (match == nullptr) {
324
8
            return false;
325
8
        }
326
2
        if (!result->fields.empty()) {
327
0
            result->child_indexes.push_back(match_index);
328
0
        }
329
2
        result->fields.push_back(match);
330
2
        result->names.push_back(parent->getFieldName(match_index));
331
2
        parent = match;
332
2
    }
333
0
    return true;
334
8
}
335
336
bool find_orc_equality_field_prefix_by_name_path(
337
        const orc::Type* root, const std::vector<const schema::external::TField*>& table_path,
338
12
        const std::string& leaf_fallback, OrcEqualityFieldPath* result) {
339
12
    DORIS_CHECK(root != nullptr);
340
12
    DORIS_CHECK(result != nullptr);
341
12
    DORIS_CHECK(!table_path.empty());
342
12
    const orc::Type* parent = root;
343
30
    for (size_t path_index = 0; path_index < table_path.size(); ++path_index) {
344
20
        const auto* table_field = table_path[path_index];
345
20
        DORIS_CHECK(table_field != nullptr);
346
20
        const auto names = equality_field_name_candidates(
347
20
                *table_field, path_index + 1 == table_path.size() ? &leaf_fallback : nullptr);
348
20
        const orc::Type* match = nullptr;
349
20
        size_t match_index = 0;
350
28
        for (const auto& name : names) {
351
52
            for (size_t child_index = 0; child_index < parent->getSubtypeCount(); ++child_index) {
352
42
                if (iequal(parent->getFieldName(child_index), name)) {
353
18
                    match = parent->getSubtype(child_index);
354
18
                    match_index = child_index;
355
18
                    break;
356
18
                }
357
42
            }
358
28
            if (match != nullptr) {
359
18
                break;
360
18
            }
361
28
        }
362
20
        if (match == nullptr) {
363
2
            return false;
364
2
        }
365
18
        if (!result->fields.empty()) {
366
8
            result->child_indexes.push_back(match_index);
367
8
        }
368
18
        result->fields.push_back(match);
369
18
        result->names.push_back(parent->getFieldName(match_index));
370
18
        parent = match;
371
18
    }
372
10
    return true;
373
12
}
374
375
} // namespace
376
377
const std::string IcebergOrcReader::ICEBERG_ORC_ATTRIBUTE = kIcebergOrcAttribute;
378
379
bool IcebergTableReader::_is_fully_dictionary_encoded(
380
16
        const tparquet::ColumnMetaData& column_metadata) {
381
28
    const auto is_dictionary_encoding = [](tparquet::Encoding::type encoding) {
382
28
        return encoding == tparquet::Encoding::PLAIN_DICTIONARY ||
383
28
               encoding == tparquet::Encoding::RLE_DICTIONARY;
384
28
    };
385
24
    const auto is_data_page = [](tparquet::PageType::type page_type) {
386
24
        return page_type == tparquet::PageType::DATA_PAGE ||
387
24
               page_type == tparquet::PageType::DATA_PAGE_V2;
388
24
    };
389
16
    const auto is_level_encoding = [](tparquet::Encoding::type encoding) {
390
4
        return encoding == tparquet::Encoding::RLE || encoding == tparquet::Encoding::BIT_PACKED;
391
4
    };
392
393
    // A column chunk may have a dictionary page but still contain plain-encoded data pages.
394
    // Only treat it as dictionary-coded when all data pages are dictionary encoded.
395
16
    if (column_metadata.__isset.encoding_stats) {
396
14
        bool has_data_page_stats = false;
397
24
        for (const tparquet::PageEncodingStats& enc_stat : column_metadata.encoding_stats) {
398
24
            if (is_data_page(enc_stat.page_type) && enc_stat.count > 0) {
399
16
                has_data_page_stats = true;
400
16
                if (!is_dictionary_encoding(enc_stat.encoding)) {
401
4
                    return false;
402
4
                }
403
16
            }
404
24
        }
405
10
        if (has_data_page_stats) {
406
8
            return true;
407
8
        }
408
10
    }
409
410
4
    bool has_dict_encoding = false;
411
4
    bool has_nondict_encoding = false;
412
6
    for (const tparquet::Encoding::type& encoding : column_metadata.encodings) {
413
6
        if (is_dictionary_encoding(encoding)) {
414
2
            has_dict_encoding = true;
415
2
        }
416
417
6
        if (!is_dictionary_encoding(encoding) && !is_level_encoding(encoding)) {
418
4
            has_nondict_encoding = true;
419
4
            break;
420
4
        }
421
6
    }
422
4
    if (!has_dict_encoding || has_nondict_encoding) {
423
4
        return false;
424
4
    }
425
426
0
    return true;
427
4
}
428
429
// ============================================================================
430
// IcebergParquetReader: on_before_init_reader (Parquet-specific schema matching)
431
// ============================================================================
432
// This format-specific setup mirrors the existing reader initialization sequence.
433
// NOLINTNEXTLINE(readability-function-cognitive-complexity,readability-function-size)
434
36
Status IcebergParquetReader::on_before_init_reader(ReaderInitContext* ctx) {
435
36
    _column_descs = ctx->column_descs;
436
36
    _fill_col_name_to_block_idx = ctx->col_name_to_block_idx;
437
36
    _file_format = Fileformat::PARQUET;
438
439
    // Get file metadata schema first (available because _open_file() already ran)
440
36
    const FieldDescriptor* field_desc = nullptr;
441
36
    RETURN_IF_ERROR(this->get_file_metadata_schema(&field_desc));
442
36
    DCHECK(field_desc != nullptr);
443
444
    // Build table_info_node by field_id or name matching.
445
    // This must happen BEFORE column classification so we can use children_column_exists
446
    // to check if a column exists in the file (by field ID, not name).
447
36
    if (!get_scan_params().__isset.history_schema_info ||
448
36
        get_scan_params().history_schema_info.empty()) [[unlikely]] {
449
2
        RETURN_IF_ERROR(BuildTableInfoUtil::by_parquet_name(ctx->tuple_descriptor, *field_desc,
450
2
                                                            ctx->table_info_node));
451
34
    } else {
452
34
        RETURN_IF_ERROR(BuildTableInfoUtil::by_parquet_field_id_with_name_mapping(
453
34
                get_scan_params().history_schema_info.front().root_field, *field_desc,
454
34
                ctx->table_info_node, supports_iceberg_scan_semantics_v1(&get_scan_params())));
455
34
    }
456
457
36
    std::unordered_set<std::string> partition_col_names;
458
36
    if (ctx->range->__isset.columns_from_path_keys) {
459
0
        partition_col_names.insert(ctx->range->columns_from_path_keys.begin(),
460
0
                                   ctx->range->columns_from_path_keys.end());
461
0
    }
462
463
    // Single pass: classify columns, detect $row_id, handle partition fallback.
464
36
    bool has_partition_from_path = false;
465
40
    for (const auto& desc : *ctx->column_descs) {
466
40
        if (desc.category == ColumnCategory::SYNTHESIZED) {
467
0
            if (desc.name == BeConsts::ICEBERG_ROWID_COL) {
468
0
                this->register_synthesized_column_handler(
469
0
                        BeConsts::ICEBERG_ROWID_COL, [this](Block* block, size_t rows) -> Status {
470
0
                            return _fill_iceberg_row_id(block, rows);
471
0
                        });
472
0
                continue;
473
0
            } else if (desc.name.starts_with(BeConsts::GLOBAL_ROWID_COL)) {
474
0
                auto topn_row_id_column_iter = _create_topn_row_id_column_iterator();
475
0
                this->register_synthesized_column_handler(
476
0
                        desc.name,
477
0
                        [iter = std::move(topn_row_id_column_iter), this, &desc](
478
0
                                Block* block, size_t rows) -> Status {
479
0
                            return fill_topn_row_id(iter, desc.name, block, rows);
480
0
                        });
481
0
                continue;
482
0
            }
483
40
        } else if (desc.category == ColumnCategory::PARTITION_KEY) {
484
0
            bool has_partition_value = partition_col_names.contains(desc.name);
485
0
            bool exists_in_file = ctx->table_info_node->children_column_exists(desc.name);
486
0
            if (!has_partition_value || exists_in_file) {
487
                // Keep PARTITION_KEY category stable for scan planning, but still read
488
                // from file when the column exists there.
489
0
                ctx->column_names.push_back(desc.name);
490
0
                continue;
491
0
            }
492
0
            has_partition_from_path = true;
493
40
        } else if (desc.category == ColumnCategory::REGULAR) {
494
40
            ctx->column_names.push_back(desc.name);
495
40
        } else if (desc.category == ColumnCategory::GENERATED) {
496
0
            _init_row_lineage_columns();
497
0
            if (desc.name == ROW_LINEAGE_ROW_ID) {
498
0
                ctx->column_names.push_back(desc.name);
499
0
                this->register_generated_column_handler(
500
0
                        ROW_LINEAGE_ROW_ID, [this](Block* block, size_t rows) -> Status {
501
0
                            return _fill_row_lineage_row_id(block, rows);
502
0
                        });
503
0
                continue;
504
0
            } else if (desc.name == ROW_LINEAGE_LAST_UPDATED_SEQ_NUMBER) {
505
0
                ctx->column_names.push_back(desc.name);
506
0
                this->register_generated_column_handler(
507
0
                        ROW_LINEAGE_LAST_UPDATED_SEQ_NUMBER,
508
0
                        [this](Block* block, size_t rows) -> Status {
509
0
                            return _fill_row_lineage_last_updated_sequence_number(block, rows);
510
0
                        });
511
0
                continue;
512
0
            }
513
0
        }
514
40
    }
515
516
    // Set up partition value extraction if any partition columns need filling from path
517
36
    if (has_partition_from_path) {
518
0
        RETURN_IF_ERROR(_extract_partition_values(*ctx->range, ctx->tuple_descriptor,
519
0
                                                  _fill_partition_values,
520
0
                                                  &_fill_partition_value_is_null));
521
0
    }
522
523
36
    _all_required_col_names = ctx->column_names;
524
525
    // Create column IDs from field descriptor
526
36
    auto column_id_result =
527
36
            _create_column_ids(field_desc, ctx->tuple_descriptor, ctx->table_info_node);
528
36
    ctx->column_ids = std::move(column_id_result.column_ids);
529
36
    ctx->filter_column_ids = std::move(column_id_result.filter_column_ids);
530
531
    // Build field_id -> block_column_name mapping for equality delete filtering.
532
    // This was previously done in init_reader() column matching (pre-CRTP refactoring).
533
40
    for (const auto* slot : ctx->tuple_descriptor->slots()) {
534
40
        _id_to_block_column_name.emplace(slot->col_unique_id(), slot->col_name());
535
40
    }
536
537
    // Process delete files (must happen before _do_init_reader so expand col IDs are included)
538
36
    RETURN_IF_ERROR(_init_row_filters());
539
540
    // Add expand column IDs for equality delete and remap expand column names
541
    // to match master's behavior:
542
    // - Use field_id to find the actual file column name in Parquet schema
543
    // - Prefix with __equality_delete_column__ to avoid name conflicts
544
    // - Correctly map table_col_name → file_col_name in table_info_node
545
36
    const static std::string EQ_DELETE_PRE = "__equality_delete_column__";
546
36
    bool all_file_columns_have_field_ids = true;
547
36
    bool any_file_column_has_field_id = false;
548
120
    for (int i = 0; i < field_desc->size(); ++i) {
549
84
        const auto* field_schema = field_desc->get_column(i);
550
84
        if (field_schema) {
551
84
            if (field_schema->field_id < 0) {
552
38
                all_file_columns_have_field_ids = false;
553
38
            }
554
84
            if (parquet_subtree_has_iceberg_id(*field_schema)) {
555
52
                any_file_column_has_field_id = true;
556
52
            }
557
84
        }
558
84
    }
559
36
    const bool use_field_ids_for_hidden_keys =
560
36
            supports_iceberg_scan_semantics_v1(&get_scan_params())
561
36
                    ? any_file_column_has_field_id
562
36
                    : all_file_columns_have_field_ids;
563
36
    const auto find_file_column_by_name = [&](const std::string& name) -> const FieldSchema* {
564
0
        for (int j = 0; j < field_desc->size(); ++j) {
565
0
            const auto* candidate = field_desc->get_column(j);
566
0
            if (candidate != nullptr && iequal(candidate->name, name)) {
567
0
                return candidate;
568
0
            }
569
0
        }
570
0
        return nullptr;
571
0
    };
572
573
    // Rebuild _expand_col_names with proper file-column-based names
574
36
    std::vector<std::string> new_expand_col_names;
575
36
    DORIS_CHECK(_expand_col_names.size() == _expand_col_field_ids.size());
576
36
    DORIS_CHECK(_expand_col_names.size() == _expand_columns.size());
577
66
    for (size_t i = 0; i < _expand_col_names.size(); ++i) {
578
30
        const auto& old_name = _expand_col_names[i];
579
30
        const int32_t field_id = _expand_col_field_ids[i];
580
581
30
        const FieldSchema* file_column = nullptr;
582
30
        ParquetEqualityFieldPath file_path;
583
30
        bool complete_file_path = false;
584
30
        if (use_field_ids_for_hidden_keys) {
585
16
            complete_file_path =
586
16
                    find_parquet_equality_field_path_by_id(field_desc, field_id, &file_path);
587
16
            if (!complete_file_path && supports_iceberg_scan_semantics_v2(&get_scan_params())) {
588
8
                const auto table_path = _find_schema_field_path(field_id);
589
8
                if (!table_path.empty()) {
590
8
                    complete_file_path = find_parquet_equality_field_prefix_by_id_path(
591
8
                            field_desc, table_path, &file_path);
592
8
                }
593
8
            }
594
16
            if (!file_path.fields.empty()) {
595
10
                file_column = file_path.fields.front();
596
10
            }
597
16
        } else {
598
14
            const auto table_path = _find_schema_field_path(field_id);
599
14
            if (!table_path.empty()) {
600
14
                complete_file_path = find_parquet_equality_field_prefix_by_name_path(
601
14
                        field_desc, table_path, old_name, &file_path);
602
14
                if (!file_path.fields.empty()) {
603
12
                    file_column = file_path.fields.front();
604
12
                }
605
14
            } else {
606
0
                file_column = find_file_column_by_name(old_name);
607
0
                complete_file_path = file_column != nullptr;
608
0
            }
609
14
        }
610
611
30
        std::string leaf_name = old_name;
612
30
        if (!file_path.fields.empty()) {
613
22
            leaf_name = file_path.fields.back()->name;
614
22
        } else if (file_column != nullptr) {
615
0
            leaf_name = file_column->name;
616
0
        }
617
30
        const std::string file_col_name = file_column == nullptr ? old_name : file_column->name;
618
30
        std::string table_col_name = EQ_DELETE_PRE + std::to_string(field_id) + "_" + leaf_name;
619
620
        // Update _id_to_block_column_name
621
30
        if (field_id >= 0) {
622
30
            _id_to_block_column_name[field_id] = table_col_name;
623
30
        }
624
625
        // Update _expand_columns name
626
30
        _expand_columns[i].name = table_col_name;
627
628
30
        if (file_column == nullptr) {
629
8
            RETURN_IF_ERROR(_register_missing_equality_delete_column(field_id, table_col_name,
630
8
                                                                     _expand_columns[i].type));
631
            // The old data file predates this equality key. Keep it in the expand block so the
632
            // synthesized-column hook can materialize its logical initial default before reader
633
            // filtering, but do not advertise it to Parquet as a physical child.
634
8
            new_expand_col_names.push_back(table_col_name);
635
8
            continue;
636
8
        }
637
638
22
        new_expand_col_names.push_back(table_col_name);
639
640
22
        if (!complete_file_path) {
641
2
            ColumnPtr missing_value;
642
2
            RETURN_IF_ERROR(_create_missing_equality_delete_value(
643
2
                    field_id, _expand_columns[i].type, file_path.fields.size(), &missing_value));
644
2
            _nested_equality_delete_columns.push_back({
645
2
                    .field_id = field_id,
646
2
                    .block_name = table_col_name,
647
2
                    .leaf_type = _expand_columns[i].type,
648
2
                    .child_indexes = file_path.child_indexes,
649
2
                    .missing_value = std::move(missing_value),
650
2
            });
651
2
            _expand_columns[i].type = make_nullable(file_column->data_type);
652
2
            _expand_columns[i].column = _expand_columns[i].type->create_column();
653
20
        } else if (!file_path.child_indexes.empty()) {
654
14
            _nested_equality_delete_columns.push_back({
655
14
                    .field_id = field_id,
656
14
                    .block_name = table_col_name,
657
14
                    .leaf_type = _expand_columns[i].type,
658
14
                    .child_indexes = file_path.child_indexes,
659
14
                    .missing_value = nullptr,
660
14
            });
661
14
            _expand_columns[i].type = make_nullable(file_column->data_type);
662
14
            _expand_columns[i].column = _expand_columns[i].type->create_column();
663
14
        }
664
665
        // A hidden nested key is read through its containing top-level struct. V1 column IDs are
666
        // pre-order ranges, so include the complete subtree before extracting the primitive leaf.
667
22
        for (uint64_t column_id = file_column->get_column_id();
668
60
             column_id <= file_column->get_max_column_id(); ++column_id) {
669
38
            ctx->column_ids.insert(column_id);
670
38
        }
671
672
        // Register in table_info_node: table_col_name → file_col_name
673
22
        ctx->column_names.push_back(table_col_name);
674
22
        ctx->table_info_node->add_children(table_col_name, file_col_name,
675
22
                                           TableSchemaChangeHelper::ConstNode::get_instance());
676
22
    }
677
36
    _expand_col_names = std::move(new_expand_col_names);
678
679
    // Enable group filtering for Iceberg
680
36
    _filter_groups = true;
681
682
36
    return Status::OK();
683
36
}
684
685
// ============================================================================
686
// IcebergParquetReader: _create_column_ids
687
// ============================================================================
688
ColumnIdResult IcebergParquetReader::_create_column_ids(
689
        const FieldDescriptor* field_desc, const TupleDescriptor* tuple_descriptor,
690
50
        const std::shared_ptr<TableSchemaChangeHelper::Node>& table_info_node) {
691
50
    auto* mutable_field_desc = const_cast<FieldDescriptor*>(field_desc);
692
50
    mutable_field_desc->assign_ids();
693
694
50
    std::unordered_map<int, const FieldSchema*> iceberg_id_to_field_schema_map;
695
234
    for (int i = 0; i < field_desc->size(); ++i) {
696
184
        const auto* field_schema = field_desc->get_column(i);
697
184
        if (!field_schema) {
698
0
            continue;
699
0
        }
700
184
        int iceberg_id = field_schema->field_id;
701
184
        iceberg_id_to_field_schema_map[iceberg_id] = field_schema;
702
184
    }
703
704
50
    std::set<uint64_t> column_ids;
705
50
    std::set<uint64_t> filter_column_ids;
706
707
50
    auto process_access_paths = [](const FieldSchema* parquet_field,
708
50
                                   const std::vector<TColumnAccessPath>& access_paths,
709
50
                                   std::set<uint64_t>& out_ids) {
710
34
        process_nested_access_paths(
711
34
                parquet_field, access_paths, out_ids,
712
34
                [](const FieldSchema* field) { return field->get_column_id(); },
713
34
                [](const FieldSchema* field) { return field->get_max_column_id(); },
714
34
                IcebergParquetNestedColumnUtils::extract_nested_column_ids);
715
34
    };
716
717
    // The Iceberg schema-mapping root is a StructNode whose registered children are the real
718
    // table columns. When present, resolve each column by name through it so the column-id set
719
    // stays consistent with the schema-mapping decision (BY_ID or BY_NAME/name-mapping);
720
    // otherwise fall back to matching by Iceberg field id.
721
50
    const auto* struct_node =
722
50
            dynamic_cast<const TableSchemaChangeHelper::StructNode*>(table_info_node.get());
723
724
70
    for (const auto* slot : tuple_descriptor->slots()) {
725
70
        const FieldSchema* field_schema = nullptr;
726
70
        if (struct_node != nullptr) {
727
            // Synthesized/metadata slots (e.g. the TopN global row-id or the $row_id column) are
728
            // never registered as children, so check membership before querying: calling
729
            // children_column_exists() on an unregistered name DCHECK-aborts in debug builds and
730
            // throws std::out_of_range from .at() in release builds.
731
44
            if (struct_node->get_children().contains(slot->col_name()) &&
732
44
                struct_node->children_column_exists(slot->col_name())) {
733
                // Use the physical child selected by the schema-mapping pass. This keeps partial-id
734
                // files in BY_NAME mode from binding a projected column through an unrelated stale
735
                // field id.
736
40
                const auto& file_column_name =
737
40
                        struct_node->children_file_column_name(slot->col_name());
738
50
                for (int i = 0; i < field_desc->size(); ++i) {
739
50
                    const auto* candidate = field_desc->get_column(i);
740
50
                    if (candidate != nullptr && candidate->name == file_column_name) {
741
40
                        field_schema = candidate;
742
40
                        break;
743
40
                    }
744
50
                }
745
40
                DORIS_CHECK(field_schema != nullptr);
746
40
            }
747
44
        } else {
748
26
            auto it = iceberg_id_to_field_schema_map.find(slot->col_unique_id());
749
26
            if (it != iceberg_id_to_field_schema_map.end()) {
750
26
                field_schema = it->second;
751
26
            }
752
26
        }
753
70
        if (field_schema == nullptr) {
754
4
            continue;
755
4
        }
756
757
66
        if ((slot->col_type() != TYPE_STRUCT && slot->col_type() != TYPE_ARRAY &&
758
66
             slot->col_type() != TYPE_MAP)) {
759
44
            column_ids.insert(field_schema->column_id);
760
44
            if (slot->is_predicate()) {
761
0
                filter_column_ids.insert(field_schema->column_id);
762
0
            }
763
44
            continue;
764
44
        }
765
766
22
        const auto& all_access_paths = slot->all_access_paths();
767
22
        process_access_paths(field_schema, all_access_paths, column_ids);
768
769
22
        const auto& predicate_access_paths = slot->predicate_access_paths();
770
22
        if (!predicate_access_paths.empty()) {
771
12
            process_access_paths(field_schema, predicate_access_paths, filter_column_ids);
772
12
        }
773
22
    }
774
50
    return {std::move(column_ids), std::move(filter_column_ids)};
775
50
}
776
777
// ============================================================================
778
// IcebergParquetReader: _read_position_delete_file
779
// ============================================================================
780
Status IcebergParquetReader::_read_position_delete_file(const TFileRangeDesc* delete_range,
781
4
                                                        DeleteFile* position_delete) {
782
4
    ParquetReader parquet_delete_reader(get_profile(), get_scan_params(), *delete_range,
783
4
                                        READ_DELETE_FILE_BATCH_SIZE, &get_state()->timezone_obj(),
784
4
                                        get_io_ctx(), get_state(), _meta_cache);
785
    // The delete file range has size=-1 (read whole file). We must disable
786
    // row group filtering before init; otherwise _do_init_reader returns EndOfFile
787
    // when _filter_groups && _range_size < 0.
788
4
    ParquetInitContext delete_ctx;
789
4
    delete_ctx.filter_groups = false;
790
4
    delete_ctx.column_names = delete_file_col_names;
791
4
    delete_ctx.col_name_to_block_idx =
792
4
            const_cast<std::unordered_map<std::string, uint32_t>*>(&DELETE_COL_NAME_TO_BLOCK_IDX);
793
4
    RETURN_IF_ERROR(parquet_delete_reader.init_reader(&delete_ctx));
794
795
0
    const tparquet::FileMetaData* meta_data = parquet_delete_reader.get_meta_data();
796
0
    bool dictionary_coded = true;
797
0
    for (const auto& row_group : meta_data->row_groups) {
798
0
        const auto& column_chunk = row_group.columns[ICEBERG_FILE_PATH_INDEX];
799
0
        if (!(column_chunk.__isset.meta_data && has_dict_page(column_chunk.meta_data))) {
800
0
            dictionary_coded = false;
801
0
            break;
802
0
        }
803
0
    }
804
0
    DataTypePtr data_type_file_path = make_nullable(std::make_shared<DataTypeString>());
805
0
    DataTypePtr data_type_pos = make_nullable(std::make_shared<DataTypeInt64>());
806
0
    bool eof = false;
807
0
    while (!eof) {
808
0
        Block block = {
809
0
                dictionary_coded
810
0
                        ? ColumnWithTypeAndName {ColumnNullable::create(ColumnDictI32::create(),
811
0
                                                                        ColumnUInt8::create()),
812
0
                                                 data_type_file_path, ICEBERG_FILE_PATH}
813
0
                        : ColumnWithTypeAndName {data_type_file_path, ICEBERG_FILE_PATH},
814
815
0
                {data_type_pos, ICEBERG_ROW_POS}};
816
0
        size_t read_rows = 0;
817
0
        RETURN_IF_ERROR(parquet_delete_reader.get_next_block(&block, &read_rows, &eof));
818
819
0
        if (read_rows <= 0) {
820
0
            break;
821
0
        }
822
0
        RETURN_IF_ERROR(_gen_position_delete_file_range(block, position_delete, read_rows,
823
0
                                                        dictionary_coded));
824
0
    }
825
0
    return Status::OK();
826
0
};
827
828
// ============================================================================
829
// IcebergOrcReader: on_before_init_reader (ORC-specific schema matching)
830
// ============================================================================
831
// This format-specific setup mirrors the existing reader initialization sequence.
832
// NOLINTNEXTLINE(readability-function-cognitive-complexity,readability-function-size)
833
32
Status IcebergOrcReader::on_before_init_reader(ReaderInitContext* ctx) {
834
32
    _column_descs = ctx->column_descs;
835
32
    _fill_col_name_to_block_idx = ctx->col_name_to_block_idx;
836
32
    _file_format = Fileformat::ORC;
837
838
    // Get ORC file type first (available because _create_file_reader() already ran)
839
32
    const orc::Type* orc_type_ptr = nullptr;
840
32
    RETURN_IF_ERROR(this->get_file_type(&orc_type_ptr));
841
842
    // Build table_info_node by field_id or name matching.
843
    // This must happen BEFORE column classification so we can use children_column_exists
844
    // to check if a column exists in the file (by field ID, not name).
845
32
    if (!get_scan_params().__isset.history_schema_info ||
846
32
        get_scan_params().history_schema_info.empty()) [[unlikely]] {
847
2
        RETURN_IF_ERROR(BuildTableInfoUtil::by_orc_name(ctx->tuple_descriptor, orc_type_ptr,
848
2
                                                        ctx->table_info_node));
849
30
    } else {
850
30
        RETURN_IF_ERROR(BuildTableInfoUtil::by_orc_field_id_with_name_mapping(
851
30
                get_scan_params().history_schema_info.front().root_field, orc_type_ptr,
852
30
                ICEBERG_ORC_ATTRIBUTE, ctx->table_info_node,
853
30
                supports_iceberg_scan_semantics_v1(&get_scan_params())));
854
30
    }
855
856
32
    std::unordered_set<std::string> partition_col_names;
857
32
    if (ctx->range->__isset.columns_from_path_keys) {
858
0
        partition_col_names.insert(ctx->range->columns_from_path_keys.begin(),
859
0
                                   ctx->range->columns_from_path_keys.end());
860
0
    }
861
862
    // Single pass: classify columns, detect $row_id, handle partition fallback.
863
32
    bool has_partition_from_path = false;
864
38
    for (const auto& desc : *ctx->column_descs) {
865
38
        if (desc.category == ColumnCategory::SYNTHESIZED) {
866
0
            if (desc.name == BeConsts::ICEBERG_ROWID_COL) {
867
0
                this->register_synthesized_column_handler(
868
0
                        BeConsts::ICEBERG_ROWID_COL, [this](Block* block, size_t rows) -> Status {
869
0
                            return _fill_iceberg_row_id(block, rows);
870
0
                        });
871
0
                continue;
872
0
            } else if (desc.name.starts_with(BeConsts::GLOBAL_ROWID_COL)) {
873
0
                auto topn_row_id_column_iter = _create_topn_row_id_column_iterator();
874
0
                this->register_synthesized_column_handler(
875
0
                        desc.name,
876
0
                        [iter = std::move(topn_row_id_column_iter), this, &desc](
877
0
                                Block* block, size_t rows) -> Status {
878
0
                            return fill_topn_row_id(iter, desc.name, block, rows);
879
0
                        });
880
0
                continue;
881
0
            }
882
38
        } else if (desc.category == ColumnCategory::PARTITION_KEY) {
883
0
            bool has_partition_value = partition_col_names.contains(desc.name);
884
0
            bool exists_in_file = ctx->table_info_node->children_column_exists(desc.name);
885
0
            if (!has_partition_value || exists_in_file) {
886
0
                ctx->column_names.push_back(desc.name);
887
0
                continue;
888
0
            }
889
0
            has_partition_from_path = true;
890
38
        } else if (desc.category == ColumnCategory::REGULAR) {
891
38
            ctx->column_names.push_back(desc.name);
892
38
        } else if (desc.category == ColumnCategory::GENERATED) {
893
0
            _init_row_lineage_columns();
894
0
            if (desc.name == ROW_LINEAGE_ROW_ID) {
895
0
                ctx->column_names.push_back(desc.name);
896
0
                this->register_generated_column_handler(
897
0
                        ROW_LINEAGE_ROW_ID, [this](Block* block, size_t rows) -> Status {
898
0
                            return _fill_row_lineage_row_id(block, rows);
899
0
                        });
900
0
                continue;
901
0
            } else if (desc.name == ROW_LINEAGE_LAST_UPDATED_SEQ_NUMBER) {
902
0
                ctx->column_names.push_back(desc.name);
903
0
                this->register_generated_column_handler(
904
0
                        ROW_LINEAGE_LAST_UPDATED_SEQ_NUMBER,
905
0
                        [this](Block* block, size_t rows) -> Status {
906
0
                            return _fill_row_lineage_last_updated_sequence_number(block, rows);
907
0
                        });
908
0
                continue;
909
0
            }
910
0
        }
911
38
    }
912
913
32
    if (has_partition_from_path) {
914
0
        RETURN_IF_ERROR(_extract_partition_values(*ctx->range, ctx->tuple_descriptor,
915
0
                                                  _fill_partition_values,
916
0
                                                  &_fill_partition_value_is_null));
917
0
    }
918
919
32
    _all_required_col_names = ctx->column_names;
920
921
    // Create column IDs from ORC type
922
32
    auto column_id_result =
923
32
            _create_column_ids(orc_type_ptr, ctx->tuple_descriptor, ctx->table_info_node);
924
32
    ctx->column_ids = std::move(column_id_result.column_ids);
925
32
    ctx->filter_column_ids = std::move(column_id_result.filter_column_ids);
926
927
    // Build field_id -> block_column_name mapping for equality delete filtering.
928
38
    for (const auto* slot : ctx->tuple_descriptor->slots()) {
929
38
        _id_to_block_column_name.emplace(slot->col_unique_id(), slot->col_name());
930
38
    }
931
932
    // Process delete files (must happen before _do_init_reader so expand col IDs are included)
933
32
    RETURN_IF_ERROR(_init_row_filters());
934
935
    // Add expand column IDs for equality delete and remap expand column names
936
    // (matching master's behavior with __equality_delete_column__ prefix)
937
32
    const static std::string EQ_DELETE_PRE = "__equality_delete_column__";
938
32
    bool all_file_columns_have_field_ids = true;
939
104
    for (uint64_t i = 0; i < orc_type_ptr->getSubtypeCount(); ++i) {
940
72
        const orc::Type* sub_type = orc_type_ptr->getSubtype(i);
941
72
        if (!sub_type->hasAttributeKey(ICEBERG_ORC_ATTRIBUTE)) {
942
32
            all_file_columns_have_field_ids = false;
943
32
        }
944
72
    }
945
32
    const bool use_field_ids_for_hidden_keys =
946
32
            supports_iceberg_scan_semantics_v1(&get_scan_params())
947
32
                    ? orc_subtree_has_iceberg_id(orc_type_ptr, ICEBERG_ORC_ATTRIBUTE)
948
32
                    : all_file_columns_have_field_ids;
949
32
    const auto find_file_column_by_name = [&](const std::string& name) -> const orc::Type* {
950
0
        for (uint64_t j = 0; j < orc_type_ptr->getSubtypeCount(); ++j) {
951
0
            if (iequal(orc_type_ptr->getFieldName(j), name)) {
952
0
                return orc_type_ptr->getSubtype(j);
953
0
            }
954
0
        }
955
0
        return nullptr;
956
0
    };
957
958
32
    std::vector<std::string> new_expand_col_names;
959
32
    DORIS_CHECK(_expand_col_names.size() == _expand_col_field_ids.size());
960
32
    DORIS_CHECK(_expand_col_names.size() == _expand_columns.size());
961
60
    for (size_t i = 0; i < _expand_col_names.size(); ++i) {
962
28
        const auto& old_name = _expand_col_names[i];
963
28
        const int32_t field_id = _expand_col_field_ids[i];
964
965
28
        const orc::Type* file_column = nullptr;
966
28
        OrcEqualityFieldPath file_path;
967
28
        bool complete_file_path = false;
968
28
        if (use_field_ids_for_hidden_keys) {
969
16
            complete_file_path =
970
16
                    find_orc_equality_field_path_by_id(orc_type_ptr, field_id, &file_path);
971
16
            if (!complete_file_path && supports_iceberg_scan_semantics_v2(&get_scan_params())) {
972
8
                const auto table_path = _find_schema_field_path(field_id);
973
8
                if (!table_path.empty()) {
974
8
                    complete_file_path = find_orc_equality_field_prefix_by_id_path(
975
8
                            orc_type_ptr, table_path, &file_path);
976
8
                }
977
8
            }
978
16
            if (!file_path.fields.empty()) {
979
8
                file_column = file_path.fields.front();
980
8
            }
981
16
        } else {
982
12
            const auto table_path = _find_schema_field_path(field_id);
983
12
            if (!table_path.empty()) {
984
12
                complete_file_path = find_orc_equality_field_prefix_by_name_path(
985
12
                        orc_type_ptr, table_path, old_name, &file_path);
986
12
                if (!file_path.fields.empty()) {
987
10
                    file_column = file_path.fields.front();
988
10
                }
989
12
            } else {
990
0
                file_column = find_file_column_by_name(old_name);
991
0
                complete_file_path = file_column != nullptr;
992
0
            }
993
12
        }
994
995
28
        std::string file_col_name = old_name;
996
28
        std::string leaf_name = old_name;
997
28
        if (!file_path.fields.empty()) {
998
18
            file_col_name = file_path.names.front();
999
18
            leaf_name = file_path.names.back();
1000
18
        } else if (file_column != nullptr) {
1001
0
            for (uint64_t j = 0; j < orc_type_ptr->getSubtypeCount(); ++j) {
1002
0
                if (orc_type_ptr->getSubtype(j) == file_column) {
1003
0
                    file_col_name = orc_type_ptr->getFieldName(j);
1004
0
                    leaf_name = file_col_name;
1005
0
                    break;
1006
0
                }
1007
0
            }
1008
0
        }
1009
28
        std::string table_col_name = EQ_DELETE_PRE + std::to_string(field_id) + "_" + leaf_name;
1010
1011
28
        if (field_id >= 0) {
1012
28
            _id_to_block_column_name[field_id] = table_col_name;
1013
28
        }
1014
28
        _expand_columns[i].name = table_col_name;
1015
28
        if (file_column == nullptr) {
1016
10
            RETURN_IF_ERROR(_register_missing_equality_delete_column(field_id, table_col_name,
1017
10
                                                                     _expand_columns[i].type));
1018
            // The old data file predates this equality key. Keep it in the expand block so the
1019
            // synthesized-column hook can materialize its logical initial default before ORC's
1020
            // block-size checks. Adding it to column_names/table_info_node would mark it as an
1021
            // existing ORC child and make OrcReader read a column that is not present in the file.
1022
10
            new_expand_col_names.push_back(table_col_name);
1023
10
            continue;
1024
10
        }
1025
18
        new_expand_col_names.push_back(table_col_name);
1026
1027
18
        if (!complete_file_path) {
1028
2
            ColumnPtr missing_value;
1029
2
            RETURN_IF_ERROR(_create_missing_equality_delete_value(
1030
2
                    field_id, _expand_columns[i].type, file_path.fields.size(), &missing_value));
1031
2
            _nested_equality_delete_columns.push_back({
1032
2
                    .field_id = field_id,
1033
2
                    .block_name = table_col_name,
1034
2
                    .leaf_type = _expand_columns[i].type,
1035
2
                    .child_indexes = file_path.child_indexes,
1036
2
                    .missing_value = std::move(missing_value),
1037
2
            });
1038
2
            _expand_columns[i].type = make_nullable(convert_to_doris_type(file_column));
1039
2
            _expand_columns[i].column = _expand_columns[i].type->create_column();
1040
16
        } else if (!file_path.child_indexes.empty()) {
1041
12
            _nested_equality_delete_columns.push_back({
1042
12
                    .field_id = field_id,
1043
12
                    .block_name = table_col_name,
1044
12
                    .leaf_type = _expand_columns[i].type,
1045
12
                    .child_indexes = file_path.child_indexes,
1046
12
                    .missing_value = nullptr,
1047
12
            });
1048
12
            _expand_columns[i].type = make_nullable(convert_to_doris_type(file_column));
1049
12
            _expand_columns[i].column = _expand_columns[i].type->create_column();
1050
12
        }
1051
1052
18
        for (uint64_t column_id = file_column->getColumnId();
1053
50
             column_id <= file_column->getMaximumColumnId(); ++column_id) {
1054
32
            ctx->column_ids.insert(column_id);
1055
32
        }
1056
1057
18
        ctx->column_names.push_back(table_col_name);
1058
18
        ctx->table_info_node->add_children(table_col_name, file_col_name,
1059
18
                                           TableSchemaChangeHelper::ConstNode::get_instance());
1060
18
    }
1061
32
    _expand_col_names = std::move(new_expand_col_names);
1062
1063
32
    return Status::OK();
1064
32
}
1065
1066
// ============================================================================
1067
// IcebergOrcReader: _create_column_ids
1068
// ============================================================================
1069
ColumnIdResult IcebergOrcReader::_create_column_ids(
1070
        const orc::Type* orc_type, const TupleDescriptor* tuple_descriptor,
1071
46
        const std::shared_ptr<TableSchemaChangeHelper::Node>& table_info_node) {
1072
46
    std::unordered_map<int, const orc::Type*> iceberg_id_to_orc_type_map;
1073
218
    for (uint64_t i = 0; i < orc_type->getSubtypeCount(); ++i) {
1074
172
        const auto* orc_sub_type = orc_type->getSubtype(i);
1075
172
        if (!orc_sub_type) {
1076
0
            continue;
1077
0
        }
1078
172
        if (!orc_sub_type->hasAttributeKey(ICEBERG_ORC_ATTRIBUTE)) {
1079
34
            continue;
1080
34
        }
1081
138
        int iceberg_id = std::stoi(orc_sub_type->getAttributeValue(ICEBERG_ORC_ATTRIBUTE));
1082
138
        iceberg_id_to_orc_type_map[iceberg_id] = orc_sub_type;
1083
138
    }
1084
1085
46
    std::set<uint64_t> column_ids;
1086
46
    std::set<uint64_t> filter_column_ids;
1087
1088
46
    auto process_access_paths = [](const orc::Type* orc_field,
1089
46
                                   const std::vector<TColumnAccessPath>& access_paths,
1090
46
                                   std::set<uint64_t>& out_ids) {
1091
34
        process_nested_access_paths(
1092
34
                orc_field, access_paths, out_ids,
1093
34
                [](const orc::Type* type) { return type->getColumnId(); },
1094
34
                [](const orc::Type* type) { return type->getMaximumColumnId(); },
1095
34
                IcebergOrcNestedColumnUtils::extract_nested_column_ids);
1096
34
    };
1097
1098
    // The Iceberg schema-mapping root is a StructNode whose registered children are the real
1099
    // table columns. When present, resolve each column by name through it so the column-id set
1100
    // stays consistent with the schema-mapping decision (BY_ID or BY_NAME/name-mapping);
1101
    // otherwise fall back to matching by Iceberg field id.
1102
46
    const auto* struct_node =
1103
46
            dynamic_cast<const TableSchemaChangeHelper::StructNode*>(table_info_node.get());
1104
1105
68
    for (const auto* slot : tuple_descriptor->slots()) {
1106
68
        const orc::Type* orc_field = nullptr;
1107
68
        if (struct_node != nullptr) {
1108
            // Synthesized/metadata slots (e.g. the TopN global row-id or the $row_id column) are
1109
            // never registered as children, so check membership before querying: calling
1110
            // children_column_exists() on an unregistered name DCHECK-aborts in debug builds and
1111
            // throws std::out_of_range from .at() in release builds.
1112
42
            if (struct_node->get_children().contains(slot->col_name()) &&
1113
42
                struct_node->children_column_exists(slot->col_name())) {
1114
                // Select the physical child resolved by the shared schema-mapping pass. Hidden
1115
                // equality keys and projected columns must obey the same BY_NAME decision for
1116
                // partial-id ORC files.
1117
38
                const auto& file_column_name =
1118
38
                        struct_node->children_file_column_name(slot->col_name());
1119
46
                for (uint64_t i = 0; i < orc_type->getSubtypeCount(); ++i) {
1120
46
                    if (orc_type->getFieldName(i) == file_column_name) {
1121
38
                        orc_field = orc_type->getSubtype(i);
1122
38
                        break;
1123
38
                    }
1124
46
                }
1125
38
                DORIS_CHECK(orc_field != nullptr);
1126
38
            }
1127
42
        } else {
1128
26
            auto it = iceberg_id_to_orc_type_map.find(slot->col_unique_id());
1129
26
            if (it != iceberg_id_to_orc_type_map.end()) {
1130
26
                orc_field = it->second;
1131
26
            }
1132
26
        }
1133
68
        if (orc_field == nullptr) {
1134
4
            continue;
1135
4
        }
1136
1137
64
        if ((slot->col_type() != TYPE_STRUCT && slot->col_type() != TYPE_ARRAY &&
1138
64
             slot->col_type() != TYPE_MAP)) {
1139
42
            column_ids.insert(orc_field->getColumnId());
1140
42
            if (slot->is_predicate()) {
1141
0
                filter_column_ids.insert(orc_field->getColumnId());
1142
0
            }
1143
42
            continue;
1144
42
        }
1145
1146
22
        const auto& all_access_paths = slot->all_access_paths();
1147
22
        process_access_paths(orc_field, all_access_paths, column_ids);
1148
1149
22
        const auto& predicate_access_paths = slot->predicate_access_paths();
1150
22
        if (!predicate_access_paths.empty()) {
1151
12
            process_access_paths(orc_field, predicate_access_paths, filter_column_ids);
1152
12
        }
1153
22
    }
1154
1155
46
    return {std::move(column_ids), std::move(filter_column_ids)};
1156
46
}
1157
1158
// ============================================================================
1159
// IcebergOrcReader: _read_position_delete_file
1160
// ============================================================================
1161
Status IcebergOrcReader::_read_position_delete_file(const TFileRangeDesc* delete_range,
1162
0
                                                    DeleteFile* position_delete) {
1163
0
    OrcReader orc_delete_reader(get_profile(), get_state(), get_scan_params(), *delete_range,
1164
0
                                READ_DELETE_FILE_BATCH_SIZE, get_state()->timezone(), get_io_ctx(),
1165
0
                                _meta_cache);
1166
0
    OrcInitContext delete_ctx;
1167
0
    delete_ctx.column_names = delete_file_col_names;
1168
0
    delete_ctx.col_name_to_block_idx =
1169
0
            const_cast<std::unordered_map<std::string, uint32_t>*>(&DELETE_COL_NAME_TO_BLOCK_IDX);
1170
0
    RETURN_IF_ERROR(orc_delete_reader.init_reader(&delete_ctx));
1171
1172
0
    bool eof = false;
1173
0
    DataTypePtr data_type_file_path {new DataTypeString};
1174
0
    DataTypePtr data_type_pos {new DataTypeInt64};
1175
0
    while (!eof) {
1176
0
        Block block = {{data_type_file_path, ICEBERG_FILE_PATH}, {data_type_pos, ICEBERG_ROW_POS}};
1177
1178
0
        size_t read_rows = 0;
1179
0
        RETURN_IF_ERROR(orc_delete_reader.get_next_block(&block, &read_rows, &eof));
1180
1181
0
        RETURN_IF_ERROR(_gen_position_delete_file_range(block, position_delete, read_rows, false));
1182
0
    }
1183
0
    return Status::OK();
1184
0
}
1185
1186
} // namespace doris