From 617afc4c1d345c1d80b635d0ad69db9e74cdf29d Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 6 Apr 2023 17:50:03 -0500 Subject: [PATCH] try again --- arrow.go | 20 +++++++------------- executor.go | 5 +++-- 2 files changed, 10 insertions(+), 15 deletions(-) diff --git a/arrow.go b/arrow.go index 730a7019d..cc74b8df9 100644 --- a/arrow.go +++ b/arrow.go @@ -435,29 +435,23 @@ func (e *executor) dataFrameExists(fname string) bool { return true } -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] + table, ok := e.arrowCache[fname] if ok { - cache.table.Retain() - return cache.table, e.pool, nil + table.Retain() + return table, e.pool, nil } // ignoring the passed in allocatorsince where caching - mem := memory.NewGoAllocator() if e.typeIsParquet() { - table, err := readTableParquetCtx(ctx, fname, mem) - e.arrowCache[fname] = &arrowCache{table: table, pool: mem} + table, err := readTableParquetCtx(ctx, fname, e.pool) + e.arrowCache[fname] = table return table, e.pool, err } - table, err := readTableArrow(fname, mem) + table, err := readTableArrow(fname, e.pool) if err != nil { return nil, nil, err } - e.arrowCache[fname] = &arrowCache{table: table, pool: mem} + e.arrowCache[fname] = table return table, e.pool, nil } diff --git a/executor.go b/executor.go index ab2b5792b..ef479ad6b 100644 --- a/executor.go +++ b/executor.go @@ -17,6 +17,7 @@ import ( "time" "unsafe" + "github.com/apache/arrow/go/v10/arrow" "github.com/apache/arrow/go/v10/arrow/memory" "github.com/featurebasedb/featurebase/v3/dax" "github.com/featurebasedb/featurebase/v3/disco" @@ -81,7 +82,7 @@ type executor struct { // Temporary flag to be removed when stablized dataframeEnabled bool datafameUseParquet bool - arrowCache map[string]*arrowCache + arrowCache map[string]arrow.Table pool memory.Allocator } @@ -146,7 +147,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]*arrowCache) + e.arrowCache = make(map[string]arrow.Table) e.pool = memory.NewGoAllocator() return e }