Skip to content
This repository was archived by the owner on Jul 24, 2024. It is now read-only.

Commit a34cc3d

Browse files
authored
lightning: fix parquet parser for decimal type (#1272) (#1275)
1 parent 1b7c88b commit a34cc3d

3 files changed

Lines changed: 276 additions & 60 deletions

File tree

2.62 KB
Binary file not shown.

pkg/lightning/mydump/parquet_parser.go

Lines changed: 200 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"context"
66
"fmt"
77
"io"
8+
"math/big"
89
"reflect"
910
"strings"
1011
"time"
@@ -191,7 +192,16 @@ func NewParquetParser(
191192
// NOTE: the SchemaElement.Name is capitalized, SchemaHandler.Infos.ExName is the raw column name
192193
// though in this context, there is no difference between these two fields
193194
columns = append(columns, strings.ToLower(c.Name))
194-
columnMetas = append(columnMetas, c)
195+
// transfer old ConvertedType to LogicalType
196+
columnMeta := c
197+
if c.ConvertedType != nil && c.LogicalType == nil {
198+
newMeta := *c
199+
columnMeta = &newMeta
200+
if err := convertToLogicType(columnMeta); err != nil {
201+
return nil, err
202+
}
203+
}
204+
columnMetas = append(columnMetas, columnMeta)
195205
}
196206
}
197207

@@ -203,6 +213,100 @@ func NewParquetParser(
203213
}, nil
204214
}
205215

216+
func convertToLogicType(se *parquet.SchemaElement) error {
217+
logicalType := &parquet.LogicalType{}
218+
switch *se.ConvertedType {
219+
case parquet.ConvertedType_UTF8:
220+
logicalType.STRING = &parquet.StringType{}
221+
case parquet.ConvertedType_ENUM:
222+
logicalType.ENUM = &parquet.EnumType{}
223+
case parquet.ConvertedType_DECIMAL:
224+
logicalType.DECIMAL = &parquet.DecimalType{
225+
Scale: *se.Scale,
226+
Precision: *se.Precision,
227+
}
228+
case parquet.ConvertedType_DATE:
229+
logicalType.DATE = &parquet.DateType{}
230+
case parquet.ConvertedType_TIME_MILLIS:
231+
logicalType.TIME = &parquet.TimeType{
232+
IsAdjustedToUTC: true,
233+
Unit: &parquet.TimeUnit{
234+
MILLIS: parquet.NewMilliSeconds(),
235+
},
236+
}
237+
case parquet.ConvertedType_TIME_MICROS:
238+
logicalType.TIME = &parquet.TimeType{
239+
IsAdjustedToUTC: true,
240+
Unit: &parquet.TimeUnit{
241+
MICROS: parquet.NewMicroSeconds(),
242+
},
243+
}
244+
case parquet.ConvertedType_TIMESTAMP_MILLIS:
245+
logicalType.TIMESTAMP = &parquet.TimestampType{
246+
IsAdjustedToUTC: true,
247+
Unit: &parquet.TimeUnit{
248+
MILLIS: parquet.NewMilliSeconds(),
249+
},
250+
}
251+
case parquet.ConvertedType_TIMESTAMP_MICROS:
252+
logicalType.TIMESTAMP = &parquet.TimestampType{
253+
IsAdjustedToUTC: true,
254+
Unit: &parquet.TimeUnit{
255+
MICROS: parquet.NewMicroSeconds(),
256+
},
257+
}
258+
case parquet.ConvertedType_UINT_8:
259+
logicalType.INTEGER = &parquet.IntType{
260+
BitWidth: 8,
261+
IsSigned: false,
262+
}
263+
case parquet.ConvertedType_UINT_16:
264+
logicalType.INTEGER = &parquet.IntType{
265+
BitWidth: 16,
266+
IsSigned: false,
267+
}
268+
case parquet.ConvertedType_UINT_32:
269+
logicalType.INTEGER = &parquet.IntType{
270+
BitWidth: 32,
271+
IsSigned: false,
272+
}
273+
case parquet.ConvertedType_UINT_64:
274+
logicalType.INTEGER = &parquet.IntType{
275+
BitWidth: 64,
276+
IsSigned: false,
277+
}
278+
case parquet.ConvertedType_INT_8:
279+
logicalType.INTEGER = &parquet.IntType{
280+
BitWidth: 8,
281+
IsSigned: true,
282+
}
283+
case parquet.ConvertedType_INT_16:
284+
logicalType.INTEGER = &parquet.IntType{
285+
BitWidth: 16,
286+
IsSigned: true,
287+
}
288+
case parquet.ConvertedType_INT_32:
289+
logicalType.INTEGER = &parquet.IntType{
290+
BitWidth: 32,
291+
IsSigned: true,
292+
}
293+
case parquet.ConvertedType_INT_64:
294+
logicalType.INTEGER = &parquet.IntType{
295+
BitWidth: 64,
296+
IsSigned: true,
297+
}
298+
case parquet.ConvertedType_JSON:
299+
logicalType.JSON = &parquet.JsonType{}
300+
case parquet.ConvertedType_BSON:
301+
logicalType.BSON = &parquet.BsonType{}
302+
// case parquet.ConvertedType_INTERVAL, parquet.ConvertedType_MAP, parquet.ConvertedType_MAP_KEY_VALUE, parquet.ConvertedType_LIST:
303+
default:
304+
return errors.Errorf("unsupported type: '%s'", *se.ConvertedType)
305+
}
306+
se.LogicalType = logicalType
307+
return nil
308+
}
309+
206310
// Pos returns the currently row number of the parquet file
207311
func (pp *ParquetParser) Pos() (pos int64, rowID int64) {
208312
return pp.curStart + int64(pp.curIndex), pp.lastRow.RowID
@@ -272,79 +376,141 @@ func (pp *ParquetParser) ReadRow() error {
272376
pp.lastRow.Row = pp.lastRow.Row[:length]
273377
}
274378
for i := 0; i < length; i++ {
275-
setDatumValue(&pp.lastRow.Row[i], v.Field(i), pp.columnMetas[i])
379+
if err := setDatumValue(&pp.lastRow.Row[i], v.Field(i), pp.columnMetas[i]); err != nil {
380+
return err
381+
}
276382
}
277383
return nil
278384
}
279385

