Coverage Report

Created: 2026-10-09 18:00

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/parquet/parquet_column_convert.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/parquet_column_convert.h"
19
20
#include <cctz/time_zone.h>
21
#include <glog/logging.h>
22
23
#include "common/cast_set.h"
24
#include "core/column/column_fixed_length_object.h"
25
#include "core/column/column_nullable.h"
26
#include "core/data_type/data_type_fixed_length_object.h"
27
#include "core/data_type/data_type_nullable.h"
28
#include "core/data_type/define_primitive_type.h"
29
#include "core/data_type/primitive_type.h"
30
#include "core/value/uuid_value.h"
31
32
namespace doris::parquet {
33
const cctz::time_zone ConvertParams::utc0 = cctz::utc_time_zone();
34
35
#define FOR_LOGICAL_DECIMAL_TYPES(M) \
36
0
    M(TYPE_DECIMAL32)                \
37
8
    M(TYPE_DECIMAL64)                \
38
8
    M(TYPE_DECIMAL128I)              \
39
0
    M(TYPE_DECIMAL256)
40
41
642
bool PhysicalToLogicalConverter::is_parquet_native_type(PrimitiveType type) {
42
642
    switch (type) {
43
24
    case TYPE_BOOLEAN:
44
264
    case TYPE_INT:
45
338
    case TYPE_BIGINT:
46
398
    case TYPE_FLOAT:
47
436
    case TYPE_DOUBLE:
48
542
    case TYPE_STRING:
49
542
    case TYPE_CHAR:
50
542
    case TYPE_VARCHAR:
51
542
        return true;
52
100
    default:
53
100
        return false;
54
642
    }
55
642
}
56
57
66
bool PhysicalToLogicalConverter::is_decimal_type(doris::PrimitiveType type) {
58
66
    switch (type) {
59
0
    case TYPE_DECIMAL32:
60
8
    case TYPE_DECIMAL64:
61
8
    case TYPE_DECIMAL128I:
62
8
    case TYPE_DECIMAL256:
63
8
    case TYPE_DECIMALV2:
64
8
        return true;
65
58
    default:
66
58
        return false;
67
66
    }
68
66
}
69
70
ColumnPtr PhysicalToLogicalConverter::get_physical_column(tparquet::Type::type src_physical_type,
71
                                                          DataTypePtr src_logical_type,
72
                                                          ColumnPtr& dst_logical_column,
73
                                                          const DataTypePtr& dst_logical_type,
74
712
                                                          bool is_dict_filter) {
75
712
    if (is_dict_filter) {
76
2
        src_physical_type = tparquet::Type::INT32;
77
2
        src_logical_type = DataTypeFactory::instance().create_data_type(
78
2
                PrimitiveType::TYPE_INT, dst_logical_type->is_nullable());
79
2
    }
80
81
712
    if (!_convert_params->is_type_compatibility && is_consistent() &&
82
712
        _logical_converter->is_consistent()) {
83
616
        if (_cached_src_physical_type == nullptr) {
84
332
            _cached_src_physical_type = dst_logical_type->is_nullable()
85
332
                                                ? make_nullable(src_logical_type)
86
332
                                                : remove_nullable(src_logical_type);
87
332
        }
88
616
        return dst_logical_column;
89
616
    }
90
91
96
    if (!_cached_src_physical_column) {
92
74
        switch (src_physical_type) {
93
0
        case tparquet::Type::type::BOOLEAN:
94
0
            _cached_src_physical_type = std::make_shared<DataTypeUInt8>();
95
0
            break;
96
46
        case tparquet::Type::type::INT32:
97
46
            _cached_src_physical_type = std::make_shared<DataTypeInt32>();
98
46
            break;
99
12
        case tparquet::Type::type::INT64:
100
12
            _cached_src_physical_type = std::make_shared<DataTypeInt64>();
101
12
            break;
102
0
        case tparquet::Type::type::FLOAT:
103
0
            _cached_src_physical_type = std::make_shared<DataTypeFloat32>();
104
0
            break;
105
0
        case tparquet::Type::type::DOUBLE:
106
0
            _cached_src_physical_type = std::make_shared<DataTypeFloat64>();
107
0
            break;
108
8
        case tparquet::Type::type::BYTE_ARRAY:
109
8
            _cached_src_physical_type = std::make_shared<DataTypeString>();
110
8
            break;
111
8
        case tparquet::Type::type::FIXED_LEN_BYTE_ARRAY:
112
8
            _cached_src_physical_type = std::make_shared<DataTypeFixedLengthObject>();
113
8
            break;
114
0
        case tparquet::Type::type::INT96:
115
0
            _cached_src_physical_type = std::make_shared<DataTypeInt8>();
116
0
            break;
117
74
        }
118
74
        const bool is_fixed_length_byte_array =
119
74
                src_physical_type == tparquet::Type::type::FIXED_LEN_BYTE_ARRAY;
120
74
        if (dst_logical_type->is_nullable()) {
121
74
            MutableColumnPtr nested_physical_column;
122
74
            if (is_fixed_length_byte_array) {
123
8
                nested_physical_column = ColumnFixedLengthObject::create(
124
8
                        _convert_params->field_schema->parquet_schema.type_length);
125
66
            } else {
126
66
                nested_physical_column = _cached_src_physical_type->create_column();
127
66
            }
128
74
            _cached_src_physical_column = ColumnNullable::create(std::move(nested_physical_column),
129
74
                                                                 ColumnUInt8::create());
130
74
            _cached_src_physical_type = make_nullable(_cached_src_physical_type);
131
74
        } else {
132
0
            if (is_fixed_length_byte_array) {
133
0
                _cached_src_physical_column = ColumnFixedLengthObject::create(
134
0
                        _convert_params->field_schema->parquet_schema.type_length);
135
0
            } else {
136
0
                _cached_src_physical_column = _cached_src_physical_type->create_column();
137
0
            }
138
0
        }
139
74
    }
140
    // remove the old cached data
141
96
    auto cached_src_physical_column = IColumn::mutate(std::move(_cached_src_physical_column));
142
96
    cached_src_physical_column->clear();
143
96
    _cached_src_physical_column = std::move(cached_src_physical_column);
144
145
96
    return _cached_src_physical_column;
146
96
}
147
148
static void get_decimal_converter(const FieldSchema* field_schema, DataTypePtr src_logical_type,
149
                                  const DataTypePtr& dst_logical_type,
150
                                  ConvertParams* convert_params,
151
8
                                  std::unique_ptr<PhysicalToLogicalConverter>& physical_converter) {
152
8
    const tparquet::SchemaElement& parquet_schema = field_schema->parquet_schema;
153
8
    if (is_decimal(dst_logical_type->get_primitive_type())) {
154
8
        src_logical_type = create_decimal(parquet_schema.precision, parquet_schema.scale, false);
155
8
    }
156
157
8
    tparquet::Type::type src_physical_type = parquet_schema.type;
158
8
    PrimitiveType src_logical_primitive = src_logical_type->get_primitive_type();
159
160
8
    if (src_physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY) {
161
4
        switch (src_logical_primitive) {
162
0
#define DISPATCH(LOGICAL_PTYPE)                                                     \
163
4
    case LOGICAL_PTYPE: {                                                           \
164
4
        physical_converter.reset(                                                   \
165
4
                new FixedSizeToDecimal<LOGICAL_PTYPE>(parquet_schema.type_length)); \
166
4
        break;                                                                      \
167
4
    }
168
4
            FOR_LOGICAL_DECIMAL_TYPES(DISPATCH)
169
0
#undef DISPATCH
170
0
        default:
171
0
            physical_converter =
172
0
                    std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
173
4
        }
174
4
    } else if (src_physical_type == tparquet::Type::BYTE_ARRAY) {
175
0
        switch (src_logical_primitive) {
176
0
#define DISPATCH(LOGICAL_PTYPE)                                         \
177
0
    case LOGICAL_PTYPE: {                                               \
178
0
        physical_converter.reset(new StringToDecimal<LOGICAL_PTYPE>()); \
179
0
        break;                                                          \
180
0
    }
181
0
            FOR_LOGICAL_DECIMAL_TYPES(DISPATCH)
182
0
#undef DISPATCH
183
0
        default:
184
0
            physical_converter =
185
0
                    std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
186
0
        }
187
4
    } else if (src_physical_type == tparquet::Type::INT32 ||
188
4
               src_physical_type == tparquet::Type::INT64) {
189
4
        switch (src_logical_primitive) {
190
0
#define DISPATCH(LOGICAL_PTYPE)                                                          \
191
4
    case LOGICAL_PTYPE: {                                                                \
192
4
        if (src_physical_type == tparquet::Type::INT32) {                                \
193
0
            physical_converter.reset(new NumberToDecimal<TYPE_INT, LOGICAL_PTYPE>());    \
194
4
        } else {                                                                         \
195
4
            physical_converter.reset(new NumberToDecimal<TYPE_BIGINT, LOGICAL_PTYPE>()); \
196
4
        }                                                                                \
197
4
        break;                                                                           \
198
4
    }
199
4
            FOR_LOGICAL_DECIMAL_TYPES(DISPATCH)
200
0
#undef DISPATCH
201
0
        default:
202
0
            physical_converter =
203
0
                    std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
204
4
        }
205
4
    } else {
206
0
        physical_converter =
207
0
                std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
208
0
    }
209
8
}
210
211
namespace {
212
213
class UUIDStringConverter final : public PhysicalToLogicalConverter {
214
public:
215
8
    Status physical_convert(ColumnPtr& src_physical_col, ColumnPtr& src_logical_column) override {
216
8
        const auto from_col = remove_nullable(src_physical_col);
217
8
        const auto src_data = get_fixed_length_physical_data(*from_col, 16);
218
8
        auto& strings = assert_cast<ColumnString&>(*get_mutable_inner_column(src_logical_column));
219
32
        for (size_t i = 0; i < src_data.rows; ++i) {
220
24
            const auto value = UUIDValue::from_big_endian(
221
24
                    reinterpret_cast<const uint8_t*>(src_data.data + i * 16));
222
24
            const auto text = UUIDValue::to_string(value);
223
24
            strings.insert_data(text.data(), text.size());
224
24
        }
225
8
        return Status::OK();
226
8
    }
227
};
228
229
} // namespace
230
231
std::unique_ptr<PhysicalToLogicalConverter> PhysicalToLogicalConverter::get_converter(
232
        const FieldSchema* field_schema, DataTypePtr src_logical_type,
233
        const DataTypePtr& dst_logical_type, const cctz::time_zone* ctz, bool is_dict_filter,
234
642
        bool preserve_binary_uuid) {
235
642
    std::unique_ptr<ConvertParams> convert_params = std::make_unique<ConvertParams>();
236
642
    const tparquet::SchemaElement& parquet_schema = field_schema->parquet_schema;
237
642
    convert_params->init(field_schema, ctz);
238
642
    tparquet::Type::type src_physical_type = parquet_schema.type;
239
642
    std::unique_ptr<PhysicalToLogicalConverter> physical_converter;
240
642
    if (is_dict_filter) {
241
2
        src_physical_type = tparquet::Type::INT32;
242
2
        src_logical_type = DataTypeFactory::instance().create_data_type(
243
2
                PrimitiveType::TYPE_INT, dst_logical_type->is_nullable());
244
2
    }
245
642
    PrimitiveType src_logical_primitive = src_logical_type->get_primitive_type();
246
247
642
    if (field_schema->is_type_compatibility) {
248
0
        if (src_logical_primitive == TYPE_SMALLINT) {
249
0
            physical_converter = std::make_unique<UnsignedIntegerConverter<TYPE_SMALLINT>>();
250
0
        } else if (src_logical_primitive == TYPE_INT) {
251
0
            physical_converter = std::make_unique<UnsignedIntegerConverter<TYPE_INT>>();
252
0
        } else if (src_logical_primitive == TYPE_BIGINT) {
253
0
            physical_converter = std::make_unique<UnsignedIntegerConverter<TYPE_BIGINT>>();
254
0
        } else if (src_logical_primitive == TYPE_LARGEINT) {
255
0
            physical_converter = std::make_unique<UnsignedIntegerConverter<TYPE_LARGEINT>>();
256
0
        } else {
257
0
            physical_converter =
258
0
                    std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
259
0
        }
260
642
    } else if (is_parquet_native_type(src_logical_primitive)) {
261
542
        bool is_string_logical_type = is_string_type(src_logical_primitive);
262
542
        if (is_string_logical_type && src_physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY) {
263
14
            if (parquet_schema.logicalType.__isset.UUID && !preserve_binary_uuid) {
264
                // Ordinary Parquet scans expose canonical text, as in File Scanner V2.
265
                // Iceberg retains its raw 16-byte STRING carrier for equality deletes/defaults.
266
8
                if (parquet_schema.type_length != UUIDValue::BINARY_LENGTH) {
267
2
                    physical_converter = std::make_unique<UnsupportedConverter>(src_physical_type,
268
2
                                                                                src_logical_type);
269
6
                } else {
270
6
                    physical_converter = std::make_unique<UUIDStringConverter>();
271
6
                }
272
8
            } else {
273
6
                physical_converter =
274
6
                        std::make_unique<FixedSizeBinaryConverter>(parquet_schema.type_length);
275
6
            }
276
528
        } else if (src_logical_primitive == TYPE_FLOAT &&
277
528
                   src_physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY &&
278
528
                   parquet_schema.logicalType.__isset.FLOAT16) {
279
0
            physical_converter =
280
0
                    std::make_unique<Float16PhysicalConverter>(parquet_schema.type_length);
281
528
        } else {
282
528
            physical_converter = std::make_unique<ConsistentPhysicalConverter>();
283
528
        }
284
542
    } else if (src_logical_primitive == TYPE_TINYINT) {
285
20
        physical_converter = std::make_unique<LittleIntPhysicalConverter<TYPE_TINYINT>>();
286
80
    } else if (src_logical_primitive == TYPE_SMALLINT) {
287
14
        physical_converter = std::make_unique<LittleIntPhysicalConverter<TYPE_SMALLINT>>();
288
66
    } else if (is_decimal_type(src_logical_primitive)) {
289
8
        get_decimal_converter(field_schema, src_logical_type, dst_logical_type,
290
8
                              convert_params.get(), physical_converter);
291
58
    } else if (src_logical_primitive == TYPE_DATEV2) {
292
14
        physical_converter = std::make_unique<Int32ToDate>();
293
44
    } else if (src_logical_primitive == TYPE_DATETIMEV2) {
294
32
        if (src_physical_type == tparquet::Type::INT96) {
295
            // int96 only stores nanoseconds in standard parquet file
296
2
            convert_params->reset_time_scale_if_missing(9);
297
2
            physical_converter = std::make_unique<Int96toTimestamp>();
298
30
        } else if (src_physical_type == tparquet::Type::INT64) {
299
30
            convert_params->reset_time_scale_if_missing(src_logical_type->get_scale());
300
30
            physical_converter = std::make_unique<Int64ToTimestamp>();
301
30
        } else {
302
0
            physical_converter =
303
0
                    std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
304
0
        }
305
32
    } else if (src_logical_primitive == TYPE_VARBINARY) {
306
8
        if (src_physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY) {
307
4
            DCHECK(parquet_schema.logicalType.__isset.UUID) << parquet_schema.name;
308
4
            physical_converter =
309
4
                    std::make_unique<UUIDVarBinaryConverter>(parquet_schema.type_length);
310
4
        } else {
311
4
            DCHECK(src_physical_type == tparquet::Type::BYTE_ARRAY) << src_physical_type;
312
4
            physical_converter = std::make_unique<ConsistentPhysicalConverter>();
313
4
        }
314
8
    } else if (src_logical_primitive == TYPE_TIMESTAMPTZ) {
315
4
        if (src_physical_type == tparquet::Type::INT96) {
316
2
            physical_converter = std::make_unique<Int96toTimestampTz>();
317
2
        } else if (src_physical_type == tparquet::Type::INT64) {
318
2
            DCHECK(src_physical_type == tparquet::Type::INT64) << src_physical_type;
319
2
            DCHECK(parquet_schema.logicalType.__isset.TIMESTAMP) << parquet_schema.name;
320
2
            physical_converter = std::make_unique<Int64ToTimestampTz>();
321
2
        } else {
322
0
            physical_converter =
323
0
                    std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
324
0
        }
325
4
    } else {
326
0
        physical_converter =
327
0
                std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
328
0
    }
329
330
642
    if (physical_converter->support()) {
331
640
        physical_converter->_convert_params = std::move(convert_params);
332
640
        physical_converter->_logical_converter = converter::ColumnTypeConverter::get_converter(
333
640
                src_logical_type, dst_logical_type, converter::FileFormat::PARQUET);
334
640
        if (!physical_converter->_logical_converter->support()) {
335
0
            physical_converter = std::make_unique<UnsupportedConverter>(
336
0
                    "Unsupported type change: " +
337
0
                    physical_converter->_logical_converter->get_error_msg());
338
0
        }
339
640
    }
340
642
    return physical_converter;
341
642
}
342
343
} // namespace doris::parquet