Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 13 additions & 5 deletions src/commands/cmd_server.cc
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ class CommandNamespace : public Commander {
WARN("New namespace: {} with token: {}, addr: {}, result: {}", args_[2], args_[3], conn->GetAddr(), s.Msg());
} else if (args_.size() == 3 && sub_command == "del") {
Status s = srv->GetNamespace()->Del(args_[2]);
if (s.IsOK()) srv->ClearNamespaceStats(args_[2]);
*output = s.IsOK() ? redis::RESP_OK : redis::Error(s);
WARN("Deleted namespace: {}, addr: {}, result: {}", args_[2], conn->GetAddr(), s.Msg());
} else if (args_.size() == 2 && sub_command == "current") {
Expand Down Expand Up @@ -1722,16 +1723,23 @@ class CommandLatency : public Commander {
return Status::OK();
}

// Histograms are per-namespace: use the caller's namespace stats, or the aggregate for the
// admin/default namespace. Hold the shared_ptr for the duration of the response build.
auto stats_holder = conn->GetNamespace() == kDefaultNamespace
? srv->AggregateNamespaceStats()
: srv->GetOrCreateNamespaceStats(conn->GetNamespace());
const Stats &cmd_stats = *stats_holder;

std::vector<const std::pair<const std::string, CommandHistogram> *> target_histograms;
if (args_.size() > 2) {
for (size_t i = 2; i < args_.size(); i++) {
auto it = srv->stats.commands_histogram.find(util::ToLower(args_[i]));
if (it != srv->stats.commands_histogram.end() && it->second.calls > 0) {
auto it = cmd_stats.commands_histogram.find(util::ToLower(args_[i]));
if (it != cmd_stats.commands_histogram.end() && it->second.calls > 0) {
target_histograms.push_back(&(*it));
}
}
} else {
for (const auto &iter : srv->stats.commands_histogram) {
for (const auto &iter : cmd_stats.commands_histogram) {
if (iter.second.calls > 0) {
target_histograms.push_back(&iter);
}
Expand All @@ -1750,8 +1758,8 @@ class CommandLatency : public Commander {
if (cumulative == 0) continue;

int64_t boundary = 0;
if (i < srv->stats.bucket_boundaries.size()) {
boundary = static_cast<int64_t>(srv->stats.bucket_boundaries[i]);
if (i < cmd_stats.bucket_boundaries.size()) {
boundary = static_cast<int64_t>(cmd_stats.bucket_boundaries[i]);
} else {
boundary = -1;
}
Expand Down
13 changes: 11 additions & 2 deletions src/server/redis_connection.cc
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,12 @@ Connection::Connection(bufferevent *bev, Worker *owner)
int64_t now = util::GetTimeStamp();
create_time_ = now;
last_interaction_ = now;
cached_ns_stats_ = srv_->GetOrCreateNamespaceStats(kDefaultNamespace);
}

void Connection::SetNamespace(std::string ns) {
ns_ = std::move(ns);
cached_ns_stats_ = srv_->GetOrCreateNamespaceStats(ns_);
}

Connection::~Connection() {
Expand Down Expand Up @@ -399,7 +405,10 @@ void Connection::RecordProfilingSampleIfNeed(const std::string &cmd, uint64_t du
Status Connection::ExecuteCommand(engine::Context &ctx, const std::string &cmd_name,
const std::vector<std::string> &cmd_tokens, Commander *current_cmd,
std::string *reply) {
srv_->stats.IncrCalls(cmd_name);
// Attribute command stats to this connection's namespace. Hold the cached pointer locally so calls
// and latency stay consistent even if the command (e.g. AUTH/SELECT/RESET) changes the namespace.
auto ns_stats = cached_ns_stats_;
ns_stats->IncrCalls(cmd_name);

auto start = std::chrono::high_resolution_clock::now();
bool is_profiling = IsProfilingEnabled(cmd_name);
Expand All @@ -409,7 +418,7 @@ Status Connection::ExecuteCommand(engine::Context &ctx, const std::string &cmd_n
if (is_profiling) RecordProfilingSampleIfNeed(cmd_name, duration);

srv_->SlowlogPushEntryIfNeeded(&cmd_tokens, duration, this);
srv_->stats.IncrLatency(static_cast<uint64_t>(duration), cmd_name);
ns_stats->IncrLatency(static_cast<uint64_t>(duration), cmd_name);
return s;
}

Expand Down
5 changes: 4 additions & 1 deletion src/server/redis_connection.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
#include "event_util.h"
#include "redis_request.h"
#include "server/redis_reply.h"
#include "stats/stats.h"

class Worker;

Expand Down Expand Up @@ -163,7 +164,7 @@ class Connection : public EvbufCallbackBase<Connection> {
void BecomeAdmin() { is_admin_ = true; }
void BecomeUser() { is_admin_ = false; }
std::string GetNamespace() const { return ns_; }
void SetNamespace(std::string ns) { ns_ = std::move(ns); }
void SetNamespace(std::string ns);

void NeedFreeBufferEvent(bool need_free = true) { need_free_bev_ = need_free; }
void NeedNotFreeBufferEvent() { NeedFreeBufferEvent(false); }
Expand Down Expand Up @@ -210,6 +211,8 @@ class Connection : public EvbufCallbackBase<Connection> {
uint64_t id_ = 0;
std::atomic<int> flags_ = 0;
std::string ns_;
// Cache of this connection's per-namespace command stats, refreshed by SetNamespace.
std::shared_ptr<Stats> cached_ns_stats_;
std::string name_;
SetInfo set_info_;
std::string ip_;
Expand Down
142 changes: 109 additions & 33 deletions src/server/server.cc
Original file line number Diff line number Diff line change
Expand Up @@ -66,21 +66,7 @@ Server::Server(engine::Storage *storage, Config *config)
config_(config),
namespace_(storage) {
// init commands stats here to prevent concurrent insert, and cause core
auto commands = redis::CommandTable::GetOriginal();

for (const auto &iter : *commands) {
stats.commands_stats[iter.first].calls = 0;
stats.commands_stats[iter.first].latency = 0;

if (stats.bucket_boundaries.size() > 0) {
// NB: Extra index for the last bucket (Inf)
for (std::size_t i{0}; i <= stats.bucket_boundaries.size(); ++i) {
stats.commands_histogram[iter.first].buckets.push_back(std::make_unique<std::atomic<uint64_t>>(0));
}
stats.commands_histogram[iter.first].calls = 0;
stats.commands_histogram[iter.first].sum = 0;
}
}
initCommandStats(&stats);

// init cursor_dict_
cursor_dict_ = std::make_unique<CursorDictType>();
Expand Down Expand Up @@ -872,7 +858,18 @@ uint64_t Server::GetClientID() { return client_id_.fetch_add(1, std::memory_orde

void Server::recordInstantaneousMetrics() {
auto rocksdb_stats = storage->GetDB()->GetDBOptions().statistics;
stats.TrackInstantaneousMetric(STATS_METRIC_COMMAND, stats.total_calls);
// Sample each namespace's command metric, and feed the sum into the global metric so the
// admin/default view reports aggregate ops/sec without keeping a global command counter on the hot path.
uint64_t total_calls = 0;
{
std::shared_lock<std::shared_mutex> lock(ns_stats_mu_);
for (const auto &[ns, ns_stats] : ns_stats_) {
auto calls = ns_stats->total_calls.load();
ns_stats->TrackInstantaneousMetric(STATS_METRIC_COMMAND, calls);
total_calls += calls;
}
}
stats.TrackInstantaneousMetric(STATS_METRIC_COMMAND, total_calls);
Comment thread
git-hulk marked this conversation as resolved.
stats.TrackInstantaneousMetric(STATS_METRIC_NET_INPUT, stats.in_bytes);
stats.TrackInstantaneousMetric(STATS_METRIC_NET_OUTPUT, stats.out_bytes);
stats.TrackInstantaneousMetric(STATS_METRIC_ROCKSDB_PUT,
Expand Down Expand Up @@ -1392,11 +1389,82 @@ int64_t Server::GetLastBgsaveTime() {
return last_bgsave_timestamp_secs_ == -1 ? start_time_secs_ : last_bgsave_timestamp_secs_;
}

Server::InfoEntries Server::GetStatsInfo() {
void Server::initCommandStats(Stats *stats) {
auto commands = redis::CommandTable::GetOriginal();
for (const auto &iter : *commands) {
stats->commands_stats[iter.first].calls = 0;
stats->commands_stats[iter.first].latency = 0;

if (stats->bucket_boundaries.size() > 0) {
// NB: Extra index for the last bucket (Inf)
for (std::size_t i{0}; i <= stats->bucket_boundaries.size(); ++i) {
stats->commands_histogram[iter.first].buckets.push_back(std::make_unique<std::atomic<uint64_t>>(0));
}
stats->commands_histogram[iter.first].calls = 0;
stats->commands_histogram[iter.first].sum = 0;
}
}
}

std::shared_ptr<Stats> Server::GetOrCreateNamespaceStats(const std::string &ns) {
{
std::shared_lock<std::shared_mutex> lock(ns_stats_mu_);
if (auto it = ns_stats_.find(ns); it != ns_stats_.end()) {
return it->second;
}
}

std::unique_lock<std::shared_mutex> lock(ns_stats_mu_);
if (auto it = ns_stats_.find(ns); it != ns_stats_.end()) {
return it->second;
}
auto ns_stats = std::make_shared<Stats>(config_->histogram_bucket_boundaries);
initCommandStats(ns_stats.get());
ns_stats_[ns] = ns_stats;
return ns_stats;
Comment thread
git-hulk marked this conversation as resolved.
}

std::shared_ptr<Stats> Server::AggregateNamespaceStats() {
auto agg = std::make_shared<Stats>(config_->histogram_bucket_boundaries);
initCommandStats(agg.get());

std::shared_lock<std::shared_mutex> lock(ns_stats_mu_);
for (const auto &[ns, ns_stats] : ns_stats_) {
agg->total_calls.fetch_add(ns_stats->total_calls.load(), std::memory_order_relaxed);
for (const auto &[cmd, stat] : ns_stats->commands_stats) {
agg->commands_stats[cmd].calls.fetch_add(stat.calls.load(), std::memory_order_relaxed);
agg->commands_stats[cmd].latency.fetch_add(stat.latency.load(), std::memory_order_relaxed);
}
for (const auto &[cmd, hist] : ns_stats->commands_histogram) {
auto &agg_hist = agg->commands_histogram[cmd];
agg_hist.calls.fetch_add(hist.calls.load(), std::memory_order_relaxed);
agg_hist.sum.fetch_add(hist.sum.load(), std::memory_order_relaxed);
for (std::size_t i = 0; i < hist.buckets.size(); ++i) {
agg_hist.buckets[i]->fetch_add(hist.buckets[i]->load(), std::memory_order_relaxed);
}
}
}
Comment thread
git-hulk marked this conversation as resolved.
return agg;
}

void Server::ClearNamespaceStats(const std::string &ns) {
std::unique_lock<std::shared_mutex> lock(ns_stats_mu_);
ns_stats_.erase(ns);
}

Server::InfoEntries Server::GetStatsInfo(const std::string &ns) {
// Command stats are per namespace; the admin/default namespace sees the aggregate across all of them.
auto cmd_stats_ptr = ns == kDefaultNamespace ? AggregateNamespaceStats() : GetOrCreateNamespaceStats(ns);
const Stats &cmd_stats = *cmd_stats_ptr;

Server::InfoEntries entries;
entries.emplace_back("total_connections_received", total_clients_.load());
entries.emplace_back("total_commands_processed", stats.total_calls.load());
entries.emplace_back("instantaneous_ops_per_sec", stats.GetInstantaneousMetric(STATS_METRIC_COMMAND));
entries.emplace_back("total_commands_processed", cmd_stats.total_calls.load());
// Per-namespace ops/sec comes from the namespace's own sampled metric; the admin/default view uses
// the global metric, which the sampler feeds with the sum across all namespaces.
auto ops_per_sec = ns == kDefaultNamespace ? stats.GetInstantaneousMetric(STATS_METRIC_COMMAND)
: cmd_stats.GetInstantaneousMetric(STATS_METRIC_COMMAND);
entries.emplace_back("instantaneous_ops_per_sec", ops_per_sec);
Comment thread
git-hulk marked this conversation as resolved.
entries.emplace_back("total_net_input_bytes", stats.in_bytes.load());
entries.emplace_back("total_net_output_bytes", stats.out_bytes.load());
entries.emplace_back("instantaneous_input_kbps",
Expand All @@ -1420,10 +1488,13 @@ Server::InfoEntries Server::GetStatsInfo() {
return entries;
}

Server::InfoEntries Server::GetCommandsStatsInfo() {
Server::InfoEntries Server::GetCommandsStatsInfo(const std::string &ns) {
auto cmd_stats_ptr = ns == kDefaultNamespace ? AggregateNamespaceStats() : GetOrCreateNamespaceStats(ns);
const Stats &cmd_stats = *cmd_stats_ptr;

InfoEntries entries;

for (const auto &cmd_stat : stats.commands_stats) {
for (const auto &cmd_stat : cmd_stats.commands_stats) {
auto calls = cmd_stat.second.calls.load();
if (calls == 0) continue;

Expand All @@ -1433,18 +1504,18 @@ Server::InfoEntries Server::GetCommandsStatsInfo() {
static_cast<double>(latency) / static_cast<double>(calls)));
}

for (const auto &cmd_hist : stats.commands_histogram) {
for (const auto &cmd_hist : cmd_stats.commands_histogram) {
auto command_name = cmd_hist.first;
auto calls = stats.commands_histogram[command_name].calls.load();
auto calls = cmd_hist.second.calls.load();
if (calls == 0) continue;

auto sum = stats.commands_histogram[command_name].sum.load();
auto sum = cmd_hist.second.sum.load();
std::string result;
for (std::size_t i{0}; i < stats.commands_histogram[command_name].buckets.size(); ++i) {
auto bucket_value = stats.commands_histogram[command_name].buckets[i]->load();
for (std::size_t i{0}; i < cmd_hist.second.buckets.size(); ++i) {
auto bucket_value = cmd_hist.second.buckets[i]->load();
auto bucket_bound = std::numeric_limits<double>::infinity();
if (i < stats.bucket_boundaries.size()) {
bucket_bound = stats.bucket_boundaries[i];
if (i < cmd_stats.bucket_boundaries.size()) {
bucket_bound = cmd_stats.bucket_boundaries[i];
}

result.append(fmt::format("{}={},", bucket_bound, bucket_value));
Expand Down Expand Up @@ -1539,11 +1610,16 @@ Server::InfoEntries Server::GetKeyspaceInfo(const std::string &ns) {
// this section can't be shown when loading(i.e. !is_loading_).
std::string Server::GetInfo(const std::string &ns, const std::vector<std::string> &sections, InfoFormat format) {
std::vector<std::pair<std::string, std::function<InfoEntries(Server *)>>> info_funcs = {
{"Server", &Server::GetServerInfo}, {"Clients", &Server::GetClientsInfo},
{"Memory", &Server::GetMemoryInfo}, {"Persistence", &Server::GetPersistenceInfo},
{"Stats", &Server::GetStatsInfo}, {"Replication", &Server::GetReplicationInfo},
{"CPU", &Server::GetCpuInfo}, {"CommandStats", &Server::GetCommandsStatsInfo},
{"Cluster", &Server::GetClusterInfo}, {"Keyspace", [&ns](Server *srv) { return srv->GetKeyspaceInfo(ns); }},
{"Server", &Server::GetServerInfo},
{"Clients", &Server::GetClientsInfo},
{"Memory", &Server::GetMemoryInfo},
{"Persistence", &Server::GetPersistenceInfo},
{"Stats", [&ns](Server *srv) { return srv->GetStatsInfo(ns); }},
{"Replication", &Server::GetReplicationInfo},
{"CPU", &Server::GetCpuInfo},
{"CommandStats", [&ns](Server *srv) { return srv->GetCommandsStatsInfo(ns); }},
{"Cluster", &Server::GetClusterInfo},
{"Keyspace", [&ns](Server *srv) { return srv->GetKeyspaceInfo(ns); }},
{"RocksDB", &Server::GetRocksDBInfo},
};

Expand Down
17 changes: 15 additions & 2 deletions src/server/server.h
Original file line number Diff line number Diff line change
Expand Up @@ -301,18 +301,24 @@ class Server {
};
using InfoEntries = std::vector<InfoEntry>;

InfoEntries GetStatsInfo();
InfoEntries GetStatsInfo(const std::string &ns);
InfoEntries GetServerInfo();
InfoEntries GetMemoryInfo();
InfoEntries GetRocksDBInfo();
InfoEntries GetClientsInfo();
InfoEntries GetReplicationInfo();
InfoEntries GetCommandsStatsInfo();
InfoEntries GetCommandsStatsInfo(const std::string &ns);
InfoEntries GetClusterInfo();
InfoEntries GetPersistenceInfo();
InfoEntries GetCpuInfo();
InfoEntries GetKeyspaceInfo(const std::string &ns);

// Per-namespace command statistics. Command calls/latency are tracked per namespace (keyed by the
// connection's namespace); the admin/default namespace view is the sum over all namespaces.
std::shared_ptr<Stats> GetOrCreateNamespaceStats(const std::string &ns);
std::shared_ptr<Stats> AggregateNamespaceStats();
void ClearNamespaceStats(const std::string &ns);

enum class InfoFormat { Text, Json };
std::string GetInfo(const std::string &ns, const std::vector<std::string> &sections,
InfoFormat format = InfoFormat::Text);
Expand Down Expand Up @@ -444,6 +450,13 @@ class Server {

std::map<std::string, DBScanInfo> db_scan_infos_;

// Per-namespace command statistics (keyed by namespace name), guarded by ns_stats_mu_. The global
// `stats` keeps the non-namespaced counters (net bytes, replication) and the sampled aggregate ops/sec.
std::unordered_map<std::string, std::shared_ptr<Stats>> ns_stats_;
Comment thread
PragmaTwice marked this conversation as resolved.
std::shared_mutex ns_stats_mu_;
// Pre-populate a Stats' per-command maps so no runtime map insertion (and thus no data race) happens.
static void initCommandStats(Stats *stats);

LogCollector<SlowEntry> slow_log_;
LogCollector<PerfEntry> perf_log_;

Expand Down
Loading
Loading