280386
// convert a parquet value to Datum
281387
//
282388
// See: https://github.com/apache/parquet-format/blob/master/LogicalTypes.md
283-
func setDatumValue(d *types.Datum, v reflect.Value, meta *parquet.SchemaElement) {
389+
func setDatumValue(d *types.Datum, v reflect.Value, meta *parquet.SchemaElement) error {
284390
switch v.Kind() {
285391
case reflect.Uint8, reflect.Uint16, reflect.Uint32, reflect.Uint64:
286392
d.SetUint64(v.Uint())
287393
case reflect.Int8, reflect.Int16:
288394
d.SetInt64(v.Int())
289395
case reflect.Int32, reflect.Int64:
290-
setDatumByInt(d, v.Int(), meta)
396+
return setDatumByInt(d, v.Int(), meta)
291397
case reflect.String:
292-
d.SetString(v.String(), "")
398+
setDatumByString(d, v.String(), meta)
293399
case reflect.Float32, reflect.Float64:
294400
d.SetFloat64(v.Float())
295401
case reflect.Ptr:
296402
if v.IsNil() {
297403
d.SetNull()
298404
} else {
299-
setDatumValue(d, v.Elem(), meta)
405+
return setDatumValue(d, v.Elem(), meta)
300406
}
301407
default:
302-
log.L().Fatal("unknown value", zap.Stringer("kind", v.Kind()),
408+
log.L().Error("unknown value", zap.Stringer("kind", v.Kind()),
303409
zap.String("type", v.Type().Name()), zap.Reflect("value", v.Interface()))
410+
return errors.Errorf("unknown value: %v", v)
411+
}
412+
return nil
413+
}
414+
415+
func setDatumByString(d *types.Datum, v string, meta *parquet.SchemaElement) {
416+
if meta.LogicalType != nil && meta.LogicalType.DECIMAL != nil {
417+
v = binaryToDecimalStr([]byte(v), int(meta.LogicalType.DECIMAL.Scale))
304418
}
419+
d.SetString(v, "")
420+
}
421+
422+
func binaryToDecimalStr(rawBytes []byte, scale int) string {
423+
negative := rawBytes[0] > 127
424+
if negative {
425+
for i := 0; i < len(rawBytes); i++ {
426+
rawBytes[i] = ^rawBytes[i]
427+
}
428+
for i := len(rawBytes) - 1; i >= 0; i-- {
429+
rawBytes[i] += 1
430+
if rawBytes[i] != 0 {
431+
break
432+
}
433+
}
434+
}
435+
436+
intValue := big.NewInt(0)
437+
intValue = intValue.SetBytes(rawBytes)
438+
val := fmt.Sprintf("%0*d", scale, intValue)
439+
dotIndex := len(val) - scale
440+
var res strings.Builder
441+
if negative {
442+
res.WriteByte('-')
443+
}
444+
if dotIndex == 0 {
445+
res.WriteByte('0')
446+
} else {
447+
res.WriteString(val[:dotIndex])
448+
}
449+
if scale > 0 {
450+
res.WriteByte('.')
451+
res.WriteString(val[dotIndex:])
452+
}
453+
return res.String()
305454
}
306455

307456
// when the value type is int32/int64, convert to value to target logical type in tidb
308-
func setDatumByInt(d *types.Datum, v int64, meta *parquet.SchemaElement) {
309-
if meta.ConvertedType == nil {
457+
func setDatumByInt(d *types.Datum, v int64, meta *parquet.SchemaElement) error {
458+
if meta.ConvertedType == nil && meta.LogicalType == nil {
310459
d.SetInt64(v)
311-
return
460+
return nil
312461
}
313-
switch *meta.ConvertedType {
314-
// decimal
315-
case parquet.ConvertedType_DECIMAL:
316-
minLen := *meta.Scale + 1
462+
463+
logicalType := meta.LogicalType
464+
switch {
465+
case logicalType.DECIMAL != nil:
466+
if logicalType.DECIMAL.Scale == 0 {
467+
d.SetInt64(v)
468+
return nil
469+
}
470+
minLen := logicalType.DECIMAL.Scale + 1
317471
if v < 0 {
318472
minLen++
319473
}
320474
val := fmt.Sprintf("%0*d", minLen, v)
321475
dotIndex := len(val) - int(*meta.Scale)
322476
d.SetString(val[:dotIndex]+"."+val[dotIndex:], "")
323-
case parquet.ConvertedType_DATE:
477+
case logicalType.DATE != nil:
324478
dateStr := time.Unix(v*86400, 0).Format("2006-01-02")
325479
d.SetString(dateStr, "")
326-
// convert all timestamp types (datetime/timestamp) to string
327-
case parquet.ConvertedType_TIMESTAMP_MICROS:
328-
dateStr := time.Unix(v/1e6, (v%1e6)*1e3).Format("2006-01-02 15:04:05.999")
329-
d.SetString(dateStr, "")
330-
case parquet.ConvertedType_TIMESTAMP_MILLIS:
331-
dateStr := time.Unix(v/1e3, (v%1e3)*1e6).Format("2006-01-02 15:04:05.999")
332-
d.SetString(dateStr, "")
333-
// covert time types to string
334-
case parquet.ConvertedType_TIME_MILLIS, parquet.ConvertedType_TIME_MICROS:
335-
if *meta.ConvertedType == parquet.ConvertedType_TIME_MICROS {
336-
v /= 1e3
337-
}
338-
millis := v % 1e3
339-
v /= 1e3
340-
sec := v % 60
341-
v /= 60
342-
min := v % 60
343-
v /= 60
344-
d.SetString(fmt.Sprintf("%d:%d:%d.%3d", v, min, sec, millis), "")
480+
case logicalType.TIMESTAMP != nil:
481+
// convert all timestamp types (datetime/timestamp) to string
482+
timeStr := formatTime(v, logicalType.TIMESTAMP.Unit, "2006-01-02 15:04:05.999999",
483+
"2006-01-02 15:04:05.999999Z", logicalType.TIMESTAMP.IsAdjustedToUTC)
484+
d.SetString(timeStr, "")
485+
case logicalType.TIME != nil:
486+
// convert all timestamp types (datetime/timestamp) to string
487+
timeStr := formatTime(v, logicalType.TIME.Unit, "15:04:05.999999", "15:04:05.999999Z",
488+
logicalType.TIME.IsAdjustedToUTC)
489+
d.SetString(timeStr, "")
345490
default:
346491
d.SetInt64(v)
347492
}
493+
return nil
494+
}
495+
496+
func formatTime(v int64, units *parquet.TimeUnit, format, utcFormat string, utc bool) string {
497+
var sec, nsec int64
498+
if units.MICROS != nil {
499+
sec = v / 1e6
500+
nsec = (v % 1e6) * 1e3
501+
} else if units.MILLIS != nil {
502+
sec = v / 1e3
503+
nsec = (v % 1e3) * 1e6
504+
} else {
505+
// nano
506+
sec = v / 1e9
507+
nsec = v % 1e9
508+
}
509+
t := time.Unix(sec, nsec).UTC()
510+
if utc {
511+
return t.Format(utcFormat)
512+
}
513+
return t.Format(format)
348514
}
349515

350516
func (pp *ParquetParser) LastRow() Row {

0 commit comments

Comments
 (0)