From 5c5260f0b62f22928f9970a0d14fcbf49a31a950 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Fri, 7 Apr 2023 11:41:19 -0500 Subject: [PATCH] singleton arrow allocator --- apply.go | 25 ++++++++++++------------- arrow.go | 29 ++++++++++++++++------------- 2 files changed, 28 insertions(+), 26 deletions(-) diff --git a/apply.go b/apply.go index a5d31944f..aa849d6de 100644 --- a/apply.go +++ b/apply.go @@ -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 } diff --git a/arrow.go b/arrow.go index cc74b8df9..1b8de9761 100644 --- a/arrow.go +++ b/arrow.go @@ -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 {