added support for querying specific range feed versions

This commit is contained in:
Evan Tschannen 2021-08-09 20:39:28 -07:00
parent 52fcf3f565
commit 42ae870c84
5 changed files with 99 additions and 26 deletions

View File

@ -3580,12 +3580,32 @@ ACTOR Future<int> cli(CLIOptions opt, LineNoise* plinenoise) {
}
}
} else if (tokencmp(tokens[1], "get")) {
if (tokens.size() != 3) {
if (tokens.size() < 3 || tokens.size() > 5) {
printUsage(tokens[0]);
is_error = true;
continue;
}
Standalone<VectorRef<MutationsAndVersionRef>> res = wait(db->getRangeFeedMutations(tokens[2]));
Version begin = 0;
Version end = std::numeric_limits<Version>::max();
if (tokens.size() > 3) {
int n = 0;
if (sscanf(tokens[3].toString().c_str(), "%ld%n", &begin, &n) != 1 ||
n != tokens[3].size()) {
printUsage(tokens[0]);
is_error = true;
continue;
}
}
if (tokens.size() > 4) {
int n = 0;
if (sscanf(tokens[4].toString().c_str(), "%ld%n", &end, &n) != 1 || n != tokens[4].size()) {
printUsage(tokens[0]);
is_error = true;
continue;
}
}
Standalone<VectorRef<MutationsAndVersionRef>> res =
wait(db->getRangeFeedMutations(tokens[2], begin, end));
for (auto& it : res) {
for (auto& it2 : it.mutations) {
printf("%lld %s\n", it.version, it2.toString().c_str());

View File

@ -252,8 +252,11 @@ public:
// Management API, create snapshot
Future<Void> createSnapshot(StringRef uid, StringRef snapshot_command);
Future<Standalone<VectorRef<MutationsAndVersionRef>>> getRangeFeedMutations(StringRef rangeID,
KeyRangeRef range = allKeys);
Future<Standalone<VectorRef<MutationsAndVersionRef>>> getRangeFeedMutations(
StringRef rangeID,
Version begin = 0,
Version end = std::numeric_limits<Version>::max(),
KeyRange range = allKeys);
Future<std::vector<std::pair<Key, KeyRange>>> getOverlappingRangeFeeds(KeyRangeRef ranges, Version minVersion);
Future<Void> popRangeFeedMutations(StringRef rangeID, Version version);

View File

@ -6519,7 +6519,9 @@ Future<Void> DatabaseContext::createSnapshot(StringRef uid, StringRef snapshot_c
ACTOR Future<Standalone<VectorRef<MutationsAndVersionRef>>> getRangeFeedMutationsActor(Reference<DatabaseContext> db,
StringRef rangeID,
KeyRangeRef range) {
Version begin,
Version end,
KeyRange range) {
state Database cx(db);
state Transaction tr(cx);
state Key rangeIDKey = rangeID.withPrefix(rangeFeedPrefix);
@ -6543,6 +6545,8 @@ ACTOR Future<Standalone<VectorRef<MutationsAndVersionRef>>> getRangeFeedMutation
state RangeFeedRequest req;
req.rangeID = rangeID;
req.begin = begin;
req.end = end;
RangeFeedReply rep = wait(loadBalance(cx.getPtr(),
locations[0].second,
@ -6555,8 +6559,10 @@ ACTOR Future<Standalone<VectorRef<MutationsAndVersionRef>>> getRangeFeedMutation
}
Future<Standalone<VectorRef<MutationsAndVersionRef>>> DatabaseContext::getRangeFeedMutations(StringRef rangeID,
KeyRangeRef range) {
return getRangeFeedMutationsActor(Reference<DatabaseContext>::addRef(this), rangeID, range);
Version begin,
Version end,
KeyRange range) {
return getRangeFeedMutationsActor(Reference<DatabaseContext>::addRef(this), rangeID, begin, end, range);
}
ACTOR Future<std::vector<std::pair<Key, KeyRange>>> getOverlappingRangeFeedsActor(Reference<DatabaseContext> db,

View File

@ -663,6 +663,8 @@ struct RangeFeedReply {
struct RangeFeedRequest {
constexpr static FileIdentifier file_identifier = 10726174;
Key rangeID;
Version begin = 0;
Version end = 0;
ReplyPromise<RangeFeedReply> reply;
RangeFeedRequest() {}
@ -670,7 +672,7 @@ struct RangeFeedRequest {
template <class Ar>
void serialize(Ar& ar) {
serializer(ar, rangeID, reply);
serializer(ar, rangeID, begin, end, reply);
}
};

View File

@ -313,6 +313,7 @@ struct FetchInjectionInfo {
struct RangeFeedInfo : ReferenceCounted<RangeFeedInfo> {
std::deque<Standalone<MutationsAndVersionRef>> mutations;
Version durableVersion = invalidVersion;
Version emptyVersion = 0;
KeyRange range;
Key id;
};
@ -1531,7 +1532,8 @@ ACTOR Future<RangeFeedReply> getRangeFeedMutations(StorageServer* data, RangeFee
state RangeFeedReply reply;
wait(delay(0));
auto& feedInfo = data->uidRangeFeed[req.rangeID];
if (feedInfo->durableVersion == invalidVersion) {
if (req.end <= feedInfo->emptyVersion + 1) {
} else if (feedInfo->durableVersion == invalidVersion || req.begin > feedInfo->durableVersion) {
for (auto& it : data->uidRangeFeed[req.rangeID]->mutations) {
reply.mutations.push_back(reply.arena, it);
}
@ -1539,10 +1541,7 @@ ACTOR Future<RangeFeedReply> getRangeFeedMutations(StorageServer* data, RangeFee
state std::deque<Standalone<MutationsAndVersionRef>> mutationsDeque =
data->uidRangeFeed[req.rangeID]->mutations;
RangeResult res = wait(data->storage.readRange(
KeyRangeRef(rangeFeedDurableKey(req.rangeID, 0), rangeFeedDurableKey(req.rangeID, data->version.get()))));
if (res.empty()) {
data->uidRangeFeed[req.rangeID]->durableVersion = invalidVersion;
}
KeyRangeRef(rangeFeedDurableKey(req.rangeID, req.begin), rangeFeedDurableKey(req.rangeID, req.end))));
Version lastVersion = invalidVersion;
for (auto& kv : res) {
Key id;
@ -1553,10 +1552,35 @@ ACTOR Future<RangeFeedReply> getRangeFeedMutations(StorageServer* data, RangeFee
lastVersion = version;
}
for (auto& it : mutationsDeque) {
if (it.version >= req.end) {
break;
}
if (it.version > lastVersion) {
reply.mutations.push_back(reply.arena, it);
}
}
if (res.empty()) {
auto& feedInfo = data->uidRangeFeed[req.rangeID];
if (req.end > feedInfo->durableVersion) {
if (req.begin == 0) {
feedInfo->durableVersion = req.end > data->storageVersion() ? invalidVersion : req.end;
} else {
RangeResult emp = wait(data->storage.readRange(
KeyRangeRef(rangeFeedDurableKey(req.rangeID, 0), rangeFeedDurableKey(req.rangeID, req.end)),
-1));
auto& feedInfo = data->uidRangeFeed[req.rangeID];
if (emp.empty()) {
feedInfo->durableVersion = req.end > data->storageVersion() ? invalidVersion : req.end;
} else {
Key id;
Version version;
std::tie(id, version) = decodeRangeFeedDurableKey(emp[0].key);
feedInfo->durableVersion = version;
}
}
}
}
}
return reply;
}
@ -2945,7 +2969,7 @@ static const KeyRangeRef persistRangeFeedKeys =
KeyRangeRef(LiteralStringRef(PERSIST_PREFIX "RF/"), LiteralStringRef(PERSIST_PREFIX "RF0"));
// data keys are unmangled (but never start with PERSIST_PREFIX because they are always in allKeys)
ACTOR Future<Void> fetchRangeFeed(StorageServer* data, Key rangeId, KeyRange range) {
ACTOR Future<Void> fetchRangeFeed(StorageServer* data, Key rangeId, KeyRange range, Version fetchVersion) {
TraceEvent("FetchRangeFeed", data->thisServerID)
.detail("RangeID", rangeId.printable())
@ -2966,9 +2990,11 @@ ACTOR Future<Void> fetchRangeFeed(StorageServer* data, Key rangeId, KeyRange ran
rangeFeedValue(range)));
state Standalone<VectorRef<MutationsAndVersionRef>> mutations =
wait(data->cx->getRangeFeedMutations(rangeId, range));
wait(data->cx->getRangeFeedMutations(rangeId, 0, fetchVersion, range));
state RangeFeedRequest req;
req.rangeID = rangeId;
req.begin = 0;
req.end = fetchVersion;
state RangeFeedReply rep = wait(getRangeFeedMutations(data, req));
state int mLoc = 0;
state int rLoc = 0;
@ -3000,7 +3026,7 @@ ACTOR Future<Void> dispatchRangeFeeds(StorageServer* data, UID fetchKeysID, KeyR
state std::vector<std::pair<Key, KeyRange>> feeds =
wait(data->cx->getOverlappingRangeFeeds(keys, fetchVersion));
for (auto& feed : feeds) {
feedFetches[feed.first] = fetchRangeFeed(data, feed.first, feed.second);
feedFetches[feed.first] = fetchRangeFeed(data, feed.first, feed.second, fetchVersion);
}
loop {
@ -3752,6 +3778,7 @@ private:
Reference<RangeFeedInfo> rangeFeedInfo(new RangeFeedInfo());
rangeFeedInfo->range = rangeFeedRange;
rangeFeedInfo->id = rangeFeedId;
rangeFeedInfo->emptyVersion = currentVersion - 1;
data->uidRangeFeed[rangeFeedId] = rangeFeedInfo;
auto rs = data->keyRangeFeed.modify(rangeFeedRange);
for (auto r = rs.begin(); r != rs.end(); ++r) {
@ -4277,11 +4304,13 @@ ACTOR Future<Void> updateStorage(StorageServer* data) {
state int curFeed = 0;
while (curFeed < updatedRangeFeeds.size()) {
auto info = data->uidRangeFeed[updatedRangeFeeds[curFeed]];
while (info->mutations.front().version < newOldestVersion) {
for (auto& it : info->mutations) {
if (it.version >= newOldestVersion) {
break;
}
data->storage.writeKeyValue(KeyValueRef(rangeFeedDurableKey(info->id, info->mutations.front().version),
rangeFeedDurableValue(info->mutations.front().mutations)));
info->durableVersion = info->mutations.front().version;
info->mutations.pop_front();
}
wait(yield(TaskPriority::UpdateStorage));
curFeed++;
@ -4320,6 +4349,16 @@ ACTOR Future<Void> updateStorage(StorageServer* data) {
throw please_reboot();
}
curFeed = 0;
while (curFeed < updatedRangeFeeds.size()) {
auto info = data->uidRangeFeed[updatedRangeFeeds[curFeed]];
while (info->mutations.front().version < newOldestVersion) {
info->mutations.pop_front();
}
wait(yield(TaskPriority::UpdateStorage));
curFeed++;
}
durableInProgress.send(Void());
wait(delay(0, TaskPriority::UpdateStorage)); // Setting durableInProgess could cause the storage server to shut
// down, so delay to check for cancellation
@ -5240,14 +5279,17 @@ ACTOR Future<Void> serveRangeFeedPopRequests(StorageServer* self, FutureStream<R
loop {
RangeFeedPopRequest req = waitNext(rangeFeedPops);
auto& feed = self->uidRangeFeed[req.rangeID];
while (!feed->mutations.empty() && feed->mutations.front().version < req.version) {
self->uidRangeFeed[req.rangeID]->mutations.pop_front();
}
if (feed->durableVersion != invalidVersion) {
self->storage.clearRange(
KeyRangeRef(rangeFeedDurableKey(feed->id, 0), rangeFeedDurableKey(feed->id, req.version)));
if (req.version > feed->durableVersion) {
feed->durableVersion = invalidVersion;
if (req.version - 1 > feed->emptyVersion) {
feed->emptyVersion = req.version - 1;
while (!feed->mutations.empty() && feed->mutations.front().version < req.version) {
self->uidRangeFeed[req.rangeID]->mutations.pop_front();
}
if (feed->durableVersion != invalidVersion) {
self->storage.clearRange(
KeyRangeRef(rangeFeedDurableKey(feed->id, 0), rangeFeedDurableKey(feed->id, req.version)));
if (req.version > feed->durableVersion) {
feed->durableVersion = invalidVersion;
}
}
}
TraceEvent("RangeFeedPopQuery", self->thisServerID)