diff --git a/fdbserver/TransactionTagCounter.cpp b/fdbserver/TransactionTagCounter.cpp index 1f0a25c2cc..a9a9847832 100644 --- a/fdbserver/TransactionTagCounter.cpp +++ b/fdbserver/TransactionTagCounter.cpp @@ -21,8 +21,39 @@ #include "fdbserver/TransactionTagCounter.h" #include "flow/Trace.h" +void TopKTags::incrementTagCount(TransactionTag tag, int previousCount, int increase) { + auto iter = std::find_if(topTags.begin(), topTags.end(), [tag](const auto& tc) { return tc.tag == tag; }); + if (iter != topTags.end()) { + ASSERT_EQ(previousCount, iter->count); + iter->count += increase; + } else if (topTags.size() < limit) { + ASSERT_EQ(previousCount, 0); + topTags.emplace_back(tag, increase); + } else { + auto toReplace = std::min_element(topTags.begin(), topTags.end()); + ASSERT_GE(toReplace->count, previousCount); + if (toReplace->count < previousCount + increase) { + toReplace->tag = tag; + toReplace->count = previousCount + increase; + } + } +} + +std::vector TopKTags::getBusiestTags(double elapsed, + double totalSampleCount) const { + std::vector result; + for (auto const& tagAndCounter : topTags) { + // FIXME: Confusing operator precedence + auto rate = tagAndCounter.count / CLIENT_KNOBS->READ_TAG_SAMPLE_RATE / elapsed; + if (rate > SERVER_KNOBS->MIN_TAG_READ_PAGES_RATE) { + result.emplace_back(tagAndCounter.tag, rate, tagAndCounter.count / totalSampleCount); + } + } + return result; +} + TransactionTagCounter::TransactionTagCounter(UID thisServerID) - : thisServerID(thisServerID), + : thisServerID(thisServerID), topTags(1), // TODO: Make this a knob busiestReadTagEventHolder(makeReference(thisServerID.toString() + "/BusiestReadTag")) {} void TransactionTagCounter::addRequest(Optional const& tags, int64_t bytes) { @@ -31,11 +62,8 @@ void TransactionTagCounter::addRequest(Optional const& tags, int64_t byt double cost = costFunction(bytes); for (auto& tag : tags.get()) { int64_t& count = intervalCounts[TransactionTag(tag, tags.get().getArena())]; + topTags.incrementTagCount(tag, count, cost); count += cost; - if (count > busiestTagCount) { - busiestTagCount = count; - busiestTag = tag; - } } intervalTotalSampledCount += cost; @@ -46,22 +74,19 @@ void TransactionTagCounter::startNewInterval() { double elapsed = now() - intervalStart; previousBusiestTags.clear(); if (intervalStart > 0 && CLIENT_KNOBS->READ_TAG_SAMPLE_RATE > 0 && elapsed > 0) { - double rate = busiestTagCount / CLIENT_KNOBS->READ_TAG_SAMPLE_RATE / elapsed; - if (rate > SERVER_KNOBS->MIN_TAG_READ_PAGES_RATE) { - previousBusiestTags.emplace_back(busiestTag, rate, (double)busiestTagCount / intervalTotalSampledCount); - } + previousBusiestTags = topTags.getBusiestTags(elapsed, intervalTotalSampledCount); TraceEvent("BusiestReadTag", thisServerID) .detail("Elapsed", elapsed) - .detail("Tag", printable(busiestTag)) - .detail("TagCost", busiestTagCount) + //.detail("Tag", printable(busiestTag)) + //.detail("TagCost", busiestTagCount) .detail("TotalSampledCost", intervalTotalSampledCount) - .detail("Reported", !previousBusiestTags.empty()) + .detail("Reported", previousBusiestTags.size()) .trackLatest(busiestReadTagEventHolder->trackingKey); } intervalCounts.clear(); intervalTotalSampledCount = 0; - busiestTagCount = 0; + topTags.clear(); intervalStart = now(); } diff --git a/fdbserver/TransactionTagCounter.h b/fdbserver/TransactionTagCounter.h index d520259c5c..aad83104bd 100644 --- a/fdbserver/TransactionTagCounter.h +++ b/fdbserver/TransactionTagCounter.h @@ -24,15 +24,39 @@ #include "fdbclient/TagThrottle.actor.h" #include "fdbserver/Knobs.h" +class TopKTags { +public: + struct TagAndCount { + TransactionTag tag; + int64_t count; + bool operator<(TagAndCount const& other) const { return count < other.count; } + explicit TagAndCount(TransactionTag tag, int64_t count) : tag(tag), count(count) {} + }; + +private: + // Because the number of tracked is expected to be small, they can be tracked + // in a simple vector. If the number of tracked tags increases, a more sophisticated + // data structure will be required. + std::vector topTags; + int limit; + +public: + explicit TopKTags(int limit) : limit(limit) { ASSERT_GT(limit, 0); } + void incrementTagCount(TransactionTag tag, int previousCount, int increase); + + std::vector getBusiestTags(double elapsed, double totalSampleCount) const; + + void clear() { topTags.clear(); } +}; + class TransactionTagCounter { + UID thisServerID; TransactionTagMap intervalCounts; int64_t intervalTotalSampledCount = 0; - TransactionTag busiestTag; - int64_t busiestTagCount = 0; + TopKTags topTags; double intervalStart = 0; std::vector previousBusiestTags; - UID thisServerID; Reference busiestReadTagEventHolder; public: