diff --git a/fdbserver/BackupWorker.actor.cpp b/fdbserver/BackupWorker.actor.cpp index 789935bd2cd..0f45f29a33e 100644 --- a/fdbserver/BackupWorker.actor.cpp +++ b/fdbserver/BackupWorker.actor.cpp @@ -46,8 +46,8 @@ struct VersionedMessage { Arena arena; // Keep a reference to the memory containing the message size_t bytes; // arena's size when inserted, which can grow afterwards - VersionedMessage(LogMessageVersion v, StringRef m, const VectorRef& t, const Arena& a) - : version(v), message(m), tags(t), arena(a), bytes(a.getSize()) {} + VersionedMessage(LogMessageVersion v, StringRef m, const VectorRef& t, const Arena& a, size_t n) + : version(v), message(m), tags(t), arena(a), bytes(n) {} Version getVersion() const { return version.version; } uint32_t getSubVersion() const { return version.sub; } @@ -924,15 +924,17 @@ ACTOR Future pullAsyncData(BackupData* self) { // Note we aggressively peek (uncommitted) messages, but only committed // messages/mutations will be flushed to disk/blob in uploadData(). while (r->hasMessage()) { + state size_t takeBytes = 0; if (!prev.sameArena(r->arena())) { TraceEvent(SevDebugMemory, "BackupWorkerMemory", self->myId) .detail("Take", r->arena().getSize()) .detail("Current", self->lock->activePermits()); - wait(self->lock->take(TaskPriority::DefaultYield, r->arena().getSize())); + takeBytes = r->arena().getSize(); // more bytes can be allocated after the wait. + wait(self->lock->take(TaskPriority::DefaultYield, takeBytes)); prev = r->arena(); } - self->messages.emplace_back(r->version(), r->getMessage(), r->getTags(), r->arena()); + self->messages.emplace_back(r->version(), r->getMessage(), r->getTags(), r->arena(), takeBytes); r->nextMessage(); }