From 38eea9b4a7cf8276afb23d5caa8751da5251cc84 Mon Sep 17 00:00:00 2001 From: Jason Aten Date: Mon, 27 Jul 2020 20:25:43 -0400 Subject: [PATCH 1/2] reparallelize view.go openFragmentsInTx() --- view.go | 37 +++++++++++++++++++++++++++++-------- 1 file changed, 29 insertions(+), 8 deletions(-) diff --git a/view.go b/view.go index 106ae2b42..682b3fb00 100644 --- a/view.go +++ b/view.go @@ -171,17 +171,38 @@ func (v *view) openFragmentsInTx() error { if err != nil { return errors.Wrap(err, "SliceOfShards") } + + eg, ctx := errgroup.WithContext(context.Background()) + var mu sync.Mutex + +shardLoop: for _, shard := range shards { - frag := v.newFragment(v.fragmentPath(shard), shard) - if err := frag.Open(); err != nil { - return fmt.Errorf("open fragment: shard=%d, err=%s", frag.shard, err) + select { + case <-ctx.Done(): + break shardLoop + default: + + workQueue <- struct{}{} + v.holder.Logger.Debugf("open index/field/view/fragment: %s/%s/%s/%d", v.index, v.field, v.name, shard) + eg.Go(func() error { + defer func() { + <-workQueue + }() + frag := v.newFragment(v.fragmentPath(shard), shard) + if err := frag.Open(); err != nil { + return fmt.Errorf("open fragment: shard=%d, err=%s", frag.shard, err) + } + frag.RowAttrStore = v.rowAttrStore + v.holder.Logger.Debugf("add index/field/view/fragment to view.fragments: %s/%s/%s/%d", v.index, v.field, v.name, shard) + mu.Lock() + v.fragments[frag.shard] = frag + v.addKnownShard(frag.shard) + mu.Unlock() + return nil + }) } - frag.RowAttrStore = v.rowAttrStore - v.holder.Logger.Debugf("add index/field/view/fragment to view.fragments: %s/%s/%s/%d", v.index, v.field, v.name, shard) - v.fragments[frag.shard] = frag - v.addKnownShard(frag.shard) } - return nil + return eg.Wait() } // close closes the view and its fragments. From 71eccd121decefb884d08e87be042375102f4139 Mon Sep 17 00:00:00 2001 From: "Jason E. Aten" Date: Tue, 28 Jul 2020 07:50:22 -0400 Subject: [PATCH 2/2] fix race in view.openFragmentsInTx --- view.go | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/view.go b/view.go index 682b3fb00..da61f1837 100644 --- a/view.go +++ b/view.go @@ -175,19 +175,32 @@ func (v *view) openFragmentsInTx() error { eg, ctx := errgroup.WithContext(context.Background()) var mu sync.Mutex + shardCh := make(chan uint64, len(shards)) + for i := range shards { + shardCh <- shards[i] + } + shardLoop: - for _, shard := range shards { + for range shards { select { case <-ctx.Done(): break shardLoop default: workQueue <- struct{}{} - v.holder.Logger.Debugf("open index/field/view/fragment: %s/%s/%s/%d", v.index, v.field, v.name, shard) eg.Go(func() error { defer func() { <-workQueue }() + + var shard uint64 + select { + case shard = <-shardCh: + default: + return nil // no more work + } + v.holder.Logger.Debugf("open index/field/view/fragment: %s/%s/%s/%d", v.index, v.field, v.name, shard) + frag := v.newFragment(v.fragmentPath(shard), shard) if err := frag.Open(); err != nil { return fmt.Errorf("open fragment: shard=%d, err=%s", frag.shard, err)