Coverage Report

Created: 2026-09-28 15:19

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/parquet/schema_desc.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/schema_desc.h"
19
20
#include <ctype.h>
21
22
#include <algorithm>
23
#include <ostream>
24
#include <utility>
25
26
#include "common/cast_set.h"
27
#include "common/logging.h"
28
#include "core/data_type/data_type_array.h"
29
#include "core/data_type/data_type_factory.hpp"
30
#include "core/data_type/data_type_map.h"
31
#include "core/data_type/data_type_struct.h"
32
#include "core/data_type/define_primitive_type.h"
33
#include "format/generic_reader.h"
34
#include "format/table/table_schema_change_helper.h"
35
#include "util/slice.h"
36
#include "util/string_util.h"
37
38
namespace doris {
39
40
5.27k
static bool is_group_node(const tparquet::SchemaElement& schema) {
41
5.27k
    return schema.num_children > 0;
42
5.27k
}
43
44
733
static bool is_list_node(const tparquet::SchemaElement& schema) {
45
733
    return schema.__isset.converted_type && schema.converted_type == tparquet::ConvertedType::LIST;
46
733
}
47
48
971
static bool is_map_node(const tparquet::SchemaElement& schema) {
49
971
    return schema.__isset.converted_type &&
50
971
           (schema.converted_type == tparquet::ConvertedType::MAP ||
51
618
            schema.converted_type == tparquet::ConvertedType::MAP_KEY_VALUE);
52
971
}
53
54
4.59k
static bool is_repeated_node(const tparquet::SchemaElement& schema) {
55
4.59k
    return schema.__isset.repetition_type &&
56
4.59k
           schema.repetition_type == tparquet::FieldRepetitionType::REPEATED;
57
4.59k
}
58
59
238
static bool is_required_node(const tparquet::SchemaElement& schema) {
60
238
    return schema.__isset.repetition_type &&
61
238
           schema.repetition_type == tparquet::FieldRepetitionType::REQUIRED;
62
238
}
63
64
4.74k
static bool is_optional_node(const tparquet::SchemaElement& schema) {
65
4.74k
    return schema.__isset.repetition_type &&
66
4.74k
           schema.repetition_type == tparquet::FieldRepetitionType::OPTIONAL;
67
4.74k
}
68
69
380
static int num_children_node(const tparquet::SchemaElement& schema) {
70
380
    return schema.__isset.num_children ? schema.num_children : 0;
71
380
}
72
73
/**
74
 * `repeated_parent_def_level` is the definition level of the first ancestor node whose repetition_type equals REPEATED.
75
 * Empty array/map values are not stored in doris columns, so have to use `repeated_parent_def_level` to skip the
76
 * empty or null values in ancestor node.
77
 *
78
 * For instance, considering an array of strings with 3 rows like the following:
79
 * null, [], [a, b, c]
80
 * We can store four elements in data column: null, a, b, c
81
 * and the offsets column is: 1, 1, 4
82
 * and the null map is: 1, 0, 0
83
 * For the i-th row in array column: range from `offsets[i - 1]` until `offsets[i]` represents the elements in this row,
84
 * so we can't store empty array/map values in doris data column.
85
 * As a comparison, spark does not require `repeated_parent_def_level`,
86
 * because the spark column stores empty array/map values , and use anther length column to indicate empty values.
87
 * Please reference: https://github.com/apache/spark/blob/master/sql/core/src/main/java/org/apache/spark/sql/execution/datasources/parquet/ParquetColumnVector.java
88
 *
89
 * Furthermore, we can also avoid store null array/map values in doris data column.
90
 * The same three rows as above, We can only store three elements in data column: a, b, c
91
 * and the offsets column is: 0, 0, 3
92
 * and the null map is: 1, 0, 0
93
 *
94
 * Inherit the repetition and definition level from parent node, if the parent node is repeated,
95
 * we should set repeated_parent_def_level = definition_level, otherwise as repeated_parent_def_level.
96
 * @param parent parent node
97
 * @param repeated_parent_def_level the first ancestor node whose repetition_type equals REPEATED
98
 */
99
971
static void set_child_node_level(FieldSchema* parent, int16_t repeated_parent_def_level) {
100
2.27k
    for (auto& child : parent->children) {
101
2.27k
        child.repetition_level = parent->repetition_level;
102
2.27k
        child.definition_level = parent->definition_level;
103
2.27k
        child.repeated_parent_def_level = repeated_parent_def_level;
104
2.27k
    }
105
971
}
106
107
380
static bool is_struct_list_node(const tparquet::SchemaElement& schema) {
108
380
    const std::string& name = schema.name;
109
380
    static const Slice array_slice("array", 5);
110
380
    static const Slice tuple_slice("_tuple", 6);
111
380
    Slice slice(name);
112
380
    return slice == array_slice || slice.ends_with(tuple_slice);
113
380
}
114
115
0
std::string FieldSchema::debug_string() const {
116
0
    std::stringstream ss;
117
0
    ss << "FieldSchema(name=" << name << ", R=" << repetition_level << ", D=" << definition_level;
118
0
    if (children.size() > 0) {
119
0
        ss << ", type=" << data_type->get_name() << ", children=[";
120
0
        for (int i = 0; i < children.size(); ++i) {
121
0
            if (i != 0) {
122
0
                ss << ", ";
123
0
            }
124
0
            ss << children[i].debug_string();
125
0
        }
126
0
        ss << "]";
127
0
    } else {
128
0
        ss << ", physical_type=" << physical_type;
129
0
        ss << " , doris_type=" << data_type->get_name();
130
0
    }
131
0
    ss << ")";
132
0
    return ss.str();
133
0
}
134
135
300
Status FieldDescriptor::parse_from_thrift(const std::vector<tparquet::SchemaElement>& t_schemas) {
136
300
    if (t_schemas.size() == 0 || !is_group_node(t_schemas[0])) {
137
0
        return Status::InvalidArgument("Wrong parquet root schema element");
138
0
    }
139
300
    const auto& root_schema = t_schemas[0];
140
300
    _fields.resize(root_schema.num_children);
141
300
    _next_schema_pos = 1;
142
143
2.76k
    for (int i = 0; i < root_schema.num_children; ++i) {
144
2.46k
        RETURN_IF_ERROR(parse_node_field(t_schemas, _next_schema_pos, &_fields[i]));
145
2.46k
        if (_name_to_field.find(_fields[i].name) != _name_to_field.end()) {
146
0
            return Status::InvalidArgument("Duplicated field name: {}", _fields[i].name);
147
0
        }
148
2.46k
        _name_to_field.emplace(_fields[i].name, &_fields[i]);
149
2.46k
    }
150
151
300
    if (_next_schema_pos != t_schemas.size()) {
152
0
        return Status::InvalidArgument("Remaining {} unparsed schema elements",
153
0
                                       t_schemas.size() - _next_schema_pos);
154
0
    }
155
156
300
    return Status::OK();
157
300
}
158
159
Status FieldDescriptor::parse_node_field(const std::vector<tparquet::SchemaElement>& t_schemas,
160
4.74k
                                         size_t curr_pos, FieldSchema* node_field) {
161
4.74k
    if (curr_pos >= t_schemas.size()) {
162
0
        return Status::InvalidArgument("Out-of-bounds index of schema elements");
163
0
    }
164
4.74k
    auto& t_schema = t_schemas[curr_pos];
165
4.74k
    if (is_group_node(t_schema)) {
166
        // nested structure or nullable list
167
971
        return parse_group_field(t_schemas, curr_pos, node_field);
168
971
    }
169
3.77k
    if (is_repeated_node(t_schema)) {
170
        // repeated <primitive-type> <name> (LIST)
171
        // produce required list<element>
172
0
        node_field->repetition_level++;
173
0
        node_field->definition_level++;
174
0
        node_field->children.resize(1);
175
0
        set_child_node_level(node_field, node_field->definition_level);
176
0
        auto child = &node_field->children[0];
177
0
        parse_physical_field(t_schema, false, child);
178
179
0
        node_field->name = t_schema.name;
180
0
        node_field->lower_case_name = to_lower(t_schema.name);
181
0
        node_field->data_type = std::make_shared<DataTypeArray>(make_nullable(child->data_type));
182
0
        _next_schema_pos = curr_pos + 1;
183
0
        node_field->field_id = t_schema.__isset.field_id ? t_schema.field_id : -1;
184
3.77k
    } else {
185
3.77k
        bool is_optional = is_optional_node(t_schema);
186
3.77k
        if (is_optional) {
187
3.27k
            node_field->definition_level++;
188
3.27k
        }
189
3.77k
        parse_physical_field(t_schema, is_optional, node_field);
190
3.77k
        _next_schema_pos = curr_pos + 1;
191
3.77k
    }
192
3.77k
    return Status::OK();
193
4.74k
}
194
195
void FieldDescriptor::parse_physical_field(const tparquet::SchemaElement& physical_schema,
196
3.77k
                                           bool is_nullable, FieldSchema* physical_field) {
197
3.77k
    physical_field->name = physical_schema.name;
198
3.77k
    physical_field->lower_case_name = to_lower(physical_field->name);
199
3.77k
    physical_field->parquet_schema = physical_schema;
200
3.77k
    physical_field->physical_type = physical_schema.type;
201
3.77k
    physical_field->column_id = UNASSIGNED_COLUMN_ID; // Initialize column_id
202
3.77k
    _physical_fields.push_back(physical_field);
203
3.77k
    physical_field->physical_column_index = cast_set<int>(_physical_fields.size() - 1);
204
3.77k
    auto type = get_doris_type(physical_schema, is_nullable);
205
3.77k
    physical_field->data_type = type.first;
206
3.77k
    physical_field->is_type_compatibility = type.second;
207
3.77k
    physical_field->field_id = physical_schema.__isset.field_id ? physical_schema.field_id : -1;
208
3.77k
}
209
210
std::pair<DataTypePtr, bool> FieldDescriptor::get_doris_type(
211
3.76k
        const tparquet::SchemaElement& physical_schema, bool nullable) {
212
3.76k
    std::pair<DataTypePtr, bool> ans = {std::make_shared<DataTypeNothing>(), false};
213
3.76k
    try {
214
3.76k
        if (physical_schema.__isset.logicalType) {
215
1.84k
            ans = convert_to_doris_type(physical_schema.logicalType, nullable);
216
1.92k
        } else if (physical_schema.__isset.converted_type) {
217
371
            ans = convert_to_doris_type(physical_schema, nullable);
218
371
        }
219
3.76k
    } catch (...) {
220
        // now the Not supported exception are ignored
221
        // so those byte_array maybe be treated as varbinary(now) : string(before)
222
0
    }
223
3.77k
    if (ans.first->get_primitive_type() == PrimitiveType::INVALID_TYPE) {
224
1.55k
        switch (physical_schema.type) {
225
194
        case tparquet::Type::BOOLEAN:
226
194
            ans.first = DataTypeFactory::instance().create_data_type(TYPE_BOOLEAN, nullable);
227
194
            break;
228
506
        case tparquet::Type::INT32:
229
506
            ans.first = DataTypeFactory::instance().create_data_type(TYPE_INT, nullable);
230
506
            break;
231
272
        case tparquet::Type::INT64:
232
272
            ans.first = DataTypeFactory::instance().create_data_type(TYPE_BIGINT, nullable);
233
272
            break;
234
52
        case tparquet::Type::INT96:
235
            // INT96 has no timezone annotation. A table's logical schema may override this
236
            // for known instant columns, but a catalog option cannot supply missing semantics.
237
52
            ans.first =
238
52
                    DataTypeFactory::instance().create_data_type(TYPE_DATETIMEV2, nullable, 0, 6);
239
52
            break;
240
180
        case tparquet::Type::FLOAT:
241
180
            ans.first = DataTypeFactory::instance().create_data_type(TYPE_FLOAT, nullable);
242
180
            break;
243
299
        case tparquet::Type::DOUBLE:
244
299
            ans.first = DataTypeFactory::instance().create_data_type(TYPE_DOUBLE, nullable);
245
299
            break;
246
49
        case tparquet::Type::BYTE_ARRAY:
247
49
            if (_enable_mapping_varbinary) {
248
                // if physical_schema not set logicalType and converted_type,
249
                // we treat BYTE_ARRAY as VARBINARY by default, so that we can read all data directly.
250
22
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_VARBINARY, nullable);
251
27
            } else {
252
27
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_STRING, nullable);
253
27
            }
254
49
            break;
255
0
        case tparquet::Type::FIXED_LEN_BYTE_ARRAY:
256
0
            ans.first = DataTypeFactory::instance().create_data_type(TYPE_STRING, nullable);
257
0
            break;
258
0
        default:
259
0
            throw Exception(Status::InternalError("Not supported parquet logicalType{}",
260
0
                                                  physical_schema.type));
261
0
            break;
262
1.55k
        }
263
1.55k
    }
