Coverage Report

Created: 2026-09-30 17:30

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