Improving memory footprint of change feeds and making it configurable

This commit is contained in:
Josh Slocum 2022-02-04 16:41:25 -06:00
parent 0a3608a33b
commit ddfc301d74
9 changed files with 66 additions and 32 deletions

View File

@ -296,6 +296,7 @@ public:
Version begin = 0,
Version end = std::numeric_limits<Version>::max(),
KeyRange range = allKeys,
int replyBufferSize = -1,
bool canReadPopped = true);
Future<std::vector<OverlappingChangeFeedEntry>> getOverlappingChangeFeeds(KeyRangeRef ranges, Version minVersion);

View File

@ -7610,6 +7610,7 @@ ACTOR Future<Void> mergeChangeFeedStream(Reference<DatabaseContext> db,
Key rangeID,
Version* begin,
Version end,
int replyBufferSize,
bool canReadPopped) {
state std::vector<Future<Void>> fetchers(interfs.size());
state std::vector<Future<Void>> onErrors(interfs.size());
@ -7624,6 +7625,7 @@ ACTOR Future<Void> mergeChangeFeedStream(Reference<DatabaseContext> 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<Void> singleChangeFeedStream(Reference<DatabaseContext> db,
Key rangeID,
Version* begin,
Version end,
int replyBufferSize,
bool canReadPopped) {
state Database cx(db);
state ChangeFeedStreamRequest req;
@ -7819,6 +7822,7 @@ ACTOR Future<Void> singleChangeFeedStream(Reference<DatabaseContext> 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<Void> getChangeFeedStreamActor(Reference<DatabaseContext> 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<Void> getChangeFeedStreamActor(Reference<DatabaseContext> 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<Void> DatabaseContext::getChangeFeedStream(Reference<ChangeFeedData> resu
Version begin,
Version end,
KeyRange range,
int replyBufferSize,
bool canReadPopped) {
return getChangeFeedStreamActor(
Reference<DatabaseContext>::addRef(this), results, rangeID, begin, end, range, canReadPopped);
Reference<DatabaseContext>::addRef(this), results, rangeID, begin, end, range, replyBufferSize, canReadPopped);
}
ACTOR Future<std::vector<OverlappingChangeFeedEntry>> singleLocationOverlappingChangeFeeds(

View File

@ -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 );

View File

@ -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;

View File

@ -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<ChangeFeedStreamReply> reply;
ChangeFeedStreamRequest() {}
template <class Ar>
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);
}
};

View File

@ -500,7 +500,7 @@ public:
Future<Void> onConnected() {
if (connected()) {
return Future<Void>(Void());
return Void();
}
if (!queue->onConnect.isValid()) {
queue->onConnect = Promise<Void>();

View File

@ -197,11 +197,6 @@ struct PeerHolder {
}
};
ACTOR template <class X>
void holdUntilConnected(Future<Void> signal, ReplyPromiseStream<X> 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<Void> signal,
Reference<Peer> peer = Reference<Peer>()) {
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());
}
}
}
}

View File

@ -177,6 +177,8 @@ struct BlobWorkerData : NonCopyable, ReferenceCounted<BlobWorkerData> {
Promise<Void> 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<Void> blobGranuleUpdateFiles(Reference<BlobWorkerData> 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<Void> blobGranuleUpdateFiles(Reference<BlobWorkerData> bwData,
Reference<ChangeFeedData> newCFData = makeReference<ChangeFeedData>();
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<Void> blobGranuleUpdateFiles(Reference<BlobWorkerData> bwData,
cfRollbackVersion + 1,
startState.changeFeedStartVersion,
metadata->keyRange,
bwData->changeFeedStreamReplyBufferSize,
false);
} else {
@ -1600,12 +1614,14 @@ ACTOR Future<Void> blobGranuleUpdateFiles(Reference<BlobWorkerData> 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<Void> 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());

View File

@ -2050,7 +2050,11 @@ ACTOR Future<Void> changeFeedStreamQ(StorageServer* data, ChangeFeedStreamReques
state UID streamUID = req.debugID;
state bool removeUID = false;
state Optional<Version> 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<Version> fetchChangeFeedApplier(StorageServer* data,
}
state Reference<ChangeFeedData> feedResults = makeReference<ChangeFeedData>();
state Future<Void> feed =
data->cx->getChangeFeedStream(feedResults, rangeId, startVersion, endVersion, range, true);
state Future<Void> feed = data->cx->getChangeFeedStream(
feedResults, rangeId, startVersion, endVersion, range, SERVER_KNOBS->CHANGEFEEDSTREAM_LIMIT_BYTES, true);
// TODO remove debugging eventually?
state Version firstVersion = invalidVersion;