diff --git a/bolt.go b/bolt.go index f2f829bde..70d05b179 100644 --- a/bolt.go +++ b/bolt.go @@ -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 diff --git a/dbshard.go b/dbshard.go index 6fc057200..ae69c8867 100644 --- a/dbshard.go +++ b/dbshard.go @@ -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" diff --git a/fragment.go b/fragment.go index 734d8157d..790b98d7a 100644 --- a/fragment.go +++ b/fragment.go @@ -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{ diff --git a/holder.go b/holder.go index 6da35f797..e02734be2 100644 --- a/holder.go +++ b/holder.go @@ -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 } diff --git a/server.go b/server.go index d543e5088..81cfeffd7 100644 --- a/server.go +++ b/server.go @@ -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)} diff --git a/txfactory.go b/txfactory.go index c0ba21f0d..17ef34c75 100644 --- a/txfactory.go +++ b/txfactory.go @@ -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, diff --git a/txfactory_internal_test.go b/txfactory_internal_test.go index 15b03b59c..26416f95a 100644 --- a/txfactory_internal_test.go +++ b/txfactory_internal_test.go @@ -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 // =============================