mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 20:37:52 +00:00
try again
This commit is contained in:
parent
0404aad84b
commit
617afc4c1d
2 changed files with 10 additions and 15 deletions
20
arrow.go
20
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
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue