From 96e83c950e9a1fafd33dd04c46145d632ad60443 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Tue, 28 Feb 2023 08:45:29 -0600 Subject: [PATCH] cache the arrow table --- apply.go | 4 +++- arrow.go | 12 +++++++++++- executor.go | 2 ++ 3 files changed, 16 insertions(+), 2 deletions(-) diff --git a/apply.go b/apply.go index dbde085fc..476857aee 100644 --- a/apply.go +++ b/apply.go @@ -407,7 +407,9 @@ func (sf *ShardFile) Process(cs *ChangesetRequest) error { if err != nil { return err } - return os.Rename(rtemp+sf.executor.TableExtension(), sf.dest+sf.executor.TableExtension()) + key := sf.dest + sf.executor.TableExtension() + delete(sf.executor.frameCache, key) + return os.Rename(rtemp+sf.executor.TableExtension(), key) } func (sf *ShardFile) LoadBlobs() error { diff --git a/arrow.go b/arrow.go index 4b416351e..f160ca166 100644 --- a/arrow.go +++ b/arrow.go @@ -436,11 +436,21 @@ func (e *executor) dataFrameExists(fname string) bool { } func (e *executor) getDataTable(ctx context.Context, fname string, mem memory.Allocator) (arrow.Table, error) { + table, ok := e.frameCache[fname] + if ok { + return table, nil + } if e.typeIsParquet() { table, err := readTableParquetCtx(ctx, fname, mem) + e.frameCache[fname] = table return table, err } - return readTableArrow(fname, mem) + table, err := readTableArrow(fname, mem) + if err != nil { + return nil, err + } + e.frameCache[fname] = table + return table, err } func (e *executor) typeIsParquet() bool { diff --git a/executor.go b/executor.go index c63771c36..bb3e2b05c 100644 --- a/executor.go +++ b/executor.go @@ -17,6 +17,7 @@ import ( "time" "unsafe" + "github.com/apache/arrow/go/v10/arrow" "github.com/featurebasedb/featurebase/v3/dax" "github.com/featurebasedb/featurebase/v3/disco" "github.com/featurebasedb/featurebase/v3/pql" @@ -80,6 +81,7 @@ type executor struct { // Temporary flag to be removed when stablized dataframeEnabled bool datafameUseParquet bool + frameCache map[string]arrow.Table } // executorOption is a functional option type for pilosa.executor