singleton arrow allocator

This commit is contained in:
Todd Gruben 2023-04-07 11:41:19 -05:00
parent 617afc4c1d
commit 5c5260f0b6
2 changed files with 28 additions and 26 deletions

View file

@ -47,7 +47,7 @@ func runIvyString(context value.Context, str string) (ok bool, err error) {
// Possibly combine all arrays together then apply some interesting
// computation at the end?
func IvyReduce(reduceCode string, opCode string, opt *ExecOptions) (func(ctx context.Context, prev, v interface{}) interface{}, func() (*dataframe.DataFrame, error)) {
func (e *executor) IvyReduce(reduceCode string, opCode string, opt *ExecOptions) (func(ctx context.Context, prev, v interface{}) interface{}, func() (*dataframe.DataFrame, error)) {
var accumulator value.Value
mu := &sync.Mutex{}
concat := value.BinaryOps[opCode]
@ -90,10 +90,9 @@ func IvyReduce(reduceCode string, opCode string, opt *ExecOptions) (func(ctx con
return nil
}
tablerFn := func() (*dataframe.DataFrame, error) {
pool := memory.NewGoAllocator() // TODO(twg) 2022/09/01 singledton?
if opt.Remote {
col := value.ToArrowColumn(accumulator, pool)
return dataframe.NewDataFrameFromColumns(pool, []arrow.Column{*col})
col := value.ToArrowColumn(accumulator, e.pool)
return dataframe.NewDataFrameFromColumns(e.pool, []arrow.Column{*col})
}
// only actually reduce on the initiating node i hate the network
// over head but oh well
@ -107,9 +106,9 @@ func IvyReduce(reduceCode string, opCode string, opt *ExecOptions) (func(ctx con
if v == nil {
return nil, errors.New("ivy reduction no result ")
}
col := value.ToArrowColumn(ctxIvy.Global("_"), pool)
col := value.ToArrowColumn(ctxIvy.Global("_"), e.pool)
return dataframe.NewDataFrameFromColumns(pool, []arrow.Column{*col})
return dataframe.NewDataFrameFromColumns(e.pool, []arrow.Column{*col})
}
return nil, errors.New("ivy reduction failed ")
}
@ -141,9 +140,9 @@ func (e *executor) executeApply(ctx context.Context, qcx *Qcx, index string, c *
if err != nil {
return nil, err
}
reduceFn, tablerFn := IvyReduce("_", ",", opt)
reduceFn, tablerFn := e.IvyReduce("_", ",", opt)
if ok {
reduceFn, tablerFn = IvyReduce(ivyReduce, ",", opt)
reduceFn, tablerFn = e.IvyReduce(ivyReduce, ",", opt)
}
_, err = e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
@ -225,12 +224,12 @@ func (e *executor) executeApplyShard(ctx context.Context, qcx *Qcx, index string
return value.NewVector([]value.Value{}), nil
}
table, pool, err := e.getDataTable(ctx, fname)
table, err := e.getDataTable(ctx, fname)
if err != nil {
return nil, err
}
defer table.Release()
df, err := dataframe.NewDataFrameFromTable(pool, table)
df, err := dataframe.NewDataFrameFromTable(e.pool, table)
if err != nil {
return nil, err
}
@ -241,7 +240,7 @@ func (e *executor) executeApplyShard(ctx context.Context, qcx *Qcx, index string
if len(ids) == 0 {
return value.NewVector([]value.Value{}), nil
}
resolver, err = filterDataframe(resolver, pool, ids)
resolver, err = filterDataframe(resolver, e.pool, ids)
if err != nil {
return nil, err
}
@ -263,7 +262,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)
table, err := e.getDataTable(ctx, name)
if err != nil {
return nil, err
}
@ -670,7 +669,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)
table, err := api.server.executor.getDataTable(ctx, fname)
if err != nil {
return nil, err
}

View file

@ -19,6 +19,7 @@ import (
"github.com/apache/arrow/go/v10/parquet/pqarrow"
"github.com/featurebasedb/featurebase/v3/pql"
"github.com/featurebasedb/featurebase/v3/tracing"
"github.com/featurebasedb/featurebase/v3/vprint"
"github.com/gomem/gomem/pkg/dataframe"
"github.com/pkg/errors"
)
@ -51,14 +52,13 @@ func (e *executor) executeArrow(ctx context.Context, qcx *Qcx, index string, c *
}
mapcounter := 0
reducecounter := 0
pool := memory.NewGoAllocator() // TODO(twg) 2022/09/01 singledton?
// Execute calls in bulk on each remote node and merge.
mu := &sync.Mutex{}
mapFn := func(ctx context.Context, shard uint64, mopt *mapOptions) (_ interface{}, err error) {
mu.Lock()
mapcounter++
mu.Unlock()
return e.executeArrowShard(ctx, qcx, index, c, shard, pool, columnFilter)
return e.executeArrowShard(ctx, qcx, index, c, shard, columnFilter)
}
tables := make([]*BasicTable, 0)
@ -79,7 +79,7 @@ func (e *executor) executeArrow(ctx context.Context, qcx *Qcx, index string, c *
}
case arrow.Table:
if t.NumRows() > 0 {
bt := BasicTableFromArrow(t, pool)
bt := BasicTableFromArrow(t, e.pool)
mu.Lock()
tables = append(tables, bt)
mu.Unlock()
@ -95,7 +95,7 @@ func (e *executor) executeArrow(ctx context.Context, qcx *Qcx, index string, c *
if len(tables) == 0 {
return &BasicTable{name: "empty"}, nil
}
tbl := Concat(tables[0].Schema(), tables, pool)
tbl := Concat(tables[0].Schema(), tables, e.pool)
r := dataframe.NewChunkResolver(tbl.Column(0))
return &BasicTable{resolver: &r, table: tbl}, nil
}
@ -363,7 +363,7 @@ func filterColumns(filters []string, table arrow.Table) arrow.Table {
return array.NewTable(filterdSchema, cols, table.NumRows())
}
func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64, pool memory.Allocator, columnFilter []string) (*BasicTable, error) {
func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64, columnFilter []string) (*BasicTable, error) {
name := fmt.Sprintf("a. %v", shard)
span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeArrowShard")
defer span.Finish()
@ -394,7 +394,7 @@ func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string
return &BasicTable{name: name}, nil
}
table, pool, err := e.getDataTable(ctx, fname)
table, err := e.getDataTable(ctx, fname)
if err != nil {
return nil, errors.Wrap(err, "arrow readTableParquet")
}
@ -402,7 +402,7 @@ func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string
if len(columnFilter) > 0 {
table = filterColumns(columnFilter, table)
}
df, err := dataframe.NewDataFrameFromTable(pool, table)
df, err := dataframe.NewDataFrameFromTable(e.pool, table)
if err != nil {
return nil, errors.Wrap(err, "arrow NewDataFromTable")
}
@ -413,7 +413,7 @@ func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string
if len(ids) == 0 {
return &BasicTable{name: name}, nil
}
resolver, err = filterDataframe(resolver, pool, ids)
resolver, err = filterDataframe(resolver, e.pool, ids)
if err != nil {
return nil, errors.Wrap(err, "filtering dataframe")
}
@ -435,24 +435,27 @@ func (e *executor) dataFrameExists(fname string) bool {
return true
}
func (e *executor) getDataTable(ctx context.Context, fname string) (arrow.Table, memory.Allocator, error) {
func (e *executor) getDataTable(ctx context.Context, fname string) (arrow.Table, error) {
table, ok := e.arrowCache[fname]
if ok {
vprint.VV("returning table from cache name:%v numcols:%v numrows:%v", fname, table.NumCols(), table.NumRows())
table.Retain()
return table, e.pool, nil
return table, nil
}
// ignoring the passed in allocatorsince where caching
if e.typeIsParquet() {
table, err := readTableParquetCtx(ctx, fname, e.pool)
e.arrowCache[fname] = table
return table, e.pool, err
return table, err
}
table, err := readTableArrow(fname, e.pool)
if err != nil {
return nil, nil, err
return nil, err
}
e.arrowCache[fname] = table
return table, e.pool, nil
vprint.VV("returning new table and cacheing at cache name:%v numcols:%v numrows:%v", fname, table.NumCols(), table.NumRows())
return table, nil
}
func (e *executor) typeIsParquet() bool {