Merge pull request #1042 from jaten-molecula/migration_logging

better migration logging
This commit is contained in:
tgruben 2020-10-29 04:38:53 -05:00 • committed by GitHub
commit ed6f59a827
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
7 changed files with 112 additions and 9 deletions

22
bolt.go
View file

@ -145,6 +145,28 @@ func (r *boltRegistrar) OpenDBWrapper(path0 string, doAllocZero bool, rbfcfg *rb
return nil, errors.Wrapf(err, fmt.Sprintf("open bolt path '%v'", path))
}
// docs on fsync from https://godoc.org/github.com/etcd-io/bbolt
//
// Setting the NoSync flag will cause the database to skip fsync()
// calls after each commit. This can be useful when bulk loading data
// into a database and you can restart the bulk load in the event of
// a system failure or database corruption. Do not set this flag for
// normal use.
//
// If the package global IgnoreNoSync constant is true, this value is
// ignored. See the comment on that constant for more details.
//
// THIS IS UNSAFE. PLEASE USE WITH CAUTION.
// NoSync bool
// When true, skips syncing freelist to disk. This improves the database
// write performance under normal operation, but requires a full database
// re-sync during recovery.
// NoFreelistSync bool
//db.NoSync = true
//db.NoFreelistSync = true
err = db.Update(func(tx *bolt.Tx) (err error) {
_, err = tx.CreateBucketIfNotExists(bucketCT)
return

View file

@ -303,6 +303,10 @@ func newShardSet() *shardSet {
func (per *DBPerShard) HasData(which int) (hasData bool, err error) {
// has to aggregate across all available DBShard for each index and shard.
if per.types[which] == roaringTxn {
return per.RoaringHasData()
}
for _, v := range per.Flatmap {
hasData, err = v.W[which].HasData()
if err != nil {
@ -315,6 +319,21 @@ func (per *DBPerShard) HasData(which int) (hasData bool, err error) {
return
}
func (per *DBPerShard) RoaringHasData() (bool, error) {
idxs := per.holder.Indexes()
const requireData = true
for _, idx := range idxs {
shards, err := per.TypedDBPerShardGetShardsForIndex(roaringTxn, idx, "", requireData)
if err != nil {
return false, err
}
if len(shards) > 0 {
return true, nil
}
}
return false, nil
}
func (per *DBPerShard) ListOpenString() (r string) {
for _, v := range per.Flatmap {
r += v.HolderPath + " -> " + v.W[per.useOpenList].OpenListString() + "\n"

View file

@ -2064,7 +2064,7 @@ func (f *fragment) Blocks() ([]FragmentBlock, error) {
// Cache checksum.
chksum := h.Sum()
f.checksums[h.blockID] = chksum
f.checksums[h.blockID] = chksum // the only place checksums is added to.
// Append block.
a = append(a, FragmentBlock{

View file

@ -203,7 +203,8 @@ type HolderConfig struct {
Txsrc string
RowcacheOff bool
RBFConfig *rbfcfg.Config
RBFConfig *rbfcfg.Config
AntiEntropyInterval time.Duration
}
func DefaultHolderConfig() *HolderConfig {
@ -651,8 +652,6 @@ func (h *Holder) Open() error {
return errors.Wrap(err, "processing foreign index fields")
}
h.Logger.Printf("open holder: complete")
h.Stats.Open()
h.opened.Close()
@ -669,6 +668,8 @@ func (h *Holder) Open() error {
}
h.txf.blueGreenOnIfRunningBlueGreen()
h.Logger.Printf("open holder: complete")
return nil
}

View file

@ -401,6 +401,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
return nil, errors.Wrap(err, "applying option")
}
}
s.holderConfig.AntiEntropyInterval = s.antiEntropyInterval
// set up executor after server opts have been processed
executorOpts := []executorOption{optExecutorInternalQueryClient(s.defaultClient)}

View file

@ -25,6 +25,7 @@ import (
"sync"
"syscall"
"text/tabwriter"
"time"
"github.com/pilosa/pilosa/v2/hash"
"github.com/pilosa/pilosa/v2/rbf"
@ -1279,6 +1280,17 @@ func (f *TxFactory) blueHasData() (hasData bool, err error) {
return f.dbPerShard.HasData(0)
}
func (f *TxFactory) greenHasData() (hasData bool, err error) {
n := len(f.types)
switch n {
case 1:
return f.dbPerShard.HasData(0)
case 2:
return f.dbPerShard.HasData(1)
}
panic(fmt.Sprintf("unsupported len(f.types): %v. Must be 1 or 2.", n))
}
// green2blue is called at the very end of Holder.Open(), so
// we know that the holder is ready to go, knowing its holder.Indexes(), fields,
// view, shards, and other metadata if any.
@ -1297,20 +1309,42 @@ func (f *TxFactory) green2blue(holder *Holder) (err error) {
blueDest := f.types[0]
greenSrc := f.types[1]
if blueDest == roaringTxn {
return fmt.Errorf("error: cannot migrate to 'roaring': not implemented.")
}
idxs := holder.Indexes()
verifyInsteadOfCopy := false
hasData, err := f.blueHasData()
blueHasData, err := f.blueHasData()
if err != nil {
return errors.Wrap(err, "TxFactory.green2blue DataSize(0)")
return errors.Wrap(err, "TxFactory.green2blue f.blueHasData()")
}
if hasData {
greenHasData, err := f.greenHasData()
if err != nil {
return errors.Wrap(err, "TxFactory.green2blue f.greenHasData()")
}
if !blueHasData && !greenHasData {
holder.Logger.Printf("no data in blue or green. No migration or verification to do.")
return nil
}
// INVAR: blue has data.
if !greenHasData {
holder.Logger.Printf("error: cannot migrate from green '%v' because it has no data in it.", greenSrc)
return fmt.Errorf("error: cannot migrate from green '%v' because it has no data in it.", greenSrc)
}
if blueHasData {
verifyInsteadOfCopy = true
} else {
holder.Logger.Printf("bitmap-backend migration starting: populating %v from %v", blueDest, greenSrc)
defer holder.Logger.Printf("bitmap-backend migration done : populated %v from %v", blueDest, greenSrc)
}
for _, idx := range idxs {
for k, idx := range idxs {
// scan directories
blueShards, err := f.dbPerShard.TypedDBPerShardGetShardsForIndex(blueDest, idx, "", false)
@ -1342,6 +1376,8 @@ func (f *TxFactory) green2blue(holder *Holder) (err error) {
}
}
lastProgress := time.Now()
progressCount := 0
for shard := range greenShards {
dbs, err := f.dbPerShard.GetDBShard(idx.name, shard, idx)
@ -1360,6 +1396,12 @@ func (f *TxFactory) green2blue(holder *Holder) (err error) {
}
} else {
// the main copy work
progressCount++
if progressCount == 1 || time.Since(lastProgress) > time.Second {
holder.Logger.Printf("migration progress on index '%v' (%v of %v): on shard %v of %v",
idx.name, k+1, len(idxs), progressCount, len(greenShards))
lastProgress = time.Now()
}
err = dbs.populateBlueFromGreen()
if err != nil {
return errors.Wrap(err,

View file

@ -118,11 +118,18 @@ func Test_TxFactory_UpdateBlueFromGreen_OnStartup(t *testing.T) {
checked := []string{"lmdb", "roaring", "rbf"}
expectError := false
for _, blue := range checked {
for _, green := range checked {
if blue == green {
continue
}
if blue == "roaring" {
// not supported
expectError = true
} else {
expectError = false
}
blue_green := blue + "_" + green
//vv("setting blue_green to '%v'", blue_green)
@ -230,7 +237,14 @@ func Test_TxFactory_UpdateBlueFromGreen_OnStartup(t *testing.T) {
h4 := NewHolder(path, nil)
//vv("about to h4.Open we should populate blue from green")
panicOn(h4.Open())
err = h4.Open()
if expectError {
if err == nil {
panic("expected error since migration to roaring not supported")
}
} else {
panicOn(err)
}
testMustHaveBit(t, h4, "i0", "f", rowID, colID)
testMustHaveBit(t, h4, "i1", "f", 100, 200)
@ -258,6 +272,10 @@ func Test_TxFactory_verifyBlueEqualsGreen(t *testing.T) {
if blue == green {
continue
}
if blue == "roaring" {
// not supported
continue
}
blue_green := blue + "_" + green
// =============================