[FB-1822] Change dataframe disk format to Arrow from parquet (#2376)

* change default backend to arrow file format instead of parquet
This commit is contained in:
tgruben 2022-12-16 14:42:39 -06:00 committed by GitHub
parent f65367ab46
commit 6d96a1474e
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
9 changed files with 243 additions and 74 deletions

View file

@ -13,8 +13,6 @@ import (
"github.com/apache/arrow/go/v10/arrow"
"github.com/apache/arrow/go/v10/arrow/array"
"github.com/apache/arrow/go/v10/arrow/memory"
"github.com/apache/arrow/go/v10/parquet/file"
"github.com/apache/arrow/go/v10/parquet/pqarrow"
"github.com/gomem/gomem/pkg/dataframe"
"github.com/molecula/featurebase/v3/pql"
"github.com/molecula/featurebase/v3/tracing"
@ -220,12 +218,14 @@ func (e *executor) executeApplyShard(ctx context.Context, qcx *Qcx, index string
if idx == nil {
return nil, newNotFoundError(ErrIndexNotFound, index)
}
fname := idx.GetDataFramePath(shard)
if _, err := os.Stat(fname + ".parquet"); os.IsNotExist(err) {
if !e.dataFrameExists(fname) {
return value.NewVector([]value.Value{}), nil
}
table, err := readTableParquet(fname)
table, err := e.getDataTable(ctx, fname, pool)
if err != nil {
return nil, err
}
@ -254,57 +254,20 @@ func (e *executor) executeApplyShard(ctx context.Context, qcx *Qcx, index string
return context.Global("_"), nil
}
func readTableParquet(filename string) (arrow.Table, error) {
r, err := os.Open(filename + ".parquet")
if err != nil {
return nil, err
}
pf, err := file.NewParquetReader(r)
if err != nil {
return nil, err
}
reader, err := pqarrow.NewFileReader(pf, pqarrow.ArrowReadProperties{}, memory.DefaultAllocator)
if err != nil {
return nil, err
}
return reader.ReadTable(context.Background())
}
func readTableParquetCtx(ctx context.Context, filename string, mem memory.Allocator) (arrow.Table, error) {
r, err := os.Open(filename + ".parquet")
if err != nil {
return nil, err
}
pf, err := file.NewParquetReader(r)
if err != nil {
return nil, err
}
reader, err := pqarrow.NewFileReader(pf, pqarrow.ArrowReadProperties{}, mem)
if err != nil {
return nil, err
}
return reader.ReadTable(ctx)
}
// ///////////////////////////////////////////////////////
// all the ingest supporting functions
// ///////////////////////////////////////////////////////
func NewShardFile(name string) (*ShardFile, error) {
if _, err := os.Stat(name + ".parquet"); os.IsNotExist(err) {
return &ShardFile{dest: name}, nil
func NewShardFile(ctx context.Context, name string, mem memory.Allocator, e *executor) (*ShardFile, error) {
if !e.dataFrameExists(name) {
return &ShardFile{dest: name, executor: e}, nil
}
// else read in existing
table, err := readTableParquet(name)
table, err := e.getDataTable(ctx, name, mem)
if err != nil {
return nil, err
}
return &ShardFile{table: table, schema: table.Schema(), dest: name}, nil
return &ShardFile{table: table, schema: table.Schema(), dest: name, executor: e}, nil
}
type NameType struct {
@ -349,6 +312,7 @@ type ShardFile struct {
added int64
columns []interface{}
dest string
executor *executor
}
func compareSchema(s1, s2 *arrow.Schema) bool {
@ -427,7 +391,7 @@ func (sf *ShardFile) Process(cs *ChangesetRequest) error {
if err != nil {
return err
}
return os.Rename(rtemp+".parquet", sf.dest+".parquet")
return os.Rename(rtemp+sf.executor.TableExtension(), sf.dest+sf.executor.TableExtension())
}
func (sf *ShardFile) process(cs *ChangesetRequest) error {
@ -521,22 +485,9 @@ func (sf *ShardFile) Save(name string) error {
}
}
rec := array.NewRecord(sf.schema, parts, sf.beforeRows+sf.added)
df, err := dataframe.NewDataFrameFromRecord(mem, rec)
if err != nil {
return err
}
// confirm change
w, err := os.Create(name + ".parquet")
if err != nil {
return err
}
table := array.NewTableFromRecords(sf.schema, []arrow.Record{rec})
err = df.ToParquet(w, 1024)
if err != nil {
return err
}
w.Close()
return nil
return sf.executor.SaveTable(name, table, mem)
}
// TODO(twg) 2022/10/03 Not a huge fan of the global variable will look at adding to executor structure
@ -573,7 +524,8 @@ func (api *API) ApplyDataframeChangeset(ctx context.Context, index string, cs *C
mu := getDataframeWritelock(shard)
mu.Lock()
defer mu.Unlock()
shardFile, err := NewShardFile(fname)
mem := memory.NewGoAllocator()
shardFile, err := NewShardFile(ctx, fname, mem, api.server.executor)
if err != nil {
return err
}
@ -600,14 +552,16 @@ func (api *API) GetDataframeSchema(ctx context.Context, indexName string) (inter
dir, _ := os.Open(base)
files, _ := dir.Readdir(0)
parts := make([]column, 0)
mem := memory.NewGoAllocator()
for i := range files {
file := files[i]
name := file.Name()
if strings.HasSuffix(name, ".parquet") {
if api.server.executor.IsDataframeFile(name) {
// strip off the parquet extenison
name = strings.TrimSuffix(name, filepath.Ext(name))
// read the parquet file and extract the schema
table, err := readTableParquet(filepath.Join(base, name))
fname := filepath.Join(base, name)
table, err := api.server.executor.getDataTable(ctx, fname, mem)
if err != nil {
return nil, err
}

151
arrow.go
View file

@ -5,12 +5,18 @@ import (
"context"
"encoding/json"
"fmt"
"io"
"os"
"strings"
"sync"
"github.com/apache/arrow/go/v10/arrow"
"github.com/apache/arrow/go/v10/arrow/array"
"github.com/apache/arrow/go/v10/arrow/ipc"
"github.com/apache/arrow/go/v10/arrow/memory"
"github.com/apache/arrow/go/v10/parquet"
"github.com/apache/arrow/go/v10/parquet/file"
"github.com/apache/arrow/go/v10/parquet/pqarrow"
"github.com/gomem/gomem/pkg/dataframe"
"github.com/molecula/featurebase/v3/pql"
"github.com/molecula/featurebase/v3/tracing"
@ -371,12 +377,14 @@ func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string
if idx == nil {
return nil, newNotFoundError(ErrIndexNotFound, index)
}
fname := idx.GetDataFramePath(shard)
if _, err := os.Stat(fname + ".parquet"); os.IsNotExist(err) {
if !e.dataFrameExists(fname) {
return &basicTable{name: name}, nil
}
table, err := readTableParquetCtx(context.TODO(), fname, pool)
table, err := e.getDataTable(ctx, fname, pool)
if err != nil {
return nil, errors.Wrap(err, "arrow readTableParquet")
}
@ -403,3 +411,142 @@ func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string
table.Retain()
return &basicTable{resolver: resolver, table: table, filtered: filter != nil, name: name}, nil
}
func (e *executor) dataFrameExists(fname string) bool {
if e.typeIsParquet() {
if _, err := os.Stat(fname + ".parquet"); os.IsNotExist(err) {
return false
}
return true
}
if _, err := os.Stat(fname + ".arrow"); os.IsNotExist(err) {
return false
}
return true
}
func (e *executor) getDataTable(ctx context.Context, fname string, mem memory.Allocator) (arrow.Table, error) {
if e.typeIsParquet() {
table, err := readTableParquetCtx(ctx, fname, mem)
return table, err
}
return readTableArrow(fname, mem)
}
func (e *executor) typeIsParquet() bool {
return e.datafameUseParquet
}
func (e *executor) IsDataframeFile(name string) bool {
if e.typeIsParquet() {
return strings.HasSuffix(name, ".parquet")
}
return strings.HasSuffix(name, ".arrow")
}
func (e *executor) SaveTable(name string, table arrow.Table, mem memory.Allocator) error {
if e.typeIsParquet() {
return writeTableParquet(table, name)
}
return writeTableArrow(table, name, mem)
}
func (e *executor) TableExtension() string {
if e.typeIsParquet() {
return ".parquet"
}
return ".arrow"
}
func readTableArrow(filename string, mem memory.Allocator) (arrow.Table, error) {
r, err := os.Open(filename + ".arrow")
if err != nil {
return nil, err
}
rr, err := ipc.NewFileReader(r, ipc.WithAllocator(mem))
if err != nil {
return nil, err
}
defer rr.Close()
records := make([]arrow.Record, rr.NumRecords(), rr.NumRecords())
i := 0
for {
rec, err := rr.Read()
if err == io.EOF {
break
} else if err != nil {
return nil, err
}
records[i] = rec
i++
}
records = records[:i]
table := array.NewTableFromRecords(rr.Schema(), records)
return table, nil
}
func readTableParquetCtx(ctx context.Context, filename string, mem memory.Allocator) (arrow.Table, error) {
r, err := os.Open(filename + ".parquet")
if err != nil {
return nil, err
}
defer r.Close()
pf, err := file.NewParquetReader(r)
if err != nil {
return nil, err
}
reader, err := pqarrow.NewFileReader(pf, pqarrow.ArrowReadProperties{}, mem)
if err != nil {
return nil, err
}
return reader.ReadTable(ctx)
}
func writeTableParquet(table arrow.Table, filename string) error {
f, err := os.Create(filename + ".parquet")
if err != nil {
return err
}
defer f.Close()
props := parquet.NewWriterProperties(parquet.WithDictionaryDefault(false))
arrProps := pqarrow.DefaultWriterProps()
chunkSize := 10 * 1024 * 1024
err = pqarrow.WriteTable(table, f, int64(chunkSize), props, arrProps)
if err != nil {
return err
}
f.Sync()
return nil
}
func writeTableArrow(table arrow.Table, filename string, mem memory.Allocator) error {
f, err := os.Create(filename + ".arrow")
if err != nil {
return err
}
defer f.Close()
writer, err := ipc.NewFileWriter(f, ipc.WithAllocator(mem), ipc.WithSchema(table.Schema()))
if err != nil {
panic(err)
}
chunkSize := int64(0)
tr := array.NewTableReader(table, chunkSize)
defer tr.Release()
n := 0
for tr.Next() {
arec := tr.Record()
err = writer.Write(arec)
if err != nil {
panic(err)
}
n++
}
err = writer.Close()
if err != nil {
panic(err)
}
f.Sync()
return nil
}

53
arrow_test.go Normal file
View file

@ -0,0 +1,53 @@
// Copyright 2021 Molecula Corp. All rights reserved.
package pilosa
import (
"context"
"encoding/hex"
"math/rand"
"os"
"path/filepath"
"testing"
"github.com/apache/arrow/go/v10/arrow"
"github.com/apache/arrow/go/v10/arrow/array"
"github.com/apache/arrow/go/v10/arrow/memory"
)
func TempFileName(prefix string) string {
randBytes := make([]byte, 16)
rand.Read(randBytes)
return filepath.Join(os.TempDir(), prefix+hex.EncodeToString(randBytes))
}
func Test_TableParquet(t *testing.T) {
// create a arrow table
schema := arrow.NewSchema(
[]arrow.Field{
{Name: "num", Type: arrow.PrimitiveTypes.Float64},
},
nil, // no metadata
)
mem := memory.NewGoAllocator()
b := array.NewRecordBuilder(mem, schema)
defer b.Release()
b.Field(0).(*array.Float64Builder).AppendValues([]float64{1.0, 1.5, 2.0}, nil)
table := array.NewTableFromRecords(schema, []arrow.Record{b.NewRecord()})
defer table.Release()
fileName := TempFileName("pq-")
// save it as a parquet file
err := writeTableParquet(table, fileName)
if err != nil {
t.Fatal(err)
}
defer os.Remove(fileName)
// read it back in and compare the result
got, err := readTableParquetCtx(context.Background(), fileName, mem)
if err != nil {
t.Fatalf("readTableParquetCtx() error = %v", err)
}
if got.NumCols() != table.NumCols() {
t.Errorf("got:%v expected:%v", got.NumCols(), table.NumCols())
}
}

View file

@ -143,6 +143,7 @@ func serverFlagSet(srv *server.Config, prefix string) *pflag.FlagSet {
flags.BoolVar(&srv.DataDog.BlockProfile, pre("datadog.block-profile"), false, "golang pprof goroutine ")
flags.BoolVar(&srv.Dataframe.Enable, pre("dataframe.enable"), false, "EXPERIMENTAL enable support for Apply and Arrow")
flags.BoolVar(&srv.Dataframe.UseParquet, pre("dataframe.use-parquet"), false, "EXPERIMENTAL use parquet for file format")
return flags
}

View file

@ -76,7 +76,8 @@ type executor struct {
maxMemory int64
// Temporary flag to be removed when stablized
dataframeEnabled bool
dataframeEnabled bool
datafameUseParquet bool
}
// executorOption is a functional option type for pilosa.executor

View file

@ -3859,7 +3859,7 @@ func (h *Handler) handlePostDataframeRestore(w http.ResponseWriter, r *http.Requ
http.Error(w, fmt.Sprintf("Index %s Not Found", indexName), http.StatusNotFound)
return
}
filename := idx.GetDataFramePath(shard) + ".parquet"
filename := idx.GetDataFramePath(shard) + h.api.server.executor.TableExtension()
dest, err := os.Create(filename)
if err != nil {
http.Error(w, fmt.Sprintf("failed to create restore dataframe shard %v %v err:%v", indexName, shard, err), http.StatusBadRequest)
@ -4183,7 +4183,7 @@ func (h *Handler) handleGetDataframe(w http.ResponseWriter, r *http.Request) {
http.Error(w, fmt.Sprintf("Index %s Not Found", indexName), http.StatusNotFound)
return
}
filename := idx.GetDataFramePath(shard) + ".parquet"
filename := idx.GetDataFramePath(shard) + h.api.server.executor.TableExtension()
http.ServeFile(w, r, filename)
}

View file

@ -98,7 +98,8 @@ type Server struct { // nolint: maligned
serverlessStorage *daxstorage.ResourceManager
dataframeEnabled bool
dataframeEnabled bool
dataframeUseParquet bool
}
type ExecutionPlannerFn func(executor Executor, api *API, sql string) sql3.CompilePlanner
@ -456,6 +457,13 @@ func OptServerIsDataframeEnabled(is bool) ServerOption {
}
}
func OptServerDataframeUseParquet(is bool) ServerOption {
return func(s *Server) error {
s.dataframeUseParquet = is
return nil
}
}
// NewServer returns a new instance of Server.
func NewServer(opts ...ServerOption) (*Server, error) {
cluster := newCluster()
@ -527,6 +535,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
}
s.executor = newExecutor(executorOpts...)
s.executor.dataframeEnabled = s.dataframeEnabled
s.executor.datafameUseParquet = s.dataframeUseParquet
path, err := expandDirName(s.dataDir)
if err != nil {

View file

@ -251,7 +251,8 @@ type Config struct {
Auth Auth
Dataframe struct {
Enable bool `toml:"enable"`
Enable bool `toml:"enable"`
UseParquet bool `toml:"use-parquet"`
} `toml:"dataframe"`
}

View file

@ -196,8 +196,10 @@ const (
// we want to set resource limits *exactly once*, and then be able
// to report on whether or not that succeeded.
var setupResourceLimitsOnce sync.Once
var setupResourceLimitsErr error
var (
setupResourceLimitsOnce sync.Once
setupResourceLimitsErr error
)
// doSetupResourceLimits is the function which actually does the
// resource limit setup, possibly yielding an error. it's a Command
@ -592,6 +594,7 @@ func (m *Command) setupServer() error {
pilosa.OptServerExecutionPlannerFn(executionPlannerFn),
pilosa.OptServerServerlessStorage(m.serverlessStorage),
pilosa.OptServerIsDataframeEnabled(m.Config.Dataframe.Enable),
pilosa.OptServerDataframeUseParquet(m.Config.Dataframe.UseParquet),
}
if m.isComputeNode {