Coverage Report

Created: 2026-09-29 15:53

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
4
    M(TYPE_DECIMAL64)                \
38
4
    M(TYPE_DECIMAL128I)              \
39
0
    M(TYPE_DECIMAL256)
40
41
321
bool PhysicalToLogicalConverter::is_parquet_native_type(PrimitiveType type) {
42
321
    switch (type) {
43
12
    case TYPE_BOOLEAN:
44
132
    case TYPE_INT:
45
169
    case TYPE_BIGINT:
46
199
    case TYPE_FLOAT:
47
218
    case TYPE_DOUBLE:
48
271
    case TYPE_STRING:
49
271
    case TYPE_CHAR:
50
271
    case TYPE_VARCHAR:
51
271
        return true;
52
50
    default:
53
50
        return false;
54
321
    }
55
321
}
56
57
33
bool PhysicalToLogicalConverter::is_decimal_type(doris::PrimitiveType type) {
58
33
    switch (type) {
59
0
    case TYPE_DECIMAL32:
60
4
    case TYPE_DECIMAL64:
61
4
    case TYPE_DECIMAL128I:
62
4
    case TYPE_DECIMAL256:
63
4
    case TYPE_DECIMALV2:
64
4
        return true;
65
29
    default:
66
29
        return false;
67
33
    }
68
33
}
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
356
                                                          bool is_dict_filter) {
75
356
    if (is_dict_filter) {
76
1
        src_physical_type = tparquet::Type::INT32;
77
1
        src_logical_type = DataTypeFactory::instance().create_data_type(
78
1
                PrimitiveType::TYPE_INT, dst_logical_type->is_nullable());
79
1
    }
80
81
356
    if (!_convert_params->is_type_compatibility && is_consistent() &&
82
356
        _logical_converter->is_consistent()) {
83
308
        if (_cached_src_physical_type == nullptr) {
84
166
            _cached_src_physical_type = dst_logical_type->is_nullable()
85
166
                                                ? make_nullable(src_logical_type)
86
166
                                                : remove_nullable(src_logical_type);
87
166
        }
88
308
        return dst_logical_column;
89
308
    }
90
91
48
    if (!_cached_src_physical_column) {
92
37
        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
23
        case tparquet::Type::type::INT32:
97
23
            _cached_src_physical_type = std::make_shared<DataTypeInt32>();
98
23
            break;
99
6
        case tparquet::Type::type::INT64:
100
6
            _cached_src_physical_type = std::make_shared<DataTypeInt64>();
101
6
            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
4
        case tparquet::Type::type::BYTE_ARRAY:
109
4
            _cached_src_physical_type = std::make_shared<DataTypeString>();
110
4
            break;
111
4
        case tparquet::Type::type::FIXED_LEN_BYTE_ARRAY:
112
4
            _cached_src_physical_type = std::make_shared<DataTypeFixedLengthObject>();
113
4
            break;
114
0
        case tparquet::Type::type::INT96:
115
0
            _cached_src_physical_type = std::make_shared<DataTypeInt8>();
116
0
            break;
117
37
        }
118
37
        const bool is_fixed_length_byte_array =
119
37
                src_physical_type == tparquet::Type::type::FIXED_LEN_BYTE_ARRAY;
120
37
        if (dst_logical_type->is_nullable()) {
121
37
            MutableColumnPtr nested_physical_column;
122
37
            if (is_fixed_length_byte_array) {
123
4
                nested_physical_column = ColumnFixedLengthObject::create(
124
4
                        _convert_params->field_schema->parquet_schema.type_length);
125
33
            } else {
126
33
                nested_physical_column = _cached_src_physical_type->create_column();
127
33
            }
128
37
            _cached_src_physical_column = ColumnNullable::create(std::move(nested_physical_column),
129
37
                                                                 ColumnUInt8::create());
130
37
            _cached_src_physical_type = make_nullable(_cached_src_physical_type);
131
37
        } 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
37
    }
140
    // remove the old cached data
141
48
    auto cached_src_physical_column = IColumn::mutate(std::move(_cached_src_physical_column));
142
48
    cached_src_physical_column->clear();
143
48
    _cached_src_physical_column = std::move(cached_src_physical_column);
144
145
48
    return _cached_src_physical_column;
