Share queue metric smoothing between storage servers and TLogs

This commit is contained in:
Trevor Clinkenbeard 2026-07-22 09:59:24 -07:00
parent b7950d9e52
commit 02add7bdf0
2 changed files with 75 additions and 60 deletions

View File

@ -1163,16 +1163,41 @@ Future<Void> ratekeeper(RatekeeperInterface rkInterf, Reference<AsyncVar<ServerD
co_await Ratekeeper::run(rkInterf, dbInfo);
}
StorageQueueInfo::StorageQueueInfo(const UID& ratekeeperID_, const UID& id_, const LocalityData& locality_)
: valid(false), ratekeeperID(ratekeeperID_), id(id_), locality(locality_), acceptingRequests(false),
smoothDurableBytes(SERVER_KNOBS->SMOOTHING_AMOUNT), smoothInputBytes(SERVER_KNOBS->SMOOTHING_AMOUNT),
verySmoothDurableBytes(SERVER_KNOBS->SLOW_SMOOTHING_AMOUNT), smoothDurableVersion(SERVER_KNOBS->SMOOTHING_AMOUNT),
smoothLatestVersion(SERVER_KNOBS->SMOOTHING_AMOUNT), smoothFreeSpace(SERVER_KNOBS->SMOOTHING_AMOUNT),
smoothTotalSpace(SERVER_KNOBS->SMOOTHING_AMOUNT), limitReason(limitReason_t::unlimited) {
// FIXME: this is a tacky workaround for a potential uninitialized use in trackStorageServerQueueInfo
lastReply.instanceID = -1;
QueueMetricsSmoother::QueueMetricsSmoother()
: smoothDurableBytes(SERVER_KNOBS->SMOOTHING_AMOUNT), smoothInputBytes(SERVER_KNOBS->SMOOTHING_AMOUNT),
verySmoothDurableBytes(SERVER_KNOBS->SLOW_SMOOTHING_AMOUNT), smoothFreeSpace(SERVER_KNOBS->SMOOTHING_AMOUNT),
smoothTotalSpace(SERVER_KNOBS->SMOOTHING_AMOUNT) {}
bool QueueMetricsSmoother::update(int64_t newInstanceID,
int64_t newDurableBytes,
int64_t inputBytes,
const StorageBytes& storageBytes,
Smoother& smoothTotalDurableBytes) {
const bool reset = !instanceID.present() || instanceID.get() != newInstanceID;
if (reset) {
smoothDurableBytes.reset(newDurableBytes);
verySmoothDurableBytes.reset(newDurableBytes);
smoothInputBytes.reset(inputBytes);
smoothFreeSpace.reset(storageBytes.available);
smoothTotalSpace.reset(storageBytes.total);
} else {
smoothTotalDurableBytes.addDelta(newDurableBytes - durableBytes);
smoothDurableBytes.setTotal(newDurableBytes);
verySmoothDurableBytes.setTotal(newDurableBytes);
smoothInputBytes.setTotal(inputBytes);
smoothFreeSpace.setTotal(storageBytes.available);
smoothTotalSpace.setTotal(storageBytes.total);
}
instanceID = newInstanceID;
durableBytes = newDurableBytes;
return reset;
}
StorageQueueInfo::StorageQueueInfo(const UID& ratekeeperID_, const UID& id_, const LocalityData& locality_)
: ratekeeperID(ratekeeperID_), smoothDurableVersion(SERVER_KNOBS->SMOOTHING_AMOUNT),
smoothLatestVersion(SERVER_KNOBS->SMOOTHING_AMOUNT), valid(false), id(id_), locality(locality_),
acceptingRequests(false), limitReason(limitReason_t::unlimited) {}
StorageQueueInfo::StorageQueueInfo(const UID& id_, const LocalityData& locality_)
: StorageQueueInfo(UID(), id_, locality_) {}
@ -1184,23 +1209,12 @@ void StorageQueueInfo::addCommitCost(TransactionTagRef tagName, TransactionCommi
void StorageQueueInfo::update(StorageQueuingMetricsReply const& reply, Smoother& smoothTotalDurableBytes) {
valid = true;
auto prevReply = std::move(lastReply);
lastReply = reply;
if (prevReply.instanceID != reply.instanceID) {
smoothDurableBytes.reset(reply.bytesDurable);
verySmoothDurableBytes.reset(reply.bytesDurable);
smoothInputBytes.reset(reply.bytesInput);
smoothFreeSpace.reset(reply.storageBytes.available);
smoothTotalSpace.reset(reply.storageBytes.total);
if (queueMetrics.update(
reply.instanceID, reply.bytesDurable, reply.bytesInput, reply.storageBytes, smoothTotalDurableBytes)) {
smoothDurableVersion.reset(reply.durableVersion);
smoothLatestVersion.reset(reply.version);
} else {
smoothTotalDurableBytes.addDelta(reply.bytesDurable - prevReply.bytesDurable);
smoothDurableBytes.setTotal(reply.bytesDurable);
verySmoothDurableBytes.setTotal(reply.bytesDurable);
smoothInputBytes.setTotal(reply.bytesInput);
smoothFreeSpace.setTotal(reply.storageBytes.available);
smoothTotalSpace.setTotal(reply.storageBytes.total);
smoothDurableVersion.setTotal(reply.durableVersion);
smoothLatestVersion.setTotal(reply.version);
}
@ -1257,31 +1271,11 @@ UpdateCommitCostRequest StorageQueueInfo::refreshCommitCost(double elapsed) {
return updateCommitCostRequest;
}
TLogQueueInfo::TLogQueueInfo(UID id)
: valid(false), id(id), smoothDurableBytes(SERVER_KNOBS->SMOOTHING_AMOUNT),
smoothInputBytes(SERVER_KNOBS->SMOOTHING_AMOUNT), verySmoothDurableBytes(SERVER_KNOBS->SLOW_SMOOTHING_AMOUNT),
smoothFreeSpace(SERVER_KNOBS->SMOOTHING_AMOUNT), smoothTotalSpace(SERVER_KNOBS->SMOOTHING_AMOUNT) {
// FIXME: this is a tacky workaround for a potential uninitialized use in trackTLogQueueInfo (copied
// from storageQueueInfO)
lastReply.instanceID = -1;
}
TLogQueueInfo::TLogQueueInfo(UID id) : valid(false), id(id) {}
void TLogQueueInfo::update(TLogQueuingMetricsReply const& reply, Smoother& smoothTotalDurableBytes) {
valid = true;
auto prevReply = lastReply;
lastReply = reply;
if (prevReply.instanceID != reply.instanceID) {
smoothDurableBytes.reset(reply.bytesDurable);
verySmoothDurableBytes.reset(reply.bytesDurable);
smoothInputBytes.reset(reply.bytesInput);
smoothFreeSpace.reset(reply.storageBytes.available);
smoothTotalSpace.reset(reply.storageBytes.total);
} else {
smoothTotalDurableBytes.addDelta(reply.bytesDurable - prevReply.bytesDurable);
smoothDurableBytes.setTotal(reply.bytesDurable);
verySmoothDurableBytes.setTotal(reply.bytesDurable);
smoothInputBytes.setTotal(reply.bytesInput);
smoothFreeSpace.setTotal(reply.storageBytes.available);
smoothTotalSpace.setTotal(reply.storageBytes.total);
}
queueMetrics.update(
reply.instanceID, reply.bytesDurable, reply.bytesInput, reply.storageBytes, smoothTotalDurableBytes);
}

View File

@ -34,6 +34,30 @@
struct ServerDBInfo;
class QueueMetricsSmoother {
Optional<int64_t> instanceID;
int64_t durableBytes{ 0 };
Smoother smoothDurableBytes;
Smoother smoothInputBytes;
Smoother verySmoothDurableBytes;
Smoother smoothFreeSpace;
Smoother smoothTotalSpace;
public:
QueueMetricsSmoother();
bool update(int64_t newInstanceID,
int64_t newDurableBytes,
int64_t inputBytes,
const StorageBytes& storageBytes,
Smoother& smoothTotalDurableBytes);
double getSmoothFreeSpace() const { return smoothFreeSpace.smoothTotal(); }
double getSmoothTotalSpace() const { return smoothTotalSpace.smoothTotal(); }
double getSmoothDurableBytes() const { return smoothDurableBytes.smoothTotal(); }
double getSmoothInputBytesRate() const { return smoothInputBytes.smoothRate(); }
double getVerySmoothDurableBytesRate() const { return verySmoothDurableBytes.smoothRate(); }
};
class StorageQueueInfo {
uint64_t totalWriteCosts{ 0 };
int totalWriteOps{ 0 };
@ -41,8 +65,7 @@ class StorageQueueInfo {
TransactionTagMap<TransactionCommitCostEstimation> tagCostEst;
UID ratekeeperID;
Smoother smoothFreeSpace, smoothTotalSpace;
Smoother smoothDurableBytes, smoothInputBytes, verySmoothDurableBytes;
QueueMetricsSmoother queueMetrics;
Smoother smoothDurableVersion, smoothLatestVersion;
public:
@ -57,35 +80,33 @@ public:
StorageQueueInfo(const UID& id, const LocalityData& locality);
StorageQueueInfo(const UID& rateKeeperID, const UID& id, const LocalityData& locality);
UpdateCommitCostRequest refreshCommitCost(double elapsed);
int64_t getStorageQueueBytes() const { return lastReply.bytesInput - smoothDurableBytes.smoothTotal(); }
int64_t getStorageQueueBytes() const { return lastReply.bytesInput - queueMetrics.getSmoothDurableBytes(); }
int64_t getDurabilityLag() const { return smoothLatestVersion.smoothTotal() - smoothDurableVersion.smoothTotal(); }
void update(StorageQueuingMetricsReply const&, Smoother& smoothTotalDurableBytes);
void addCommitCost(TransactionTagRef tagName, TransactionCommitCostEstimation const& cost);
double getSmoothFreeSpace() const { return smoothFreeSpace.smoothTotal(); }
double getSmoothTotalSpace() const { return smoothTotalSpace.smoothTotal(); }
double getSmoothDurableBytes() const { return smoothDurableBytes.smoothTotal(); }
double getSmoothInputBytesRate() const { return smoothInputBytes.smoothRate(); }
double getVerySmoothDurableBytesRate() const { return verySmoothDurableBytes.smoothRate(); }
double getSmoothFreeSpace() const { return queueMetrics.getSmoothFreeSpace(); }
double getSmoothTotalSpace() const { return queueMetrics.getSmoothTotalSpace(); }
double getSmoothDurableBytes() const { return queueMetrics.getSmoothDurableBytes(); }
double getSmoothInputBytesRate() const { return queueMetrics.getSmoothInputBytesRate(); }
double getVerySmoothDurableBytesRate() const { return queueMetrics.getVerySmoothDurableBytesRate(); }
Version getLatestVersion() const { return lastReply.version; }
};
class TLogQueueInfo {
Smoother smoothDurableBytes, smoothInputBytes, verySmoothDurableBytes;
Smoother smoothFreeSpace;
Smoother smoothTotalSpace;
QueueMetricsSmoother queueMetrics;
public:
TLogQueuingMetricsReply lastReply;
bool valid;
UID id;
double getSmoothFreeSpace() const { return smoothFreeSpace.smoothTotal(); }
double getSmoothTotalSpace() const { return smoothTotalSpace.smoothTotal(); }
double getSmoothDurableBytes() const { return smoothDurableBytes.smoothTotal(); }
double getSmoothInputBytesRate() const { return smoothInputBytes.smoothRate(); }
double getVerySmoothDurableBytesRate() const { return verySmoothDurableBytes.smoothRate(); }
double getSmoothFreeSpace() const { return queueMetrics.getSmoothFreeSpace(); }
double getSmoothTotalSpace() const { return queueMetrics.getSmoothTotalSpace(); }
double getSmoothDurableBytes() const { return queueMetrics.getSmoothDurableBytes(); }
double getSmoothInputBytesRate() const { return queueMetrics.getSmoothInputBytesRate(); }
double getVerySmoothDurableBytesRate() const { return queueMetrics.getVerySmoothDurableBytesRate(); }
explicit TLogQueueInfo(UID id);
Version getLastCommittedVersion() const { return lastReply.v; }