refactor table cache to include allocator

This commit is contained in:
Todd Gruben 2023-04-06 14:16:48 -05:00
parent d1332e0c67
commit 7525f46295
3 changed files with 19 additions and 17 deletions

View file

@ -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
}

View file

@ -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 {

View file

@ -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
}