Coverage Report

Created: 2026-07-25 12:34

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/parquet/vparquet_column_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/parquet/vparquet_column_reader.h"
19
20
#include <gen_cpp/parquet_types.h>
21
#include <limits.h>
22
#include <sys/types.h>
23
24
#include <algorithm>
25
#include <utility>
26
27
#include "common/status.h"
28
#include "core/column/column.h"
29
#include "core/column/column_array.h"
30
#include "core/column/column_map.h"
31
#include "core/column/column_nullable.h"
32
#include "core/column/column_struct.h"
33
#include "core/data_type/data_type_array.h"
34
#include "core/data_type/data_type_map.h"
35
#include "core/data_type/data_type_nullable.h"
36
#include "core/data_type/data_type_struct.h"
37
#include "core/data_type/define_primitive_type.h"
38
#include "format/parquet/level_decoder.h"
39
#include "format/parquet/schema_desc.h"
40
#include "format/parquet/vparquet_column_chunk_reader.h"
41
#include "io/fs/tracing_file_reader.h"
42
#include "runtime/runtime_profile.h"
43
#include "util/url_coding.h"
44
45
namespace doris {
46
static Status build_initial_default_column(
47
        const std::optional<TableSchemaChangeHelper::InitialDefaultValue>& metadata,
48
1
        const DataTypePtr& type, size_t rows, ColumnPtr* column) {
49
1
    DORIS_CHECK(column != nullptr);
50
1
    if (!metadata.has_value()) {
51
0
        *column = nullptr;
52
0
        return Status::OK();
53
0
    }
54
1
    const auto nested_type = remove_nullable(type);
55
1
    Field value;
56
1
    if (metadata->is_base64 || nested_type->get_primitive_type() == TYPE_VARBINARY) {
57
0
        std::string decoded;
58
0
        if (!base64_decode(metadata->value, &decoded)) {
59
0
            return Status::InvalidArgument("Invalid Base64 Iceberg nested initial default");
60
0
        }
61
0
        value = nested_type->get_primitive_type() == TYPE_VARBINARY
62
0
                        ? Field::create_field<TYPE_VARBINARY>(StringView(decoded))
63
0
                        : Field::create_field<TYPE_STRING>(decoded);
64
        // StringView borrows decoded for payloads longer than 12 bytes. Build the owning column
65
        // before the decode buffer leaves scope (UUID/FIXED defaults routinely exceed that size).
66
0
        *column = type->create_column_const(rows, value)->convert_to_full_column_if_const();
67
0
        return Status::OK();
68
1
    } else {
69
1
        RETURN_IF_ERROR(nested_type->get_serde()->from_fe_string(metadata->value, value));
70
1
    }
71
1
    *column = type->create_column_const(rows, value)->convert_to_full_column_if_const();
72
1
    return Status::OK();
73
1
}
74
75
static void fill_struct_null_map(FieldSchema* field, NullMap& null_map,
76
                                 const std::vector<level_t>& rep_levels,
77
13
                                 const std::vector<level_t>& def_levels) {
78
13
    size_t num_levels = def_levels.size();
79
13
    DCHECK_EQ(num_levels, rep_levels.size());
80
13
    size_t origin_size = null_map.size();
81
13
    null_map.resize(origin_size + num_levels);
82
13
    size_t pos = origin_size;
83
30
    for (size_t i = 0; i < num_levels; ++i) {
84
        // skip the levels affect its ancestor or its descendants
85
17
        if (def_levels[i] < field->repeated_parent_def_level ||
86
17
            rep_levels[i] > field->repetition_level) {
87
0
            continue;
88
0
        }
89
17
        if (def_levels[i] >= field->definition_level) {
90
17
            null_map[pos++] = 0;
91
17
        } else {
92
0
            null_map[pos++] = 1;
93
0
        }
94
17
    }
95
13
    null_map.resize(pos);
96
13
}
97
98
static void fill_array_offset(FieldSchema* field, ColumnArray::Offsets64& offsets_data,
99
                              NullMap* null_map_ptr, const std::vector<level_t>& rep_levels,
100
2
                              const std::vector<level_t>& def_levels) {
101
2
    size_t num_levels = rep_levels.size();
102
2
    DCHECK_EQ(num_levels, def_levels.size());
103
2
    size_t origin_size = offsets_data.size();
104
2
    offsets_data.resize(origin_size + num_levels);
105
2
    if (null_map_ptr != nullptr) {
106
2
        null_map_ptr->resize(origin_size + num_levels);
107
2
    }
108
2
    size_t offset_pos = origin_size - 1;
109
8
    for (size_t i = 0; i < num_levels; ++i) {
110
        // skip the levels affect its ancestor or its descendants
111
6
        if (def_levels[i] < field->repeated_parent_def_level ||
112
6
            rep_levels[i] > field->repetition_level) {
113
0
            continue;
114
0
        }
115
6
        if (rep_levels[i] == field->repetition_level) {
116
4
            offsets_data[offset_pos]++;
117
4
            continue;
118
4
        }
119
2
        offset_pos++;
120
2
        offsets_data[offset_pos] = offsets_data[offset_pos - 1];
121
2
        if (def_levels[i] >= field->definition_level) {
122
2
            offsets_data[offset_pos]++;
123
2
        }
124
2
        if (def_levels[i] >= field->definition_level - 1) {
125
2
            (*null_map_ptr)[offset_pos] = 0;
126
2
        } else {
127
0
            (*null_map_ptr)[offset_pos] = 1;
128
0
        }
129
2
    }
130
2
    offsets_data.resize(offset_pos + 1);
131
2
    if (null_map_ptr != nullptr) {
132
2
        null_map_ptr->resize(offset_pos + 1);
133
2
    }
134
2
}
135
136
Status ParquetColumnReader::create(io::FileReaderSPtr file, FieldSchema* field,
137
                                   const tparquet::RowGroup& row_group, const RowRanges& row_ranges,
138
                                   const cctz::time_zone* ctz, io::IOContext* io_ctx,
139
                                   std::unique_ptr<ParquetColumnReader>& reader,
140
                                   size_t max_buf_size,
141
                                   std::unordered_map<int, tparquet::OffsetIndex>& col_offsets,
142
                                   RuntimeState* state, bool in_collection,
143
                                   const std::set<uint64_t>& column_ids,
144
140
                                   const std::set<uint64_t>& filter_column_ids) {
145
140
    size_t total_rows = row_group.num_rows;
146
140
    if (field->data_type->get_primitive_type() == TYPE_ARRAY) {
147
2
        std::unique_ptr<ParquetColumnReader> element_reader;
148
2
        RETURN_IF_ERROR(create(file, &field->children[0], row_group, row_ranges, ctz, io_ctx,
149
2
                               element_reader, max_buf_size, col_offsets, state, true, column_ids,
150
2
                               filter_column_ids));
151
2
        auto array_reader = ArrayColumnReader::create_unique(row_ranges, total_rows, ctz, io_ctx);
152
2
        element_reader->set_column_in_nested();
153
2
        RETURN_IF_ERROR(array_reader->init(std::move(element_reader), field));
154
2
        array_reader->_filter_column_ids = filter_column_ids;
155
2
        reader.reset(array_reader.release());
156
138
    } else if (field->data_type->get_primitive_type() == TYPE_MAP) {
157
0
        std::unique_ptr<ParquetColumnReader> key_reader;
158
0
        std::unique_ptr<ParquetColumnReader> value_reader;
159
160
0
        if (column_ids.empty() ||
161
0
            column_ids.find(field->children[0].get_column_id()) != column_ids.end()) {
162
            // Create key reader
163
0
            RETURN_IF_ERROR(create(file, &field->children[0], row_group, row_ranges, ctz, io_ctx,
164
0
                                   key_reader, max_buf_size, col_offsets, state, true, column_ids,
165
0
                                   filter_column_ids));
166
0
        } else {
167
0
            auto skip_reader = std::make_unique<SkipReadingReader>(row_ranges, total_rows, ctz,
168
0
                                                                   io_ctx, &field->children[0]);
169
0
            key_reader = std::move(skip_reader);
170
0
        }
171
172
0
        if (column_ids.empty() ||
173
0
            column_ids.find(field->children[1].get_column_id()) != column_ids.end()) {
174
            // Create value reader
175
0
            RETURN_IF_ERROR(create(file, &field->children[1], row_group, row_ranges, ctz, io_ctx,
176
0
                                   value_reader, max_buf_size, col_offsets, state, true, column_ids,
177
0
                                   filter_column_ids));
178
0
        } else {
179
0
            auto skip_reader = std::make_unique<SkipReadingReader>(row_ranges, total_rows, ctz,
180
0
                                                                   io_ctx, &field->children[0]);
181
0
            value_reader = std::move(skip_reader);
182
0
        }
183
184
0
        auto map_reader = MapColumnReader::create_unique(row_ranges, total_rows, ctz, io_ctx);
185
0
        key_reader->set_column_in_nested();
186
0
        value_reader->set_column_in_nested();
187
0
        RETURN_IF_ERROR(map_reader->init(std::move(key_reader), std::move(value_reader), field));
188
0
        map_reader->_filter_column_ids = filter_column_ids;
189
0
        reader.reset(map_reader.release());
190
138
    } else if (field->data_type->get_primitive_type() == TYPE_STRUCT) {
191
13
        std::unordered_map<std::string, std::unique_ptr<ParquetColumnReader>> child_readers;
192
13
        child_readers.reserve(field->children.size());
193
13
        int non_skip_reader_idx = -1;
194
41
        for (int i = 0; i < field->children.size(); ++i) {
195
28
            auto& child = field->children[i];
196
28
            std::unique_ptr<ParquetColumnReader> child_reader;
197
28
            if (column_ids.empty() || column_ids.find(child.get_column_id()) != column_ids.end()) {
198
24
                RETURN_IF_ERROR(create(file, &child, row_group, row_ranges, ctz, io_ctx,
199
24
                                       child_reader, max_buf_size, col_offsets, state,
200
24
                                       in_collection, column_ids, filter_column_ids));
201
24
                child_readers[child.name] = std::move(child_reader);
202
                // Record the first non-SkippingReader
203
24
                if (non_skip_reader_idx == -1) {
204
13
                    non_skip_reader_idx = i;
205
13
                }
206
24
            } else {
207
4
                auto skip_reader = std::make_unique<SkipReadingReader>(row_ranges, total_rows, ctz,
208
4
                                                                       io_ctx, &child);
209
4
                skip_reader->_filter_column_ids = filter_column_ids;
210
4
                child_readers[child.name] = std::move(skip_reader);
211
4
            }
212
28
            child_readers[child.name]->set_column_in_nested();
213
28
        }
214
        // If all children are SkipReadingReader, force the first child to call create
215
13
        if (non_skip_reader_idx == -1) {
216
0
            std::unique_ptr<ParquetColumnReader> child_reader;
217
0
            RETURN_IF_ERROR(create(file, &field->children[0], row_group, row_ranges, ctz, io_ctx,
218
0
                                   child_reader, max_buf_size, col_offsets, state, in_collection,
219
0
                                   column_ids, filter_column_ids));
220
0
            child_reader->set_column_in_nested();
221
0
            child_readers[field->children[0].name] = std::move(child_reader);
222
0
        }
223
13
        auto struct_reader = StructColumnReader::create_unique(row_ranges, total_rows, ctz, io_ctx);
224
13
        RETURN_IF_ERROR(struct_reader->init(std::move(child_readers), field));
225
13
        struct_reader->_filter_column_ids = filter_column_ids;
226
13
        reader.reset(struct_reader.release());
227
125
    } else {
228
125
        auto physical_index = field->physical_column_index;
229
125
        const tparquet::OffsetIndex* offset_index =
230
125
                col_offsets.find(physical_index) != col_offsets.end() ? &col_offsets[physical_index]
231
125
                                                                      : nullptr;
232
233
125
        const tparquet::ColumnChunk& chunk = row_group.columns[physical_index];
234
125
        if (in_collection) {
235
3
            if (offset_index == nullptr) {
236
3
                auto scalar_reader = ScalarColumnReader<true, false>::create_unique(
237
3
                        row_ranges, total_rows, chunk, offset_index, ctz, io_ctx);
238
239
3
                RETURN_IF_ERROR(scalar_reader->init(file, field, max_buf_size, state));
240
3
                scalar_reader->_filter_column_ids = filter_column_ids;
241
3
                reader.reset(scalar_reader.release());
242
3
            } else {
243
0
                auto scalar_reader = ScalarColumnReader<true, true>::create_unique(
244
0
                        row_ranges, total_rows, chunk, offset_index, ctz, io_ctx);
245
246
0
                RETURN_IF_ERROR(scalar_reader->init(file, field, max_buf_size, state));
247
0
                scalar_reader->_filter_column_ids = filter_column_ids;
248
0
                reader.reset(scalar_reader.release());
249
0
            }
250
122
        } else {
251
122
            if (offset_index == nullptr) {
252
122
                auto scalar_reader = ScalarColumnReader<false, false>::create_unique(
253
122
                        row_ranges, total_rows, chunk, offset_index, ctz, io_ctx);
254
255
122
                RETURN_IF_ERROR(scalar_reader->init(file, field, max_buf_size, state));
256
122
                scalar_reader->_filter_column_ids = filter_column_ids;
257
122
                reader.reset(scalar_reader.release());
258
122
            } else {
259
0
                auto scalar_reader = ScalarColumnReader<false, true>::create_unique(
260
0
                        row_ranges, total_rows, chunk, offset_index, ctz, io_ctx);
261
262
0
                RETURN_IF_ERROR(scalar_reader->init(file, field, max_buf_size, state));
263
0
                scalar_reader->_filter_column_ids = filter_column_ids;
264
0
                reader.reset(scalar_reader.release());
265
0
            }
266
122
        }
267
125
    }
268
140
    return Status::OK();
269
140
}
270
271
void ParquetColumnReader::_generate_read_ranges(RowRange page_row_range,
272
253
                                                RowRanges* result_ranges) const {
273
253
    result_ranges->add(page_row_range);
274
253
    RowRanges::ranges_intersection(*result_ranges, _row_ranges, result_ranges);
275
253
}
276
277
template <bool IN_COLLECTION, bool OFFSET_INDEX>
278
Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::init(io::FileReaderSPtr file,
279
                                                             FieldSchema* field,
280
                                                             size_t max_buf_size,
281
125
                                                             RuntimeState* state) {
282
125
    _field_schema = field;
283
125
    auto& chunk_meta = _chunk_meta.meta_data;
284
125
    int64_t chunk_start = has_dict_page(chunk_meta) ? chunk_meta.dictionary_page_offset
285
125
                                                    : chunk_meta.data_page_offset;
286
125
    size_t chunk_len = chunk_meta.total_compressed_size;
287
125
    size_t prefetch_buffer_size = std::min(chunk_len, max_buf_size);
288
125
    if ((typeid_cast<doris::io::TracingFileReader*>(file.get()) &&
289
125
         typeid_cast<io::MergeRangeFileReader*>(
290
59
                 ((doris::io::TracingFileReader*)(file.get()))->inner_reader().get())) ||
291
125
        typeid_cast<io::MergeRangeFileReader*>(file.get())) {
292
        // turn off prefetch data when using MergeRangeFileReader
293
125
        prefetch_buffer_size = 0;
294
125
    }
295
125
    _stream_reader = std::make_unique<io::BufferedFileStreamReader>(file, chunk_start, chunk_len,
296
125
                                                                    prefetch_buffer_size);
297
125
    ParquetPageReadContext ctx(
298
125
            (state == nullptr) ? true : state->query_options().enable_parquet_file_page_cache);
299
300
125
    _chunk_reader = std::make_unique<ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>>(
301
125
            _stream_reader.get(), &_chunk_meta, field, _offset_index, _total_rows, _io_ctx, ctx);
302
125
    RETURN_IF_ERROR(_chunk_reader->init());
303
125
    return Status::OK();
304
125
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb1EE4initESt10shared_ptrINS_2io10FileReaderEEPNS_11FieldSchemaEmPNS_12RuntimeStateE
_ZN5doris18ScalarColumnReaderILb1ELb0EE4initESt10shared_ptrINS_2io10FileReaderEEPNS_11FieldSchemaEmPNS_12RuntimeStateE
Line
Count
Source
281
3
                                                             RuntimeState* state) {
282
3
    _field_schema = field;
283
3
    auto& chunk_meta = _chunk_meta.meta_data;
284
3
    int64_t chunk_start = has_dict_page(chunk_meta) ? chunk_meta.dictionary_page_offset
285
3
                                                    : chunk_meta.data_page_offset;
286
3
    size_t chunk_len = chunk_meta.total_compressed_size;
287
3
    size_t prefetch_buffer_size = std::min(chunk_len, max_buf_size);
288
3
    if ((typeid_cast<doris::io::TracingFileReader*>(file.get()) &&
289
3
         typeid_cast<io::MergeRangeFileReader*>(
290
0
                 ((doris::io::TracingFileReader*)(file.get()))->inner_reader().get())) ||
291
3
        typeid_cast<io::MergeRangeFileReader*>(file.get())) {
292
        // turn off prefetch data when using MergeRangeFileReader
293
3
        prefetch_buffer_size = 0;
294
3
    }
295
3
    _stream_reader = std::make_unique<io::BufferedFileStreamReader>(file, chunk_start, chunk_len,
296
3
                                                                    prefetch_buffer_size);
297
3
    ParquetPageReadContext ctx(
298
3
            (state == nullptr) ? true : state->query_options().enable_parquet_file_page_cache);
299
300
3
    _chunk_reader = std::make_unique<ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>>(
301
3
            _stream_reader.get(), &_chunk_meta, field, _offset_index, _total_rows, _io_ctx, ctx);
302
3
    RETURN_IF_ERROR(_chunk_reader->init());
303
3
    return Status::OK();
304
3
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb1EE4initESt10shared_ptrINS_2io10FileReaderEEPNS_11FieldSchemaEmPNS_12RuntimeStateE
_ZN5doris18ScalarColumnReaderILb0ELb0EE4initESt10shared_ptrINS_2io10FileReaderEEPNS_11FieldSchemaEmPNS_12RuntimeStateE
Line
Count
Source
281
122
                                                             RuntimeState* state) {
282
122
    _field_schema = field;
283
122
    auto& chunk_meta = _chunk_meta.meta_data;
284
122
    int64_t chunk_start = has_dict_page(chunk_meta) ? chunk_meta.dictionary_page_offset
285
122
                                                    : chunk_meta.data_page_offset;
286
122
    size_t chunk_len = chunk_meta.total_compressed_size;
287
122
    size_t prefetch_buffer_size = std::min(chunk_len, max_buf_size);
288
122
    if ((typeid_cast<doris::io::TracingFileReader*>(file.get()) &&
289
122
         typeid_cast<io::MergeRangeFileReader*>(
290
59
                 ((doris::io::TracingFileReader*)(file.get()))->inner_reader().get())) ||
291
122
        typeid_cast<io::MergeRangeFileReader*>(file.get())) {
292
        // turn off prefetch data when using MergeRangeFileReader
293
122
        prefetch_buffer_size = 0;
294
122
    }
295
122
    _stream_reader = std::make_unique<io::BufferedFileStreamReader>(file, chunk_start, chunk_len,
296
122
                                                                    prefetch_buffer_size);
297
122
    ParquetPageReadContext ctx(
298
122
            (state == nullptr) ? true : state->query_options().enable_parquet_file_page_cache);
299
300
122
    _chunk_reader = std::make_unique<ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>>(
301
122
            _stream_reader.get(), &_chunk_meta, field, _offset_index, _total_rows, _io_ctx, ctx);
302
122
    RETURN_IF_ERROR(_chunk_reader->init());
303
122
    return Status::OK();
304
122
}
305
306
template <bool IN_COLLECTION, bool OFFSET_INDEX>
307
223
Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::_skip_values(size_t num_values) {
308
223
    if (num_values == 0) {
309
121
        return Status::OK();
310
121
    }
311
102
    if (_chunk_reader->max_def_level() > 0) {
312
102
        LevelDecoder& def_decoder = _chunk_reader->def_level_decoder();
313
102
        size_t skipped = 0;
314
102
        size_t null_size = 0;
315
102
        size_t nonnull_size = 0;
316
217
        while (skipped < num_values) {
317
115
            level_t def_level = -1;
318
115
            size_t loop_skip = def_decoder.get_next_run(&def_level, num_values - skipped);
319
115
            if (loop_skip == 0) {
320
0
                std::stringstream ss;
321
0
                auto& bit_reader = def_decoder.rle_decoder().bit_reader();
322
0
                ss << "def_decoder buffer (hex): ";
323
0
                for (size_t i = 0; i < bit_reader.max_bytes(); ++i) {
324
0
                    ss << std::hex << std::setw(2) << std::setfill('0')
325
0
                       << static_cast<int>(bit_reader.buffer()[i]) << " ";
326
0
                }
327
0
                LOG(WARNING) << ss.str();
328
0
                return Status::InternalError("Failed to decode definition level.");
329
0
            }
330
115
            if (def_level < _field_schema->definition_level) {
331
8
                null_size += loop_skip;
332
107
            } else {
333
107
                nonnull_size += loop_skip;
334
107
            }
335
115
            skipped += loop_skip;
336
115
        }
337
102
        if (null_size > 0) {
338
5
            RETURN_IF_ERROR(_chunk_reader->skip_values(null_size, false));
339
5
        }
340
102
        if (nonnull_size > 0) {
341
101
            RETURN_IF_ERROR(_chunk_reader->skip_values(nonnull_size, true));
342
101
        }
343
102
    } else {
344
0
        RETURN_IF_ERROR(_chunk_reader->skip_values(num_values));
345
0
    }
346
102
    return Status::OK();
347
102
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb1EE12_skip_valuesEm
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb0EE12_skip_valuesEm
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb1EE12_skip_valuesEm
_ZN5doris18ScalarColumnReaderILb0ELb0EE12_skip_valuesEm
Line
Count
Source
307
223
Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::_skip_values(size_t num_values) {
308
223
    if (num_values == 0) {
309
121
        return Status::OK();
310
121
    }
311
102
    if (_chunk_reader->max_def_level() > 0) {
312
102
        LevelDecoder& def_decoder = _chunk_reader->def_level_decoder();
313
102
        size_t skipped = 0;
314
102
        size_t null_size = 0;
315
102
        size_t nonnull_size = 0;
316
217
        while (skipped < num_values) {
317
115
            level_t def_level = -1;
318
115
            size_t loop_skip = def_decoder.get_next_run(&def_level, num_values - skipped);
319
115
            if (loop_skip == 0) {
320
0
                std::stringstream ss;
321
0
                auto& bit_reader = def_decoder.rle_decoder().bit_reader();
322
0
                ss << "def_decoder buffer (hex): ";
323
0
                for (size_t i = 0; i < bit_reader.max_bytes(); ++i) {
324
0
                    ss << std::hex << std::setw(2) << std::setfill('0')
325
0
                       << static_cast<int>(bit_reader.buffer()[i]) << " ";
326
0
                }
327
0
                LOG(WARNING) << ss.str();
328
0
                return Status::InternalError("Failed to decode definition level.");
329
0
            }
330
115
            if (def_level < _field_schema->definition_level) {
331
8
                null_size += loop_skip;
332
107
            } else {
333
107
                nonnull_size += loop_skip;
334
107
            }
335
115
            skipped += loop_skip;
336
115
        }
337
102
        if (null_size > 0) {
338
5
            RETURN_IF_ERROR(_chunk_reader->skip_values(null_size, false));
339
5
        }
340
102
        if (nonnull_size > 0) {
341
101
            RETURN_IF_ERROR(_chunk_reader->skip_values(nonnull_size, true));
342
101
        }
343
102
    } else {
344
0
        RETURN_IF_ERROR(_chunk_reader->skip_values(num_values));
345
0
    }
346
102
    return Status::OK();
347
102
}
348
349
template <bool IN_COLLECTION, bool OFFSET_INDEX>
350
Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::_read_values(size_t num_values,
351
                                                                     ColumnPtr& doris_column,
352
                                                                     DataTypePtr& type,
353
                                                                     FilterMap& filter_map,
354
223
                                                                     bool is_dict_filter) {
355
223
    if (num_values == 0) {
356
0
        return Status::OK();
357
0
    }
358
223
    MutableColumnPtr data_column;
359
223
    std::vector<uint16_t> null_map;
360
223
    NullMap* map_data_column = nullptr;
361
223
    doris_column = IColumn::mutate(std::move(doris_column));
362
223
    if (is_column_nullable(*doris_column)) {
363
221
        SCOPED_RAW_TIMER(&_decode_null_map_time);
364
221
        auto mutable_column = doris_column->assert_mutable();
365
221
        auto* nullable_column = assert_cast<ColumnNullable*>(mutable_column.get());
366
367
221
        data_column = nullable_column->get_nested_column_ptr();
368
221
        map_data_column = &(nullable_column->get_null_map_data());
369
221
        if (_chunk_reader->max_def_level() > 0) {
370
167
            LevelDecoder& def_decoder = _chunk_reader->def_level_decoder();
371
167
            size_t has_read = 0;
372
167
            bool prev_is_null = true;
373
335
            while (has_read < num_values) {
374
168
                level_t def_level;
375
168
                size_t loop_read = def_decoder.get_next_run(&def_level, num_values - has_read);
376
168
                if (loop_read == 0) {
377
0
                    std::stringstream ss;
378
0
                    auto& bit_reader = def_decoder.rle_decoder().bit_reader();
379
0
                    ss << "def_decoder buffer (hex): ";
380
0
                    for (size_t i = 0; i < bit_reader.max_bytes(); ++i) {
381
0
                        ss << std::hex << std::setw(2) << std::setfill('0')
382
0
                           << static_cast<int>(bit_reader.buffer()[i]) << " ";
383
0
                    }
384
0
                    LOG(WARNING) << ss.str();
385
0
                    return Status::InternalError("Failed to decode definition level.");
386
0
                }
387
388
168
                bool is_null = def_level < _field_schema->definition_level;
389
168
                if (!(prev_is_null ^ is_null)) {
390
46
                    null_map.emplace_back(0);
391
46
                }
392
168
                size_t remaining = loop_read;
393
168
                while (remaining > USHRT_MAX) {
394
0
                    null_map.emplace_back(USHRT_MAX);
395
0
                    null_map.emplace_back(0);
396
0
                    remaining -= USHRT_MAX;
397
0
                }
398
168
                null_map.emplace_back((u_short)remaining);
399
168
                prev_is_null = is_null;
400
168
                has_read += loop_read;
401
168
            }
402
167
        }
403
221
    } else {
404
2
        if (_chunk_reader->max_def_level() > 0) {
405
0
            return Status::Corruption("Not nullable column has null values in parquet file");
406
0
        }
407
2
        data_column = doris_column->assert_mutable();
408
2
    }
409
223
    if (null_map.size() == 0) {
410
56
        size_t remaining = num_values;
411
56
        while (remaining > USHRT_MAX) {
412
0
            null_map.emplace_back(USHRT_MAX);
413
0
            null_map.emplace_back(0);
414
0
            remaining -= USHRT_MAX;
415
0
        }
416
56
        null_map.emplace_back((u_short)remaining);
417
56
    }
418
223
    ColumnSelectVector select_vector;
419
223
    {
420
223
        SCOPED_RAW_TIMER(&_decode_null_map_time);
421
223
        RETURN_IF_ERROR(select_vector.init(null_map, num_values, map_data_column, &filter_map,
422
223
                                           _filter_map_index));
423
223
        _filter_map_index += num_values;
424
223
    }
425
0
    return _chunk_reader->decode_values(data_column, type, select_vector, is_dict_filter);
426
223
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb1EE12_read_valuesEmRNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEb
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb0EE12_read_valuesEmRNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEb
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb1EE12_read_valuesEmRNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEb
_ZN5doris18ScalarColumnReaderILb0ELb0EE12_read_valuesEmRNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEb
Line
Count
Source
354
223
                                                                     bool is_dict_filter) {
355
223
    if (num_values == 0) {
356
0
        return Status::OK();
357
0
    }
358
223
    MutableColumnPtr data_column;
359
223
    std::vector<uint16_t> null_map;
360
223
    NullMap* map_data_column = nullptr;
361
223
    doris_column = IColumn::mutate(std::move(doris_column));
362
223
    if (is_column_nullable(*doris_column)) {
363
221
        SCOPED_RAW_TIMER(&_decode_null_map_time);
364
221
        auto mutable_column = doris_column->assert_mutable();
365
221
        auto* nullable_column = assert_cast<ColumnNullable*>(mutable_column.get());
366
367
221
        data_column = nullable_column->get_nested_column_ptr();
368
221
        map_data_column = &(nullable_column->get_null_map_data());
369
221
        if (_chunk_reader->max_def_level() > 0) {
370
167
            LevelDecoder& def_decoder = _chunk_reader->def_level_decoder();
371
167
            size_t has_read = 0;
372
167
            bool prev_is_null = true;
373
335
            while (has_read < num_values) {
374
168
                level_t def_level;
375
168
                size_t loop_read = def_decoder.get_next_run(&def_level, num_values - has_read);
376
168
                if (loop_read == 0) {
377
0
                    std::stringstream ss;
378
0
                    auto& bit_reader = def_decoder.rle_decoder().bit_reader();
379
0
                    ss << "def_decoder buffer (hex): ";
380
0
                    for (size_t i = 0; i < bit_reader.max_bytes(); ++i) {
381
0
                        ss << std::hex << std::setw(2) << std::setfill('0')
382
0
                           << static_cast<int>(bit_reader.buffer()[i]) << " ";
383
0
                    }
384
0
                    LOG(WARNING) << ss.str();
385
0
                    return Status::InternalError("Failed to decode definition level.");
386
0
                }
387
388
168
                bool is_null = def_level < _field_schema->definition_level;
389
168
                if (!(prev_is_null ^ is_null)) {
390
46
                    null_map.emplace_back(0);
391
46
                }
392
168
                size_t remaining = loop_read;
393
168
                while (remaining > USHRT_MAX) {
394
0
                    null_map.emplace_back(USHRT_MAX);
395
0
                    null_map.emplace_back(0);
396
0
                    remaining -= USHRT_MAX;
397
0
                }
398
168
                null_map.emplace_back((u_short)remaining);
399
168
                prev_is_null = is_null;
400
168
                has_read += loop_read;
401
168
            }
402
167
        }
403
221
    } else {
404
2
        if (_chunk_reader->max_def_level() > 0) {
405
0
            return Status::Corruption("Not nullable column has null values in parquet file");
406
0
        }
407
2
        data_column = doris_column->assert_mutable();
408
2
    }
409
223
    if (null_map.size() == 0) {
410
56
        size_t remaining = num_values;
411
56
        while (remaining > USHRT_MAX) {
412
0
            null_map.emplace_back(USHRT_MAX);
413
0
            null_map.emplace_back(0);
414
0
            remaining -= USHRT_MAX;
415
0
        }
416
56
        null_map.emplace_back((u_short)remaining);
417
56
    }
418
223
    ColumnSelectVector select_vector;
419
223
    {
420
223
        SCOPED_RAW_TIMER(&_decode_null_map_time);
421
223
        RETURN_IF_ERROR(select_vector.init(null_map, num_values, map_data_column, &filter_map,
422
223
                                           _filter_map_index));
423
223
        _filter_map_index += num_values;
424
223
    }
425
0
    return _chunk_reader->decode_values(data_column, type, select_vector, is_dict_filter);
426
223
}
427
428
/**
429
 * Load the nested column data of complex type.
430
 * A row of complex type may be stored across two(or more) pages, and the parameter `align_rows` indicates that
431
 * whether the reader should read the remaining value of the last row in previous page.
432
 */
433
template <bool IN_COLLECTION, bool OFFSET_INDEX>
434
Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::_read_nested_column(
435
        ColumnPtr& doris_column, DataTypePtr& type, FilterMap& filter_map, size_t batch_size,
436
15
        size_t* read_rows, bool* eof, bool is_dict_filter) {
437
15
    _rep_levels.clear();
438
15
    _def_levels.clear();
439
440
    // Handle nullable columns
441
15
    MutableColumnPtr data_column;
442
15
    NullMap* map_data_column = nullptr;
443
15
    doris_column = IColumn::mutate(std::move(doris_column));
444
15
    if (is_column_nullable(*doris_column)) {
445
14
        SCOPED_RAW_TIMER(&_decode_null_map_time);
446
14
        auto mutable_column = doris_column->assert_mutable();
447
14
        auto* nullable_column = assert_cast<ColumnNullable*>(mutable_column.get());
448
14
        data_column = nullable_column->get_nested_column_ptr();
449
14
        map_data_column = &(nullable_column->get_null_map_data());
450
14
    } else {
451
1
        if (_field_schema->data_type->is_nullable()) {
452
0
            return Status::Corruption("Not nullable column has null values in parquet file");
453
0
        }
454
1
        data_column = doris_column->assert_mutable();
455
1
    }
456
457
15
    std::vector<uint16_t> null_map;
458
15
    std::unordered_set<size_t> ancestor_null_indices;
459
15
    std::vector<uint8_t> nested_filter_map_data;
460
461
15
    auto read_and_fill_data = [&](size_t before_rep_level_sz, size_t filter_map_index) {
462
15
        RETURN_IF_ERROR(_chunk_reader->fill_def(_def_levels));
463
15
        std::unique_ptr<FilterMap> nested_filter_map = std::make_unique<FilterMap>();
464
15
        if (filter_map.has_filter()) {
465
0
            RETURN_IF_ERROR(gen_filter_map(filter_map, filter_map_index, before_rep_level_sz,
466
0
                                           _rep_levels.size(), nested_filter_map_data,
467
0
                                           &nested_filter_map));
468
0
        }
469
470
15
        null_map.clear();
471
15
        ancestor_null_indices.clear();
472
15
        RETURN_IF_ERROR(gen_nested_null_map(before_rep_level_sz, _rep_levels.size(), null_map,
473
15
                                            ancestor_null_indices));
474
475
15
        ColumnSelectVector select_vector;
476
15
        {
477
15
            SCOPED_RAW_TIMER(&_decode_null_map_time);
478
15
            RETURN_IF_ERROR(select_vector.init(
479
15
                    null_map,
480
15
                    _rep_levels.size() - before_rep_level_sz - ancestor_null_indices.size(),
481
15
                    map_data_column, nested_filter_map.get(), 0, &ancestor_null_indices));
482
15
        }
483
484
15
        RETURN_IF_ERROR(
485
15
                _chunk_reader->decode_values(data_column, type, select_vector, is_dict_filter));
486
15
        if (ancestor_null_indices.size() != 0) {
487
0
            RETURN_IF_ERROR(_chunk_reader->skip_values(ancestor_null_indices.size(), false));
488
0
        }
489
15
        if (filter_map.has_filter()) {
490
0
            auto new_rep_sz = before_rep_level_sz;
491
0
            for (size_t idx = before_rep_level_sz; idx < _rep_levels.size(); idx++) {
492
0
                if (nested_filter_map_data[idx - before_rep_level_sz]) {
493
0
                    _rep_levels[new_rep_sz] = _rep_levels[idx];
494
0
                    _def_levels[new_rep_sz] = _def_levels[idx];
495
0
                    new_rep_sz++;
496
0
                }
497
0
            }
498
0
            _rep_levels.resize(new_rep_sz);
499
0
            _def_levels.resize(new_rep_sz);
500
0
        }
501
15
        return Status::OK();
502
15
    };
Unexecuted instantiation: _ZZN5doris18ScalarColumnReaderILb1ELb1EE19_read_nested_columnERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEmPmPbbENKUlmmE_clEmm
_ZZN5doris18ScalarColumnReaderILb1ELb0EE19_read_nested_columnERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEmPmPbbENKUlmmE_clEmm
Line
Count
Source
461
3
    auto read_and_fill_data = [&](size_t before_rep_level_sz, size_t filter_map_index) {
462
3
        RETURN_IF_ERROR(_chunk_reader->fill_def(_def_levels));
463
3
        std::unique_ptr<FilterMap> nested_filter_map = std::make_unique<FilterMap>();
464
3
        if (filter_map.has_filter()) {
465
0
            RETURN_IF_ERROR(gen_filter_map(filter_map, filter_map_index, before_rep_level_sz,
466
0
                                           _rep_levels.size(), nested_filter_map_data,
467
0
                                           &nested_filter_map));
468
0
        }
469
470
3
        null_map.clear();
471
3
        ancestor_null_indices.clear();
472
3
        RETURN_IF_ERROR(gen_nested_null_map(before_rep_level_sz, _rep_levels.size(), null_map,
473
3
                                            ancestor_null_indices));
474
475
3
        ColumnSelectVector select_vector;
476
3
        {
477
3
            SCOPED_RAW_TIMER(&_decode_null_map_time);
478
3
            RETURN_IF_ERROR(select_vector.init(
479
3
                    null_map,
480
3
                    _rep_levels.size() - before_rep_level_sz - ancestor_null_indices.size(),
481
3
                    map_data_column, nested_filter_map.get(), 0, &ancestor_null_indices));
482
3
        }
483
484
3
        RETURN_IF_ERROR(
485
3
                _chunk_reader->decode_values(data_column, type, select_vector, is_dict_filter));
486
3
        if (ancestor_null_indices.size() != 0) {
487
0
            RETURN_IF_ERROR(_chunk_reader->skip_values(ancestor_null_indices.size(), false));
488
0
        }
489
3
        if (filter_map.has_filter()) {
490
0
            auto new_rep_sz = before_rep_level_sz;
491
0
            for (size_t idx = before_rep_level_sz; idx < _rep_levels.size(); idx++) {
492
0
                if (nested_filter_map_data[idx - before_rep_level_sz]) {
493
0
                    _rep_levels[new_rep_sz] = _rep_levels[idx];
494
0
                    _def_levels[new_rep_sz] = _def_levels[idx];
495
0
                    new_rep_sz++;
496
0
                }
497
0
            }
498
0
            _rep_levels.resize(new_rep_sz);
499
0
            _def_levels.resize(new_rep_sz);
500
0
        }
501
3
        return Status::OK();
502
3
    };
Unexecuted instantiation: _ZZN5doris18ScalarColumnReaderILb0ELb1EE19_read_nested_columnERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEmPmPbbENKUlmmE_clEmm
_ZZN5doris18ScalarColumnReaderILb0ELb0EE19_read_nested_columnERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEmPmPbbENKUlmmE_clEmm
Line
Count
Source
461
12
    auto read_and_fill_data = [&](size_t before_rep_level_sz, size_t filter_map_index) {
462
12
        RETURN_IF_ERROR(_chunk_reader->fill_def(_def_levels));
463
12
        std::unique_ptr<FilterMap> nested_filter_map = std::make_unique<FilterMap>();
464
12
        if (filter_map.has_filter()) {
465
0
            RETURN_IF_ERROR(gen_filter_map(filter_map, filter_map_index, before_rep_level_sz,
466
0
                                           _rep_levels.size(), nested_filter_map_data,
467
0
                                           &nested_filter_map));
468
0
        }
469
470
12
        null_map.clear();
471
12
        ancestor_null_indices.clear();
472
12
        RETURN_IF_ERROR(gen_nested_null_map(before_rep_level_sz, _rep_levels.size(), null_map,
473
12
                                            ancestor_null_indices));
474
475
12
        ColumnSelectVector select_vector;
476
12
        {
477
12
            SCOPED_RAW_TIMER(&_decode_null_map_time);
478
12
            RETURN_IF_ERROR(select_vector.init(
479
12
                    null_map,
480
12
                    _rep_levels.size() - before_rep_level_sz - ancestor_null_indices.size(),
481
12
                    map_data_column, nested_filter_map.get(), 0, &ancestor_null_indices));
482
12
        }
483
484
12
        RETURN_IF_ERROR(
485
12
                _chunk_reader->decode_values(data_column, type, select_vector, is_dict_filter));
486
12
        if (ancestor_null_indices.size() != 0) {
487
0
            RETURN_IF_ERROR(_chunk_reader->skip_values(ancestor_null_indices.size(), false));
488
0
        }
489
12
        if (filter_map.has_filter()) {
490
0
            auto new_rep_sz = before_rep_level_sz;
491
0
            for (size_t idx = before_rep_level_sz; idx < _rep_levels.size(); idx++) {
492
0
                if (nested_filter_map_data[idx - before_rep_level_sz]) {
493
0
                    _rep_levels[new_rep_sz] = _rep_levels[idx];
494
0
                    _def_levels[new_rep_sz] = _def_levels[idx];
495
0
                    new_rep_sz++;
496
0
                }
497
0
            }
498
0
            _rep_levels.resize(new_rep_sz);
499
0
            _def_levels.resize(new_rep_sz);
500
0
        }
501
12
        return Status::OK();
502
12
    };
503
504
19
    while (_current_range_idx < _row_ranges.range_size()) {
505
15
        size_t left_row =
506
15
                std::max(_current_row_index, _row_ranges.get_range_from(_current_range_idx));
507
15
        size_t right_row = std::min(left_row + batch_size - *read_rows,
508
15
                                    (size_t)_row_ranges.get_range_to(_current_range_idx));
509
15
        _current_row_index = left_row;
510
15
        RETURN_IF_ERROR(_chunk_reader->seek_to_nested_row(left_row));
511
15
        size_t load_rows = 0;
512
15
        bool cross_page = false;
513
15
        size_t before_rep_level_sz = _rep_levels.size();
514
15
        RETURN_IF_ERROR(_chunk_reader->load_page_nested_rows(_rep_levels, right_row - left_row,
515
15
                                                             &load_rows, &cross_page));
516
15
        RETURN_IF_ERROR(read_and_fill_data(before_rep_level_sz, _filter_map_index));
517
15
        _filter_map_index += load_rows;
518
15
        while (cross_page) {
519
0
            before_rep_level_sz = _rep_levels.size();
520
0
            RETURN_IF_ERROR(_chunk_reader->load_cross_page_nested_row(_rep_levels, &cross_page));
521
0
            RETURN_IF_ERROR(read_and_fill_data(before_rep_level_sz, _filter_map_index - 1));
522
0
        }
523
15
        *read_rows += load_rows;
524
15
        _current_row_index += load_rows;
525
15
        _current_range_idx += (_current_row_index == _row_ranges.get_range_to(_current_range_idx));
526
15
        if (*read_rows == batch_size) {
527
11
            break;
528
11
        }
529
15
    }
530
15
    *eof = _current_range_idx == _row_ranges.range_size();
531
15
    return Status::OK();
532
15
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb1EE19_read_nested_columnERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEmPmPbb
_ZN5doris18ScalarColumnReaderILb1ELb0EE19_read_nested_columnERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEmPmPbb
Line
Count
Source
436
3
        size_t* read_rows, bool* eof, bool is_dict_filter) {
437
3
    _rep_levels.clear();
438
3
    _def_levels.clear();
439
440
    // Handle nullable columns
441
3
    MutableColumnPtr data_column;
442
3
    NullMap* map_data_column = nullptr;
443
3
    doris_column = IColumn::mutate(std::move(doris_column));
444
3
    if (is_column_nullable(*doris_column)) {
445
3
        SCOPED_RAW_TIMER(&_decode_null_map_time);
446
3
        auto mutable_column = doris_column->assert_mutable();
447
3
        auto* nullable_column = assert_cast<ColumnNullable*>(mutable_column.get());
448
3
        data_column = nullable_column->get_nested_column_ptr();
449
3
        map_data_column = &(nullable_column->get_null_map_data());
450
3
    } else {
451
0
        if (_field_schema->data_type->is_nullable()) {
452
0
            return Status::Corruption("Not nullable column has null values in parquet file");
453
0
        }
454
0
        data_column = doris_column->assert_mutable();
455
0
    }
456
457
3
    std::vector<uint16_t> null_map;
458
3
    std::unordered_set<size_t> ancestor_null_indices;
459
3
    std::vector<uint8_t> nested_filter_map_data;
460
461
3
    auto read_and_fill_data = [&](size_t before_rep_level_sz, size_t filter_map_index) {
462
3
        RETURN_IF_ERROR(_chunk_reader->fill_def(_def_levels));
463
3
        std::unique_ptr<FilterMap> nested_filter_map = std::make_unique<FilterMap>();
464
3
        if (filter_map.has_filter()) {
465
3
            RETURN_IF_ERROR(gen_filter_map(filter_map, filter_map_index, before_rep_level_sz,
466
3
                                           _rep_levels.size(), nested_filter_map_data,
467
3
                                           &nested_filter_map));
468
3
        }
469
470
3
        null_map.clear();
471
3
        ancestor_null_indices.clear();
472
3
        RETURN_IF_ERROR(gen_nested_null_map(before_rep_level_sz, _rep_levels.size(), null_map,
473
3
                                            ancestor_null_indices));
474
475
3
        ColumnSelectVector select_vector;
476
3
        {
477
3
            SCOPED_RAW_TIMER(&_decode_null_map_time);
478
3
            RETURN_IF_ERROR(select_vector.init(
479
3
                    null_map,
480
3
                    _rep_levels.size() - before_rep_level_sz - ancestor_null_indices.size(),
481
3
                    map_data_column, nested_filter_map.get(), 0, &ancestor_null_indices));
482
3
        }
483
484
3
        RETURN_IF_ERROR(
485
3
                _chunk_reader->decode_values(data_column, type, select_vector, is_dict_filter));
486
3
        if (ancestor_null_indices.size() != 0) {
487
3
            RETURN_IF_ERROR(_chunk_reader->skip_values(ancestor_null_indices.size(), false));
488
3
        }
489
3
        if (filter_map.has_filter()) {
490
3
            auto new_rep_sz = before_rep_level_sz;
491
3
            for (size_t idx = before_rep_level_sz; idx < _rep_levels.size(); idx++) {
492
3
                if (nested_filter_map_data[idx - before_rep_level_sz]) {
493
3
                    _rep_levels[new_rep_sz] = _rep_levels[idx];
494
3
                    _def_levels[new_rep_sz] = _def_levels[idx];
495
3
                    new_rep_sz++;
496
3
                }
497
3
            }
498
3
            _rep_levels.resize(new_rep_sz);
499
3
            _def_levels.resize(new_rep_sz);
500
3
        }
501
3
        return Status::OK();
502
3
    };
503
504
3
    while (_current_range_idx < _row_ranges.range_size()) {
505
3
        size_t left_row =
506
3
                std::max(_current_row_index, _row_ranges.get_range_from(_current_range_idx));
507
3
        size_t right_row = std::min(left_row + batch_size - *read_rows,
508
3
                                    (size_t)_row_ranges.get_range_to(_current_range_idx));
509
3
        _current_row_index = left_row;
510
3
        RETURN_IF_ERROR(_chunk_reader->seek_to_nested_row(left_row));
511
3
        size_t load_rows = 0;
512
3
        bool cross_page = false;
513
3
        size_t before_rep_level_sz = _rep_levels.size();
514
3
        RETURN_IF_ERROR(_chunk_reader->load_page_nested_rows(_rep_levels, right_row - left_row,
515
3
                                                             &load_rows, &cross_page));
516
3
        RETURN_IF_ERROR(read_and_fill_data(before_rep_level_sz, _filter_map_index));
517
3
        _filter_map_index += load_rows;
518
3
        while (cross_page) {
519
0
            before_rep_level_sz = _rep_levels.size();
520
0
            RETURN_IF_ERROR(_chunk_reader->load_cross_page_nested_row(_rep_levels, &cross_page));
521
0
            RETURN_IF_ERROR(read_and_fill_data(before_rep_level_sz, _filter_map_index - 1));
522
0
        }
523
3
        *read_rows += load_rows;
524
3
        _current_row_index += load_rows;
525
3
        _current_range_idx += (_current_row_index == _row_ranges.get_range_to(_current_range_idx));
526
3
        if (*read_rows == batch_size) {
527
3
            break;
528
3
        }
529
3
    }
530
3
    *eof = _current_range_idx == _row_ranges.range_size();
531
3
    return Status::OK();
532
3
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb1EE19_read_nested_columnERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEmPmPbb
_ZN5doris18ScalarColumnReaderILb0ELb0EE19_read_nested_columnERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERSt10shared_ptrIKNS_9IDataTypeEERNS_9FilterMapEmPmPbb
Line
Count
Source
436
12
        size_t* read_rows, bool* eof, bool is_dict_filter) {
437
12
    _rep_levels.clear();
438
12
    _def_levels.clear();
439
440
    // Handle nullable columns
441
12
    MutableColumnPtr data_column;
442
12
    NullMap* map_data_column = nullptr;
443
12
    doris_column = IColumn::mutate(std::move(doris_column));
444
12
    if (is_column_nullable(*doris_column)) {
445
11
        SCOPED_RAW_TIMER(&_decode_null_map_time);
446
11
        auto mutable_column = doris_column->assert_mutable();
447
11
        auto* nullable_column = assert_cast<ColumnNullable*>(mutable_column.get());
448
11
        data_column = nullable_column->get_nested_column_ptr();
449
11
        map_data_column = &(nullable_column->get_null_map_data());
450
11
    } else {
451
1
        if (_field_schema->data_type->is_nullable()) {
452
0
            return Status::Corruption("Not nullable column has null values in parquet file");
453
0
        }
454
1
        data_column = doris_column->assert_mutable();
455
1
    }
456
457
12
    std::vector<uint16_t> null_map;
458
12
    std::unordered_set<size_t> ancestor_null_indices;
459
12
    std::vector<uint8_t> nested_filter_map_data;
460
461
12
    auto read_and_fill_data = [&](size_t before_rep_level_sz, size_t filter_map_index) {
462
12
        RETURN_IF_ERROR(_chunk_reader->fill_def(_def_levels));
463
12
        std::unique_ptr<FilterMap> nested_filter_map = std::make_unique<FilterMap>();
464
12
        if (filter_map.has_filter()) {
465
12
            RETURN_IF_ERROR(gen_filter_map(filter_map, filter_map_index, before_rep_level_sz,
466
12
                                           _rep_levels.size(), nested_filter_map_data,
467
12
                                           &nested_filter_map));
468
12
        }
469
470
12
        null_map.clear();
471
12
        ancestor_null_indices.clear();
472
12
        RETURN_IF_ERROR(gen_nested_null_map(before_rep_level_sz, _rep_levels.size(), null_map,
473
12
                                            ancestor_null_indices));
474
475
12
        ColumnSelectVector select_vector;
476
12
        {
477
12
            SCOPED_RAW_TIMER(&_decode_null_map_time);
478
12
            RETURN_IF_ERROR(select_vector.init(
479
12
                    null_map,
480
12
                    _rep_levels.size() - before_rep_level_sz - ancestor_null_indices.size(),
481
12
                    map_data_column, nested_filter_map.get(), 0, &ancestor_null_indices));
482
12
        }
483
484
12
        RETURN_IF_ERROR(
485
12
                _chunk_reader->decode_values(data_column, type, select_vector, is_dict_filter));
486
12
        if (ancestor_null_indices.size() != 0) {
487
12
            RETURN_IF_ERROR(_chunk_reader->skip_values(ancestor_null_indices.size(), false));
488
12
        }
489
12
        if (filter_map.has_filter()) {
490
12
            auto new_rep_sz = before_rep_level_sz;
491
12
            for (size_t idx = before_rep_level_sz; idx < _rep_levels.size(); idx++) {
492
12
                if (nested_filter_map_data[idx - before_rep_level_sz]) {
493
12
                    _rep_levels[new_rep_sz] = _rep_levels[idx];
494
12
                    _def_levels[new_rep_sz] = _def_levels[idx];
495
12
                    new_rep_sz++;
496
12
                }
497
12
            }
498
12
            _rep_levels.resize(new_rep_sz);
499
12
            _def_levels.resize(new_rep_sz);
500
12
        }
501
12
        return Status::OK();
502
12
    };
503
504
16
    while (_current_range_idx < _row_ranges.range_size()) {
505
12
        size_t left_row =
506
12
                std::max(_current_row_index, _row_ranges.get_range_from(_current_range_idx));
507
12
        size_t right_row = std::min(left_row + batch_size - *read_rows,
508
12
                                    (size_t)_row_ranges.get_range_to(_current_range_idx));
509
12
        _current_row_index = left_row;
510
12
        RETURN_IF_ERROR(_chunk_reader->seek_to_nested_row(left_row));
511
12
        size_t load_rows = 0;
512
12
        bool cross_page = false;
513
12
        size_t before_rep_level_sz = _rep_levels.size();
514
12
        RETURN_IF_ERROR(_chunk_reader->load_page_nested_rows(_rep_levels, right_row - left_row,
515
12
                                                             &load_rows, &cross_page));
516
12
        RETURN_IF_ERROR(read_and_fill_data(before_rep_level_sz, _filter_map_index));
517
12
        _filter_map_index += load_rows;
518
12
        while (cross_page) {
519
0
            before_rep_level_sz = _rep_levels.size();
520
0
            RETURN_IF_ERROR(_chunk_reader->load_cross_page_nested_row(_rep_levels, &cross_page));
521
0
            RETURN_IF_ERROR(read_and_fill_data(before_rep_level_sz, _filter_map_index - 1));
522
0
        }
523
12
        *read_rows += load_rows;
524
12
        _current_row_index += load_rows;
525
12
        _current_range_idx += (_current_row_index == _row_ranges.get_range_to(_current_range_idx));
526
12
        if (*read_rows == batch_size) {
527
8
            break;
528
8
        }
529
12
    }
530
12
    *eof = _current_range_idx == _row_ranges.range_size();
531
12
    return Status::OK();
532
12
}
533
534
template <bool IN_COLLECTION, bool OFFSET_INDEX>
535
Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::read_dict_values_to_column(
536
0
        MutableColumnPtr& doris_column, bool* has_dict) {
537
0
    bool loaded;
538
0
    RETURN_IF_ERROR(_try_load_dict_page(&loaded, has_dict));
539
0
    if (loaded && *has_dict) {
540
0
        return _chunk_reader->read_dict_values_to_column(doris_column);
541
0
    }
542
0
    return Status::OK();
543
0
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb1EE26read_dict_values_to_columnERNS_3COWINS_7IColumnEE11mutable_ptrIS3_EEPb
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb0EE26read_dict_values_to_columnERNS_3COWINS_7IColumnEE11mutable_ptrIS3_EEPb
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb1EE26read_dict_values_to_columnERNS_3COWINS_7IColumnEE11mutable_ptrIS3_EEPb
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb0EE26read_dict_values_to_columnERNS_3COWINS_7IColumnEE11mutable_ptrIS3_EEPb
544
template <bool IN_COLLECTION, bool OFFSET_INDEX>
545
Result<MutableColumnPtr>
546
ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::convert_dict_column_to_string_column(
547
0
        const ColumnInt32* dict_column) {
548
0
    return _chunk_reader->convert_dict_column_to_string_column(dict_column);
549
0
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb1EE36convert_dict_column_to_string_columnEPKNS_12ColumnVectorILNS_13PrimitiveTypeE5EEE
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb0EE36convert_dict_column_to_string_columnEPKNS_12ColumnVectorILNS_13PrimitiveTypeE5EEE
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb1EE36convert_dict_column_to_string_columnEPKNS_12ColumnVectorILNS_13PrimitiveTypeE5EEE
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb0EE36convert_dict_column_to_string_columnEPKNS_12ColumnVectorILNS_13PrimitiveTypeE5EEE
550
551
template <bool IN_COLLECTION, bool OFFSET_INDEX>
552
Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::_try_load_dict_page(bool* loaded,
553
0
                                                                            bool* has_dict) {
554
    // _chunk_reader init will load first page header to check whether has dict page
555
0
    *loaded = true;
556
0
    *has_dict = _chunk_reader->has_dict();
557
0
    return Status::OK();
558
0
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb1EE19_try_load_dict_pageEPbS2_
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb0EE19_try_load_dict_pageEPbS2_
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb1EE19_try_load_dict_pageEPbS2_
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb0EE19_try_load_dict_pageEPbS2_
559
560
template <bool IN_COLLECTION, bool OFFSET_INDEX>
561
Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::read_column_data(
562
        ColumnPtr& doris_column, const DataTypePtr& type,
563
        const std::shared_ptr<TableSchemaChangeHelper::Node>& root_node, FilterMap& filter_map,
564
        size_t batch_size, size_t* read_rows, bool* eof, bool is_dict_filter,
565
268
        int64_t real_column_size) {
566
268
    if (_converter == nullptr) {
567
125
        _converter = parquet::PhysicalToLogicalConverter::get_converter(
568
125
                _field_schema, _field_schema->data_type, type, _ctz, is_dict_filter);
569
125
        if (!_converter->support()) {
570
0
            return Status::InternalError(
571
0
                    "The column type of '{}' is not supported: {}, is_dict_filter: {}, "
572
0
                    "src_logical_type: {}, dst_logical_type: {}",
573
0
                    _field_schema->name, _converter->get_error_msg(), is_dict_filter,
574
0
                    _field_schema->data_type->get_name(), type->get_name());
575
0
        }
576
125
    }
577
    // !FIXME: We should verify whether the get_physical_column logic is correct, why do we return a doris_column?
578
268
    ColumnPtr resolved_column =
579
268
            _converter->get_physical_column(_field_schema->physical_type, _field_schema->data_type,
580
268
                                            doris_column, type, is_dict_filter);
581
268
    if (_converter->read_directly_into_dst_logical_column()) {
582
239
        DCHECK_EQ(resolved_column.get(), doris_column.get());
583
239
        resolved_column = std::move(doris_column);
584
239
    }
585
268
    DataTypePtr& resolved_type = _converter->get_physical_type();
586
587
268
    _def_levels.clear();
588
268
    _rep_levels.clear();
589
268
    *read_rows = 0;
590
591
268
    if (_in_nested) {
592
15
        RETURN_IF_ERROR(_read_nested_column(resolved_column, resolved_type, filter_map, batch_size,
593
15
                                            read_rows, eof, is_dict_filter));
594
15
        return _converter->convert(resolved_column, _field_schema->data_type, type, doris_column,
595
15
                                   is_dict_filter);
596
15
    }
597
598
253
    int64_t right_row = 0;
599
253
    if constexpr (OFFSET_INDEX == false) {
600
253
        RETURN_IF_ERROR(_chunk_reader->parse_page_header());
601
253
        right_row = _chunk_reader->page_end_row();
602
253
    } else {
603
0
        right_row = _chunk_reader->page_end_row();
604
0
    }
605
606
253
    do {
607
        // generate the row ranges that should be read
608
253
        RowRanges read_ranges;
609
253
        _generate_read_ranges(RowRange {_current_row_index, right_row}, &read_ranges);
610
253
        if (read_ranges.count() == 0) {
611
            // skip the whole page
612
63
            _current_row_index = right_row;
613
190
        } else {
614
190
            bool skip_whole_batch = false;
615
            // Determining whether to skip page or batch will increase the calculation time.
616
            // When the filtering effect is greater than 60%, it is possible to skip the page or batch.
617
190
            if (filter_map.has_filter() && filter_map.filter_ratio() > 0.6) {
618
                // lazy read
619
0
                size_t remaining_num_values = read_ranges.count();
620
0
                if (batch_size >= remaining_num_values &&
621
0
                    filter_map.can_filter_all(remaining_num_values, _filter_map_index)) {
622
                    // We can skip the whole page if the remaining values are filtered by predicate columns
623
0
                    _filter_map_index += remaining_num_values;
624
0
                    _current_row_index = right_row;
625
0
                    *read_rows = remaining_num_values;
626
0
                    break;
627
0
                }
628
0
                skip_whole_batch = batch_size <= remaining_num_values &&
629
0
                                   filter_map.can_filter_all(batch_size, _filter_map_index);
630
0
                if (skip_whole_batch) {
631
0
                    _filter_map_index += batch_size;
632
0
                }
633
0
            }
634
            // load page data to decode or skip values
635
190
            RETURN_IF_ERROR(_chunk_reader->parse_page_header());
636
190
            RETURN_IF_ERROR(_chunk_reader->load_page_data_idempotent());
637
190
            size_t has_read = 0;
638
332
            for (size_t idx = 0; idx < read_ranges.range_size(); idx++) {
639
223
                auto range = read_ranges.get_range(idx);
640
                // generate the skipped values
641
223
                size_t skip_values = range.from() - _current_row_index;
642
223
                RETURN_IF_ERROR(_skip_values(skip_values));
643
223
                _current_row_index += skip_values;
644
                // generate the read values
645
223
                size_t read_values =
646
223
                        std::min((size_t)(range.to() - range.from()), batch_size - has_read);
647
223
                if (skip_whole_batch) {
648
0
                    RETURN_IF_ERROR(_skip_values(read_values));
649
223
                } else {
650
223
                    RETURN_IF_ERROR(_read_values(read_values, resolved_column, resolved_type,
651
223
                                                 filter_map, is_dict_filter));
652
223
                }
653
223
                has_read += read_values;
654
223
                *read_rows += read_values;
655
223
                _current_row_index += read_values;
656
223
                if (has_read == batch_size) {
657
81
                    break;
658
81
                }
659
223
            }
660
190
        }
661
253
    } while (false);
662
663
253
    if (right_row == _current_row_index) {
664
110
        if (!_chunk_reader->has_next_page()) {
665
110
            *eof = true;
666
110
        } else {
667
0
            RETURN_IF_ERROR(_chunk_reader->next_page());
668
0
        }
669
110
    }
670
671
253
    {
672
253
        SCOPED_RAW_TIMER(&_convert_time);
673
253
        RETURN_IF_ERROR(_converter->convert(resolved_column, _field_schema->data_type, type,
674
253
                                            doris_column, is_dict_filter));
675
253
    }
676
253
    return Status::OK();
677
253
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb1ELb1EE16read_column_dataERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERKSt10shared_ptrIKNS_9IDataTypeEERKS8_INS_23TableSchemaChangeHelper4NodeEERNS_9FilterMapEmPmPbbl
_ZN5doris18ScalarColumnReaderILb1ELb0EE16read_column_dataERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERKSt10shared_ptrIKNS_9IDataTypeEERKS8_INS_23TableSchemaChangeHelper4NodeEERNS_9FilterMapEmPmPbbl
Line
Count
Source
565
3
        int64_t real_column_size) {
566
3
    if (_converter == nullptr) {
567
3
        _converter = parquet::PhysicalToLogicalConverter::get_converter(
568
3
                _field_schema, _field_schema->data_type, type, _ctz, is_dict_filter);
569
3
        if (!_converter->support()) {
570
0
            return Status::InternalError(
571
0
                    "The column type of '{}' is not supported: {}, is_dict_filter: {}, "
572
0
                    "src_logical_type: {}, dst_logical_type: {}",
573
0
                    _field_schema->name, _converter->get_error_msg(), is_dict_filter,
574
0
                    _field_schema->data_type->get_name(), type->get_name());
575
0
        }
576
3
    }
577
    // !FIXME: We should verify whether the get_physical_column logic is correct, why do we return a doris_column?
578
3
    ColumnPtr resolved_column =
579
3
            _converter->get_physical_column(_field_schema->physical_type, _field_schema->data_type,
580
3
                                            doris_column, type, is_dict_filter);
581
3
    if (_converter->read_directly_into_dst_logical_column()) {
582
3
        DCHECK_EQ(resolved_column.get(), doris_column.get());
583
3
        resolved_column = std::move(doris_column);
584
3
    }
585
3
    DataTypePtr& resolved_type = _converter->get_physical_type();
586
587
3
    _def_levels.clear();
588
3
    _rep_levels.clear();
589
3
    *read_rows = 0;
590
591
3
    if (_in_nested) {
592
3
        RETURN_IF_ERROR(_read_nested_column(resolved_column, resolved_type, filter_map, batch_size,
593
3
                                            read_rows, eof, is_dict_filter));
594
3
        return _converter->convert(resolved_column, _field_schema->data_type, type, doris_column,
595
3
                                   is_dict_filter);
596
3
    }
597
598
0
    int64_t right_row = 0;
599
0
    if constexpr (OFFSET_INDEX == false) {
600
0
        RETURN_IF_ERROR(_chunk_reader->parse_page_header());
601
0
        right_row = _chunk_reader->page_end_row();
602
    } else {
603
        right_row = _chunk_reader->page_end_row();
604
    }
605
606
0
    do {
607
        // generate the row ranges that should be read
608
0
        RowRanges read_ranges;
609
0
        _generate_read_ranges(RowRange {_current_row_index, right_row}, &read_ranges);
610
0
        if (read_ranges.count() == 0) {
611
            // skip the whole page
612
0
            _current_row_index = right_row;
613
0
        } else {
614
0
            bool skip_whole_batch = false;
615
            // Determining whether to skip page or batch will increase the calculation time.
616
            // When the filtering effect is greater than 60%, it is possible to skip the page or batch.
617
0
            if (filter_map.has_filter() && filter_map.filter_ratio() > 0.6) {
618
                // lazy read
619
0
                size_t remaining_num_values = read_ranges.count();
620
0
                if (batch_size >= remaining_num_values &&
621
0
                    filter_map.can_filter_all(remaining_num_values, _filter_map_index)) {
622
                    // We can skip the whole page if the remaining values are filtered by predicate columns
623
0
                    _filter_map_index += remaining_num_values;
624
0
                    _current_row_index = right_row;
625
0
                    *read_rows = remaining_num_values;
626
0
                    break;
627
0
                }
628
0
                skip_whole_batch = batch_size <= remaining_num_values &&
629
0
                                   filter_map.can_filter_all(batch_size, _filter_map_index);
630
0
                if (skip_whole_batch) {
631
0
                    _filter_map_index += batch_size;
632
0
                }
633
0
            }
634
            // load page data to decode or skip values
635
0
            RETURN_IF_ERROR(_chunk_reader->parse_page_header());
636
0
            RETURN_IF_ERROR(_chunk_reader->load_page_data_idempotent());
637
0
            size_t has_read = 0;
638
0
            for (size_t idx = 0; idx < read_ranges.range_size(); idx++) {
639
0
                auto range = read_ranges.get_range(idx);
640
                // generate the skipped values
641
0
                size_t skip_values = range.from() - _current_row_index;
642
0
                RETURN_IF_ERROR(_skip_values(skip_values));
643
0
                _current_row_index += skip_values;
644
                // generate the read values
645
0
                size_t read_values =
646
0
                        std::min((size_t)(range.to() - range.from()), batch_size - has_read);
647
0
                if (skip_whole_batch) {
648
0
                    RETURN_IF_ERROR(_skip_values(read_values));
649
0
                } else {
650
0
                    RETURN_IF_ERROR(_read_values(read_values, resolved_column, resolved_type,
651
0
                                                 filter_map, is_dict_filter));
652
0
                }
653
0
                has_read += read_values;
654
0
                *read_rows += read_values;
655
0
                _current_row_index += read_values;
656
0
                if (has_read == batch_size) {
657
0
                    break;
658
0
                }
659
0
            }
660
0
        }
661
0
    } while (false);
662
663
0
    if (right_row == _current_row_index) {
664
0
        if (!_chunk_reader->has_next_page()) {
665
0
            *eof = true;
666
0
        } else {
667
0
            RETURN_IF_ERROR(_chunk_reader->next_page());
668
0
        }
669
0
    }
670
671
0
    {
672
0
        SCOPED_RAW_TIMER(&_convert_time);
673
0
        RETURN_IF_ERROR(_converter->convert(resolved_column, _field_schema->data_type, type,
674
0
                                            doris_column, is_dict_filter));
675
0
    }
676
0
    return Status::OK();
677
0
}
Unexecuted instantiation: _ZN5doris18ScalarColumnReaderILb0ELb1EE16read_column_dataERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERKSt10shared_ptrIKNS_9IDataTypeEERKS8_INS_23TableSchemaChangeHelper4NodeEERNS_9FilterMapEmPmPbbl
_ZN5doris18ScalarColumnReaderILb0ELb0EE16read_column_dataERNS_3COWINS_7IColumnEE13immutable_ptrIS3_EERKSt10shared_ptrIKNS_9IDataTypeEERKS8_INS_23TableSchemaChangeHelper4NodeEERNS_9FilterMapEmPmPbbl
Line
Count
Source
565
265
        int64_t real_column_size) {
566
265
    if (_converter == nullptr) {
567
122
        _converter = parquet::PhysicalToLogicalConverter::get_converter(
568
122
                _field_schema, _field_schema->data_type, type, _ctz, is_dict_filter);
569
122
        if (!_converter->support()) {
570
0
            return Status::InternalError(
571
0
                    "The column type of '{}' is not supported: {}, is_dict_filter: {}, "
572
0
                    "src_logical_type: {}, dst_logical_type: {}",
573
0
                    _field_schema->name, _converter->get_error_msg(), is_dict_filter,
574
0
                    _field_schema->data_type->get_name(), type->get_name());
575
0
        }
576
122
    }
577
    // !FIXME: We should verify whether the get_physical_column logic is correct, why do we return a doris_column?
578
265
    ColumnPtr resolved_column =
579
265
            _converter->get_physical_column(_field_schema->physical_type, _field_schema->data_type,
580
265
                                            doris_column, type, is_dict_filter);
581
265
    if (_converter->read_directly_into_dst_logical_column()) {
582
236
        DCHECK_EQ(resolved_column.get(), doris_column.get());
583
236
        resolved_column = std::move(doris_column);
584
236
    }
585
265
    DataTypePtr& resolved_type = _converter->get_physical_type();
586
587
265
    _def_levels.clear();
588
265
    _rep_levels.clear();
589
265
    *read_rows = 0;
590
591
265
    if (_in_nested) {
592
12
        RETURN_IF_ERROR(_read_nested_column(resolved_column, resolved_type, filter_map, batch_size,
593
12
                                            read_rows, eof, is_dict_filter));
594
12
        return _converter->convert(resolved_column, _field_schema->data_type, type, doris_column,
595
12
                                   is_dict_filter);
596
12
    }
597
598
253
    int64_t right_row = 0;
599
253
    if constexpr (OFFSET_INDEX == false) {
600
253
        RETURN_IF_ERROR(_chunk_reader->parse_page_header());
601
253
        right_row = _chunk_reader->page_end_row();
602
    } else {
603
        right_row = _chunk_reader->page_end_row();
604
    }
605
606
253
    do {
607
        // generate the row ranges that should be read
608
253
        RowRanges read_ranges;
609
253
        _generate_read_ranges(RowRange {_current_row_index, right_row}, &read_ranges);
610
253
        if (read_ranges.count() == 0) {
611
            // skip the whole page
612
63
            _current_row_index = right_row;
613
190
        } else {
614
190
            bool skip_whole_batch = false;
615
            // Determining whether to skip page or batch will increase the calculation time.
616
            // When the filtering effect is greater than 60%, it is possible to skip the page or batch.
617
190
            if (filter_map.has_filter() && filter_map.filter_ratio() > 0.6) {
618
                // lazy read
619
0
                size_t remaining_num_values = read_ranges.count();
620
0
                if (batch_size >= remaining_num_values &&
621
0
                    filter_map.can_filter_all(remaining_num_values, _filter_map_index)) {
622
                    // We can skip the whole page if the remaining values are filtered by predicate columns
623
0
                    _filter_map_index += remaining_num_values;
624
0
                    _current_row_index = right_row;
625
0
                    *read_rows = remaining_num_values;
626
0
                    break;
627
0
                }
628
0
                skip_whole_batch = batch_size <= remaining_num_values &&
629
0
                                   filter_map.can_filter_all(batch_size, _filter_map_index);
630
0
                if (skip_whole_batch) {
631
0
                    _filter_map_index += batch_size;
632
0
                }
633
0
            }
634
            // load page data to decode or skip values
635
190
            RETURN_IF_ERROR(_chunk_reader->parse_page_header());
636
190
            RETURN_IF_ERROR(_chunk_reader->load_page_data_idempotent());
637
190
            size_t has_read = 0;
638
332
            for (size_t idx = 0; idx < read_ranges.range_size(); idx++) {
639
223
                auto range = read_ranges.get_range(idx);
640
                // generate the skipped values
641
223
                size_t skip_values = range.from() - _current_row_index;
642
223
                RETURN_IF_ERROR(_skip_values(skip_values));
643
223
                _current_row_index += skip_values;
644
                // generate the read values
645
223
                size_t read_values =
646
223
                        std::min((size_t)(range.to() - range.from()), batch_size - has_read);
647
223
                if (skip_whole_batch) {
648
0
                    RETURN_IF_ERROR(_skip_values(read_values));
649
223
                } else {
650
223
                    RETURN_IF_ERROR(_read_values(read_values, resolved_column, resolved_type,
651
223
                                                 filter_map, is_dict_filter));
652
223
                }
653
223
                has_read += read_values;
654
223
                *read_rows += read_values;
655
223
                _current_row_index += read_values;
656
223
                if (has_read == batch_size) {
657
81
                    break;
658
81
                }
659
223
            }
660
190
        }
661
253
    } while (false);
662
663
253
    if (right_row == _current_row_index) {
664
110
        if (!_chunk_reader->has_next_page()) {
665
110
            *eof = true;
666
110
        } else {
667
0
            RETURN_IF_ERROR(_chunk_reader->next_page());
668
0
        }
669
110
    }
670
671
253
    {
672
253
        SCOPED_RAW_TIMER(&_convert_time);
673
253
        RETURN_IF_ERROR(_converter->convert(resolved_column, _field_schema->data_type, type,
674
253
                                            doris_column, is_dict_filter));
675
253
    }
676
253
    return Status::OK();
677
253
}
678
679
Status ArrayColumnReader::init(std::unique_ptr<ParquetColumnReader> element_reader,
680
2
                               FieldSchema* field) {
681
2
    _field_schema = field;
682
2
    _element_reader = std::move(element_reader);
683
2
    return Status::OK();
684
2
}
685
686
Status ArrayColumnReader::read_column_data(
687
        ColumnPtr& doris_column, const DataTypePtr& type,
688
        const std::shared_ptr<TableSchemaChangeHelper::Node>& root_node, FilterMap& filter_map,
689
        size_t batch_size, size_t* read_rows, bool* eof, bool is_dict_filter,
690
2
        int64_t real_column_size) {
691
2
    MutableColumnPtr data_column;
692
2
    NullMap* null_map_ptr = nullptr;
693
2
    doris_column = IColumn::mutate(std::move(doris_column));
694
2
    if (is_column_nullable(*doris_column)) {
695
2
        auto mutable_column = doris_column->assert_mutable();
696
2
        auto* nullable_column = assert_cast<ColumnNullable*>(mutable_column.get());
697
2
        null_map_ptr = &nullable_column->get_null_map_data();
698
2
        data_column = nullable_column->get_nested_column_ptr();
699
2
    } else {
700
0
        if (_field_schema->data_type->is_nullable()) {
701
0
            return Status::Corruption("Not nullable column has null values in parquet file");
702
0
        }
703
0
        data_column = doris_column->assert_mutable();
704
0
    }
705
2
    if (type->get_primitive_type() != PrimitiveType::TYPE_ARRAY) {
706
0
        return Status::Corruption(
707
0
                "Wrong data type for column '{}', expected Array type, actual type: {}.",
708
0
                _field_schema->name, type->get_name());
709
0
    }
710
711
2
    ColumnPtr& element_column = assert_cast<ColumnArray&>(*data_column).get_data_ptr();
712
2
    const DataTypePtr& element_type =
713
2
            (assert_cast<const DataTypeArray*>(remove_nullable(type).get()))->get_nested_type();
714
    // read nested column
715
2
    RETURN_IF_ERROR(_element_reader->read_column_data(element_column, element_type,
716
2
                                                      root_node->get_element_node(), filter_map,
717
2
                                                      batch_size, read_rows, eof, is_dict_filter));
718
2
    if (*read_rows == 0) {
719
0
        return Status::OK();
720
0
    }
721
722
2
    ColumnArray::Offsets64& offsets_data = assert_cast<ColumnArray&>(*data_column).get_offsets();
723
    // fill offset and null map
724
2
    fill_array_offset(_field_schema, offsets_data, null_map_ptr, _element_reader->get_rep_level(),
725
2
                      _element_reader->get_def_level());
726
2
    DCHECK_EQ(element_column->size(), offsets_data.back());
727
2
#ifndef NDEBUG
728
2
    doris_column->sanity_check();
729
2
#endif
730
2
    return Status::OK();
731
2
}
732
733
Status MapColumnReader::init(std::unique_ptr<ParquetColumnReader> key_reader,
734
                             std::unique_ptr<ParquetColumnReader> value_reader,
735
0
                             FieldSchema* field) {
736
0
    _field_schema = field;
737
0
    _key_reader = std::move(key_reader);
738
0
    _value_reader = std::move(value_reader);
739
0
    return Status::OK();
740
0
}
741
742
Status MapColumnReader::read_column_data(
743
        ColumnPtr& doris_column, const DataTypePtr& type,
744
        const std::shared_ptr<TableSchemaChangeHelper::Node>& root_node, FilterMap& filter_map,
745
        size_t batch_size, size_t* read_rows, bool* eof, bool is_dict_filter,
746
0
        int64_t real_column_size) {
747
0
    MutableColumnPtr data_column;
748
0
    NullMap* null_map_ptr = nullptr;
749
0
    doris_column = IColumn::mutate(std::move(doris_column));
750
0
    if (is_column_nullable(*doris_column)) {
751
0
        auto mutable_column = doris_column->assert_mutable();
752
0
        auto* nullable_column = assert_cast<ColumnNullable*>(mutable_column.get());
753
0
        null_map_ptr = &nullable_column->get_null_map_data();
754
0
        data_column = nullable_column->get_nested_column_ptr();
755
0
    } else {
756
0
        if (_field_schema->data_type->is_nullable()) {
757
0
            return Status::Corruption("Not nullable column has null values in parquet file");
758
0
        }
759
0
        data_column = doris_column->assert_mutable();
760
0
    }
761
0
    if (remove_nullable(type)->get_primitive_type() != PrimitiveType::TYPE_MAP) {
762
0
        return Status::Corruption(
763
0
                "Wrong data type for column '{}', expected Map type, actual type id {}.",
764
0
                _field_schema->name, type->get_name());
765
0
    }
766
767
0
    auto& map = assert_cast<ColumnMap&>(*data_column);
768
0
    const DataTypePtr& key_type =
769
0
            assert_cast<const DataTypeMap*>(remove_nullable(type).get())->get_key_type();
770
0
    const DataTypePtr& value_type =
771
0
            assert_cast<const DataTypeMap*>(remove_nullable(type).get())->get_value_type();
772
0
    ColumnPtr& key_column = map.get_keys_ptr();
773
0
    ColumnPtr& value_column = map.get_values_ptr();
774
775
0
    size_t key_rows = 0;
776
0
    size_t value_rows = 0;
777
0
    bool key_eof = false;
778
0
    bool value_eof = false;
779
0
    int64_t orig_col_column_size = key_column->size();
780
781
0
    RETURN_IF_ERROR(_key_reader->read_column_data(key_column, key_type, root_node->get_key_node(),
782
0
                                                  filter_map, batch_size, &key_rows, &key_eof,
783
0
                                                  is_dict_filter));
784
785
0
    while (value_rows < key_rows && !value_eof) {
786
0
        size_t loop_rows = 0;
787
0
        RETURN_IF_ERROR(_value_reader->read_column_data(
788
0
                value_column, value_type, root_node->get_value_node(), filter_map,
789
0
                key_rows - value_rows, &loop_rows, &value_eof, is_dict_filter,
790
0
                key_column->size() - orig_col_column_size));
791
0
        value_rows += loop_rows;
792
0
    }
793
0
    DCHECK_EQ(key_rows, value_rows);
794
0
    *read_rows = key_rows;
795
0
    *eof = key_eof;
796
797
0
    if (*read_rows == 0) {
798
0
        return Status::OK();
799
0
    }
800
801
0
    DCHECK_EQ(key_column->size(), value_column->size());
802
    // fill offset and null map
803
0
    fill_array_offset(_field_schema, map.get_offsets(), null_map_ptr, _key_reader->get_rep_level(),
804
0
                      _key_reader->get_def_level());
805
0
    DCHECK_EQ(key_column->size(), map.get_offsets().back());
806
0
#ifndef NDEBUG
807
0
    doris_column->sanity_check();
808
0
#endif
809
0
    return Status::OK();
810
0
}
811
812
Status StructColumnReader::init(
813
        std::unordered_map<std::string, std::unique_ptr<ParquetColumnReader>>&& child_readers,
814
13
        FieldSchema* field) {
815
13
    _field_schema = field;
816
13
    _child_readers = std::move(child_readers);
817
13
    return Status::OK();
818
13
}
819
Status StructColumnReader::read_column_data(
820
        ColumnPtr& doris_column, const DataTypePtr& type,
821
        const std::shared_ptr<TableSchemaChangeHelper::Node>& root_node, FilterMap& filter_map,
822
        size_t batch_size, size_t* read_rows, bool* eof, bool is_dict_filter,
823
13
        int64_t real_column_size) {
824
13
    MutableColumnPtr data_column;
825
13
    NullMap* null_map_ptr = nullptr;
826
13
    doris_column = IColumn::mutate(std::move(doris_column));
827
13
    if (is_column_nullable(*doris_column)) {
828
13
        auto mutable_column = doris_column->assert_mutable();
829
13
        auto* nullable_column = assert_cast<ColumnNullable*>(mutable_column.get());
830
13
        null_map_ptr = &nullable_column->get_null_map_data();
831
13
        data_column = nullable_column->get_nested_column_ptr();
832
13
    } else {
833
0
        if (_field_schema->data_type->is_nullable()) {
834
0
            return Status::Corruption("Not nullable column has null values in parquet file");
835
0
        }
836
0
        data_column = doris_column->assert_mutable();
837
0
    }
838
13
    if (type->get_primitive_type() != PrimitiveType::TYPE_STRUCT) {
839
0
        return Status::Corruption(
840
0
                "Wrong data type for column '{}', expected Struct type, actual type id {}.",
841
0
                _field_schema->name, type->get_name());
842
0
    }
843
844
13
    auto& doris_struct = assert_cast<ColumnStruct&>(*data_column);
845
13
    const auto* doris_struct_type = assert_cast<const DataTypeStruct*>(remove_nullable(type).get());
846
847
13
    int64_t not_missing_column_id = -1;
848
13
    size_t not_missing_orig_column_size = 0;
849
13
    std::vector<size_t> missing_column_idxs {};
850
13
    std::vector<size_t> skip_reading_column_idxs {};
851
852
13
    _read_column_names.clear();
853
854
41
    for (size_t i = 0; i < doris_struct.tuple_size(); ++i) {
855
28
        ColumnPtr& doris_field = doris_struct.get_column_ptr(i);
856
28
        auto& doris_type = doris_struct_type->get_element(i);
857
28
        auto& doris_name = doris_struct_type->get_element_name(i);
858
28
        if (!root_node->children_column_exists(doris_name)) {
859
1
            missing_column_idxs.push_back(i);
860
1
            VLOG_DEBUG << "[ParquetReader] Missing column in schema: column_idx[" << i
861
0
                       << "], doris_name: " << doris_name << " (column not exists in root node)";
862
1
            continue;
863
1
        }
864
27
        auto file_name = root_node->children_file_column_name(doris_name);
865
866
        // Check if this is a SkipReadingReader - we should skip it when choosing reference column
867
        // because SkipReadingReader doesn't know the actual data size in nested context
868
27
        bool is_skip_reader =
869
27
                dynamic_cast<SkipReadingReader*>(_child_readers[file_name].get()) != nullptr;
870
871
27
        if (is_skip_reader) {
872
            // Store SkipReadingReader columns to fill them later based on reference column size
873
4
            skip_reading_column_idxs.push_back(i);
874
4
            continue;
875
4
        }
876
877
        // Only add non-SkipReadingReader columns to _read_column_names
878
        // This ensures get_rep_level() and get_def_level() return valid levels
879
23
        _read_column_names.emplace_back(file_name);
880
881
23
        size_t field_rows = 0;
882
23
        bool field_eof = false;
883
23
        if (not_missing_column_id == -1) {
884
12
            not_missing_column_id = i;
885
12
            not_missing_orig_column_size = doris_field->size();
886
12
            RETURN_IF_ERROR(_child_readers[file_name]->read_column_data(
887
12
                    doris_field, doris_type, root_node->get_children_node(doris_name), filter_map,
888
12
                    batch_size, &field_rows, &field_eof, is_dict_filter));
889
12
            *read_rows = field_rows;
890
12
            *eof = field_eof;
891
            /*
892
             * Considering the issue in the `_read_nested_column` function where data may span across pages, leading
893
             * to missing definition and repetition levels, when filling the null_map of the struct later, it is
894
             * crucial to use the definition and repetition levels from the first read column
895
             * (since `_read_nested_column` is not called repeatedly).
896
             *
897
             *  It is worth mentioning that, theoretically, any sub-column can be chosen to fill the null_map,
898
             *  and selecting the shortest one will offer better performance
899
             */
900
12
        } else {
901
22
            while (field_rows < *read_rows && !field_eof) {
902
11
                size_t loop_rows = 0;
903
11
                RETURN_IF_ERROR(_child_readers[file_name]->read_column_data(
904
11
                        doris_field, doris_type, root_node->get_children_node(doris_name),
905
11
                        filter_map, *read_rows - field_rows, &loop_rows, &field_eof,
906
11
                        is_dict_filter));
907
11
                field_rows += loop_rows;
908
11
            }
909
11
            DCHECK_EQ(*read_rows, field_rows);
910
            //            DCHECK_EQ(*eof, field_eof);
911
11
        }
912
23
    }
913
914
13
    int64_t missing_column_sz = -1;
915
916
13
    if (not_missing_column_id == -1) {
917
        // All queried columns are missing in the file (e.g., all added after schema change)
918
        // We need to pick a column from _field_schema children that exists in the file for RL/DL reference
919
1
        std::string reference_file_column_name;
920
1
        std::unique_ptr<ParquetColumnReader>* reference_reader = nullptr;
921
922
1
        for (const auto& child : _field_schema->children) {
923
1
            auto it = _child_readers.find(child.name);
924
1
            if (it != _child_readers.end()) {
925
                // Skip SkipReadingReader as they don't have valid RL/DL
926
1
                bool is_skip_reader = dynamic_cast<SkipReadingReader*>(it->second.get()) != nullptr;
927
1
                if (!is_skip_reader) {
928
1
                    reference_file_column_name = child.name;
929
1
                    reference_reader = &(it->second);
930
1
                    break;
931
1
                }
932
1
            }
933
1
        }
934
935
1
        if (reference_reader != nullptr) {
936
            // Read the reference column to get correct RL/DL information
937
            // TODO: Optimize by only reading RL/DL without actual data decoding
938
939
            // We need to find the FieldSchema for the reference column from _field_schema children
940
1
            FieldSchema* ref_field_schema = nullptr;
941
1
            for (auto& child : _field_schema->children) {
942
1
                if (child.name == reference_file_column_name) {
943
1
                    ref_field_schema = &child;
944
1
                    break;
945
1
                }
946
1
            }
947
948
1
            if (ref_field_schema == nullptr) {
949
0
                return Status::InternalError(
950
0
                        "Cannot find field schema for reference column '{}' in struct '{}'",
951
0
                        reference_file_column_name, _field_schema->name);
952
0
            }
953
954
            // Create a temporary column to hold the data (we'll use its size for missing_column_sz)
955
1
            ColumnPtr temp_column = ref_field_schema->data_type->create_column();
956
1
            auto temp_type = ref_field_schema->data_type;
957
958
1
            size_t field_rows = 0;
959
1
            bool field_eof = false;
960
961
            // Use ConstNode for the reference column instead of looking up from root_node.
962
            // The reference column is only used to get RL/DL information for determining the number
963
            // of elements in the struct. It may be a column that has been dropped from the table
964
            // schema (e.g., 'removed' field), but still exists in older parquet files.
965
            // Since we don't need schema mapping for this column (we just need its RL/DL levels),
966
            // using ConstNode is safe and avoids the issue where the reference column doesn't exist
967
            // in root_node (because it was dropped from table schema).
968
1
            auto ref_child_node = TableSchemaChangeHelper::ConstNode::get_instance();
969
1
            not_missing_orig_column_size = temp_column->size();
970
971
1
            RETURN_IF_ERROR((*reference_reader)
972
1
                                    ->read_column_data(temp_column, temp_type, ref_child_node,
973
1
                                                       filter_map, batch_size, &field_rows,
974
1
                                                       &field_eof, is_dict_filter));
975
976
1
            *read_rows = field_rows;
977
1
            *eof = field_eof;
978
979
            // Store this reference column name for get_rep_level/get_def_level to use
980
1
            _read_column_names.emplace_back(reference_file_column_name);
981
982
1
            missing_column_sz = temp_column->size() - not_missing_orig_column_size;
983
1
        } else {
984
0
            return Status::Corruption(
985
0
                    "Cannot read struct '{}': all queried columns are missing and no reference "
986
0
                    "column found in file",
987
0
                    _field_schema->name);
988
0
        }
989
1
    }
990
991
    //  This missing_column_sz is not *read_rows. Because read_rows returns the number of rows.
992
    //  For example: suppose we have a column array<struct<a:int,b:string>>,
993
    //  where b is a newly added column, that is, a missing column.
994
    //  There are two rows of data in this column,
995
    //      [{1,null},{2,null},{3,null}]
996
    //      [{4,null},{5,null}]
997
    //  When you first read subcolumn a, you read 5 data items and the value of *read_rows is 2.
998
    //  You should insert 5 records into subcolumn b instead of 2.
999
13
    if (missing_column_sz == -1) {
1000
12
        missing_column_sz = doris_struct.get_column(not_missing_column_id).size() -
1001
12
                            not_missing_orig_column_size;
1002
12
    }
1003
1004
    // Fill SkipReadingReader columns with the correct amount of data based on reference column
1005
    // Let SkipReadingReader handle the data filling through its read_column_data method
1006
13
    for (auto idx : skip_reading_column_idxs) {
1007
4
        auto& doris_field = doris_struct.get_column_ptr(idx);
1008
4
        auto& doris_type = const_cast<DataTypePtr&>(doris_struct_type->get_element(idx));
1009
4
        auto& doris_name = const_cast<String&>(doris_struct_type->get_element_name(idx));
1010
4
        auto file_name = root_node->children_file_column_name(doris_name);
1011
1012
4
        size_t field_rows = 0;
1013
4
        bool field_eof = false;
1014
4
        RETURN_IF_ERROR(_child_readers[file_name]->read_column_data(
1015
4
                doris_field, doris_type, root_node->get_children_node(doris_name), filter_map,
1016
4
                missing_column_sz, &field_rows, &field_eof, is_dict_filter, missing_column_sz));
1017
4
    }
1018
1019
    // Fill truly missing columns (not in root_node) with null or default value
1020
13
    for (auto idx : missing_column_idxs) {
1021
1
        auto& doris_field = doris_struct.get_column_ptr(idx);
1022
1
        auto& doris_type = doris_struct_type->get_element(idx);
1023
1
        auto mutable_field = IColumn::mutate(std::move(doris_field));
1024
1
        ColumnPtr initial_default;
1025
1
        RETURN_IF_ERROR(build_initial_default_column(
1026
1
                root_node->children_initial_default_value(doris_struct_type->get_element_name(idx)),
1027
1
                doris_type, missing_column_sz, &initial_default));
1028
1
        if (initial_default.get() != nullptr) {
1029
            // Iceberg initial defaults are logical row values, including for nested fields absent
1030
            // from the physical file; append them instead of the type's generic NULL/default.
1031
1
            mutable_field->insert_range_from(*initial_default, 0, missing_column_sz);
1032
1
        } else {
1033
0
            DCHECK(doris_type->is_nullable());
1034
0
            static_cast<ColumnNullable*>(mutable_field.get())
1035
0
                    ->insert_many_defaults(missing_column_sz);
1036
0
        }
1037
1
        doris_field = std::move(mutable_field);
1038
1
    }
1039
1040
13
    if (null_map_ptr != nullptr) {
1041
13
        fill_struct_null_map(_field_schema, *null_map_ptr, this->get_rep_level(),
1042
13
                             this->get_def_level());
1043
13
    }
1044
13
#ifndef NDEBUG
1045
13
    doris_column->sanity_check();
1046
13
#endif
1047
13
    return Status::OK();
1048
13
}
1049
1050
template class ScalarColumnReader<true, true>;
1051
template class ScalarColumnReader<true, false>;
1052
template class ScalarColumnReader<false, true>;
1053
template class ScalarColumnReader<false, false>;
1054
1055
}; // namespace doris