From cc05f5e9dbe6bf55c516bfbd2ddbbefe72437e2c Mon Sep 17 00:00:00 2001 From: Xiaoxi Wang Date: Sun, 24 Apr 2022 17:10:58 -0700 Subject: [PATCH] fix getMetrics keys bug --- fdbserver/DataDistribution.actor.h | 4 ++-- fdbserver/DataDistributionQueue.actor.cpp | 19 ++++++++++++------- fdbserver/DataDistributionTracker.actor.cpp | 5 +++-- 3 files changed, 17 insertions(+), 11 deletions(-) diff --git a/fdbserver/DataDistribution.actor.h b/fdbserver/DataDistribution.actor.h index 16c42f23c2..66cfef3c48 100644 --- a/fdbserver/DataDistribution.actor.h +++ b/fdbserver/DataDistribution.actor.h @@ -152,13 +152,13 @@ struct GetTeamRequest { }; struct GetMetricsRequest { - // whether a < b + // whether a > b typedef std::function MetricsComparator; std::vector keys; int topK = 1; // default only return the top 1 shard based on the comparator Promise> reply; // topK storage metrics Optional - comparator; // if comparator is assigned, return the largest topK in keys, otherwise return the sum of metrics + comparator; // Return true if a.score > b.score.if comparator is assigned, return the largest topK in keys, otherwise return the sum of metrics GetMetricsRequest() {} GetMetricsRequest(KeyRange const& keys, int topK = 1) : keys({ keys }), topK(topK) {} diff --git a/fdbserver/DataDistributionQueue.actor.cpp b/fdbserver/DataDistributionQueue.actor.cpp index 9a2fe3e1a5..1d61812a37 100644 --- a/fdbserver/DataDistributionQueue.actor.cpp +++ b/fdbserver/DataDistributionQueue.actor.cpp @@ -1081,7 +1081,7 @@ struct DDQueueData { bool timeThrottle(const std::vector& ids) const { return std::any_of(ids.begin(), ids.end(), [this](const UID& id) { if (this->lastAsSource.count(id)) { - return (now() - this->lastAsSource.at(id)) * 3.0 < SERVER_KNOBS->STORAGE_METRICS_AVERAGE_INTERVAL; + return (now() - this->lastAsSource.at(id)) * 5.0 < SERVER_KNOBS->STORAGE_METRICS_AVERAGE_INTERVAL; } return false; }); @@ -1504,10 +1504,15 @@ ACTOR Future dataDistributionRelocator(DDQueueData* self, RelocateData rd, } } -inline double getWorstCpu(const HealthMetrics& metrics) { +inline double getWorstCpu(const HealthMetrics& metrics, const std::vector& ids) { double cpu = 0; - for (auto p : metrics.storageStats) { - cpu = std::max(cpu, p.second.cpuUsage); + for (auto& id : ids) { + if (metrics.storageStats.count(id)) { + cpu = std::max(cpu, metrics.storageStats.at(id).cpuUsage); + } else { + // assume the server is too busy to report its stats + cpu = std::max(cpu, 100.0); + } } return cpu; } @@ -1546,12 +1551,12 @@ ACTOR Future rebalanceReadLoad(DDQueueData* self, state Future healthMetrics = self->cx->getHealthMetrics(true); state GetMetricsRequest req(shards, 10); req.comparator = [](const StorageMetrics& a, const StorageMetrics& b) { - return a.bytesReadPerKSecond / std::max(a.bytes * 1.0, 1.0 * SERVER_KNOBS->MIN_SHARD_BYTES) < + return a.bytesReadPerKSecond / std::max(a.bytes * 1.0, 1.0 * SERVER_KNOBS->MIN_SHARD_BYTES) > b.bytesReadPerKSecond / std::max(b.bytes * 1.0, 1.0 * SERVER_KNOBS->MIN_SHARD_BYTES); }; state std::vector metricsList = wait(brokenPromiseToNever(self->getShardMetrics.getReply(req))); wait(ready(healthMetrics)); - if (getWorstCpu(healthMetrics.get()) < 25.0) { // 25% + if (getWorstCpu(healthMetrics.get(), sourceTeam->getServerIDs()) < 25.0) { // 25% traceEvent->detail("SkipReason", "LowReadLoad"); return false; } @@ -1559,7 +1564,7 @@ ACTOR Future rebalanceReadLoad(DDQueueData* self, deterministicRandom()->randomShuffle(metricsList); int chosenIdx = -1; for (int i = 0; i < metricsList.size(); ++i) { - if (metricsList[i].keys.present() && metricsList[i].bytes > 0) { + if (metricsList[i].keys.present() && metricsList[i].bytesReadPerKSecond > 0) { chosenIdx = i; break; } diff --git a/fdbserver/DataDistributionTracker.actor.cpp b/fdbserver/DataDistributionTracker.actor.cpp index 0f8883cccc..3c578b3f84 100644 --- a/fdbserver/DataDistributionTracker.actor.cpp +++ b/fdbserver/DataDistributionTracker.actor.cpp @@ -856,6 +856,7 @@ ACTOR Future fetchShardMetrics_impl(DataDistributionTracker* self, GetMetr } if (req.comparator.present()) { + metrics.keys = range; returnMetrics.push_back(metrics); } else { returnMetrics[0] += metrics; @@ -867,11 +868,11 @@ ACTOR Future fetchShardMetrics_impl(DataDistributionTracker* self, GetMetr req.reply.send(returnMetrics); else if (req.comparator.present()) { std::nth_element(returnMetrics.begin(), - returnMetrics.end() - req.topK, + returnMetrics.begin() + req.topK - 1, returnMetrics.end(), req.comparator.get()); req.reply.send( - std::vector(returnMetrics.rbegin(), returnMetrics.rbegin() + req.topK)); + std::vector(returnMetrics.begin(), returnMetrics.begin() + req.topK)); } return Void(); }