diff --git a/apply.go b/apply.go index dbde085fc..198010cdd 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()) + fname := sf.dest + sf.executor.TableExtension() + delete(sf.executor.arrowCache, fname) + return os.Rename(rtemp+sf.executor.TableExtension(), fname) } func (sf *ShardFile) LoadBlobs() error { diff --git a/arrow.go b/arrow.go index 348dec94d..e3c015db3 100644 --- a/arrow.go +++ b/arrow.go @@ -435,12 +435,24 @@ func (e *executor) dataFrameExists(fname string) bool { return true } -func (e *executor) getDataTable(ctx context.Context, fname string, mem memory.Allocator) (arrow.Table, error) { +func (e *executor) getDataTable(ctx context.Context, fname string, memignore memory.Allocator) (arrow.Table, error) { + table, ok := e.arrowCache[fname] + if ok { + return table, nil + } + // ignoring the passed in allocatorsince where caching + mem := memory.NewGoAllocator() if e.typeIsParquet() { table, err := readTableParquetCtx(ctx, fname, mem) + e.arrowCache[fname] = table return table, err } - return readTableArrow(fname, mem) + table, err := readTableArrow(fname, mem) + if err != nil { + return nil, err + } + e.arrowCache[fname] = table + return table, nil } func (e *executor) typeIsParquet() bool { diff --git a/executor.go b/executor.go index 5b804c4c2..f3c6ab352 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 + arrowCache map[string]arrow.Table } // executorOption is a functional option type for pilosa.executor @@ -143,6 +145,7 @@ func newExecutor(opts ...executorOption) *executor { e.work = make(chan job, e.workerPoolSize) _ = testhook.Opened(NewAuditor(), e, nil) e.workers = task.NewPool(e.workerPoolSize, e.doOneJob, e) + e.arrowCache = make(map[string]arrow.Table) return e }