From ddfc301d745c31a15e4074e8cd2492179fc49690 Mon Sep 17 00:00:00 2001 From: Josh Slocum Date: Fri, 4 Feb 2022 16:41:25 -0600 Subject: [PATCH] Improving memory footprint of change feeds and making it configurable --- fdbclient/DatabaseContext.h | 1 + fdbclient/NativeAPI.actor.cpp | 16 +++++++++---- fdbclient/ServerKnobs.cpp | 2 ++ fdbclient/ServerKnobs.h | 2 ++ fdbclient/StorageServerInterface.h | 4 +++- fdbrpc/fdbrpc.h | 2 +- fdbrpc/genericactors.actor.h | 23 +++++++++--------- fdbserver/BlobWorker.actor.cpp | 38 +++++++++++++++++++++--------- fdbserver/storageserver.actor.cpp | 10 +++++--- 9 files changed, 66 insertions(+), 32 deletions(-) diff --git a/fdbclient/DatabaseContext.h b/fdbclient/DatabaseContext.h index f9ad1a63a8..92a9004c3a 100644 --- a/fdbclient/DatabaseContext.h +++ b/fdbclient/DatabaseContext.h @@ -296,6 +296,7 @@ public: Version begin = 0, Version end = std::numeric_limits::max(), KeyRange range = allKeys, + int replyBufferSize = -1, bool canReadPopped = true); Future> getOverlappingChangeFeeds(KeyRangeRef ranges, Version minVersion); diff --git a/fdbclient/NativeAPI.actor.cpp b/fdbclient/NativeAPI.actor.cpp index 2ddf1cd674..6d7e84489c 100644 --- a/fdbclient/NativeAPI.actor.cpp +++ b/fdbclient/NativeAPI.actor.cpp @@ -7610,6 +7610,7 @@ ACTOR Future mergeChangeFeedStream(Reference db, Key rangeID, Version* begin, Version end, + int replyBufferSize, bool canReadPopped) { state std::vector> fetchers(interfs.size()); state std::vector> onErrors(interfs.size()); @@ -7624,6 +7625,7 @@ ACTOR Future mergeChangeFeedStream(Reference db, req.end = end; req.range = it.second; req.canReadPopped = canReadPopped; + req.replyBufferSize = replyBufferSize / interfs.size(); UID debugID = deterministicRandom()->randomUniqueID(); debugIDs.push_back(debugID); req.debugID = debugID; @@ -7810,6 +7812,7 @@ ACTOR Future singleChangeFeedStream(Reference db, Key rangeID, Version* begin, Version end, + int replyBufferSize, bool canReadPopped) { state Database cx(db); state ChangeFeedStreamRequest req; @@ -7819,6 +7822,7 @@ ACTOR Future singleChangeFeedStream(Reference db, req.end = end; req.range = range; req.canReadPopped = canReadPopped; + req.replyBufferSize = replyBufferSize; req.debugID = debugID; results->streams.clear(); @@ -7859,6 +7863,7 @@ ACTOR Future getChangeFeedStreamActor(Reference db, Version begin, Version end, KeyRange range, + int replyBufferSize, bool canReadPopped) { state Database cx(db); state Span span("NAPI:GetChangeFeedStream"_loc); @@ -7938,11 +7943,13 @@ ACTOR Future getChangeFeedStreamActor(Reference db, interfs.push_back(std::make_pair(locations[i].second->getInterface(chosenLocations[i]), locations[i].first & range)); } - wait(mergeChangeFeedStream(db, interfs, results, rangeID, &begin, end, canReadPopped) || - cx->connectionFileChanged()); + wait( + mergeChangeFeedStream(db, interfs, results, rangeID, &begin, end, replyBufferSize, canReadPopped) || + cx->connectionFileChanged()); } else { StorageServerInterface interf = locations[0].second->getInterface(chosenLocations[0]); - wait(singleChangeFeedStream(db, interf, range, results, rangeID, &begin, end, canReadPopped) || + wait(singleChangeFeedStream( + db, interf, range, results, rangeID, &begin, end, replyBufferSize, canReadPopped) || cx->connectionFileChanged()); } } catch (Error& e) { @@ -7991,9 +7998,10 @@ Future DatabaseContext::getChangeFeedStream(Reference resu Version begin, Version end, KeyRange range, + int replyBufferSize, bool canReadPopped) { return getChangeFeedStreamActor( - Reference::addRef(this), results, rangeID, begin, end, range, canReadPopped); + Reference::addRef(this), results, rangeID, begin, end, range, replyBufferSize, canReadPopped); } ACTOR Future> singleLocationOverlappingChangeFeeds( diff --git a/fdbclient/ServerKnobs.cpp b/fdbclient/ServerKnobs.cpp index fde9d49203..2e0d7658cf 100644 --- a/fdbclient/ServerKnobs.cpp +++ b/fdbclient/ServerKnobs.cpp @@ -649,6 +649,8 @@ void ServerKnobs::initialize(Randomize randomize, ClientKnobs* clientKnobs, IsSi init( FETCH_KEYS_TOO_LONG_TIME_CRITERIA, 300.0 ); init( MAX_STORAGE_COMMIT_TIME, 120.0 ); //The max fsync stall time on the storage server and tlog before marking a disk as failed init( RANGESTREAM_LIMIT_BYTES, 2e6 ); if( randomize && BUGGIFY ) RANGESTREAM_LIMIT_BYTES = 1; + init( CHANGEFEEDSTREAM_LIMIT_BYTES, 1e6 ); if( randomize && BUGGIFY ) CHANGEFEEDSTREAM_LIMIT_BYTES = 1; + init( BLOBWORKERSTATUSSTREAM_LIMIT_BYTES, 1e4 ); if( randomize && BUGGIFY ) BLOBWORKERSTATUSSTREAM_LIMIT_BYTES = 1; init( ENABLE_CLEAR_RANGE_EAGER_READS, true ); init( QUICK_GET_VALUE_FALLBACK, true ); init( QUICK_GET_KEY_VALUES_FALLBACK, true ); diff --git a/fdbclient/ServerKnobs.h b/fdbclient/ServerKnobs.h index 6f1f4a9971..9e02ca63d5 100644 --- a/fdbclient/ServerKnobs.h +++ b/fdbclient/ServerKnobs.h @@ -590,6 +590,8 @@ public: double FETCH_KEYS_TOO_LONG_TIME_CRITERIA; double MAX_STORAGE_COMMIT_TIME; int64_t RANGESTREAM_LIMIT_BYTES; + int64_t CHANGEFEEDSTREAM_LIMIT_BYTES; + int64_t BLOBWORKERSTATUSSTREAM_LIMIT_BYTES; bool ENABLE_CLEAR_RANGE_EAGER_READS; bool QUICK_GET_VALUE_FALLBACK; bool QUICK_GET_KEY_VALUES_FALLBACK; diff --git a/fdbclient/StorageServerInterface.h b/fdbclient/StorageServerInterface.h index 556c8548e0..db6fb9f9e4 100644 --- a/fdbclient/StorageServerInterface.h +++ b/fdbclient/StorageServerInterface.h @@ -715,15 +715,17 @@ struct ChangeFeedStreamRequest { Version begin = 0; Version end = 0; KeyRange range; + int replyBufferSize = -1; bool canReadPopped = true; // TODO REMOVE once BG is correctness clean!! Useful for debugging UID debugID; + ReplyPromiseStream reply; ChangeFeedStreamRequest() {} template void serialize(Ar& ar) { - serializer(ar, rangeID, begin, end, range, reply, spanContext, canReadPopped, debugID, arena); + serializer(ar, rangeID, begin, end, range, reply, spanContext, replyBufferSize, canReadPopped, debugID, arena); } }; diff --git a/fdbrpc/fdbrpc.h b/fdbrpc/fdbrpc.h index 666454ed25..fb726309cb 100644 --- a/fdbrpc/fdbrpc.h +++ b/fdbrpc/fdbrpc.h @@ -500,7 +500,7 @@ public: Future onConnected() { if (connected()) { - return Future(Void()); + return Void(); } if (!queue->onConnect.isValid()) { queue->onConnect = Promise(); diff --git a/fdbrpc/genericactors.actor.h b/fdbrpc/genericactors.actor.h index 11c9c38395..e93825b43b 100644 --- a/fdbrpc/genericactors.actor.h +++ b/fdbrpc/genericactors.actor.h @@ -197,11 +197,6 @@ struct PeerHolder { } }; -ACTOR template -void holdUntilConnected(Future signal, ReplyPromiseStream stream) { - wait(stream.onConnected() || signal); -} - // Implements getReplyStream, this a void actor with the same lifetime as the input ReplyPromiseStream. // Because this actor holds a reference to the stream, normally it would be impossible to know when there are no other // references. To get around this, there is a SAV inside the stream that has one less promise reference than it should @@ -215,13 +210,17 @@ void endStreamOnDisconnect(Future signal, Reference peer = Reference()) { state PeerHolder holder = PeerHolder(peer); stream.setRequestStreamEndpoint(endpoint); - choose { - when(wait(signal)) { stream.sendError(connection_failed()); } - when(wait(stream.getErrorFutureAndDelPromiseRef())) { - // Wait for a response from the server - /*if (!stream.connected()) { - // TODO WANT TO DO holdAndConnected ACTOR HERE INSTEAD! - }*/ + try { + choose { + when(wait(signal)) { stream.sendError(connection_failed()); } + when(wait(stream.getErrorFutureAndDelPromiseRef())) {} + } + } catch (Error& e) { + if (e.code() == error_code_broken_promise) { + // getErrorFutureAndDelPromiseRef returned, wait on stream connect or error + if (!stream.connected()) { + wait(signal || stream.onConnected()); + } } } } diff --git a/fdbserver/BlobWorker.actor.cpp b/fdbserver/BlobWorker.actor.cpp index 6806a7cac7..e64e293f98 100644 --- a/fdbserver/BlobWorker.actor.cpp +++ b/fdbserver/BlobWorker.actor.cpp @@ -177,6 +177,8 @@ struct BlobWorkerData : NonCopyable, ReferenceCounted { Promise doGRVCheck; NotifiedVersion grvVersion; + int changeFeedStreamReplyBufferSize = SERVER_KNOBS->BG_DELTA_FILE_TARGET_BYTES / 2; + BlobWorkerData(UID id, Database db) : id(id), db(db), stats(id, SERVER_KNOBS->WORKER_LOGGING_INTERVAL) {} ~BlobWorkerData() { if (BW_DEBUG) { @@ -1338,12 +1340,18 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, startVersion + 1, startState.changeFeedStartVersion, metadata->keyRange, + bwData->changeFeedStreamReplyBufferSize, false); } else { readOldChangeFeed = false; - changeFeedFuture = bwData->db->getChangeFeedStream( - newCFData, cfKey, startVersion + 1, MAX_VERSION, metadata->keyRange, false); + changeFeedFuture = bwData->db->getChangeFeedStream(newCFData, + cfKey, + startVersion + 1, + MAX_VERSION, + metadata->keyRange, + bwData->changeFeedStreamReplyBufferSize, + false); } // Start actors BEFORE setting new change feed data to ensure the change feed data is properly initialized by @@ -1469,8 +1477,13 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, Reference newCFData = makeReference(); - changeFeedFuture = bwData->db->getChangeFeedStream( - newCFData, cfKey, startState.changeFeedStartVersion, MAX_VERSION, metadata->keyRange, false); + changeFeedFuture = bwData->db->getChangeFeedStream(newCFData, + cfKey, + startState.changeFeedStartVersion, + MAX_VERSION, + metadata->keyRange, + bwData->changeFeedStreamReplyBufferSize, + false); // Start actors BEFORE setting new change feed data to ensure the change feed data is properly // initialized by the client @@ -1590,6 +1603,7 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, cfRollbackVersion + 1, startState.changeFeedStartVersion, metadata->keyRange, + bwData->changeFeedStreamReplyBufferSize, false); } else { @@ -1600,12 +1614,14 @@ ACTOR Future blobGranuleUpdateFiles(Reference bwData, } ASSERT(cfRollbackVersion >= startState.changeFeedStartVersion); - changeFeedFuture = bwData->db->getChangeFeedStream(newCFData, - cfKey, - cfRollbackVersion + 1, - MAX_VERSION, - metadata->keyRange, - false); + changeFeedFuture = + bwData->db->getChangeFeedStream(newCFData, + cfKey, + cfRollbackVersion + 1, + MAX_VERSION, + metadata->keyRange, + bwData->changeFeedStreamReplyBufferSize, + false); } // Start actors BEFORE setting new change feed data to ensure the change feed data @@ -3094,7 +3110,7 @@ ACTOR Future blobWorker(BlobWorkerInterface bwInterf, self->currentManagerStatusStream.get().sendError(connection_failed()); // TODO: pick a reasonable byte limit instead of just piggy-backing - req.reply.setByteLimit(SERVER_KNOBS->RANGESTREAM_LIMIT_BYTES); + req.reply.setByteLimit(SERVER_KNOBS->BLOBWORKERSTATUSSTREAM_LIMIT_BYTES); self->currentManagerStatusStream.set(req.reply); } else { req.reply.sendError(blob_manager_replaced()); diff --git a/fdbserver/storageserver.actor.cpp b/fdbserver/storageserver.actor.cpp index 7f4d352b50..2334ca6ed5 100644 --- a/fdbserver/storageserver.actor.cpp +++ b/fdbserver/storageserver.actor.cpp @@ -2050,7 +2050,11 @@ ACTOR Future changeFeedStreamQ(StorageServer* data, ChangeFeedStreamReques state UID streamUID = req.debugID; state bool removeUID = false; state Optional blockedVersion; - req.reply.setByteLimit(SERVER_KNOBS->RANGESTREAM_LIMIT_BYTES); + if (req.replyBufferSize <= 0) { + req.reply.setByteLimit(SERVER_KNOBS->CHANGEFEEDSTREAM_LIMIT_BYTES); + } else { + req.reply.setByteLimit(req.replyBufferSize); + } wait(delay(0, TaskPriority::DefaultEndpoint)); @@ -4145,8 +4149,8 @@ ACTOR Future fetchChangeFeedApplier(StorageServer* data, } state Reference feedResults = makeReference(); - state Future feed = - data->cx->getChangeFeedStream(feedResults, rangeId, startVersion, endVersion, range, true); + state Future feed = data->cx->getChangeFeedStream( + feedResults, rangeId, startVersion, endVersion, range, SERVER_KNOBS->CHANGEFEEDSTREAM_LIMIT_BYTES, true); // TODO remove debugging eventually? state Version firstVersion = invalidVersion;