From ff2cd691cdfacc076f5fc3e967c91a34d2dad46e Mon Sep 17 00:00:00 2001 From: Josh Slocum Date: Fri, 10 Dec 2021 15:27:25 -0600 Subject: [PATCH 1/2] Switching back to GRV for committed version checking, with proper rollback checking --- fdbclient/NativeAPI.actor.cpp | 28 +++--- fdbclient/Notified.h | 2 + fdbserver/BlobWorker.actor.cpp | 161 ++++++++++++++++++++++----------- 3 files changed, 123 insertions(+), 68 deletions(-) diff --git a/fdbclient/NativeAPI.actor.cpp b/fdbclient/NativeAPI.actor.cpp index 6d2642da62..23107e027c 100644 --- a/fdbclient/NativeAPI.actor.cpp +++ b/fdbclient/NativeAPI.actor.cpp @@ -7163,26 +7163,28 @@ ACTOR Future changeFeedWaitLatest(ChangeFeedData* self, Version version) { } ACTOR Future changeFeedWhenAtLatest(ChangeFeedData* self, Version version) { - state Future lastReturned = self->lastReturnedVersion.whenAtLeast(version); if (DEBUG_CF_WAIT_VERSION == version) { fmt::print("CFW {0}) WhenAtLeast: LR={1}\n", version, self->lastReturnedVersion.get()); } + if (version <= self->getVersion()) { + if (DEBUG_CF_WAIT_VERSION == version) { + fmt::print("CFW {0}) WhenAtLeast: Already done\n", version, self->lastReturnedVersion.get()); + } + return Void(); + } + state Future lastReturned = self->lastReturnedVersion.whenAtLeast(version); + loop { if (DEBUG_CF_WAIT_VERSION == version) { fmt::print("CFW {0}) WhenAtLeast: NotAtLatest={1}\n", version, self->notAtLatest.get()); } - if (self->notAtLatest.get() == 0) { - choose { - when(wait(changeFeedWaitLatest(self, version))) { break; } - when(wait(self->refresh.getFuture())) {} - when(wait(self->notAtLatest.onChange())) {} - } - } else { - choose { - when(wait(lastReturned)) { break; } - when(wait(self->notAtLatest.onChange())) {} - when(wait(self->refresh.getFuture())) {} - } + // only allowed to use empty versions if you're caught up + Future waitEmptyVersion = (self->notAtLatest.get() == 0) ? changeFeedWaitLatest(self, version) : Never(); + choose { + when(wait(waitEmptyVersion)) { break; } + when(wait(lastReturned)) { break; } + when(wait(self->refresh.getFuture())) {} + when(wait(self->notAtLatest.onChange())) {} } } diff --git a/fdbclient/Notified.h b/fdbclient/Notified.h index 94a1bb2144..dc18cda123 100644 --- a/fdbclient/Notified.h +++ b/fdbclient/Notified.h @@ -80,6 +80,8 @@ struct Notified { val = std::move(r.val); } + int numWaiting() { return waiting.size(); } + private: using Item = std::pair>; struct ItemCompare { diff --git a/fdbserver/BlobWorker.actor.cpp b/fdbserver/BlobWorker.actor.cpp index d24993db11..aba8332dce 100644 --- a/fdbserver/BlobWorker.actor.cpp +++ b/fdbserver/BlobWorker.actor.cpp @@ -200,6 +200,9 @@ struct BlobWorkerData : NonCopyable, ReferenceCounted { PromiseStream granuleUpdateErrors; + Promise doGRVCheck; + NotifiedVersion grvVersion; + BlobWorkerData(UID id, Database db) : id(id), db(db), stats(id, SERVER_KNOBS->WORKER_LOGGING_INTERVAL) {} bool managerEpochOk(int64_t epoch) { @@ -497,7 +500,7 @@ ACTOR Future writeDeltaFile(Reference bwData, GranuleDeltas deltasToWrite, Version currentDeltaVersion, Future previousDeltaFileFuture, - NotifiedVersion* granuleCommittedVersion, + Future waitCommitted, Optional> oldGranuleComplete) { wait(delay(0, TaskPriority::BlobWorkerUpdateStorage)); @@ -522,9 +525,8 @@ ACTOR Future writeDeltaFile(Reference bwData, state int numIterations = 0; try { // before updating FDB, wait for the delta file version to be committed and previous delta files to finish - if (currentDeltaVersion > granuleCommittedVersion->get()) { - wait(granuleCommittedVersion->whenAtLeast(currentDeltaVersion)); - } + // TODO fix file leak here on error pre-transaction. + wait(waitCommitted); BlobFileIndex prev = wait(previousDeltaFileFuture); wait(delay(0, TaskPriority::BlobWorkerUpdateFDB)); @@ -1139,6 +1141,50 @@ static Version doGranuleRollback(Reference metadata, return cfRollbackVersion; } +ACTOR Future waitVersionCommitted(Reference bwData, + Reference metadata, + Version version) { + // TODO REMOVE debugs + if (version > bwData->grvVersion.get()) { + /*if (BW_DEBUG) { + fmt::print("waitVersionCommitted waiting {0}\n", version); + }*/ + // this order is important, since we need to register a waiter on the notified version before waking the GRV + // actor + Future grvAtLeast = bwData->grvVersion.whenAtLeast(version); + if (bwData->doGRVCheck.canBeSet()) { + bwData->doGRVCheck.send(Void()); + } + wait(grvAtLeast); + } + state Version grvVersion = bwData->grvVersion.get(); + /*if (BW_DEBUG) { + fmt::print("waitVersionCommitted got {0} < {1}, waiting on CF (currently {2})\n", + version, + grvVersion, + metadata->activeCFData.get()->getVersion()); + }*/ + // make sure the change feed has consumed mutations up through grvVersion to ensure none of them are rollbacks + + loop { + state Future atLeast = metadata->activeCFData.get()->whenAtLeast(grvVersion); + choose { + when(wait(atLeast)) { break; } + when(wait(metadata->activeCFData.onChange())) {} + } + } + // sanity check to make sure whenAtLeast didn't return early + if (grvVersion > metadata->waitForVersionReturned) { + metadata->waitForVersionReturned = grvVersion; + } + /*if (BW_DEBUG) { + fmt::print( + "waitVersionCommitted CF whenAtLeast {0}: {1}\n", grvVersion, metadata->activeCFData.get()->getVersion()); + }*/ + + return Void(); +} + // TODO REMOVE once correctness clean #define DEBUG_BW_START_VERSION invalidVersion #define DEBUG_BW_END_VERSION invalidVersion @@ -1159,7 +1205,6 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, state Optional> oldChangeFeedDataComplete; state Key cfKey; state Optional oldCFKey; - state NotifiedVersion committedVersion; state int pendingSnapshots = 0; state std::deque> rollbacksInProgress; @@ -1248,7 +1293,6 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, metadata->durableDeltaVersion.set(startVersion); metadata->pendingDeltaVersion = startVersion; metadata->bufferedDeltaVersion = startVersion; - committedVersion.set(startVersion); metadata->activeCFData.set(newChangeFeedData(startVersion)); if (startState.parentGranule.present() && startVersion < startState.changeFeedStartVersion) { @@ -1383,7 +1427,6 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, // process mutations if (!mutations.empty()) { bool processedAnyMutations = false; - Version knownNoRollbacksPast = invalidVersion; Version lastDeltaVersion = invalidVersion; for (MutationsAndVersionRef deltas : mutations) { @@ -1404,7 +1447,6 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, br >> rollbackVersion; ASSERT(rollbackVersion >= metadata->durableDeltaVersion.get()); - ASSERT(rollbackVersion >= committedVersion.get()); if (!rollbacksInProgress.empty()) { ASSERT(rollbacksInProgress.front().first == rollbackVersion); @@ -1492,10 +1534,7 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, .detail("OldChangeFeed", readOldChangeFeed ? "T" : "F"); } if (DEBUG_BW_VERSION(deltas.version)) { - fmt::print("BW {0}: ({1}), KCV={2}\n", - deltas.version, - deltas.mutations.size(), - deltas.knownCommittedVersion); + fmt::print("BW {0}: ({1})\n", deltas.version, deltas.mutations.size()); } metadata->currentDeltas.push_back_deep(metadata->deltaArena, deltas); @@ -1503,14 +1542,6 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, ASSERT(deltas.version != invalidVersion); ASSERT(deltas.version > lastDeltaVersion); lastDeltaVersion = deltas.version; - - Version nextKnownNoRollbacksPast = std::min(deltas.version, deltas.knownCommittedVersion); - ASSERT(nextKnownNoRollbacksPast >= knownNoRollbacksPast || - nextKnownNoRollbacksPast == invalidVersion); - - if (nextKnownNoRollbacksPast != invalidVersion) { - knownNoRollbacksPast = nextKnownNoRollbacksPast; - } } } if (justDidRollback) { @@ -1518,24 +1549,12 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, } } if (!justDidRollback && processedAnyMutations) { - // update buffered version and committed version + // update buffered version ASSERT(lastDeltaVersion != invalidVersion); ASSERT(lastDeltaVersion > metadata->bufferedDeltaVersion); // Update buffered delta version so new waitForVersion checks can bypass waiting entirely metadata->bufferedDeltaVersion = lastDeltaVersion; - - // This is the only place it is safe to set committedVersion, as it has to come from the - // mutation stream, or we could have a situation where the blob worker has consumed an - // uncommitted mutation, but not its rollback, fro m the change feed, and could thus - // think the uncommitted mutation is committed because it saw a higher committed version - // than the mutation's version. - // We also can only set it after consuming all of the mutations from the vector from the promise - // stream, as yielding when consuming from a change feed can cause bugs if this wakes up one of the - // file writers - if (knownNoRollbacksPast > committedVersion.get()) { - committedVersion.set(knownNoRollbacksPast); - } } justDidRollback = false; @@ -1564,17 +1583,18 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, } else { previousFuture = Future(BlobFileIndex()); } - Future dfFuture = writeDeltaFile(bwData, - metadata->keyRange, - startState.granuleID, - metadata->originalEpoch, - metadata->originalSeqno, - metadata->deltaArena, - metadata->currentDeltas, - lastDeltaVersion, - previousFuture, - &committedVersion, - oldChangeFeedDataComplete); + Future dfFuture = + writeDeltaFile(bwData, + metadata->keyRange, + startState.granuleID, + metadata->originalEpoch, + metadata->originalSeqno, + metadata->deltaArena, + metadata->currentDeltas, + lastDeltaVersion, + previousFuture, + waitVersionCommitted(bwData, metadata, lastDeltaVersion), + oldChangeFeedDataComplete); inFlightFiles.push_back( InFlightFile(dfFuture, lastDeltaVersion, metadata->bufferedDeltaBytes, false)); @@ -1641,15 +1661,12 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, state int waitIdx = 0; int idx = 0; for (auto& f : inFlightFiles) { - if (f.snapshot && f.version < metadata->pendingSnapshotVersion && - f.version <= committedVersion.get()) { + if (f.snapshot && f.version < metadata->pendingSnapshotVersion) { if (BW_DEBUG) { - fmt::print("[{0} - {1}) Waiting on previous snapshot file @ {2} <= known " - "committed {3}\n", + fmt::print("[{0} - {1}) Waiting on previous snapshot file @ {2}\n", metadata->keyRange.begin.printable(), metadata->keyRange.end.printable(), - f.version, - committedVersion.get()); + f.version); } waitIdx = idx + 1; } @@ -1912,6 +1929,10 @@ ACTOR Future waitForVersion(Reference metadata, Version v if (v == DEBUG_BW_WAIT_VERSION) { fmt::print("{0}) got CF version {1}\n", v, metadata->activeCFData.get()->getVersion()); } + // TODO REMOVE debugging + if (v > metadata->waitForVersionReturned) { + metadata->waitForVersionReturned = v; + } } // wait for any pending delta and snapshot files as of the moment the change feed version caught up. @@ -1969,11 +1990,6 @@ ACTOR Future waitForVersion(Reference metadata, Version v fmt::print("{0}) done\n", v); } - // TODO REMOVE debugging - if (v > metadata->waitForVersionReturned) { - metadata->waitForVersionReturned = v; - } - return Void(); } @@ -2776,6 +2792,40 @@ ACTOR Future monitorRemoval(Reference bwData) { } } +// Because change feeds send uncommitted data and explicit rollback messages, we speculatively buffer/write +// uncommitted data. This means we must ensure the data is actually committed before "committing" those writes in +// the blob granule. The simplest way to do this is to have the blob worker do a periodic GRV, which is guaranteed +// to be an earlier committed version. Then, once the change feed has consumed up through the GRV's data, we can +// guarantee nothing will roll back the in-memory mutations +ACTOR Future runGRVChecks(Reference bwData) { + state Transaction tr(bwData->db); + loop { + // only do grvs to get committed version if we need it to persist delta files + while (bwData->grvVersion.numWaiting() == 0) { + // printf("GRV checker sleeping\n"); + wait(bwData->doGRVCheck.getFuture()); + bwData->doGRVCheck.reset(); + // printf("GRV checker waking: %d pending\n", bwData->grvVersion.numWaiting()); + } + + // batch potentially multiple delta files into one GRV, and also rate limit GRVs for this worker + wait(delay(0.1)); // TODO KNOB? + // printf("GRV checker doing grv @ %.2f\n", now()); + + tr.reset(); + try { + Version readVersion = wait(tr.getReadVersion()); + ASSERT(readVersion >= bwData->grvVersion.get()); + // printf("GRV checker got GRV %lld\n", readVersion); + bwData->grvVersion.set(readVersion); + + ++bwData->stats.commitVersionChecks; + } catch (Error& e) { + wait(tr.onError(e)); + } + } +} + ACTOR Future blobWorker(BlobWorkerInterface bwInterf, ReplyPromise recruitReply, Reference const> dbInfo) { @@ -2828,6 +2878,7 @@ ACTOR Future blobWorker(BlobWorkerInterface bwInterf, recruitReply.send(rep); self->addActor.send(waitFailureServer(bwInterf.waitFailure.getFuture())); + self->addActor.send(runGRVChecks(self)); state Future selfRemoved = monitorRemoval(self); TraceEvent("BlobWorkerInit", self->id); From 307d049c9d713aecab5ffd9fc05295562d9a1715 Mon Sep 17 00:00:00 2001 From: Josh Slocum Date: Fri, 10 Dec 2021 16:12:06 -0600 Subject: [PATCH 2/2] Cleaning up some memory lifetime issues --- fdbclient/NativeAPI.actor.cpp | 6 +++--- fdbserver/BlobWorker.actor.cpp | 8 +++++--- 2 files changed, 8 insertions(+), 6 deletions(-) diff --git a/fdbclient/NativeAPI.actor.cpp b/fdbclient/NativeAPI.actor.cpp index 23107e027c..9d6e1136b0 100644 --- a/fdbclient/NativeAPI.actor.cpp +++ b/fdbclient/NativeAPI.actor.cpp @@ -7085,7 +7085,7 @@ Version ChangeFeedData::getVersion() { #define DEBUG_CF_WAIT_VERSION invalidVersion #define DEBUG_CF_VERSION(v) DEBUG_CF_START_VERSION <= v&& v <= DEBUG_CF_END_VERSION -ACTOR Future changeFeedWaitLatest(ChangeFeedData* self, Version version) { +ACTOR Future changeFeedWaitLatest(Reference self, Version version) { // first, wait on SS to have sent up through version int desired = 0; int waiting = 0; @@ -7162,7 +7162,7 @@ ACTOR Future changeFeedWaitLatest(ChangeFeedData* self, Version version) { return Void(); } -ACTOR Future changeFeedWhenAtLatest(ChangeFeedData* self, Version version) { +ACTOR Future changeFeedWhenAtLatest(Reference self, Version version) { if (DEBUG_CF_WAIT_VERSION == version) { fmt::print("CFW {0}) WhenAtLeast: LR={1}\n", version, self->lastReturnedVersion.get()); } @@ -7200,7 +7200,7 @@ ACTOR Future changeFeedWhenAtLatest(ChangeFeedData* self, Version version) } Future ChangeFeedData::whenAtLeast(Version version) { - return changeFeedWhenAtLatest(this, version); + return changeFeedWhenAtLatest(Reference::addRef(this), version); } ACTOR Future singleChangeFeedStream(StorageServerInterface interf, diff --git a/fdbserver/BlobWorker.actor.cpp b/fdbserver/BlobWorker.actor.cpp index aba8332dce..91de0695f6 100644 --- a/fdbserver/BlobWorker.actor.cpp +++ b/fdbserver/BlobWorker.actor.cpp @@ -1167,7 +1167,9 @@ ACTOR Future waitVersionCommitted(Reference bwData, // make sure the change feed has consumed mutations up through grvVersion to ensure none of them are rollbacks loop { - state Future atLeast = metadata->activeCFData.get()->whenAtLeast(grvVersion); + // if not valid, we're about to be cancelled anyway + state Future atLeast = + metadata->activeCFData.get().isValid() ? metadata->activeCFData.get()->whenAtLeast(grvVersion) : Never(); choose { when(wait(atLeast)) { break; } when(wait(metadata->activeCFData.onChange())) {} @@ -1885,6 +1887,8 @@ ACTOR Future waitForVersion(Reference metadata, Version v // if we don't have to wait for change feed version to catch up or wait for any pending file writes to complete, // nothing to do + ASSERT(metadata->activeCFData.get().isValid()); + if (v == DEBUG_BW_WAIT_VERSION) { fmt::print("{0}) [{1} - {2}) waiting for {3}\n readable:{4}\n bufferedDelta={5}\n pendingDelta={6}\n " "durableDelta={7}\n pendingSnapshot={8}\n durableSnapshot={9}\n", @@ -1900,8 +1904,6 @@ ACTOR Future waitForVersion(Reference metadata, Version v metadata->durableSnapshotVersion.get()); } - ASSERT(metadata->activeCFData.get().isValid()); - if (v <= metadata->activeCFData.get()->getVersion() && (v <= metadata->durableDeltaVersion.get() || metadata->durableDeltaVersion.get() == metadata->pendingDeltaVersion) &&