Merge pull request #1024 from jaten-molecula/avoid_excessive_directory_scans

pilosa: avoid re-scanning shards during holder open
This commit is contained in:
jaten-molecula 2020-10-23 20:31:05 -05:00 committed by GitHub
commit 14cb29ea7d
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
6 changed files with 185 additions and 73 deletions

View file

@ -228,6 +228,69 @@ type DBPerShard struct {
// roaring doesn't keep a list of open Tx sn.
// or default to the 2nd.
useOpenList int
// cache the shards per index to avoid excessive
// directory scans of the index directory. Keep per
// txtype to allow blue-green migrate open to be fast too.
// Keep it up-to-date as we add shards to avoid doing
// a filesystem rescan on new shard creation.
//
// txtype -> index -> *shardSet
index2shards map[txtype]map[string]*shardSet
isBlueGreen bool
}
func newIndex2Shards() (r map[txtype]map[string]*shardSet) {
r = make(map[txtype]map[string]*shardSet)
return
}
type shardSet struct {
shards map[uint64]bool
shardsVer int64 // increment with each change.
// give out readonly to repeated consumers if
// readonlyVer == shardsVer
readonly map[uint64]bool
readonlyVer int64
}
func (ss *shardSet) add(shard uint64) {
_, already := ss.shards[shard]
if !already {
ss.shards[shard] = true
ss.shardsVer++
}
}
// CloneMaybe maintains a re-usable readonly version
// ss.shards that can be returned to multiple goroutine
// reads as it will never change. A copy is only made
// once for each change in the shard set.
func (ss *shardSet) CloneMaybe() map[uint64]bool {
if ss.readonlyVer == ss.shardsVer {
return ss.readonly
}
// readonlyVer is out of date.
// readonly needs update. We cannot
// modify the readonly map in place;
// must make a fully new copy here.
ss.readonly = make(map[uint64]bool)
for k, v := range ss.shards {
ss.readonly[k] = v
}
ss.readonlyVer = ss.shardsVer
return ss.readonly
}
func newShardSet() *shardSet {
return &shardSet{
shards: make(map[uint64]bool),
}
}
// HasData returns true if the database has at least one key.
@ -261,11 +324,11 @@ func (per *DBPerShard) LoadExistingDBs() (err error) {
for _, idx := range idxs {
sos, err := per.txf.GetShardsForIndex(idx, "", true)
shardset, err := per.txf.GetShardsForIndex(idx, "", true)
if err != nil {
return err
}
for _, shard := range sos {
for shard := range shardset {
_, err := per.GetDBShard(idx.name, shard, idx)
if err != nil {
return errors.Wrap(err, "DBPerShard.LoadExistingDBs GetDBShard()")
@ -294,14 +357,16 @@ func (txf *TxFactory) NewDBPerShard(types []txtype, holderDir string, holder *Ho
}
d = &DBPerShard{
types: types,
HolderDir: holderDir,
holder: holder,
dbh: NewDBHolder(),
Flatmap: make(map[flatkey]*DBShard),
txf: txf,
useOpenList: useOpenList,
hasRoaring: hasRoaring,
types: types,
HolderDir: holderDir,
holder: holder,
dbh: NewDBHolder(),
Flatmap: make(map[flatkey]*DBShard),
txf: txf,
useOpenList: useOpenList,
hasRoaring: hasRoaring,
isBlueGreen: len(types) > 1,
index2shards: newIndex2Shards(),
}
return
}
@ -439,6 +504,31 @@ func (per *DBPerShard) prefixForType(idx *Index, ty txtype) string {
var ErrNoData = fmt.Errorf("no data")
// keep our cache of shards up-to-date in memory; after the initial
// directory scan, this is all we should we need. Prevents us from
// doing additional, expensive, directory scans.
//
// Caller must hold per.Mu.Lock() already.
func (per *DBPerShard) updateIndex2ShardCacheWithNewShard(dbs *DBShard) {
for _, ty := range dbs.types {
mapIndex2shardSet, ok := per.index2shards[ty]
if !ok {
mapIndex2shardSet = make(map[string]*shardSet)
per.index2shards[ty] = mapIndex2shardSet
}
// INVAR: mapIndex2shardSet is good, but may be an empty map
shardset, ok := mapIndex2shardSet[dbs.Index]
if !ok {
shardset = newShardSet()
mapIndex2shardSet[dbs.Index] = shardset
}
// INVAR: shardset is present, not nil; a map that can be added to.
shardset.add(dbs.Shard)
}
}
func (per *DBPerShard) GetDBShard(index string, shard uint64, idx *Index) (dbs *DBShard, err error) {
per.Mu.Lock()
@ -478,6 +568,7 @@ func (per *DBPerShard) GetDBShard(index string, shard uint64, idx *Index) (dbs *
}
dbi.Shard[shard] = dbs
per.updateIndex2ShardCacheWithNewShard(dbs)
}
if !dbs.Open {
var registry DBRegistry
@ -541,22 +632,22 @@ func (per *DBPerShard) Close() (err error) {
// DBPerShardGetShardsForIndex returns the shards for idx.
// If requireData, we open the database and see that it has a key, rather
// than assume that the database file presence is enough.
func (f *TxFactory) GetShardsForIndex(idx *Index, roaringViewPath string, requireData bool) (sliceOfShards []uint64, err error) {
func (f *TxFactory) GetShardsForIndex(idx *Index, roaringViewPath string, requireData bool) (map[uint64]bool, error) {
var shards [][]uint64
for _, ty := range f.types {
var slc []uint64
slc, err = f.dbPerShard.TypedDBPerShardGetShardsForIndex(ty, idx, roaringViewPath, requireData)
if err != nil {
return
}
shards = append(shards, slc)
}
n := len(f.types)
if n != 1 && n != 2 {
panic(fmt.Sprintf("internal error. only green or blue/green supported. we see types len %v", n))
}
var shards []map[uint64]bool
for _, ty := range f.types {
ss, err := f.dbPerShard.TypedDBPerShardGetShardsForIndex(ty, idx, roaringViewPath, requireData)
if err != nil {
return nil, err
}
shards = append(shards, ss)
}
// Note: we don't actually know when the blue call and when the green call comes
// through here. So if we are deleting a shard, we will see a difference earlier
// in one than the other. TestAPI_ClearFlagForImportAndImportValues for example.
@ -569,7 +660,36 @@ func (f *TxFactory) GetShardsForIndex(idx *Index, roaringViewPath string, requir
// if roaringViewPath is "" then for ty == roaringTxn we go to disk to discover
// all the view paths under idx for type ty.
// requireData means open the database file and verify that at least one key is set.
func (per *DBPerShard) TypedDBPerShardGetShardsForIndex(ty txtype, idx *Index, roaringViewPath string, requireData bool) (sliceOfShards []uint64, err error) {
// The returned sliceOfShards should not be modified. We will cache it for subsequent
// queries.
//
// when a new DBShard is made, we will update the list of shards then. Thus
// the per.index2shard should always be up to date AFTER the first call here.
func (per *DBPerShard) TypedDBPerShardGetShardsForIndex(ty txtype, idx *Index, roaringViewPath string, requireData bool) (shardMap map[uint64]bool, err error) {
// use the cache, always
per.Mu.Lock()
defer per.Mu.Unlock()
i2ss, ok := per.index2shards[ty]
if !ok {
// index -> shardSet
i2ss = make(map[string]*shardSet)
per.index2shards[ty] = i2ss
}
// INVAR: i2ss is good, but may be an empty map
ss, ok := i2ss[idx.name]
if ok {
return ss.CloneMaybe(), nil
}
// INVAR: cache miss, and index2shards[ty] exists.
// gotta read shards from disk directory layout.
setOfShards := newShardSet()
per.index2shards[ty][idx.name] = setOfShards
// Upon return, cache the setOfShards value and reuse it next time
if ty == roaringTxn {
rx := &RoaringTx{
@ -585,13 +705,24 @@ func (per *DBPerShard) TypedDBPerShardGetShardsForIndex(ty txtype, idx *Index, r
errors.Wrap(err, fmt.Sprintf(
"TypedDBPerShardGetLocalShardsForIndex roaringTxn view.path='%v'", view.path))
}
sliceOfShards = append(sliceOfShards, sos...)
for _, shard := range sos {
setOfShards.add(shard)
}
}
}
return dedupShardSlice(sliceOfShards), nil
return setOfShards.CloneMaybe(), nil
}
return rx.SliceOfShards("", "", "", roaringViewPath)
sos, err := rx.SliceOfShards("", "", "", roaringViewPath)
if err != nil {
return nil, err
}
for _, shard := range sos {
setOfShards.add(shard)
}
return setOfShards.CloneMaybe(), nil
}
// INVAR: not-roaring.
requiredSuffix := ty.FileSuffix()
path := per.prefixForType(idx, ty)
@ -631,14 +762,14 @@ func (per *DBPerShard) TypedDBPerShardGetShardsForIndex(ty txtype, idx *Index, r
return nil, err
}
if hasData {
sliceOfShards = append(sliceOfShards, shard)
setOfShards.add(shard)
}
} else {
// file presence is enough
sliceOfShards = append(sliceOfShards, shard)
setOfShards.add(shard)
}
}
return
return setOfShards.CloneMaybe(), nil
}
func (per *DBPerShard) TypedIndexShardHasData(ty txtype, idx *Index, shard uint64) (hasData bool, err error) {
@ -663,10 +794,10 @@ func (per *DBPerShard) TypedIndexShardHasData(ty txtype, idx *Index, shard uint6
}
func listDirUnderDir(root string, includeRoot bool, requiredSuffix string, ignoreEmpty bool) (files []string, err error) {
if !dirExists(root) {
return
}
n := len(root) + 1
if includeRoot {
n = 0
@ -881,14 +1012,3 @@ func (dbs *DBShard) verifyBlueEqualsGreen() (err error) {
return nil
}
func dedupShardSlice(sos []uint64) (r []uint64) {
m := make(map[uint64]struct{})
for _, s := range sos {
m[s] = struct{}{}
}
for k := range m {
r = append(r, k)
}
return
}

View file

@ -88,21 +88,21 @@ func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) {
estd := "rick/_exists/views/standard"
std := "rick/f/views/standard"
sos, err := holder.txf.GetShardsForIndex(idx, tmpdir+sep+std, false)
shards, err := holder.txf.GetShardsForIndex(idx, tmpdir+sep+std, false)
panicOn(err)
for _, shard := range []uint64{93, 223, 221, 215, 219, 217} {
if !inSlice(sos, shard) {
panic(fmt.Sprintf("missing shard=%v from sos='%#v'", shard, sos))
if !shards[shard] {
panic(fmt.Sprintf("missing shard=%v from shards='%#v'", shard, shards))
}
}
if src == "roaring" {
// check estd too
sos, err = holder.txf.GetShardsForIndex(idx, tmpdir+sep+estd, false)
shards, err = holder.txf.GetShardsForIndex(idx, tmpdir+sep+estd, false)
panicOn(err)
for _, shard := range []uint64{93, 223, 221, 215, 219, 217} {
if !inSlice(sos, shard) {
panic(fmt.Sprintf("missing shard=%v from sos='%#v'", shard, sos))
if !shards[shard] {
panic(fmt.Sprintf("missing shard=%v from shards='%#v'", shard, shards))
}
}
}
@ -110,15 +110,6 @@ func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) {
}
}
func inSlice(sos []uint64, shard uint64) bool {
for i := range sos {
if shard == sos[i] {
return true
}
}
return false
}
// data for Test_DBPerShard_GetShardsForIndex
//
var sampleRoaringDirList = map[string]string{"roaring": `

View file

@ -462,7 +462,7 @@ func (s *Server) InternalClient() InternalClient {
// UpAndDown brings the server up minimally and shuts it down
// again; basically, it exists for testing holder open and close.
func (s *Server) UpAndDown() error {
s.logger.Printf("open server")
s.logger.Printf("open server. PID %v", os.Getpid())
// Log startup
err := s.holder.logStartup()
@ -485,7 +485,7 @@ func (s *Server) UpAndDown() error {
// Open opens and initializes the server.
func (s *Server) Open() error {
s.logger.Printf("open server")
s.logger.Printf("open server. PID %v", os.Getpid())
if s.holder.NeedsSnapshot() {
// Start background monitoring.

View file

@ -1305,20 +1305,24 @@ func (f *TxFactory) green2blue(holder *Holder) (err error) {
}
if verifyInsteadOfCopy {
diff := f.shardSliceDiff(blueShards, greenShards)
diff := f.shardSetDiff(blueShards, greenShards)
if diff != "" {
return fmt.Errorf("verifyInsteadOfCopy true, blue[%v]=%#v and green[%v]=%#v have different shards for index '%v': '%v'; stack=\n%v", blueDest, blueShards, greenSrc, greenShards, idx.name, diff, stack())
}
// can also check against meta data
shards := idx.AvailableShards(localOnly).Slice()
diff2 := f.shardSliceDiff(greenShards, shards)
meta := make(map[uint64]bool)
for _, shard := range shards {
meta[shard] = true
}
diff2 := f.shardSetDiff(greenShards, meta)
if diff2 != "" {
return fmt.Errorf("green[%v] = '%#v' and meta data '%#v' have different shards for index '%v': %v", greenSrc, greenShards, shards, idx.name, diff2)
}
}
for _, shard := range greenShards {
for shard := range greenShards {
dbs, err := f.dbPerShard.GetDBShard(idx.name, shard, idx)
if err != nil {
@ -1349,22 +1353,14 @@ func (f *TxFactory) green2blue(holder *Holder) (err error) {
return nil
}
func (f *TxFactory) shardSliceDiff(blueShards, greenShards []uint64) (diff string) {
func (f *TxFactory) shardSetDiff(blueShards, greenShards map[uint64]bool) (diff string) {
nb := len(blueShards)
ng := len(greenShards)
if nb != ng {
diff = fmt.Sprintf("blueShard[%v] count = %v; greenShard[%v] count = %v; ", f.types[0], nb, f.types[1], ng)
}
b := make(map[uint64]bool)
g := make(map[uint64]bool)
for _, bs := range blueShards {
b[bs] = true
}
for _, gs := range greenShards {
g[gs] = true
}
bmg := mapDiff(b, g) // get blue - green
gmb := mapDiff(g, b) // get green - blue
bmg := mapDiff(blueShards, greenShards) // get blue - green
gmb := mapDiff(greenShards, blueShards) // get green - blue
if len(bmg) == 0 && len(gmb) == 0 {
return ""

View file

@ -14,13 +14,17 @@
package pilosa
import "time"
import (
"runtime"
"time"
)
var Version string
var Commit string
var Variant string
var BuildTime string
var LatticeCommit string
var GoVersion string = runtime.Version()
func VersionInfo() string {
var prefix string
@ -48,6 +52,7 @@ func VersionInfo() string {
case buildTime != "":
suffix += " (" + buildTime + ")"
}
suffix += " " + GoVersion
return prefix + "Pilosa" + suffix
}

View file

@ -181,8 +181,8 @@ func (v *view) openFragmentsInTx() error {
var mu sync.Mutex
shardCh := make(chan uint64, len(shards))
for i := range shards {
shardCh <- shards[i]
for shard := range shards {
shardCh <- shard
}
shardLoop: