From 24372a94056bb90b404a93450ab1ee627f504b7f Mon Sep 17 00:00:00 2001 From: tgruben Date: Fri, 16 Dec 2022 14:42:39 -0600 Subject: [PATCH] [FB-1822] Change dataframe disk format to Arrow from parquet (#2376) * change default backend to arrow file format instead of parquet (cherry picked from commit 6d96a1474ee929a28c8921e71f33cfb739b1dfd4) --- apply.go | 84 ++++++-------------------- arrow.go | 151 ++++++++++++++++++++++++++++++++++++++++++++++- arrow_test.go | 53 +++++++++++++++++ ctl/server.go | 1 + executor.go | 3 +- http_handler.go | 4 +- server.go | 11 +++- server/config.go | 3 +- server/server.go | 7 ++- 9 files changed, 243 insertions(+), 74 deletions(-) create mode 100644 arrow_test.go diff --git a/apply.go b/apply.go index ecbee7735..12f89aa25 100644 --- a/apply.go +++ b/apply.go @@ -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/featurebasedb/featurebase/v3/pql" "github.com/featurebasedb/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 } diff --git a/arrow.go b/arrow.go index d5b75e820..1d009b285 100644 --- a/arrow.go +++ b/arrow.go @@ -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/featurebasedb/featurebase/v3/pql" "github.com/featurebasedb/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 +} diff --git a/arrow_test.go b/arrow_test.go new file mode 100644 index 000000000..8142997db --- /dev/null +++ b/arrow_test.go @@ -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()) + } +} diff --git a/ctl/server.go b/ctl/server.go index f4f3be907..b1320a73b 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -144,6 +144,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 } diff --git a/executor.go b/executor.go index 2a33aba24..2172f39b3 100644 --- a/executor.go +++ b/executor.go @@ -77,7 +77,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 diff --git a/http_handler.go b/http_handler.go index 797f31125..ecdcb623d 100644 --- a/http_handler.go +++ b/http_handler.go @@ -3860,7 +3860,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) @@ -4184,7 +4184,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) } diff --git a/server.go b/server.go index 39d5e7995..eb253add3 100644 --- a/server.go +++ b/server.go @@ -99,7 +99,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 @@ -457,6 +458,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() @@ -528,6 +536,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 { diff --git a/server/config.go b/server/config.go index 20e53848f..43ac93383 100644 --- a/server/config.go +++ b/server/config.go @@ -252,7 +252,8 @@ type Config struct { Auth Auth Dataframe struct { - Enable bool `toml:"enable"` + Enable bool `toml:"enable"` + UseParquet bool `toml:"use-parquet"` } `toml:"dataframe"` } diff --git a/server/server.go b/server/server.go index b192d380f..8f2b5b152 100644 --- a/server/server.go +++ b/server/server.go @@ -195,8 +195,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 @@ -591,6 +593,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 {