264
3.77k
    return ans;
265
3.77k
}
266
267
std::pair<DataTypePtr, bool> FieldDescriptor::convert_to_doris_type(
268
1.84k
        tparquet::LogicalType logicalType, bool nullable) {
269
1.84k
    std::pair<DataTypePtr, bool> ans = {std::make_shared<DataTypeNothing>(), false};
270
1.84k
    bool& is_type_compatibility = ans.second;
271
1.84k
    if (logicalType.__isset.STRING) {
272
1.17k
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_STRING, nullable);
273
1.17k
    } else if (logicalType.__isset.DECIMAL) {
274
262
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_DECIMAL128I, nullable,
275
262
                                                                 logicalType.DECIMAL.precision,
276
262
                                                                 logicalType.DECIMAL.scale);
277
409
    } else if (logicalType.__isset.DATE) {
278
106
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_DATEV2, nullable);
279
303
    } else if (logicalType.__isset.INTEGER) {
280
187
        if (logicalType.INTEGER.isSigned) {
281
187
            if (logicalType.INTEGER.bitWidth <= 8) {
282
62
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_TINYINT, nullable);
283
125
            } else if (logicalType.INTEGER.bitWidth <= 16) {
284
125
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_SMALLINT, nullable);
285
125
            } else if (logicalType.INTEGER.bitWidth <= 32) {
286
0
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_INT, nullable);
287
0
            } else {
288
0
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_BIGINT, nullable);
289
0
            }
