Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .circleci/config.yml
Original file line number Diff line number Diff line change
Expand Up @@ -9,13 +9,13 @@ orbs:
jobs:
build:
docker:
- image: circleci/golang:1.17.2
- image: cimg/go:1.18
environment:
PARQUET_COMPATIBILITY_REPO_ROOT: /tmp/parquet-compatibility
PARQUET_TESTING_ROOT: /tmp/parquet-testing
steps:
- checkout
- run: curl -sSfL https://raw.githubusercontent.com/golangci/golangci-lint/master/install.sh | sh -s -- -b $(go env GOPATH)/bin v1.25.0
- run: curl -sSfL https://raw.githubusercontent.com/golangci/golangci-lint/master/install.sh | sh -s -- -b $(go env GOPATH)/bin v1.45.0
- run: golangci-lint run
- run: git clone https://github.com/Parquet/parquet-compatibility.git ${PARQUET_COMPATIBILITY_REPO_ROOT}
- run: git clone https://github.com/apache/parquet-testing.git ${PARQUET_TESTING_ROOT}
Expand Down
34 changes: 15 additions & 19 deletions .golangci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -25,25 +25,21 @@ linters-settings:
- wrapperFunc

linters:
enable-all: true
disable:
- lll
- scopelint
- gochecknoglobals
- goconst
- gocyclo
- funlen
- godox
- wsl
- unparam
- dupl
- gocritic
- gocognit
- gochecknoinits
- testpackage
- nestif
- gomnd
- godot
disable-all: true
enable:
- deadcode
- errcheck
# - gosimple
- govet
- ineffassign
# - structcheck
- typecheck
# - unused
- varcheck
- gofmt
- gosec
# - revive

run:
# timeout for analysis, e.g. 30s, 5m, default is 1m
deadline: 5m
Expand Down
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

- Switched number encoding and decoding (int32, int64, float32, float64) to Go 1.18 generics.

## [v0.10.0] - 2022-02-18

- Updated to parquet-format 2.9.0.
Expand Down
36 changes: 18 additions & 18 deletions chunk_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,13 +24,13 @@ func getDictValuesDecoder(typ *parquet.SchemaElement) (valuesDecoder, error) {
}
return &byteArrayPlainDecoder{length: int(*typ.TypeLength)}, nil
case parquet.Type_FLOAT:
return &floatPlainDecoder{}, nil
return &numberPlainDecoder[float32, internalFloat32]{}, nil
case parquet.Type_DOUBLE:
return &doublePlainDecoder{}, nil
return &numberPlainDecoder[float64, internalFloat64]{}, nil
case parquet.Type_INT32:
return &int32PlainDecoder{}, nil
return &numberPlainDecoder[int32, internalInt32]{}, nil
case parquet.Type_INT64:
return &int64PlainDecoder{}, nil
return &numberPlainDecoder[int64, internalInt64]{}, nil
case parquet.Type_INT96:
return &int96PlainDecoder{}, nil
}
Expand All @@ -49,7 +49,7 @@ func getBooleanValuesDecoder(pageEncoding parquet.Encoding) (valuesDecoder, erro
}
}

func getByteArrayValuesDecoder(pageEncoding parquet.Encoding, dictValues []interface{}) (valuesDecoder, error) {
func getByteArrayValuesDecoder(pageEncoding parquet.Encoding, dictValues []any) (valuesDecoder, error) {
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &byteArrayPlainDecoder{}, nil
Expand All @@ -64,7 +64,7 @@ func getByteArrayValuesDecoder(pageEncoding parquet.Encoding, dictValues []inter
}
}

func getFixedLenByteArrayValuesDecoder(pageEncoding parquet.Encoding, len int, dictValues []interface{}) (valuesDecoder, error) {
func getFixedLenByteArrayValuesDecoder(pageEncoding parquet.Encoding, len int, dictValues []any) (valuesDecoder, error) {
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &byteArrayPlainDecoder{length: len}, nil
Expand All @@ -77,33 +77,33 @@ func getFixedLenByteArrayValuesDecoder(pageEncoding parquet.Encoding, len int, d
}
}

func getInt32ValuesDecoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []interface{}) (valuesDecoder, error) {
func getInt32ValuesDecoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []any) (valuesDecoder, error) {
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &int32PlainDecoder{}, nil
return &numberPlainDecoder[int32, internalInt32]{}, nil
case parquet.Encoding_DELTA_BINARY_PACKED:
return &int32DeltaBPDecoder{}, nil
return &deltaBitPackDecoder[int32, internalInt32]{}, nil
case parquet.Encoding_RLE_DICTIONARY:
return &dictDecoder{uniqueValues: dictValues}, nil
default:
return nil, fmt.Errorf("unsupported encoding %s for int32", pageEncoding)
}
}