146
48
}
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
4
                                  std::unique_ptr<PhysicalToLogicalConverter>& physical_converter) {
152
4
    const tparquet::SchemaElement& parquet_schema = field_schema->parquet_schema;
153
4
    if (is_decimal(dst_logical_type->get_primitive_type())) {
154
4
        src_logical_type = create_decimal(parquet_schema.precision, parquet_schema.scale, false);
155
4
    }
156
157
4
    tparquet::Type::type src_physical_type = parquet_schema.type;
158
4
    PrimitiveType src_logical_primitive = src_logical_type->get_primitive_type();
159
160
4
    if (src_physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY) {
161
2
        switch (src_logical_primitive) {
162
0
#define DISPATCH(LOGICAL_PTYPE)                                                     \
163
2
    case LOGICAL_PTYPE: {                                                           \
164
2
        physical_converter.reset(                                                   \
165
2
                new FixedSizeToDecimal<LOGICAL_PTYPE>(parquet_schema.type_length)); \
166
2
        break;                                                                      \
167
2
    }
168
2
            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
2
        }
174
2
    } 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
2
    } else if (src_physical_type == tparquet::Type::INT32 ||
188
2
               src_physical_type == tparquet::Type::INT64) {
189
2
        switch (src_logical_primitive) {
190
0
#define DISPATCH(LOGICAL_PTYPE)                                                          \
191
2
    case LOGICAL_PTYPE: {                                                                \
192
2
        if (src_physical_type == tparquet::Type::INT32) {                                \
193
0
            physical_converter.reset(new NumberToDecimal<TYPE_INT, LOGICAL_PTYPE>());    \
194
2
        } else {                                                                         \
195
2
            physical_converter.reset(new NumberToDecimal<TYPE_BIGINT, LOGICAL_PTYPE>()); \
196
2
        }                                                                                \
197
2
        break;                                                                           \
198
2
    }
199
2
            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
2
        }
205
2
    } else {
206
0
        physical_converter =
207
0
                std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
208
0
    }