290
187
        } else {
291
0
            is_type_compatibility = true;
292
0
            if (logicalType.INTEGER.bitWidth <= 8) {
293
0
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_SMALLINT, nullable);
294
0
            } else if (logicalType.INTEGER.bitWidth <= 16) {
295
0
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_INT, nullable);
296
0
            } else if (logicalType.INTEGER.bitWidth <= 32) {
297
0
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_BIGINT, nullable);
298
0
            } else {
299
0
                ans.first = DataTypeFactory::instance().create_data_type(TYPE_LARGEINT, nullable);
300
0
            }
301
0
        }
302
187
    } else if (logicalType.__isset.TIME) {
303
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_TIMEV2, nullable);
304
116
    } else if (logicalType.__isset.TIMESTAMP) {
305
112
        if (_enable_mapping_timestamp_tz) {
306
0
            if (logicalType.TIMESTAMP.isAdjustedToUTC) {
307
                // treat TIMESTAMP with isAdjustedToUTC as TIMESTAMPTZ
308
0
                ans.first = DataTypeFactory::instance().create_data_type(
309
0
                        TYPE_TIMESTAMPTZ, nullable, 0,
310
0
                        logicalType.TIMESTAMP.unit.__isset.MILLIS ? 3 : 6);
311
0
                return ans;
312
0
            }
313
0
        }
314
112
        ans.first = DataTypeFactory::instance().create_data_type(
315
112
                TYPE_DATETIMEV2, nullable, 0, logicalType.TIMESTAMP.unit.__isset.MILLIS ? 3 : 6);
316
112
    } else if (logicalType.__isset.JSON) {
317
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_STRING, nullable);
318
4
    } else if (logicalType.__isset.UUID) {
319
4
        if (_enable_mapping_varbinary) {
320
3
            ans.first = DataTypeFactory::instance().create_data_type(TYPE_VARBINARY, nullable, -1,
321
3
                                                                     -1, 16);
322
3
        } else {
323
1
            ans.first = DataTypeFactory::instance().create_data_type(TYPE_STRING, nullable);
324
1
        }
325
4
    } else if (logicalType.__isset.FLOAT16) {
326
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_FLOAT, nullable);
327
0
    } else {
328
0
        throw Exception(Status::InternalError("Not supported parquet logicalType"));
329
0
    }