func getInt64ValuesDecoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []interface{}) (valuesDecoder, error) {
func getInt64ValuesDecoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []any) (valuesDecoder, error) {
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &int64PlainDecoder{}, nil
return &numberPlainDecoder[int64, internalInt64]{}, nil
case parquet.Encoding_DELTA_BINARY_PACKED:
return &int64DeltaBPDecoder{}, nil
return &deltaBitPackDecoder[int64, internalInt64]{}, nil
case parquet.Encoding_RLE_DICTIONARY:
return &dictDecoder{uniqueValues: dictValues}, nil
default:
return nil, fmt.Errorf("unsupported encoding %s for int64", pageEncoding)
}
}

func getValuesDecoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []interface{}) (valuesDecoder, error) {
func getValuesDecoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []any) (valuesDecoder, error) {
// Change the deprecated value
if pageEncoding == parquet.Encoding_PLAIN_DICTIONARY {
pageEncoding = parquet.Encoding_RLE_DICTIONARY
Expand All @@ -124,15 +124,15 @@ func getValuesDecoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement,
case parquet.Type_FLOAT:
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &floatPlainDecoder{}, nil
return &numberPlainDecoder[float32, internalFloat32]{}, nil
case parquet.Encoding_RLE_DICTIONARY:
return &dictDecoder{uniqueValues: dictValues}, nil
}

case parquet.Type_DOUBLE:
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &doublePlainDecoder{}, nil
return &numberPlainDecoder[float64, internalFloat64]{}, nil
case parquet.Encoding_RLE_DICTIONARY:
return &dictDecoder{uniqueValues: dictValues}, nil
}
Expand Down Expand Up @@ -235,7 +235,7 @@ func readPages(ctx context.Context, sch *schema, r *offsetReader, col *Column, c
default:
return nil, false, fmt.Errorf("DATA_PAGE or DATA_PAGE_V2 type supported, but was %s", ph.Type)
}
var dictValue []interface{}
var dictValue []any
if dictPage != nil {
dictValue = dictPage.values
}
Expand All @@ -255,8 +255,8 @@ func readPages(ctx context.Context, sch *schema, r *offsetReader, col *Column, c
return pages, dictPage != nil, nil
}

func clone(in []interface{}) []interface{} {
out := make([]interface{}, len(in))
func clone(in []any) []any {
out := make([]any, len(in))
copy(out, in)
return out
}
Expand Down
61 changes: 20 additions & 41 deletions chunk_writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ func getBooleanValuesEncoder(pageEncoding parquet.Encoding) (valuesEncoder, erro
}
}

