diff --git a/field.go b/field.go index 341695816..b6047e3a8 100644 --- a/field.go +++ b/field.go @@ -91,8 +91,7 @@ type Field struct { logger logger.Logger - snapshotQueue chan *fragment - + snapshotQueue snapshotQueue // Instantiates new translation store on open. OpenTranslateStore OpenTranslateStoreFunc } @@ -924,7 +923,9 @@ func (f *Field) newView(path, name string) *view { view.rowAttrStore = f.rowAttrStore view.stats = f.Stats view.broadcaster = f.broadcaster - view.snapshotQueue = f.snapshotQueue + if f.snapshotQueue != nil { + view.snapshotQueue = f.snapshotQueue + } return view } diff --git a/fragment.go b/fragment.go index 852b11389..794a5b43b 100644 --- a/fragment.go +++ b/fragment.go @@ -107,20 +107,16 @@ type fragment struct { shard uint64 // File-backed storage - path string - flags byte // user-defined flags passed to roaring - gen generation - storage *roaring.Bitmap - totalOpN int64 // total opN values - totalOps int64 // total ops (across all snapshots) - opN int // number of ops since snapshot (may be approximate for imports) - ops int // number of higher-level operations, as opposed to bit changes - snapshotsRequested int // number of times we've requested a snapshot - snapshotsTaken int // number of actual snapshot operations - snapshotting bool // set to true when requesting a snapshot, set to false after snapshot completes - snapshotCond sync.Cond - snapshotDelays int - snapshotDelayTime time.Duration + path string + flags byte // user-defined flags passed to roaring + gen generation + storage *roaring.Bitmap + opN int // number of ops since snapshot (may be approximate for imports) + ops int // number of higher-level operations, as opposed to bit changes + snapshotPending bool // set to true when requesting a snapshot, set to false after snapshot completes + snapshotCond sync.Cond + snapshotErr error // error yielded by the last snapshot operation + snapshotStamp time.Time // timestamp of last snapshot // Cache for row counts. CacheType string // passed in by field @@ -154,7 +150,7 @@ type fragment struct { stats stats.StatsClient - snapshotQueue chan *fragment + snapshotQueue snapshotQueue } // newFragment returns a new instance of Fragment. @@ -172,7 +168,8 @@ func newFragment(path, index, field, view string, shard uint64, flags byte) *fra Logger: logger.NopLogger, MaxOpN: defaultFragmentMaxOpN, - stats: stats.NopStatsClient, + stats: stats.NopStatsClient, + snapshotQueue: defaultSnapshotQueue, } f.snapshotCond = sync.Cond{L: &f.mu} return f @@ -181,62 +178,6 @@ func newFragment(path, index, field, view string, shard uint64, flags byte) *fra // cachePath returns the path to the fragment's cache data. func (f *fragment) cachePath() string { return f.path + cacheExt } -// newSnapshotQueue makes a new snapshot queue, of depth N, and spawns a -// goroutine for it. -func newSnapshotQueue(n int, w int, l logger.Logger) chan *fragment { - ch := make(chan *fragment, n) - for i := 0; i < w; i++ { - go snapshotQueueWorker(ch, l) - } - return ch -} - -func snapshotQueueWorker(snapshotQueue chan *fragment, l logger.Logger) { - for f := range snapshotQueue { - err := f.protectedSnapshot(true) - if err != nil { - l.Printf("snapshot error: %v", err) - } - f.snapshotCond.Broadcast() - } -} - -// enqueueSnapshot requests that the fragment be snapshotted at some point -// in the future, if this has not already been requested. Call this only when -// the mutex is held. -func (f *fragment) enqueueSnapshot() { - f.snapshotsRequested++ - if f.snapshotting { - return - } - f.snapshotting = true - if f.snapshotQueue != nil { - select { - case f.snapshotQueue <- f: - default: - before := time.Now() - // wait forever, but notice that we're waiting - f.snapshotQueue <- f - f.snapshotDelays++ - f.snapshotDelayTime += time.Since(before) - if f.snapshotDelays >= 10 { - f.Logger.Printf("snapshotting %s: last ten enqueue delays took %v", f.path, f.snapshotDelayTime) - f.snapshotDelays = 0 - f.snapshotDelayTime = 0 - } - } - } else { - // in testing, for instance, there may be no holder, thus no one - // to handle these snapshots. - err := f.snapshot() - if err != nil { - f.Logger.Printf("snapshot failed: %v", err) - } - f.snapshotting = false - f.snapshotCond.Broadcast() - } -} - // Open opens the underlying storage. func (f *fragment) Open() error { f.mu.Lock() @@ -445,31 +386,12 @@ func (f *fragment) openCache() error { func (f *fragment) Close() error { f.mu.Lock() defer f.mu.Unlock() - for f.snapshotting { + for f.snapshotPending { f.snapshotCond.Wait() } return f.close() } -// awaitSnapshot lets us delay until the snapshot gets written, preventing tests -// from misleadingly showing amazingly fast performance because the snapshots they -// trigger haven't happened yet. -func (f *fragment) awaitSnapshot() { - f.mu.Lock() - defer f.mu.Unlock() - for f.snapshotting { - f.snapshotCond.Wait() - } -} - -// unprotectedAwaitSnapshot assumes you already hold the lock, and waits for -// the snapshot fairy to come along. -func (f *fragment) unprotectedAwaitSnapshot() { - for f.snapshotting { - f.snapshotCond.Wait() - } -} - func (f *fragment) close() error { // Flush cache if closing gracefully. if err := f.flushCache(); err != nil { @@ -727,7 +649,7 @@ func (f *fragment) unprotectedSetRow(row *Row, rowID uint64) (changed bool, err f.rowCache.Add(rowID, nil) // Snapshot storage. - f.enqueueSnapshot() + f.snapshotQueue.Enqueue(f) f.stats.Count("setRow", 1, 1.0) return changed, nil @@ -768,7 +690,7 @@ func (f *fragment) unprotectedClearRow(rowID uint64) (changed bool, err error) { f.rowCache.Add(rowID, nil) // Snapshot storage. - f.enqueueSnapshot() + f.snapshotQueue.Enqueue(f) f.stats.Count("clearRow", 1, 1.0) @@ -2126,19 +2048,20 @@ func (f *fragment) importValue(columnIDs []uint64, values []int64, bitDepth uint _ = f.openStorage(true) return err } - // We don't actually care, except we want our stats to be accurate. - f.incrementOpN(totalChanges) + // Keep stats accurate. We don't call incrementOpN here because it may + // or may not enqueue a request, which would then be in the queue + // taking up space and otherwise being a possible nuisance, when we're + // about to force a snapshot anyway. + f.opN += totalChanges + f.ops++ // Reset the rowCache. f.rowCache = &simpleCache{make(map[uint64]*Row)} - // in theory, this should probably have happened anyway, but if enough + // in theory, this should probably have been queued anyway, but if enough // of the bits matched existing bits, we'll be under our opN estimate, and // we want to ensure that the snapshot happens. - f.enqueueSnapshot() - f.unprotectedAwaitSnapshot() - - return nil + return f.snapshotQueue.Immediate(f) } // importRoaring imports from the official roaring data format defined at @@ -2198,14 +2121,14 @@ func (f *fragment) incrementOpN(changed int) { f.opN += changed f.ops++ if f.opN > f.MaxOpN { - f.enqueueSnapshot() + f.snapshotQueue.Enqueue(f) } } // Snapshot writes the storage bitmap to disk and reopens it. This may // coexist with existing background-queue snapshotting; it does not remove // things from the queue. You probably don't want to do this; use -// enqueueSnapshot/awaitSnapshot. +// the snapshotQueue's Enqueue/Await. func (f *fragment) Snapshot() error { f.mu.Lock() defer f.mu.Unlock() @@ -2218,25 +2141,13 @@ func track(start time.Time, message string, stats stats.StatsClient, logger logg stats.Histogram("snapshot", elapsed.Seconds(), 1.0) } -// protectedSnapshot grabs the lock and unconditionally calls snapshot(). If -// fromQueue is true, the snapshotting state is also cleared. -func (f *fragment) protectedSnapshot(fromQueue bool) error { - f.mu.Lock() - defer f.mu.Unlock() - err := f.snapshot() - if fromQueue { - f.snapshotting = false - } - return err -} - // snapshot does the actual snapshot operation. it does not check or care -// about f.snapshotting. +// about f.snapshotPending. func (f *fragment) snapshot() error { - f.totalOpN += int64(f.opN) - f.totalOps += int64(f.ops) - f.snapshotsTaken++ _, err := unprotectedWriteToFragment(f, f.storage) + if err == nil { + f.snapshotStamp = time.Now() + } return err } diff --git a/fragment_internal_test.go b/fragment_internal_test.go index 13d47181b..229021907 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -2070,7 +2070,11 @@ func BenchmarkImportRoaring(b *testing.B) { b.StartTimer() err := f.importRoaringT(data, false) if err != nil { - f.awaitSnapshot() + // we don't actually particularly + // care whether this succeeds, + // but if it's happening we want + // it to be done. + _ = f.snapshotQueue.Await(f) f.Clean(b) b.Fatalf("import error: %v", err) } @@ -2109,7 +2113,9 @@ func BenchmarkImportRoaringConcurrent(b *testing.B) { j := j eg.Go(func() error { err := frags[j].importRoaringT(data[j], false) - frags[j].awaitSnapshot() + // error unimportant if it happened, but we want + // any snapshots to have finished. + _ = frags[j].snapshotQueue.Await(frags[j]) return err }) } @@ -2147,11 +2153,13 @@ func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) { // is excessive. force storage into snapshotted state, then use import // to generate an op log and/or snapshot. _, _, err := frags[j].storage.ImportRoaringBits(data, false, false, 0) - frags[j].enqueueSnapshot() - frags[j].awaitSnapshot() if err != nil { b.Fatalf("importing roaring: %v", err) } + err = frags[j].snapshotQueue.Immediate(frags[j]) + if err != nil { + b.Fatalf("snapshot after import: %v", err) + } } eg := errgroup.Group{} b.StartTimer() @@ -2159,7 +2167,10 @@ func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) { j := j eg.Go(func() error { err := frags[j].importRoaringT(updata, false) - frags[j].awaitSnapshot() + err2 := frags[j].snapshotQueue.Await(frags[j]) + if err == nil { + err = err2 + } return err }) } @@ -2221,18 +2232,23 @@ func BenchmarkImportRoaringUpdate(b *testing.B) { // is excessive. force storage into snapshotted state, then use import // to generate an op log and/or snapshot. _, _, err := f.storage.ImportRoaringBits(data, false, false, 0) - f.enqueueSnapshot() - f.awaitSnapshot() if err != nil { b.Errorf("import error: %v", err) } + err = f.snapshotQueue.Immediate(f) + if err != nil { + b.Errorf("snapshot after import error: %v", err) + } b.StartTimer() err = f.importRoaringT(updata, false) - f.awaitSnapshot() if err != nil { f.Clean(b) b.Errorf("import error: %v", err) } + err = f.snapshotQueue.Await(f) + if err != nil { + b.Errorf("snapshot after import error: %v", err) + } b.StopTimer() var stat os.FileInfo var statTarget io.Writer @@ -2542,7 +2558,12 @@ func (f *fragment) sanityCheck(t testing.TB) { } func (f *fragment) Clean(t testing.TB) { - f.awaitSnapshot() + f.mu.Lock() + err := f.snapshotQueue.Await(f) + f.mu.Unlock() + if err != nil { + t.Fatalf("snapshot failed before sanity check: %v", err) + } f.sanityCheck(t) if f.storage != nil && f.storage.Source != nil { if f.storage.Source.Dead() { @@ -2558,7 +2579,7 @@ func (f *fragment) Clean(t testing.TB) { t.Fatal("cleaning up fragment: ", errc, errf, errp) } if f.snapshotQueue != nil { - close(f.snapshotQueue) + f.snapshotQueue.Stop() f.snapshotQueue = nil } // not all fragments have cache files @@ -2581,7 +2602,7 @@ func (f *fragment) CleanKeep(t testing.TB) { t.Fatal("closing fragment: ", errc, errp) } if f.snapshotQueue != nil { - close(f.snapshotQueue) + f.snapshotQueue.Stop() f.snapshotQueue = nil } // not all fragments have cache files @@ -3069,10 +3090,10 @@ func TestFragmentRowIterator(t *testing.T) { func TestUnionInPlaceMapped(t *testing.T) { f := mustOpenFragment("i", "f", "v", 0, CacheTypeNone) + // note: clean has to be deferred first, because it has to run with + // the lock *not* held, because it is sometimes so it has to grab the + // lock... defer f.Clean(t) - // I know this doesn't actually matter in our current context, but - // strictly speaking, we do say you have to hold the lock while calling - // unprotectedWriteToFragment... f.mu.Lock() defer f.mu.Unlock() r0 := rand.New(rand.NewSource(2)) @@ -3105,8 +3126,15 @@ func TestUnionInPlaceMapped(t *testing.T) { f.storage.UnionInPlace(setBM1) countUnion := f.storage.Count() // UnionInPlace produces no ops log, we have to make it snapshot, to - // ensure that the on-disk representation is correct. - f.enqueueSnapshot() + // ensure that the on-disk representation is correct. Note, UIP is + // not used for things that are modifying real fragments, usually; + // it's used only in computation of things that usually don't go to + // disk, which is why we handle this specially in testing and not + // generically. + err = f.snapshotQueue.Immediate(f) + if err != nil { + t.Fatalf("snapshot after union-in-place: %v", err) + } if count0 != countF { t.Fatalf("writing bitmap to storage changed count: %d => %d", count0, countF) diff --git a/holder.go b/holder.go index d54fae43e..5045bff4d 100644 --- a/holder.go +++ b/holder.go @@ -75,7 +75,7 @@ type Holder struct { Logger logger.Logger - snapshotQueue chan *fragment + snapshotQueue snapshotQueue // Manages replication from the primary node. primaryTranslateNode *Node @@ -167,7 +167,7 @@ func (h *Holder) Open() error { // Run snapshots asynchronously. The snapshotQueue will have a background // task associated with it which flushes it and waits until this channel // is closed, so we should always close this channel when done. - h.snapshotQueue = newSnapshotQueue(100, 2, h.Logger) + h.snapshotQueue = newSnapshotQueue(10, 2, h.Logger) for _, fi := range fis { // Skip files or hidden directories. @@ -203,6 +203,7 @@ func (h *Holder) Open() error { go func() { defer h.wg.Done(); h.monitorCacheFlush() }() h.Stats.Open() + h.snapshotQueue.ScanHolder(h) h.opened.Close() return nil @@ -222,8 +223,7 @@ func (h *Holder) Close() error { } } if h.snapshotQueue != nil { - close(h.snapshotQueue) - // assuming the snapshotQueueWorker has already started, this is safe. + h.snapshotQueue.Stop() h.snapshotQueue = nil } diff --git a/index.go b/index.go index 5c73809c1..c96e3bd5f 100644 --- a/index.go +++ b/index.go @@ -58,7 +58,7 @@ type Index struct { Stats stats.StatsClient logger logger.Logger - snapshotQueue chan *fragment + snapshotQueue snapshotQueue // Used for notifying holder when a field is added. holder *Holder @@ -462,7 +462,9 @@ func (i *Index) newField(path, name string) (*Field, error) { f.Stats = i.Stats f.broadcaster = i.broadcaster f.rowAttrStore = i.newAttrStore(filepath.Join(f.path, ".data")) - f.snapshotQueue = i.snapshotQueue + if i.snapshotQueue != nil { + f.snapshotQueue = i.snapshotQueue + } f.OpenTranslateStore = i.OpenTranslateStore return f, nil } diff --git a/snapshotqueue.go b/snapshotqueue.go new file mode 100644 index 000000000..210c2b772 --- /dev/null +++ b/snapshotqueue.go @@ -0,0 +1,405 @@ +// Copyright 2019 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package pilosa + +import ( + "fmt" + "os" + "sync" + "time" + + "github.com/pilosa/pilosa/v2/logger" + "github.com/pkg/errors" +) + +// snapshotQueue is a thing which can handle enqueuing snapshots. A snapshot +// queue distinguishes between high-priority requests, which get satisfied +// by the next available worker, and regular requests, which get enqueued +// if there's space in the queue, and otherwise dropped. There's also a +// separate background task to scan a holder for fragments which may need +// snapshots, but which is processed only when the queue is empty, and only +// slowly. "Await" awaits an existing snapshot if one is already enqueued. +// "Immediate" tries to do one right away. (If one's already enqueued, this +// can leave it in the queue, which will ignore anything that shows up with +// the request flag cleared.) +// +// Await, Enqueue, and Immediate should be called only with the fragment lock +// held. +// +// ScanHolder spawns a new goroutine. You don't need to use `go` on it. +type snapshotQueue interface { + Immediate(*fragment) error + Enqueue(*fragment) + Await(*fragment) error + ScanHolder(*Holder) + Stop() +} + +// queuelessSnapshotQueue isn't a snapshot queue, but it satisfies the +// interface. +type queuelessSnapshotQueue struct{} + +func (q *queuelessSnapshotQueue) Enqueue(f *fragment) { + _ = f.snapshot() +} + +func (q *queuelessSnapshotQueue) Await(f *fragment) error { + return nil +} + +func (q *queuelessSnapshotQueue) Immediate(f *fragment) error { + return f.snapshot() +} + +func (q *queuelessSnapshotQueue) ScanHolder(h *Holder) { +} + +func (q *queuelessSnapshotQueue) Stop() { +} + +// defaultSnapshotQueue is the fallback to use if none is available, +// and currently uses queueless -- it runs all snapshots immediately. +var defaultSnapshotQueue *queuelessSnapshotQueue + +// newSnapshotQueue makes a new snapshot queue, of depth N, with +// w worker threads. +func newSnapshotQueue(n int, w int, l logger.Logger) snapshotQueue { + sq := prioritySnapshotQueue{normal: make(chan snapshotRequest, n), urgent: make(chan snapshotRequest), background: make(chan snapshotRequest), done: make(chan struct{}), logger: l} + if sq.logger == nil { + sq.logger = logger.NewStandardLogger(os.Stderr) + } + sq.spawnWorkers(w) + return &sq +} + +type snapshotRequest struct { + frag *fragment + when time.Time +} + +// prioritySnapshotQueue gives preference to "immediate" requests, and +// dispreference to "background" requests from ScanHolder. It timestamps +// requests, so it can discard a request if the most recent snapshot is +// newer than the request. The snapshotPending flag in the fragment is +// used to track that a given fragment thinks it has been successfully +// enqueued. Background requests are not considered enqueued, since +// they'll never get processed if there's anything else. In normal workloads, +// immediate/urgent snapshots should be rare, but we'll happily drop +// most requests on the floor; the scanner should pick them up once things +// are quiet. +type prioritySnapshotQueue struct { + logger logger.Logger + urgent chan snapshotRequest + normal chan snapshotRequest + background chan snapshotRequest + done chan struct{} + mu sync.RWMutex + scanWG, workerWG sync.WaitGroup + stats struct { + enqueued int64 + skipped int64 + } +} + +func (sq *prioritySnapshotQueue) spawnWorkers(w int) { + sq.mu.Lock() + defer sq.mu.Unlock() + if sq.done == nil { + sq.logger.Printf("prioritySnapshotQueue worker: no done channel, already done?") + return + } + sq.workerWG.Add(w) + for i := 0; i < w; i++ { + go sq.worker(sq.urgent, sq.normal, sq.background, sq.done) + } +} + +func (sq *prioritySnapshotQueue) worker(urgent, normal, background chan snapshotRequest, done chan struct{}) { + // We don't want a race condition on these. If they're non-nil when + // we get them, they should get closed at some point. If done is + // already nil, we shouldn't do anything. + defer sq.workerWG.Done() + ok := true + var req snapshotRequest + for ok { + req.frag = nil + + select { + case req, ok = <-urgent: + default: + select { + case req, ok = <-urgent: + case req, ok = <-normal: + default: + select { + case req, ok = <-urgent: + case req, ok = <-normal: + case req, ok = <-background: + case _, ok = <-done: + } + } + } + if req.frag != nil { + sq.process(req) + } + } +} + +// process actually runs a fragment. it will do this if either the fragment +// has a pending snapshot, or the force flag is set. +func (sq *prioritySnapshotQueue) process(req snapshotRequest) { + f := req.frag + f.mu.Lock() + defer f.mu.Unlock() + if f.snapshotStamp.Before(req.when) { + f.snapshotErr = f.snapshot() + if f.snapshotErr != nil { + fmt.Printf("snapshot error: %v\n", f.snapshotErr) + sq.logger.Printf("snapshot error: %v", f.snapshotErr) + } + f.snapshotPending = false + f.snapshotCond.Broadcast() + } +} + +// Stop shuts down the snapshot queue. It first marks it as done, causing +// the background scanner(s), if any, to shut down, then waits for them, then +// closes and nils the queues. The background scanner has to get stopped +// because otherwise it might try to write to those closed queues. +func (sq *prioritySnapshotQueue) Stop() { + sq.mu.Lock() + defer sq.mu.Unlock() + close(sq.done) + // scanners need to be done before we close the other channels. + sq.scanWG.Wait() + sq.done = nil + close(sq.normal) + sq.normal = nil + close(sq.urgent) + sq.urgent = nil + close(sq.background) + sq.background = nil + sq.logger.Printf("snapshot queue: enqueued %d, skipped %d\n", sq.stats.enqueued, sq.stats.skipped) +} + +// Enqueue tries to add a fragment to the queue, if the fragment is not already +// enqueued. You should hold a lock on the fragment when calling this. +func (sq *prioritySnapshotQueue) Enqueue(f *fragment) { + if f.snapshotPending { + return + } + sq.mu.Lock() + defer sq.mu.Unlock() + if sq.normal == nil { + sq.logger.Printf("requested snapshot after snapshot queue was closed") + return + } + // we have to set this before enqueing, because it's + // otherwise possible that we're at the head of the queue, + // and the recipient gets the fragment before we execute the + // line after the send. + f.snapshotPending = true + // try to enqueue snapshot + select { + case sq.normal <- snapshotRequest{frag: f, when: time.Now()}: + sq.stats.enqueued++ + return + default: + sq.stats.skipped++ + f.snapshotPending = false + return + } +} + +// Await returns when f is not pending a snapshot. Call with the fragment lock +// held. Await waits on a condition variable inside f, associated with the +// fragment's lock, so this does not conflict with the lock being used for +// snapshots. +func (sq *prioritySnapshotQueue) Await(f *fragment) (err error) { + for f.snapshotPending { + f.snapshotCond.Wait() + } + err, f.snapshotErr = f.snapshotErr, nil + return err +} + +// Immediate forces an immediate snapshot of the given fragment. Call with +// the fragment locked. If the queue is already closing, the fragment does +// not get snapshotted. +func (sq *prioritySnapshotQueue) Immediate(f *fragment) error { + sq.mu.RLock() + // no deferred unlock, because we want to unlock this before calling Await. + // Not because that needs this lock, but because once we're that far, we + // *don't* need this lock anymore so someone else should have it. + if sq.urgent == nil { + sq.mu.RUnlock() + sq.logger.Printf("requested immediate snapshot after snapshot queue was closed") + return errors.New("requested immediate snapshot after snapshot queue was closed") + } + f.snapshotPending = true + req := snapshotRequest{frag: f, when: time.Now()} + // if the fragment was already in the work queue, it's *possible* + // that the only available worker just picked it off the queue, and + // is now waiting on getting the fragment's lock, so it can run + // a snapshot. So we let go of the lock on the fragment, send the + // request, then request the fragment lock again, because Await will + // be sleeping on the condition variable associated with the lock, + // which means it needs to hold the lock so it can let it go during + // the wait... No, really, this made sense. + f.mu.Unlock() + sq.urgent <- req + sq.mu.RUnlock() + f.mu.Lock() + return sq.Await(f) +} + +// needsSnapshot determines whether a fragment probably wants snapshotting. +// Specifically, it looks for fragments not already marked to receive +// snapshots, but which have a high enough opN to justify a snapshot. This +// is only used from the background scan. +func (sq *prioritySnapshotQueue) needsSnapshot(f *fragment) bool { + if f == nil { + return false + } + f.mu.Lock() + defer f.mu.Unlock() + if f.snapshotPending { + return false + } + if f.opN > f.MaxOpN { + return true + } + return false +} + +// ScanHolder spawns a goroutine which iterates through the holder's +// indexes/fields/views/fragments, looking for fragments which have OpN +// high enough to justify a snapshot but don't seem to have one pending. +// It then dumps these in the low priority background queue. +func (sq *prioritySnapshotQueue) ScanHolder(h *Holder) { + sq.mu.Lock() + sq.scanWG.Add(1) + go sq.scanHolderWorker(h, sq.background, sq.done) + sq.mu.Unlock() +} + +// scanHolderWorker is a background task that scans a holder looking for +// fragments which need snapshots taken. It's the cleanup task for snapshots +// that would have been requested by Enqueue, but the queue was full. +func (sq *prioritySnapshotQueue) scanHolderWorker(h *Holder, background chan snapshotRequest, done chan struct{}) { + defer sq.scanWG.Done() + var indexNames, fieldNames, viewNames []string + var fragNums []uint64 + for { + // To avoid abusing things, cap activity rate; every time we finish + // the holder, or every couple hundred fragments considered, we + // pause for a bit. + counter := 0 + hits := 0 + h.mu.Lock() + indexNames = indexNames[:0] + for indexName := range h.indexes { + indexNames = append(indexNames, indexName) + } + h.mu.Unlock() + for _, indexName := range indexNames { + h.mu.Lock() + index := h.indexes[indexName] + h.mu.Unlock() + if index == nil { + continue + } + fieldNames = fieldNames[:0] + index.mu.Lock() + for fieldName := range index.fields { + fieldNames = append(fieldNames, fieldName) + } + index.mu.Unlock() + for _, fieldName := range fieldNames { + index.mu.Lock() + field := index.fields[fieldName] + index.mu.Unlock() + if field == nil { + continue + } + viewNames = viewNames[:0] + field.mu.Lock() + for viewName := range field.viewMap { + viewNames = append(viewNames, viewName) + } + field.mu.Unlock() + for _, viewName := range viewNames { + field.mu.Lock() + view := field.viewMap[viewName] + field.mu.Unlock() + if view == nil { + continue + } + fragNums := fragNums[:0] + view.mu.Lock() + for fragNum := range view.fragments { + fragNums = append(fragNums, fragNum) + } + view.mu.Unlock() + for _, fragNum := range fragNums { + view.mu.Lock() + frag := view.fragments[fragNum] + view.mu.Unlock() + if sq.needsSnapshot(frag) { + hits++ + select { + case background <- snapshotRequest{frag: frag, when: time.Now()}: + sq.logger.Debugf("found fragment needing snapshot: %s\n", frag.path) + case <-done: + return + } + } else { + // Count fragments examined *without* finding anything that + // needed a snapshot. When we find things that need snapshots, + // the time it takes the workers to respond to us is enough + // of a delay to keep us from eating every CPU. So, if a lot + // of things need snapshots, and the workers aren't doing + // anything else, ScanHolder will mostly keep them saturated. + // If they're busy, we'll block forever in the write to the + // background queue. If there's nothing that needs snapshots, + // we pause frequently for a second or so at a time. + counter++ + if counter == 100 { + select { + case <-time.After(1 * time.Second): + case <-done: + return + } + counter = 0 + } + } + } + } + } + } + if hits > 0 { + sq.logger.Printf("background scan: %d fragments needed snapshots\n", hits) + hits = 0 + } else { + sq.logger.Printf("background scan: no fragments needed snapshots, waiting\n") + // No reason to be active if we're not finding anything. + select { + case <-time.After(60 * time.Second): + case <-done: + return + } + } + } +} diff --git a/view.go b/view.go index 69ae1b22f..e1648a838 100644 --- a/view.go +++ b/view.go @@ -59,7 +59,7 @@ type view struct { stats stats.StatsClient rowAttrStore AttrStore logger logger.Logger - snapshotQueue chan *fragment + snapshotQueue snapshotQueue } // newView returns a new instance of View. @@ -309,7 +309,9 @@ func (v *view) newFragment(path string, shard uint64) *fragment { frag.CacheSize = v.cacheSize frag.Logger = v.logger frag.stats = v.stats - frag.snapshotQueue = v.snapshotQueue + if v.snapshotQueue != nil { + frag.snapshotQueue = v.snapshotQueue + } if v.fieldType == FieldTypeMutex { frag.mutexVector = newRowsVector(frag) } else if v.fieldType == FieldTypeBool {