S3 select switch to new parquet library and reduce locking (#14731)

- This change switches to a new parquet library
- SelectObjectContent now takes a single lock at the beginning and holds it
during the operation. Previously the operation took a lock every time the
parquet library performed a Seek on the underlying object stream.
- Add basic support for LogicalType annotations for timestamps.
This commit is contained in:
Aditya Manthramurthy
2022-04-14 06:54:47 -07:00
committed by GitHub
parent 67e17ed3f8
commit e8e48e4c4a
7 changed files with 336 additions and 348 deletions
+88 -83
View File
@@ -18,88 +18,55 @@
package parquet
import (
"fmt"
"errors"
"io"
"time"
"github.com/bcicen/jstream"
parquetgo "github.com/fraugster/parquet-go"
parquettypes "github.com/fraugster/parquet-go/parquet"
jsonfmt "github.com/minio/minio/internal/s3select/json"
"github.com/minio/minio/internal/s3select/sql"
parquetgo "github.com/minio/parquet-go"
parquetgen "github.com/minio/parquet-go/gen-go/parquet"
)
// Reader - Parquet record reader for S3Select.
// Reader implements reading records from parquet input.
type Reader struct {
args *ReaderArgs
reader *parquetgo.Reader
io.Closer
r *parquetgo.FileReader
}
// Read - reads single record.
func (r *Reader) Read(dst sql.Record) (rec sql.Record, rerr error) {
defer func() {
if rec := recover(); rec != nil {
rerr = fmt.Errorf("panic reading parquet record: %v", rec)
}
}()
parquetRecord, err := r.reader.Read()
// NewParquetReader creates a Reader2 from a io.ReadSeekCloser.
func NewParquetReader(rsc io.ReadSeekCloser, _ *ReaderArgs) (r *Reader, err error) {
fr, err := parquetgo.NewFileReader(rsc)
if err != nil {
if err != io.EOF {
return nil, errParquetParsingError(err)
}
return nil, errParquetParsingError(err)
}
return nil, err
return &Reader{Closer: rsc, r: fr}, nil
}
func (pr *Reader) Read(dst sql.Record) (rec sql.Record, rerr error) {
nextRow, err := pr.r.NextRow()
if err != nil {
if err == io.EOF {
return nil, err
}
return nil, errParquetParsingError(err)
}
kvs := jstream.KVS{}
f := func(name string, v parquetgo.Value) bool {
if v.Value == nil {
kvs = append(kvs, jstream.KV{Key: name, Value: nil})
return true
}
for _, col := range pr.r.Columns() {
var value interface{}
switch v.Type {
case parquetgen.Type_BOOLEAN:
value = v.Value.(bool)
case parquetgen.Type_INT32:
value = int64(v.Value.(int32))
if v.Schema != nil && v.Schema.ConvertedType != nil {
switch *v.Schema.ConvertedType {
case parquetgen.ConvertedType_DATE:
value = sql.FormatSQLTimestamp(time.Unix(60*60*24*int64(v.Value.(int32)), 0).UTC())
}
if v, ok := nextRow[col.FlatName()]; ok {
value, err = convertFromAnnotation(col.Element(), v)
if err != nil {
return nil, errParquetParsingError(err)
}
case parquetgen.Type_INT64:
value = v.Value.(int64)
if v.Schema != nil && v.Schema.ConvertedType != nil {
switch *v.Schema.ConvertedType {
// Only UTC supported, add one NS to never be exactly midnight.
case parquetgen.ConvertedType_TIMESTAMP_MILLIS:
value = sql.FormatSQLTimestamp(time.Unix(0, 0).Add(time.Duration(v.Value.(int64)) * time.Millisecond).UTC())
case parquetgen.ConvertedType_TIMESTAMP_MICROS:
value = sql.FormatSQLTimestamp(time.Unix(0, 0).Add(time.Duration(v.Value.(int64)) * time.Microsecond).UTC())
}
}
case parquetgen.Type_FLOAT:
value = float64(v.Value.(float32))
case parquetgen.Type_DOUBLE:
value = v.Value.(float64)
case parquetgen.Type_INT96, parquetgen.Type_BYTE_ARRAY, parquetgen.Type_FIXED_LEN_BYTE_ARRAY:
value = string(v.Value.([]byte))
default:
rerr = errParquetParsingError(nil)
return false
}
kvs = append(kvs, jstream.KV{Key: name, Value: value})
return true
kvs = append(kvs, jstream.KV{Key: col.FlatName(), Value: value})
}
// Apply our range
parquetRecord.Range(f)
// Reuse destination if we can.
dstRec, ok := dst.(*jsonfmt.Record)
if !ok {
@@ -110,29 +77,67 @@ func (r *Reader) Read(dst sql.Record) (rec sql.Record, rerr error) {
return dstRec, nil
}
// Close - closes underlying readers.
func (r *Reader) Close() error {
return r.reader.Close()
}
// NewReader - creates new Parquet reader using readerFunc callback.
func NewReader(getReaderFunc func(offset, length int64) (io.ReadCloser, error), args *ReaderArgs) (r *Reader, err error) {
defer func() {
if rec := recover(); rec != nil {
err = fmt.Errorf("panic reading parquet header: %v", rec)
}
}()
reader, err := parquetgo.NewReader(getReaderFunc, nil)
if err != nil {
if err != io.EOF {
return nil, errParquetParsingError(err)
}
return nil, err
// convertFromAnnotation - converts values based on the Parquet column's type
// annotations. LogicalType annotations if present override the deprecated
// ConvertedType annotations. Ref:
// https://github.com/apache/parquet-format/blob/master/LogicalTypes.md
func convertFromAnnotation(se *parquettypes.SchemaElement, v interface{}) (interface{}, error) {
if se == nil {
return v, nil
}
return &Reader{
args: args,
reader: reader,
}, nil
var value interface{}
switch val := v.(type) {
case []byte:
// TODO: only strings are supported in s3select output (not
// binary arrays) - perhaps we need to check the annotation to
// ensure it's UTF8 encoded.
value = string(val)
case [12]byte:
// TODO: This is returned for the parquet INT96 type. We just
// treat it same as []byte (but AWS S3 treats it as a large int)
// - fix this later.
value = string(val[:])
case int32:
value = int64(val)
if logicalType := se.GetLogicalType(); logicalType != nil {
if logicalType.IsSetDATE() {
value = sql.FormatSQLTimestamp(time.Unix(60*60*24*int64(val), 0).UTC())
}
} else if se.GetConvertedType() == parquettypes.ConvertedType_DATE {
value = sql.FormatSQLTimestamp(time.Unix(60*60*24*int64(val), 0).UTC())
}
case int64:
value = val
if logicalType := se.GetLogicalType(); logicalType != nil {
if ts := logicalType.GetTIMESTAMP(); ts != nil {
var duration time.Duration
// Only support UTC normalized timestamps.
if ts.IsAdjustedToUTC {
switch {
case ts.Unit.IsSetNANOS():
duration = time.Duration(val) * time.Nanosecond
case ts.Unit.IsSetMILLIS():
duration = time.Duration(val) * time.Millisecond
case ts.Unit.IsSetMICROS():
duration = time.Duration(val) * time.Microsecond
default:
return nil, errors.New("Invalid LogicalType annotation found")
}
value = sql.FormatSQLTimestamp(time.Unix(0, 0).Add(duration))
}
} else if se.GetConvertedType() == parquettypes.ConvertedType_TIMESTAMP_MILLIS {
duration := time.Duration(val) * time.Millisecond
value = sql.FormatSQLTimestamp(time.Unix(0, 0).Add(duration))
} else if se.GetConvertedType() == parquettypes.ConvertedType_TIMESTAMP_MICROS {
duration := time.Duration(val) * time.Microsecond
value = sql.FormatSQLTimestamp(time.Unix(0, 0).Add(duration))
}
}
case float32:
value = float64(val)
default:
value = v
}
return value, nil
}
+113 -11
View File
@@ -214,10 +214,10 @@ type RequestProgress struct {
// ScanRange represents the ScanRange parameter.
type ScanRange struct {
// Start if byte offset form the start off the file.
// Start is the byte offset to read from (from the start of the file).
Start *uint64 `xml:"Start"`
// End is the last byte that should be returned, if Start is set,
// or the offset from EOF to start reading if start is not present.
// End is the offset of the last byte that should be returned when Start
// is set, otherwise it is the offset from EOF to start reading.
End *uint64 `xml:"End"`
}
@@ -362,21 +362,29 @@ func (s3Select *S3Select) getProgress() (bytesScanned, bytesProcessed int64) {
// Open - opens S3 object by using callback for SQL selection query.
// Currently CSV, JSON and Apache Parquet formats are supported.
func (s3Select *S3Select) Open(getReader func(offset, length int64) (io.ReadCloser, error)) error {
offset, end, err := s3Select.ScanRange.StartLen()
func (s3Select *S3Select) Open(rsc io.ReadSeekCloser) error {
offset, length, err := s3Select.ScanRange.StartLen()
if err != nil {
return err
}
seekDirection := io.SeekStart
if offset < 0 {
seekDirection = io.SeekEnd
}
switch s3Select.Input.format {
case csvFormat:
rc, err := getReader(offset, end)
_, err = rsc.Seek(offset, seekDirection)
if err != nil {
return err
}
var rc io.ReadCloser = rsc
if length != -1 {
rc = newLimitedReadCloser(rsc, length)
}
s3Select.progressReader, err = newProgressReader(rc, s3Select.Input.CompressionType)
if err != nil {
rc.Close()
rsc.Close()
return err
}
@@ -405,14 +413,18 @@ func (s3Select *S3Select) Open(getReader func(offset, length int64) (io.ReadClos
}
return nil
case jsonFormat:
rc, err := getReader(offset, end)
_, err = rsc.Seek(offset, seekDirection)
if err != nil {
return err
}
var rc io.ReadCloser = rsc
if length != -1 {
rc = newLimitedReadCloser(rsc, length)
}
s3Select.progressReader, err = newProgressReader(rc, s3Select.Input.CompressionType)
if err != nil {
rc.Close()
rsc.Close()
return err
}
@@ -431,12 +443,12 @@ func (s3Select *S3Select) Open(getReader func(offset, length int64) (io.ReadClos
if !strings.EqualFold(os.Getenv("MINIO_API_SELECT_PARQUET"), "on") {
return errors.New("parquet format parsing not enabled on server")
}
if offset != 0 || end != -1 {
if offset != 0 || length != -1 {
// Offsets do not make sense in parquet files.
return errors.New("parquet format does not support offsets")
}
var err error
s3Select.recordReader, err = parquet.NewReader(getReader, &s3Select.Input.ParquetArgs)
s3Select.recordReader, err = parquet.NewParquetReader(rsc, &s3Select.Input.ParquetArgs)
return err
}
@@ -657,3 +669,93 @@ func NewS3Select(r io.Reader) (*S3Select, error) {
return s3Select, nil
}
//////////////////
// Helpers
/////////////////
// limitedReadCloser is like io.LimitedReader, but also implements io.Closer.
type limitedReadCloser struct {
io.LimitedReader
io.Closer
}
func newLimitedReadCloser(r io.ReadCloser, n int64) *limitedReadCloser {
return &limitedReadCloser{
LimitedReader: io.LimitedReader{R: r, N: n},
Closer: r,
}
}
// ObjectSegmentReaderFn is a function that returns a reader for a contiguous
// suffix segment of an object starting at the given (non-negative) offset.
type ObjectSegmentReaderFn func(offset int64) (io.ReadCloser, error)
// ObjectReadSeekCloser implements ReadSeekCloser interface for reading objects.
// It uses a function that returns a io.ReadCloser for the object.
type ObjectReadSeekCloser struct {
segmentReader ObjectSegmentReaderFn
size int64 // actual object size regardless of compression/encryption
offset int64
reader io.ReadCloser
}
// NewObjectReadSeekCloser creates a new ObjectReadSeekCloser.
func NewObjectReadSeekCloser(segmentReader ObjectSegmentReaderFn, actualSize int64) *ObjectReadSeekCloser {
return &ObjectReadSeekCloser{
segmentReader: segmentReader,
size: actualSize,
offset: 0,
reader: nil,
}
}
// Seek call to implement io.Seeker
func (rsc *ObjectReadSeekCloser) Seek(offset int64, whence int) (int64, error) {
// fmt.Printf("actual: %v offset: %v (%v) whence: %v\n", rsc.size, offset, rsc.offset, whence)
switch whence {
case io.SeekStart:
rsc.offset = offset
case io.SeekCurrent:
rsc.offset += offset
case io.SeekEnd:
rsc.offset = rsc.size + offset
}
if rsc.offset < 0 {
return rsc.offset, errors.New("seek to invalid negative offset")
}
if rsc.offset >= rsc.size {
return rsc.offset, errors.New("seek past end of object")
}
if rsc.reader != nil {
_ = rsc.reader.Close()
rsc.reader = nil
}
return rsc.offset, nil
}
// Read call to implement io.Reader
func (rsc *ObjectReadSeekCloser) Read(p []byte) (n int, err error) {
if rsc.reader == nil {
rsc.reader, err = rsc.segmentReader(rsc.offset)
if err != nil {
return 0, err
}
}
return rsc.reader.Read(p)
}
// Close call to implement io.Closer. Calling Read/Seek after Close reopens the
// object for reading and a subsequent Close call is required to ensure
// resources are freed.
func (rsc *ObjectReadSeekCloser) Close() error {
if rsc.reader != nil {
err := rsc.reader.Close()
if err == nil {
rsc.reader = nil
}
return err
}
return nil
}
+1 -5
View File
@@ -20,8 +20,6 @@ package s3select
import (
"bytes"
"encoding/csv"
"io"
"io/ioutil"
"math/rand"
"net/http"
"strconv"
@@ -112,9 +110,7 @@ func benchmarkSelect(b *testing.B, count int, query string) {
b.Fatal(err)
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
return ioutil.NopCloser(bytes.NewReader(csvData)), nil
}); err != nil {
if err = s3Select.Open(newBytesRSC(csvData)); err != nil {
b.Fatal(err)
}
+49 -200
View File
@@ -20,7 +20,6 @@ package s3select
import (
"bytes"
"encoding/xml"
"errors"
"fmt"
"io"
"io/ioutil"
@@ -35,6 +34,22 @@ import (
"github.com/minio/simdjson-go"
)
func newStringRSC(s string) io.ReadSeekCloser {
return newBytesRSC([]byte(s))
}
func newBytesRSC(b []byte) io.ReadSeekCloser {
r := bytes.NewReader(b)
segmentReader := func(offset int64) (io.ReadCloser, error) {
_, err := r.Seek(offset, io.SeekStart)
if err != nil {
return nil, err
}
return io.NopCloser(r), nil
}
return NewObjectReadSeekCloser(segmentReader, int64(len(b)))
}
type testResponseWriter struct {
statusCode int
response []byte
@@ -608,13 +623,11 @@ func TestJSONQueries(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
in := input
if len(testCase.withJSON) > 0 {
in = testCase.withJSON
}
return ioutil.NopCloser(bytes.NewBufferString(in)), nil
}); err != nil {
in := input
if len(testCase.withJSON) > 0 {
in = testCase.withJSON
}
if err = s3Select.Open(newStringRSC(in)); err != nil {
t.Fatal(err)
}
@@ -656,13 +669,11 @@ func TestJSONQueries(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
in := input
if len(testCase.withJSON) > 0 {
in = testCase.withJSON
}
return ioutil.NopCloser(bytes.NewBufferString(in)), nil
}); err != nil {
in := input
if len(testCase.withJSON) > 0 {
in = testCase.withJSON
}
if err = s3Select.Open(newStringRSC(in)); err != nil {
t.Fatal(err)
}
@@ -743,9 +754,7 @@ func TestCSVQueries(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
return ioutil.NopCloser(bytes.NewBufferString(input)), nil
}); err != nil {
if err = s3Select.Open(newStringRSC(input)); err != nil {
t.Fatal(err)
}
@@ -928,9 +937,7 @@ func TestCSVQueries2(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
return ioutil.NopCloser(bytes.NewBuffer(testCase.input)), nil
}); err != nil {
if err = s3Select.Open(newBytesRSC(testCase.input)); err != nil {
t.Fatal(err)
}
@@ -1074,9 +1081,7 @@ true`,
t.Fatal(err)
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
return ioutil.NopCloser(bytes.NewBufferString(input)), nil
}); err != nil {
if err = s3Select.Open(newStringRSC(input)); err != nil {
t.Fatal(err)
}
@@ -1220,9 +1225,7 @@ func TestCSVInput(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
return ioutil.NopCloser(bytes.NewReader(csvData)), nil
}); err != nil {
if err = s3Select.Open(newBytesRSC(csvData)); err != nil {
t.Fatal(err)
}
@@ -1342,9 +1345,7 @@ func TestJSONInput(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
return ioutil.NopCloser(bytes.NewReader(jsonData)), nil
}); err != nil {
if err = s3Select.Open(newBytesRSC(jsonData)); err != nil {
t.Fatal(err)
}
@@ -1616,37 +1617,7 @@ func TestCSVRanges(t *testing.T) {
return
}
if err = s3Select.Open(func(offset, length int64) (io.ReadCloser, error) {
in := testCase.input
if offset != 0 || length != -1 {
// Copy from SelectObjectContentHandler
isSuffixLength := false
if offset < 0 {
isSuffixLength = true
}
if length > 0 {
length--
}
rs := &httpRangeSpec{
IsSuffixLength: isSuffixLength,
Start: offset,
End: offset + length,
}
if length == -1 {
rs.End = -1
}
t.Log("input, offset:", offset, "length:", length, "size:", len(in))
offset, length, err = rs.GetOffsetLength(int64(len(in)))
if err != nil {
return nil, err
}
t.Log("rs:", *rs, "offset:", offset, "length:", length)
in = in[offset : offset+length]
}
return ioutil.NopCloser(bytes.NewBuffer(in)), nil
}); err != nil {
if err = s3Select.Open(newBytesRSC(testCase.input)); err != nil {
if !testCase.wantErr {
t.Fatal(err)
}
@@ -1741,27 +1712,10 @@ func TestParquetInput(t *testing.T) {
for i, testCase := range testTable {
t.Run(fmt.Sprint(i), func(t *testing.T) {
getReader := func(offset int64, length int64) (io.ReadCloser, error) {
testdataFile := "testdata/testdata.parquet"
file, err := os.Open(testdataFile)
if err != nil {
return nil, err
}
fi, err := file.Stat()
if err != nil {
return nil, err
}
if offset < 0 {
offset = fi.Size() + offset
}
if _, err = file.Seek(offset, io.SeekStart); err != nil {
return nil, err
}
return file, nil
testdataFile := "testdata/testdata.parquet"
file, err := os.Open(testdataFile)
if err != nil {
t.Fatal(err)
}
s3Select, err := NewS3Select(bytes.NewReader(testCase.requestXML))
@@ -1769,10 +1723,12 @@ func TestParquetInput(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(getReader); err != nil {
if err = s3Select.Open(file); err != nil {
t.Fatal(err)
}
fmt.Printf("R: \nE: %s\n" /* string(w.response), */, string(testCase.expectedResult))
w := &testResponseWriter{}
s3Select.Evaluate(w)
s3Select.Close()
@@ -1862,27 +1818,10 @@ func TestParquetInputSchema(t *testing.T) {
for i, testCase := range testTable {
t.Run(fmt.Sprint(i), func(t *testing.T) {
getReader := func(offset int64, length int64) (io.ReadCloser, error) {
testdataFile := "testdata/lineitem_shipdate.parquet"
file, err := os.Open(testdataFile)
if err != nil {
return nil, err
}
fi, err := file.Stat()
if err != nil {
return nil, err
}
if offset < 0 {
offset = fi.Size() + offset
}
if _, err = file.Seek(offset, io.SeekStart); err != nil {
return nil, err
}
return file, nil
testdataFile := "testdata/lineitem_shipdate.parquet"
file, err := os.Open(testdataFile)
if err != nil {
t.Fatal(err)
}
s3Select, err := NewS3Select(bytes.NewReader(testCase.requestXML))
@@ -1890,7 +1829,7 @@ func TestParquetInputSchema(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(getReader); err != nil {
if err = s3Select.Open(file); err != nil {
t.Fatal(err)
}
@@ -1980,27 +1919,10 @@ func TestParquetInputSchemaCSV(t *testing.T) {
for i, testCase := range testTable {
t.Run(fmt.Sprint(i), func(t *testing.T) {
getReader := func(offset int64, length int64) (io.ReadCloser, error) {
testdataFile := "testdata/lineitem_shipdate.parquet"
file, err := os.Open(testdataFile)
if err != nil {
return nil, err
}
fi, err := file.Stat()
if err != nil {
return nil, err
}
if offset < 0 {
offset = fi.Size() + offset
}
if _, err = file.Seek(offset, io.SeekStart); err != nil {
return nil, err
}
return file, nil
testdataFile := "testdata/lineitem_shipdate.parquet"
file, err := os.Open(testdataFile)
if err != nil {
t.Fatal(err)
}
s3Select, err := NewS3Select(bytes.NewReader(testCase.requestXML))
@@ -2008,7 +1930,7 @@ func TestParquetInputSchemaCSV(t *testing.T) {
t.Fatal(err)
}
if err = s3Select.Open(getReader); err != nil {
if err = s3Select.Open(file); err != nil {
t.Fatal(err)
}
@@ -2037,76 +1959,3 @@ func TestParquetInputSchemaCSV(t *testing.T) {
})
}
}
// httpRangeSpec represents a range specification as supported by S3 GET
// object request.
//
// Case 1: Not present -> represented by a nil RangeSpec
// Case 2: bytes=1-10 (absolute start and end offsets) -> RangeSpec{false, 1, 10}
// Case 3: bytes=10- (absolute start offset with end offset unspecified) -> RangeSpec{false, 10, -1}
// Case 4: bytes=-30 (suffix length specification) -> RangeSpec{true, -30, -1}
type httpRangeSpec struct {
// Does the range spec refer to a suffix of the object?
IsSuffixLength bool
// Start and end offset specified in range spec
Start, End int64
}
func (h *httpRangeSpec) GetLength(resourceSize int64) (rangeLength int64, err error) {
switch {
case resourceSize < 0:
return 0, errors.New("Resource size cannot be negative")
case h == nil:
rangeLength = resourceSize
case h.IsSuffixLength:
specifiedLen := -h.Start
rangeLength = specifiedLen
if specifiedLen > resourceSize {
rangeLength = resourceSize
}
case h.Start >= resourceSize:
return 0, errors.New("errInvalidRange")
case h.End > -1:
end := h.End
if resourceSize <= end {
end = resourceSize - 1
}
rangeLength = end - h.Start + 1
case h.End == -1:
rangeLength = resourceSize - h.Start
default:
return 0, errors.New("Unexpected range specification case")
}
return rangeLength, nil
}
// GetOffsetLength computes the start offset and length of the range
// given the size of the resource
func (h *httpRangeSpec) GetOffsetLength(resourceSize int64) (start, length int64, err error) {
if h == nil {
// No range specified, implies whole object.
return 0, resourceSize, nil
}
length, err = h.GetLength(resourceSize)
if err != nil {
return 0, 0, err
}
start = h.Start
if h.IsSuffixLength {
start = resourceSize + h.Start
if start < 0 {
start = 0
}
}
return start, length, nil
}