func getByteArrayValuesEncoder(pageEncoding parquet.Encoding, dictValues []interface{}) (valuesEncoder, error) {
func getByteArrayValuesEncoder(pageEncoding parquet.Encoding, dictValues []any) (valuesEncoder, error) {
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &byteArrayPlainEncoder{}, nil
Expand All @@ -36,7 +36,7 @@ func getByteArrayValuesEncoder(pageEncoding parquet.Encoding, dictValues []inter
}
}

func getFixedLenByteArrayValuesEncoder(pageEncoding parquet.Encoding, len int, dictValues []interface{}) (valuesEncoder, error) {
func getFixedLenByteArrayValuesEncoder(pageEncoding parquet.Encoding, len int, dictValues []any) (valuesEncoder, error) {
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &byteArrayPlainEncoder{length: len}, nil
Expand All @@ -49,45 +49,24 @@ func getFixedLenByteArrayValuesEncoder(pageEncoding parquet.Encoding, len int, d
}
}

func getInt32ValuesEncoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []interface{}) (valuesEncoder, error) {
func getIntValuesEncoder[T intType, I internalIntType[T]](pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []any) (valuesEncoder, error) {
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &int32PlainEncoder{}, nil
return &numberPlainEncoder[T, I]{}, nil
case parquet.Encoding_DELTA_BINARY_PACKED:
return &int32DeltaBPEncoder{
deltaBitPackEncoder32: deltaBitPackEncoder32{
blockSize: 128,
miniBlockCount: 4,
},
return &deltaBitPackEncoder[T, I]{
blockSize: 128,
miniBlockCount: 4,
}, nil
case parquet.Encoding_RLE_DICTIONARY:
return &dictEncoder{dictValues: dictValues}, nil
default:
return nil, fmt.Errorf("unsupported encoding %s for int32", pageEncoding)
var t T
return nil, fmt.Errorf("unsupported encoding %s for %T", pageEncoding, t)
}
}

func getInt64ValuesEncoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []interface{}) (valuesEncoder, error) {
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &int64PlainEncoder{}, nil
case parquet.Encoding_DELTA_BINARY_PACKED:
return &int64DeltaBPEncoder{
deltaBitPackEncoder64: deltaBitPackEncoder64{
blockSize: 128,
miniBlockCount: 4,
},
}, nil
case parquet.Encoding_RLE_DICTIONARY:
return &dictEncoder{
dictValues: dictValues,
}, nil
default:
return nil, fmt.Errorf("unsupported encoding %s for int64", pageEncoding)
}
}

func getValuesEncoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []interface{}) (valuesEncoder, error) {
func getValuesEncoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement, dictValues []any) (valuesEncoder, error) {
// Change the deprecated value
if pageEncoding == parquet.Encoding_PLAIN_DICTIONARY {
pageEncoding = parquet.Encoding_RLE_DICTIONARY
Expand All @@ -109,7 +88,7 @@ func getValuesEncoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement,
case parquet.Type_FLOAT:
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &floatPlainEncoder{}, nil
return &numberPlainEncoder[float32, internalFloat32]{}, nil
case parquet.Encoding_RLE_DICTIONARY:
return &dictEncoder{
dictValues: dictValues,
Expand All @@ -119,18 +98,18 @@ func getValuesEncoder(pageEncoding parquet.Encoding, typ *parquet.SchemaElement,
case parquet.Type_DOUBLE:
switch pageEncoding {
case parquet.Encoding_PLAIN:
return &doublePlainEncoder{}, nil
return &numberPlainEncoder[float64, internalFloat64]{}, nil
case parquet.Encoding_RLE_DICTIONARY:
return &dictEncoder{
dictValues: dictValues,
}, nil
}

case parquet.Type_INT32:
return getInt32ValuesEncoder(pageEncoding, typ, dictValues)
return getIntValuesEncoder[int32, internalInt32](pageEncoding, typ, dictValues)

case parquet.Type_INT64:
return getInt64ValuesEncoder(pageEncoding, typ, dictValues)
return getIntValuesEncoder[int64, internalInt64](pageEncoding, typ, dictValues)

case parquet.Type_INT96:
switch pageEncoding {
Expand Down Expand Up @@ -159,13 +138,13 @@ func getDictValuesEncoder(typ *parquet.SchemaElement) (valuesEncoder, error) {
}
return &byteArrayPlainEncoder{length: int(*typ.TypeLength)}, nil
case parquet.Type_FLOAT:
return &floatPlainEncoder{}, nil
return &numberPlainEncoder[float32, internalFloat32]{}, nil
case parquet.Type_DOUBLE:
return &doublePlainEncoder{}, nil
return &numberPlainEncoder[float64, internalFloat64]{}, nil
case parquet.Type_INT32:
return &int32PlainEncoder{}, nil
return &numberPlainEncoder[int32, internalInt32]{}, nil
case parquet.Type_INT64:
return &int64PlainEncoder{}, nil
return &numberPlainEncoder[int64, internalInt64]{}, nil
case parquet.Type_INT96:
return &int96PlainEncoder{}, nil
}
Expand Down Expand Up @@ -193,8 +172,8 @@ func writeChunk(ctx context.Context, w writePos, sch *schema, col *Column, codec
return nil, err
}

dictValues := []interface{}{}
indices := map[interface{}]int32{}
dictValues := []any{}
indices := map[any]int32{}

for _, page := range col.data.dataPages {
for _, v := range page.values {
Expand Down
Loading