From 7525f46295eb3249388351630c3d37cf14c37f9d Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 6 Apr 2023 14:16:48 -0500 Subject: [PATCH] refactor table cache to include allocator --- apply.go | 8 +++----- arrow.go | 23 ++++++++++++++--------- executor.go | 5 ++--- 3 files changed, 19 insertions(+), 17 deletions(-) diff --git a/apply.go b/apply.go index 198010cdd..a5d31944f 100644 --- a/apply.go +++ b/apply.go @@ -211,7 +211,6 @@ func (e *executor) executeApplyShard(ctx context.Context, qcx *Qcx, index string } } // - pool := memory.NewGoAllocator() // TODO(twg) 2022/09/01 singledton? ids := filter.ShardColumns() // needs to be shard columns // Fetch index. @@ -226,7 +225,7 @@ func (e *executor) executeApplyShard(ctx context.Context, qcx *Qcx, index string return value.NewVector([]value.Value{}), nil } - table, err := e.getDataTable(ctx, fname, pool) + table, pool, err := e.getDataTable(ctx, fname) if err != nil { return nil, err } @@ -264,7 +263,7 @@ func NewShardFile(ctx context.Context, name string, mem memory.Allocator, e *exe return &ShardFile{dest: name, executor: e, strings: make(map[key][]string)}, nil } // else read in existing - table, err := e.getDataTable(ctx, name, mem) + table, _, err := e.getDataTable(ctx, name) if err != nil { return nil, err } @@ -663,7 +662,6 @@ 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() @@ -672,7 +670,7 @@ func (api *API) GetDataframeSchema(ctx context.Context, indexName string) (inter name = strings.TrimSuffix(name, filepath.Ext(name)) // read the parquet file and extract the schema fname := filepath.Join(base, name) - table, err := api.server.executor.getDataTable(ctx, fname, mem) + table, _, err := api.server.executor.getDataTable(ctx, fname) if err != nil { return nil, err } diff --git a/arrow.go b/arrow.go index e3c015db3..ee2477fd6 100644 --- a/arrow.go +++ b/arrow.go @@ -394,7 +394,7 @@ func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string return &BasicTable{name: name}, nil } - table, err := e.getDataTable(ctx, fname, pool) + table, pool, err := e.getDataTable(ctx, fname) if err != nil { return nil, errors.Wrap(err, "arrow readTableParquet") } @@ -435,24 +435,29 @@ func (e *executor) dataFrameExists(fname string) bool { return true } -func (e *executor) getDataTable(ctx context.Context, fname string, memignore memory.Allocator) (arrow.Table, error) { - table, ok := e.arrowCache[fname] +type arrowCache struct { + table arrow.Table + pool memory.Allocator +} + +func (e *executor) getDataTable(ctx context.Context, fname string) (arrow.Table, memory.Allocator, error) { + cache, ok := e.arrowCache[fname] if ok { - return table, nil + return cache.table, cache.pool, 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 + e.arrowCache[fname] = &arrowCache{table: table, pool: mem} + return table, mem, err } table, err := readTableArrow(fname, mem) if err != nil { - return nil, err + return nil, nil, err } - e.arrowCache[fname] = table - return table, nil + e.arrowCache[fname] = &arrowCache{table: table, pool: mem} + return table, mem, nil } func (e *executor) typeIsParquet() bool { diff --git a/executor.go b/executor.go index f3c6ab352..2f37d0ac9 100644 --- a/executor.go +++ b/executor.go @@ -17,7 +17,6 @@ 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" @@ -81,7 +80,7 @@ type executor struct { // Temporary flag to be removed when stablized dataframeEnabled bool datafameUseParquet bool - arrowCache map[string]arrow.Table + arrowCache map[string]*arrowCache } // executorOption is a functional option type for pilosa.executor @@ -145,7 +144,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) + e.arrowCache = make(map[string]*arrowCache) return e }