330
1.84k
    return ans;
331
1.84k
}
332
333
std::pair<DataTypePtr, bool> FieldDescriptor::convert_to_doris_type(
334
371
        const tparquet::SchemaElement& physical_schema, bool nullable) {
335
371
    std::pair<DataTypePtr, bool> ans = {std::make_shared<DataTypeNothing>(), false};
336
371
    bool& is_type_compatibility = ans.second;
337
371
    switch (physical_schema.converted_type) {
338
155
    case tparquet::ConvertedType::type::UTF8:
339
155
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_STRING, nullable);
340
155
        break;
341
62
    case tparquet::ConvertedType::type::DECIMAL:
342
62
        ans.first = DataTypeFactory::instance().create_data_type(
343
62
                TYPE_DECIMAL128I, nullable, physical_schema.precision, physical_schema.scale);
344
62
        break;
345
56
    case tparquet::ConvertedType::type::DATE:
346
56
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_DATEV2, nullable);
347
56
        break;
348
0
    case tparquet::ConvertedType::type::TIME_MILLIS:
349
0
        [[fallthrough]];
350
0
    case tparquet::ConvertedType::type::TIME_MICROS:
351
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_TIMEV2, nullable);
352
0
        break;
353
0
    case tparquet::ConvertedType::type::TIMESTAMP_MILLIS:
