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 |