209
4
}
210
211
namespace {
212
213
class UUIDStringConverter final : public PhysicalToLogicalConverter {
214
public:
215
4
    Status physical_convert(ColumnPtr& src_physical_col, ColumnPtr& src_logical_column) override {
216
4
        const auto from_col = remove_nullable(src_physical_col);
217
4
        const auto src_data = get_fixed_length_physical_data(*from_col, 16);
218
4
        auto& strings = assert_cast<ColumnString&>(*get_mutable_inner_column(src_logical_column));
219
16
        for (size_t i = 0; i < src_data.rows; ++i) {
220
12
            const auto value = UUIDValue::from_big_endian(
221
12
                    reinterpret_cast<const uint8_t*>(src_data.data + i * 16));
222
12
            const auto text = UUIDValue::to_string(value);
223
12
            strings.insert_data(text.data(), text.size());
224
12
        }
225
4
        return Status::OK();
226
4
    }
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
321
        bool preserve_binary_uuid) {
235
321
    std::unique_ptr<ConvertParams> convert_params = std::make_unique<ConvertParams>();
236
321
    const tparquet::SchemaElement& parquet_schema = field_schema->parquet_schema;
237
321
    convert_params->init(field_schema, ctz);
238
321
    tparquet::Type::type src_physical_type = parquet_schema.type;
239
321
    std::unique_ptr<PhysicalToLogicalConverter> physical_converter;
240
321
    if (is_dict_filter) {
241
1
        src_physical_type = tparquet::Type::INT32;
242
1
        src_logical_type = DataTypeFactory::instance().create_data_type(
243
1
                PrimitiveType::TYPE_INT, dst_logical_type->is_nullable());
244
1
    }
245
321
    PrimitiveType src_logical_primitive = src_logical_type->get_primitive_type();
246
247
321
    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
321
    } else if (is_parquet_native_type(src_logical_primitive)) {
261
271
        bool is_string_logical_type = is_string_type(src_logical_primitive);
262
271
        if (is_string_logical_type && src_physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY) {
263
7
            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
4
                if (parquet_schema.type_length != UUIDValue::BINARY_LENGTH) {
267
1
                    physical_converter = std::make_unique<UnsupportedConverter>(src_physical_type,
268
1
                                                                                src_logical_type);
269
3
                } else {
270
3
                    physical_converter = std::make_unique<UUIDStringConverter>();
271
3
                }
272
4
            } else {
273
3
                physical_converter =
274
3
                        std::make_unique<FixedSizeBinaryConverter>(parquet_schema.type_length);
275
3
            }
276
264
        } else if (src_logical_primitive == TYPE_FLOAT &&
277
264
                   src_physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY &&
278
264
                   parquet_schema.logicalType.__isset.FLOAT16) {
279
0
            physical_converter =
280
0
                    std::make_unique<Float16PhysicalConverter>(parquet_schema.type_length);
281
264
        } else {
282
264
            physical_converter = std::make_unique<ConsistentPhysicalConverter>();
283
264
        }
284
271
    } else if (src_logical_primitive == TYPE_TINYINT) {
285
10
        physical_converter = std::make_unique<LittleIntPhysicalConverter<TYPE_TINYINT>>();
286
40
    } else if (src_logical_primitive == TYPE_SMALLINT) {
287
7
        physical_converter = std::make_unique<LittleIntPhysicalConverter<TYPE_SMALLINT>>();
288
33
    } else if (is_decimal_type(src_logical_primitive)) {
289
4
        get_decimal_converter(field_schema, src_logical_type, dst_logical_type,
290
4
                              convert_params.get(), physical_converter);
291
29
    } else if (src_logical_primitive == TYPE_DATEV2) {
292
7
        physical_converter = std::make_unique<Int32ToDate>();
293
22
    } else if (src_logical_primitive == TYPE_DATETIMEV2) {
294
16
        if (src_physical_type == tparquet::Type::INT96) {
295
            // int96 only stores nanoseconds in standard parquet file
296
1
            convert_params->reset_time_scale_if_missing(9);
297
1
            physical_converter = std::make_unique<Int96toTimestamp>();
298
15
        } else if (src_physical_type == tparquet::Type::INT64) {
299
15
            convert_params->reset_time_scale_if_missing(src_logical_type->get_scale());
300
15
            physical_converter = std::make_unique<Int64ToTimestamp>();
301
15
        } else {
302
0
            physical_converter =
303
0
                    std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
304
0
        }
305
16
    } else if (src_logical_primitive == TYPE_VARBINARY) {
306
4
        if (src_physical_type == tparquet::Type::FIXED_LEN_BYTE_ARRAY) {
307
2
            DCHECK(parquet_schema.logicalType.__isset.UUID) << parquet_schema.name;
308
2
            physical_converter =
309
2
                    std::make_unique<UUIDVarBinaryConverter>(parquet_schema.type_length);
310
2
        } else {
311
2
            DCHECK(src_physical_type == tparquet::Type::BYTE_ARRAY) << src_physical_type;
312
2
            physical_converter = std::make_unique<ConsistentPhysicalConverter>();
313
2
        }
314
4
    } else if (src_logical_primitive == TYPE_TIMESTAMPTZ) {
315
2
        if (src_physical_type == tparquet::Type::INT96) {
316
1
            physical_converter = std::make_unique<Int96toTimestampTz>();
317
1
        } else if (src_physical_type == tparquet::Type::INT64) {
318
1
            DCHECK(src_physical_type == tparquet::Type::INT64) << src_physical_type;
319
1
            DCHECK(parquet_schema.logicalType.__isset.TIMESTAMP) << parquet_schema.name;
320
1
            physical_converter = std::make_unique<Int64ToTimestampTz>();
321
1
        } else {
322
0
            physical_converter =
323
0
                    std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
324
0
        }
325
2
    } else {
326
0
        physical_converter =
327
0
                std::make_unique<UnsupportedConverter>(src_physical_type, src_logical_type);
328
0
    }
329
330
321
    if (physical_converter->support()) {
331
320
        physical_converter->_convert_params = std::move(convert_params);
332
320
        physical_converter->_logical_converter = converter::ColumnTypeConverter::get_converter(
333
320
                src_logical_type, dst_logical_type, converter::FileFormat::PARQUET);
334
320
        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
320
    }
340
321
    return physical_converter;
341
321
}
342
343
} // namespace doris::parquet