354
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_DATETIMEV2, nullable, 0, 3);
355
0
        break;
356
24
    case tparquet::ConvertedType::type::TIMESTAMP_MICROS:
357
24
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_DATETIMEV2, nullable, 0, 6);
358
24
        break;
359
36
    case tparquet::ConvertedType::type::INT_8:
360
36
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_TINYINT, nullable);
361
36
        break;
362
0
    case tparquet::ConvertedType::type::UINT_8:
363
0
        is_type_compatibility = true;
364
0
        [[fallthrough]];
365
38
    case tparquet::ConvertedType::type::INT_16:
366
38
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_SMALLINT, nullable);
367
38
        break;
368
0
    case tparquet::ConvertedType::type::UINT_16:
369
0
        is_type_compatibility = true;
370
0
        [[fallthrough]];
371
0
    case tparquet::ConvertedType::type::INT_32:
372
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_INT, nullable);
373
0
        break;
374
0
    case tparquet::ConvertedType::type::UINT_32:
375
0
        is_type_compatibility = true;
376
0
        [[fallthrough]];
377
0
    case tparquet::ConvertedType::type::INT_64:
378
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_BIGINT, nullable);
379
0
        break;
380
0
    case tparquet::ConvertedType::type::UINT_64:
381
0
        is_type_compatibility = true;
382
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_LARGEINT, nullable);
383
0
        break;
384
0
    case tparquet::ConvertedType::type::JSON:
385
0
        ans.first = DataTypeFactory::instance().create_data_type(TYPE_STRING, nullable);
386
0
        break;
387
0
    default:
388
0
        throw Exception(Status::InternalError("Not supported parquet ConvertedType: {}",
389
0
                                              physical_schema.converted_type));
390
371
    }
391
371
    return ans;
