Coverage Report

Created: 2026-09-17 22:12

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format_v2/table/iceberg_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 <memory>
21
#include <optional>
22
#include <string>
23
#include <unordered_map>
24
#include <vector>
25
26
#include "common/status.h"
27
#include "core/block/block.h"
28
#include "format/table/iceberg_delete_file_reader_helper.h"
29
#include "format_v2/file_reader.h"
30
#include "format_v2/table/iceberg_schema_utils.h"
31
#include "format_v2/table_reader.h"
32
#include "gen_cpp/PlanNodes_types.h"
33
34
namespace doris {
35
class Block;
36
class EqualityDeleteHashIndex;
37
struct DeleteFileDesc;
38
namespace io {
39
struct FileDescription;
40
struct FileSystemProperties;
41
} // namespace io
42
} // namespace doris
43
44
namespace doris::format::iceberg {
45
46
Status prepare_iceberg_initial_default_exprs(format::ColumnDefinition* column);
47
48
// Iceberg table-level reader.
49
// It reuses TableReader for split orchestration, dynamic partition pruning and table-block
50
// finalization, while composing a FileReader for physical data-file reads instead of inheriting
51
// from a concrete file-format reader.
52
class IcebergTableReader : public format::TableReader {
53
public:
54
142
    ~IcebergTableReader() override = default;
55
    static Status validate_variant_file_mappings(
56
            FileFormat format, const std::vector<format::ColumnMapping>& mappings);
57
117
    Status init(format::TableReadOptions&& options) override {
58
117
        RETURN_IF_ERROR(format::TableReader::init(std::move(options)));
59
117
        _mapper_options.mode = format::TableColumnMappingMode::BY_FIELD_ID;
60
117
        _mapper_options.reject_missing_required_field =
61
117
                supports_iceberg_scan_semantics_v2(_scan_params);
62
117
        return Status::OK();
63
117
    }
64
65
    Status prepare_split(const format::SplitReadOptions& options) override;
66
    Status annotate_projected_column(const TFileScanSlotInfo& slot_info,
67
                                     format::ProjectedColumnBuildContext* context,
68
                                     format::ColumnDefinition* column) const override;
69
    std::string debug_string() const override;
70
299
    format::TableColumnMappingMode mapping_mode() const override {
71
299
        const bool has_field_ids = supports_iceberg_scan_semantics_v1(_scan_params)
72
299
                                           ? schema_has_any_field_id(_data_reader.file_schema)
73
299
                                           : schema_has_all_field_ids(_data_reader.file_schema);
74
299
        if (!_data_reader.file_schema.empty() && has_field_ids) {
75
254
            return format::TableColumnMappingMode::BY_FIELD_ID;
76
254
        }
77
45
        if (!_data_reader.file_schema.empty() && supports_iceberg_scan_semantics_v2(_scan_params) &&
78
45
            !_scan_has_any_authoritative_name_mapping()) {
79
            // ID-less migrated files are name-readable only while Iceberg's explicit default name
80
            // mapping exists; current names must not resurrect file fields after it is removed.
81
2
            return format::TableColumnMappingMode::BY_FIELD_ID;
82
2
        }
83
43
        return format::TableColumnMappingMode::BY_NAME;
84
45
    }
85
86
protected:
87
    // Iceberg UUID uses the same 16-byte STRING/VARBINARY carrier in data, defaults and deletes.
88
161
    bool preserve_binary_uuid() const override { return true; }
89
90
    Status validate_file_mapping(const format::TableColumnMapper& mapper) const override;
91
92
144
    void configure_mapper_options(format::TableColumnMapperOptions* options) const override {
93
144
        options->enable_row_lineage_virtual_columns = true;
94
144
        options->enable_iceberg_metadata_virtual_columns = true;
95
144
        options->reject_missing_required_field = supports_iceberg_scan_semantics_v2(_scan_params);
96
144
        options->allow_idless_complex_wrapper_projection =
97
144
                supports_iceberg_scan_semantics_v1(_scan_params) && _format == FileFormat::PARQUET;
98
144
    }
99
100
    Status materialize_virtual_columns(Block* table_block) override;
101
102
    Status customize_file_scan_request(format::FileScanRequest* file_request) override;
103
104
    bool _supports_aggregate_pushdown(TPushAggOp::type agg_type) const override;
105
106
    Status _parse_deletion_vector_file(const TTableFormatFileDesc& t_desc, DeleteFileDesc* desc,
107
                                       bool* has_delete_file) override;
108
109
    Status _init_delete_predicates(const TTableFormatFileDesc& t_desc);
110
111
private:
112
    struct EqualityDeleteFilter;
113
    static constexpr int MIN_SUPPORT_DELETE_FILES_VERSION = 2;
114
    static constexpr int POSITION_DELETE = 1;
115
    static constexpr int EQUALITY_DELETE = 2;
116
    static constexpr int DELETION_VECTOR = 3;
117
118
    struct RowLineageColumns {
119
        int64_t first_row_id = -1;
120
        int64_t last_updated_sequence_number = -1;
121
    };
122
123
    static constexpr const char* ICEBERG_FILE_PATH = "file_path";
124
    static constexpr const char* ICEBERG_ROW_POS = "pos";
125
    static constexpr size_t ICEBERG_FILE_PATH_BLOCK_POSITION = 0;
126
    static constexpr size_t ICEBERG_ROW_POS_BLOCK_POSITION = 1;
127
128
    bool _scan_has_any_authoritative_name_mapping() const;
129
130
    class PositionDeleteRowsCollector final {
131
    public:
132
        using PositionDeleteFile = std::unordered_map<std::string, format::DeleteRows>;
133
134
        explicit PositionDeleteRowsCollector(PositionDeleteFile* rows_by_data_file);
135
136
        Status collect(const Block& block, size_t read_rows);
137
138
    private:
139
        PositionDeleteFile* _rows_by_data_file = nullptr;
140
    };
141
142
    static std::shared_ptr<io::FileSystemProperties> _delete_file_system_properties(
143
            const TFileScanRangeParams& scan_params);
144
145
    static std::unique_ptr<io::FileDescription> _delete_file_description(
146
            const TFileRangeDesc& range);
147
148
    std::string _data_file_path() const;
149
150
    // Append row position column to file scan request for position delete handling.
151
    Status _append_row_position_output_column(format::FileScanRequest* request);
152
    // Append equality delete predicates to file scan request based on the delete files in iceberg
153
    // params. DeleteVector and position delete files use the common DeleteRows path in TableReader.
154
    using EqualityDeleteColumnPath = std::vector<const format::ColumnDefinition*>;
155
    Status _append_equality_delete_predicates(format::FileScanRequest* request);
156
    Status _build_missing_equality_delete_key_expr(const EqualityDeleteFilter& filter,
157
                                                   size_t key_idx,
158
                                                   const EqualityDeleteColumnPath& data_path,
159
                                                   format::FileScanRequest* request,
160
                                                   VExprSPtr* key_expr);
161
    Status _find_equality_delete_data_field(const EqualityDeleteFilter& filter, size_t key_idx,
162
                                            EqualityDeleteColumnPath* data_path,
163
                                            bool* complete_path) const;
164
    Status _find_equality_delete_table_field(const EqualityDeleteFilter& filter, size_t key_idx,
165
                                             format::ColumnDefinition* table_field) const;
166
    void _append_equality_delete_row_count_carrier(format::FileScanRequest* request);
167
    std::string _delete_file_cache_key(const char* prefix, const std::string& path) const;
168
169
    Status _init_equality_delete_predicates(
170
            const std::vector<TIcebergDeleteFileDesc>& delete_files);
171
172
    // Read equality/position delete files.
173
    Status _create_delete_file_reader(const TIcebergDeleteFileDesc& delete_file,
174
                                      const TFileScanRangeParams& scan_params,
175
                                      IcebergDeleteFileIOContext* delete_io_ctx,
176
                                      std::unique_ptr<format::FileReader>* reader);
177
    Status _read_equality_delete_file(const TIcebergDeleteFileDesc& delete_file,
178
                                      const TFileScanRangeParams& scan_params,
179
                                      IcebergDeleteFileIOContext* delete_io_ctx);
180
    Status _load_equality_delete_file(const TIcebergDeleteFileDesc& delete_file,
181
                                      const TFileScanRangeParams& scan_params,
182
                                      IcebergDeleteFileIOContext* delete_io_ctx,
183
                                      EqualityDeleteFilter* result);
184
    Status _resolve_equality_delete_fields(const TIcebergDeleteFileDesc& delete_file,
185
                                           const std::vector<format::ColumnDefinition>& schema,
186
                                           std::vector<EqualityDeleteColumnPath>* delete_paths,
187
                                           EqualityDeleteFilter* result) const;
188
    Status _read_position_delete_file(const TIcebergDeleteFileDesc& delete_file,
189
                                      const TFileScanRangeParams& scan_params,
190
                                      IcebergDeleteFileIOContext* delete_io_ctx,
191
                                      PositionDeleteRowsCollector* collector);
192
193
    // Read position delete files and collect deleted row positions to update DeletePredicate.
194
    Status _init_position_delete_rows(const std::vector<TIcebergDeleteFileDesc>& delete_files);
195
196
    // Materialize row lineage virtual columns based on the position delete file.
197
    Status _materialize_iceberg_rowid(Block* table_block, size_t column_idx);
198
    Status _materialize_iceberg_file_path(Block* table_block, size_t column_idx);
199
    Status _materialize_iceberg_row_position(Block* table_block, size_t column_idx);
200
    Status _materialize_row_lineage_row_id(Block* table_block, size_t column_idx);
201
    Status _materialize_row_lineage_last_updated_sequence_number(Block* table_block,
202
                                                                 size_t column_idx);
203
204
    RowLineageColumns _row_lineage_columns;
205
    size_t _row_position_block_position = 0;
206
    std::optional<TIcebergFileDesc> _iceberg_params;
207
    bool _delete_predicates_initialized = false;
208
    format::DeleteRows _position_delete_rows_storage;
209
    struct EqualityDeleteFilter {
210
        std::vector<int> field_ids;
211
        // Delete-file names are retained for Iceberg tables imported from formats that did not
212
        // persist field ids. In BY_NAME mode they are the fallback binding key.
213
        std::vector<std::string> field_names;
214
        std::vector<DataTypePtr> key_types;
215
        Block delete_block;
216
        std::shared_ptr<const EqualityDeleteHashIndex> hash_index;
217
    };
218
    std::vector<EqualityDeleteFilter> _equality_delete_filters;
219
    // Scanner-shared cache supplied in SplitReadOptions. Parsed delete files outlive one data-file
220
    // split and can be reused by every split referencing the same delete file.
221
    ShardedKVCache* _split_cache = nullptr;
222
223
    bool _need_row_lineage_row_id() const;
224
    bool _need_iceberg_rowid() const;
225
    bool _need_iceberg_metadata() const;
226
};
227
228
} // namespace doris::format::iceberg