Merge pull request #615 from molecula/parallelized_open_frag

Parallelize view.OpenFragmentsInTx
This commit is contained in:
jaten-molecula 2020-07-28 11:18:13 -04:00 committed by GitHub
commit 4cdf62ab89
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23

54
view.go
View file

@ -171,17 +171,51 @@ func (v *view) openFragmentsInTx() error {
if err != nil {
return errors.Wrap(err, "SliceOfShards")
}
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)
}
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)
eg, ctx := errgroup.WithContext(context.Background())
var mu sync.Mutex
shardCh := make(chan uint64, len(shards))
for i := range shards {
shardCh <- shards[i]
}
return nil
shardLoop:
for range shards {
select {
case <-ctx.Done():
break shardLoop
default:
workQueue <- struct{}{}
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)
}
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
})
}
}
return eg.Wait()
}
// close closes the view and its fragments.