392
371
}
393
394
Status FieldDescriptor::parse_group_field(const std::vector<tparquet::SchemaElement>& t_schemas,
395
971
                                          size_t curr_pos, FieldSchema* group_field) {
396
971
    auto& group_schema = t_schemas[curr_pos];
397
971
    if (is_map_node(group_schema)) {
398
        // the map definition:
399
        // optional group <name> (MAP) {
400
        //   repeated group map (MAP_KEY_VALUE) {
401
        //     required <type> key;
402
        //     optional <type> value;
403
        //   }
404
        // }
405
238
        return parse_map_field(t_schemas, curr_pos, group_field);
406
238
    }
407
733
    if (is_list_node(group_schema)) {
408
        // the list definition:
409
        // optional group <name> (LIST) {
410
        //   repeated group [bag | list] { // hive or spark
411
        //     optional <type> [array_element | element]; // hive or spark
412
        //   }
413
        // }
414
380
        return parse_list_field(t_schemas, curr_pos, group_field);
415
380
    }
416
417
353
    if (is_repeated_node(group_schema)) {
418
0
        group_field->repetition_level++;
419
0
        group_field->definition_level++;
420
0
        group_field->children.resize(1);
421
0
        set_child_node_level(group_field, group_field->definition_level);
422
0
        auto struct_field = &group_field->children[0];
423
        // the list of struct:
424
        // repeated group <name> (LIST) {
425
        //   optional/required <type> <name>;
426
        //   ...
427
        // }
428
        // produce a non-null list<struct>
429
0
        RETURN_IF_ERROR(parse_struct_field(t_schemas, curr_pos, struct_field));
430
431
0
        group_field->name = group_schema.name;
432
0
        group_field->lower_case_name = to_lower(group_field->name);
433
0
        group_field->column_id = UNASSIGNED_COLUMN_ID; // Initialize column_id
434
0
        group_field->data_type =
435
0
                std::make_shared<DataTypeArray>(make_nullable(struct_field->data_type));
436
0
        group_field->field_id = group_schema.__isset.field_id ? group_schema.field_id : -1;
437
353
    } else {
438
353
        RETURN_IF_ERROR(parse_struct_field(t_schemas, curr_pos, group_field));
439
353
    }
440
441
353
    return Status::OK();
442
353
}
443
444
Status FieldDescriptor::parse_list_field(const std::vector<tparquet::SchemaElement>& t_schemas,
445
380
                                         size_t curr_pos, FieldSchema* list_field) {
446
    // the list definition:
447
    // spark and hive have three level schemas but with different schema name
448
    // spark: <column-name> - "list" - "element"
449
    // hive: <column-name> - "bag" - "array_element"
450
    // parse three level schemas to two level primitive like: LIST<INT>,
451
    // or nested structure like: LIST<MAP<INT, INT>>
452
380
    auto& first_level = t_schemas[curr_pos];
453
380
    if (first_level.num_children != 1) {
454
0
        return Status::InvalidArgument("List element should have only one child");
455
0
    }
456
457
380
    if (curr_pos + 1 >= t_schemas.size()) {
458
0
        return Status::InvalidArgument("List element should have the second level schema");
459
0
    }
460
461
380
    if (first_level.repetition_type == tparquet::FieldRepetitionType::REPEATED) {
462
0
        return Status::InvalidArgument("List element can't be a repeated schema");
463
0
    }
464
465
    // the repeated schema element
466
380
    auto& second_level = t_schemas[curr_pos + 1];
467
380
    if (second_level.repetition_type != tparquet::FieldRepetitionType::REPEATED) {
468
0
        return Status::InvalidArgument("The second level of list element should be repeated");
469
0
    }
470
471
    // This indicates if this list is nullable.
472
380
    bool is_optional = is_optional_node(first_level);
473
380
    if (is_optional) {
474
328
        list_field->definition_level++;
475
328
    }
476
380
    list_field->repetition_level++;
477
380
    list_field->definition_level++;
478
380
    list_field->children.resize(1);
479
380
    FieldSchema* list_child = &list_field->children[0];
480
481
380
    size_t num_children = num_children_node(second_level);
482
380
    if (num_children > 0) {
483
380
        if (num_children == 1 && !is_struct_list_node(second_level)) {
484
            // optional field, and the third level element is the nested structure in list
485
            // produce nested structure like: LIST<INT>, LIST<MAP>, LIST<LIST<...>>
486
            // skip bag/list, it's a repeated element.
487
380
            set_child_node_level(list_field, list_field->definition_level);
488
380
            RETURN_IF_ERROR(parse_node_field(t_schemas, curr_pos + 2, list_child));
489
380
        } else {
490
            // required field, produce the list of struct
491
0
            set_child_node_level(list_field, list_field->definition_level);
492
0
            RETURN_IF_ERROR(parse_struct_field(t_schemas, curr_pos + 1, list_child));
493
0
        }
494
380
    } else if (num_children == 0) {
495
        // required two level list, for compatibility reason.
496
0
        set_child_node_level(list_field, list_field->definition_level);
497
0
        parse_physical_field(second_level, false, list_child);
498
0
        _next_schema_pos = curr_pos + 2;
499
0
    }
500
501
380
    list_field->name = first_level.name;
502
380
    list_field->lower_case_name = to_lower(first_level.name);
503
380
    list_field->column_id = UNASSIGNED_COLUMN_ID; // Initialize column_id
504
380
    list_field->data_type =
505
380
            std::make_shared<DataTypeArray>(make_nullable(list_field->children[0].data_type));
506
380
    if (is_optional) {
507
328
        list_field->data_type = make_nullable(list_field->data_type);
508
328
    }
509
380
    list_field->field_id = first_level.__isset.field_id ? first_level.field_id : -1;
510
511
380
    return Status::OK();
512
380
}
513
514
Status FieldDescriptor::parse_map_field(const std::vector<tparquet::SchemaElement>& t_schemas,
515
238
                                        size_t curr_pos, FieldSchema* map_field) {
516
    // the map definition in parquet:
517
    // optional group <name> (MAP) {
518
    //   repeated group map (MAP_KEY_VALUE) {
519
    //     required <type> key;
520
    //     optional <type> value;
521
    //   }
522
    // }
523
    // Map value can be optional, the map without values is a SET
524
238
    if (curr_pos + 2 >= t_schemas.size()) {
525
0
        return Status::InvalidArgument("Map element should have at least three levels");
526
0
    }
527
238
    auto& map_schema = t_schemas[curr_pos];
528
238
    if (map_schema.num_children != 1) {
529
0
        return Status::InvalidArgument(
530
0
                "Map element should have only one child(name='map', type='MAP_KEY_VALUE')");
531
0
    }
532
238
    if (is_repeated_node(map_schema)) {
533
0
        return Status::InvalidArgument("Map element can't be a repeated schema");
534
0
    }
535
238
    auto& map_key_value = t_schemas[curr_pos + 1];
536
238
    if (!is_group_node(map_key_value) || !is_repeated_node(map_key_value)) {
537
0
        return Status::InvalidArgument(
538
0
                "the second level in map must be a repeated group(key and value)");
539
0
    }
540
238
    auto& map_key = t_schemas[curr_pos + 2];
541
238
    if (!is_required_node(map_key)) {
542
0
        LOG(WARNING) << "Filed " << map_schema.name << " is map type, but with nullable key column";
543
0
    }
544
545
238
    if (map_key_value.num_children == 1) {
546
        // The map with three levels is a SET
547
0
        return parse_list_field(t_schemas, curr_pos, map_field);
548
0
    }
549
238
    if (map_key_value.num_children != 2) {
550
        // A standard map should have four levels
551
0
        return Status::InvalidArgument(
552
0
                "the second level in map(MAP_KEY_VALUE) should have two children");
553
0
    }
554
    // standard map
555
238
    bool is_optional = is_optional_node(map_schema);
556
238
    if (is_optional) {
557
212
        map_field->definition_level++;
558
212
    }
559
238
    map_field->repetition_level++;
560
238
    map_field->definition_level++;
561
562
    // Directly create key and value children instead of intermediate key_value node
563
238
    map_field->children.resize(2);
564
    // map is a repeated node, we should set the `repeated_parent_def_level` of its children as `definition_level`
565
238
    set_child_node_level(map_field, map_field->definition_level);
566
567
238
    auto key_field = &map_field->children[0];
568
238
    auto value_field = &map_field->children[1];
569
570
    // Parse key and value fields directly from the key_value group's children
571
238
    _next_schema_pos = curr_pos + 2; // Skip key_value group, go directly to key
572
238
    RETURN_IF_ERROR(parse_node_field(t_schemas, _next_schema_pos, key_field));
573
238
    RETURN_IF_ERROR(parse_node_field(t_schemas, _next_schema_pos, value_field));
574
575
238
    map_field->name = map_schema.name;
576
238
    map_field->lower_case_name = to_lower(map_field->name);
577
238
    map_field->column_id = UNASSIGNED_COLUMN_ID; // Initialize column_id
578
238
    map_field->data_type = std::make_shared<DataTypeMap>(make_nullable(key_field->data_type),
579
238
                                                         make_nullable(value_field->data_type));
580
238
    if (is_optional) {
581
212
        map_field->data_type = make_nullable(map_field->data_type);
582
212
    }
583
238
    map_field->field_id = map_schema.__isset.field_id ? map_schema.field_id : -1;
584
585
238
    return Status::OK();
586
238
}
587
588
Status FieldDescriptor::parse_struct_field(const std::vector<tparquet::SchemaElement>& t_schemas,
589
353
                                           size_t curr_pos, FieldSchema* struct_field) {
590
    // the nested column in parquet, parse group to struct.
591
353
    auto& struct_schema = t_schemas[curr_pos];
592
353
    bool is_optional = is_optional_node(struct_schema);
593
353
    if (is_optional) {
594
334
        struct_field->definition_level++;
595
334
    }
596
353
    auto num_children = struct_schema.num_children;
597
353
    struct_field->children.resize(num_children);
598
353
    set_child_node_level(struct_field, struct_field->repeated_parent_def_level);
599
353
    _next_schema_pos = curr_pos + 1;
600
1.77k
    for (int i = 0; i < num_children; ++i) {
601
1.42k
        RETURN_IF_ERROR(parse_node_field(t_schemas, _next_schema_pos, &struct_field->children[i]));
602
1.42k
    }
603
353
    struct_field->name = struct_schema.name;
604
353
    struct_field->lower_case_name = to_lower(struct_field->name);
605
353
    struct_field->column_id = UNASSIGNED_COLUMN_ID; // Initialize column_id
606
607
353
    struct_field->field_id = struct_schema.__isset.field_id ? struct_schema.field_id : -1;
608
353
    DataTypes res_data_types;
609
353
    std::vector<String> names;
610
1.77k
    for (int i = 0; i < num_children; ++i) {
611
1.42k
        res_data_types.push_back(make_nullable(struct_field->children[i].data_type));
612
1.42k
        names.push_back(struct_field->children[i].name);
613
1.42k
    }
614
353
    struct_field->data_type = std::make_shared<DataTypeStruct>(res_data_types, names);
615
353
    if (is_optional) {
616
334
        struct_field->data_type = make_nullable(struct_field->data_type);
617
334
    }
618
353
    return Status::OK();
619
353
}
620
621
5
int FieldDescriptor::get_column_index(const std::string& column) const {
622
15
    for (int32_t i = 0; i < _fields.size(); i++) {
623
15
        if (_fields[i].name == column) {
624
5
            return i;
625
5
        }
626
15
    }
627
0
    return -1;
628
5
}
629
630
3.42k
FieldSchema* FieldDescriptor::get_column(const std::string& name) const {
631
3.42k
    auto it = _name_to_field.find(name);
632
3.42k
    if (it != _name_to_field.end()) {
633
3.42k
        return it->second;
634
3.42k
    }
635
18.4E
    throw Exception(Status::InternalError("Name {} not found in FieldDescriptor!", name));
636
0
    return nullptr;
637
3.42k
}
638
639
37
void FieldDescriptor::get_column_names(std::unordered_set<std::string>* names) const {
640
37
    names->clear();
641
650
    for (const FieldSchema& f : _fields) {
642
650
        names->emplace(f.name);
643
650
    }
644
37
}
645
646
0
std::string FieldDescriptor::debug_string() const {
647
0
    std::stringstream ss;
648
0
    ss << "fields=[";
649
0
    for (int i = 0; i < _fields.size(); ++i) {
650
0
        if (i != 0) {
651
0
            ss << ", ";
652
0
        }
653
0
        ss << _fields[i].debug_string();
654
0
    }
655
0
    ss << "]";
656
0
    return ss.str();
657
0
}
658
659
51
void FieldDescriptor::assign_ids() {
660
51
    uint64_t next_id = 1;
661
398
    for (auto& field : _fields) {
662
398
        field.assign_ids(next_id);
663
398
    }
664
51
}
665
666
0
const FieldSchema* FieldDescriptor::find_column_by_id(uint64_t column_id) const {
667
0
    for (const auto& field : _fields) {
668
0
        if (auto result = field.find_column_by_id(column_id)) {
669
0
            return result;
670
0
        }
671
0
    }
672
0
    return nullptr;
673
0
}
674
675
1.95k
void FieldSchema::assign_ids(uint64_t& next_id) {
676
1.95k
    column_id = next_id++;
677
678
1.95k
    for (auto& child : children) {
679
1.55k
        child.assign_ids(next_id);
680
1.55k
    }
681
682
1.95k
    max_column_id = next_id - 1;
683
1.95k
}
684
685
0
const FieldSchema* FieldSchema::find_column_by_id(uint64_t target_id) const {
686
0
    if (column_id == target_id) {
687
0
        return this;
688
0
    }
689
690
0
    for (const auto& child : children) {
691
0
        if (auto result = child.find_column_by_id(target_id)) {
692
0
            return result;
693
0
        }
694
0
    }
695
696
0
    return nullptr;
697
0
}
698
699
399
uint64_t FieldSchema::get_column_id() const {
700
399
    return column_id;
701
399
}
702
703
0
void FieldSchema::set_column_id(uint64_t id) {
704
0
    column_id = id;
705
0
}
706
707
121
uint64_t FieldSchema::get_max_column_id() const {
708
121
    return max_column_id;
709
121
}
710
711
} // namespace doris