From 5a10f3572674416893ff45ee63ed7e4f20a0080a Mon Sep 17 00:00:00 2001 From: yuanhao Date: Fri, 14 Aug 2026 19:07:46 +0800 Subject: [PATCH 1/6] =?UTF-8?q?fix(tracing):=20=E6=B7=BB=E5=8A=A0=E7=BC=BA?= =?UTF-8?q?=E5=A4=B1=E7=9A=84=20span=20=E7=BB=93=E6=9D=9F=E8=B0=83?= =?UTF-8?q?=E7=94=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- mooncake-store/src/client_service.cpp | 29 +++++++++++++++++++++++++-- mooncake-store/src/master_service.cpp | 1 + 2 files changed, 28 insertions(+), 2 deletions(-) diff --git a/mooncake-store/src/client_service.cpp b/mooncake-store/src/client_service.cpp index 254eb7bf0d..501e5567fc 100644 --- a/mooncake-store/src/client_service.cpp +++ b/mooncake-store/src/client_service.cpp @@ -1693,6 +1693,7 @@ tl::expected Client::Put(const ObjectKey& key, auto checksum_result = ComputeObjectChecksumForSlices( key, slices, CalculateSliceSize(slices)); if (!checksum_result) { + pt_full.End(-1); return tl::unexpected(checksum_result.error()); } object_checksum = *checksum_result; @@ -1719,6 +1720,7 @@ tl::expected Client::Put(const ObjectKey& key, ErrorCode err = start_result.error(); if (err == ErrorCode::OBJECT_ALREADY_EXISTS) { VLOG(1) << "object_already_exists key=" << key; + pt_full.End(0); return {}; } if (err == ErrorCode::NO_AVAILABLE_HANDLE) { @@ -1728,6 +1730,7 @@ tl::expected Client::Put(const ObjectKey& key, LOG(ERROR) << "Failed to start put operation for key=" << key << ": " << toString(err); } + pt_full.End(-1); return tl::unexpected(err); } @@ -1789,6 +1792,7 @@ tl::expected Client::Put(const ObjectKey& key, if (!end_result) { ErrorCode err = end_result.error(); LOG(ERROR) << "Failed to end put operation: " << err; + pt_full.End(-1); return tl::unexpected(err); } } @@ -1798,14 +1802,17 @@ tl::expected Client::Put(const ObjectKey& key, master_client_.PutRevoke(key, *finalize_decision.revoke_type); if (!revoke_result) { LOG(ERROR) << "Failed to revoke put operation"; + pt_full.End(-1); return tl::unexpected(revoke_result.error()); } } if (!finalize_decision.success) { + pt_full.End(-1); return tl::unexpected(finalize_decision.error); } + pt_full.End(0); return {}; } @@ -2138,9 +2145,15 @@ void Client::StartBatchPut(std::vector& ops, op.SetError(ErrorCode::RPC_FAIL, "BatchPutStart response size mismatch"); } + pt_batch_start.End(-1); return; } + const bool any_start_succeeded = + std::any_of(start_responses.begin(), start_responses.end(), + [](const auto& response) { return response.has_value(); }); + pt_batch_start.End(any_start_succeeded ? 0 : -1); + // Process individual responses with robust error handling for (size_t i = 0; i < active_indices.size(); ++i) { auto& op = ops[active_indices[i]]; @@ -2828,17 +2841,24 @@ std::vector> Client::BatchPut( if (client_cfg.nof_replica_num > 0) { LOG(ERROR) << "prefer_alloc_in_same_node is not supported with " "NoF replicas"; + pt_full.End(-1); return std::vector>( keys.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); } if (client_cfg.replica_num != 1) { LOG(ERROR) << "prefer_alloc_in_same_node is not supported with " "replica_num != 1"; + pt_full.End(-1); return std::vector>( keys.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); } StartBatchPut(ops, client_cfg); - return BatchPutWhenPreferSameNode(ops); + auto results = BatchPutWhenPreferSameNode(ops); + const bool any_succeeded = + std::any_of(results.begin(), results.end(), + [](const auto& result) { return result.has_value(); }); + pt_full.End(any_succeeded ? 0 : -1); + return results; } StartBatchPut(ops, client_cfg); @@ -2853,7 +2873,12 @@ std::vector> Client::BatchPut( } FinalizeBatchPut(ops); - return CollectResults(ops); + auto results = CollectResults(ops); + const bool any_succeeded = + std::any_of(results.begin(), results.end(), + [](const auto& result) { return result.has_value(); }); + pt_full.End(any_succeeded ? 0 : -1); + return results; } tl::expected Client::Remove(const ObjectKey& key, bool force) { diff --git a/mooncake-store/src/master_service.cpp b/mooncake-store/src/master_service.cpp index a319d1aef3..d5660e8e74 100644 --- a/mooncake-store/src/master_service.cpp +++ b/mooncake-store/src/master_service.cpp @@ -9316,6 +9316,7 @@ void MasterService::ClientMonitorFunc() { } } RecomputeTenantEffectiveQuotas(); + pt_unmount.End(0); } pt_monitor.End(0); From b874ea58320b08087efb33f16e3b7209b49e8928 Mon Sep 17 00:00:00 2001 From: ZhiningPan Date: Thu, 27 Aug 2026 19:09:07 +0800 Subject: [PATCH 2/6] =?UTF-8?q?feat(spdiag):=20=E4=B8=BA=20vLLM=20?= =?UTF-8?q?=E8=B0=83=E7=94=A8=E8=B7=AF=E5=BE=84=E6=B7=BB=E5=8A=A0=20PerfPo?= =?UTF-8?q?int=20=E6=89=93=E7=82=B9=E4=B8=8E=20MC=5FLOG=20=E6=97=A5?= =?UTF-8?q?=E5=BF=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 依据 vllm_spdiag_logging_plan.md,按 Q1b/Q2A/Q3b 方案为 Mooncake Store vLLM 调用链添加 SpDiag 性能打点与日志输出: - mooncake_perf_points.def: 追加 25 个新 PerfKey,覆盖 store_py 入口 层(STORE_PY_*)、RealClient 核心层(RC_*)、Client 服务层 (CLIENT_BATCH_QUERY)。 - store_py.cpp: 为 setup/register_buffer/batch_put_from_multi_buffers/ batch_get_into_multi_buffers/batchIsExist/remove_all/close 等 8 个 Python 绑定方法添加 PerfPoint(KEY_MODULE),remove_all 与 close 额外输出 MC_LOG 汇总行。 - real_client.cpp: 为 setup_real/batchIsExist/batch_put_from_multi_ buffers/batch_get_into_multi_buffers 等 9 个方法添加 PerfPoint; 下沉层 _internal 方法用 MODULE 级别,并按 Q1b 输出 MC_LOG 汇总行 + per-key 多行(包含 success/key/size/replica/endpoint 字段, 缺失字段不输出)。 - client_service.cpp: 在 L1117 BatchQuery 实际实现中添加 CLIENT_BATCH_QUERY PerfPoint(Q3b,仅实际实现,避免重复统计)。 --- .../store/mooncake_perf_points.def | 35 +++++ mooncake-integration/store/store_py.cpp | 85 ++++++++-- mooncake-store/src/client_service.cpp | 5 + mooncake-store/src/real_client.cpp | 145 +++++++++++++++++- 4 files changed, 254 insertions(+), 16 deletions(-) diff --git a/mooncake-integration/store/mooncake_perf_points.def b/mooncake-integration/store/mooncake_perf_points.def index e1b45a7e78..eaf6dd992d 100644 --- a/mooncake-integration/store/mooncake_perf_points.def +++ b/mooncake-integration/store/mooncake_perf_points.def @@ -179,3 +179,38 @@ PERF_KEY_DEF(MASTER_BG_DISCARD_EXPIRED, "master_service.cpp::EvictionThreadFun PERF_KEY_DEF(MASTER_BG_SNAPSHOT_PERSIST, "master_service.cpp::SnapshotThreadFunc", "SnapshotPersist") PERF_KEY_DEF(MASTER_BG_CLIENT_MONITOR, "master_service.cpp::ClientMonitorFunc", "ClientMonitorScan") PERF_KEY_DEF(MASTER_BG_CLIENT_UNMOUNT, "master_service.cpp::ClientMonitorFunc", "ExpiredClientUnmount") + +// ============================================================ +// === vLLM MooncakeStoreConnector 路径打点(vllm 0.26.1rc0)=== +// === 覆盖 S2-S9 入口 + T1-T8 下沉 + Client 服务层 === +// ============================================================ + +// === store_py.cpp Python 绑定层(vllm 直接调用入口,S2-S9)=== +PERF_KEY_DEF(STORE_PY_SETUP, "store_py.cpp::setup", "Setup") +PERF_KEY_DEF(STORE_PY_REGISTER_BUFFER, "store_py.cpp::register_buffer", "RegisterBuffer") +PERF_KEY_DEF(STORE_PY_BATCH_PUT_MULTI, "store_py.cpp::batch_put_from_multi_buffers","BatchPutMultiBuf") +PERF_KEY_DEF(STORE_PY_BATCH_GET_INTO_MULTI, "store_py.cpp::batch_get_into_multi_buffers","BatchGetIntoMultiBuf") +PERF_KEY_DEF(STORE_PY_BATCH_IS_EXIST, "store_py.cpp::batch_is_exist", "BatchIsExist") +PERF_KEY_DEF(STORE_PY_BATCH_GET_REPLICA_DESC, "store_py.cpp::batch_get_replica_desc", "BatchGetReplicaDesc") +PERF_KEY_DEF(STORE_PY_REMOVE_ALL, "store_py.cpp::remove_all", "RemoveAll") +PERF_KEY_DEF(STORE_PY_CLOSE, "store_py.cpp::close", "Close") + +// === RealClient 核心逻辑层(下沉路径,T1-T8)=== +PERF_KEY_DEF(RC_SETUP_REAL, "real_client.cpp::setup_real", "SetupReal") +PERF_KEY_DEF(RC_SETUP_INTERNAL, "real_client.cpp::setup_internal", "SetupInternal") +PERF_KEY_DEF(RC_TEARDOWN_ALL, "real_client.cpp::tearDownAll", "TeardownAll") +PERF_KEY_DEF(RC_TEARDOWN_ALL_INTERNAL, "real_client.cpp::tearDownAll_internal", "TeardownAllInternal") +PERF_KEY_DEF(RC_REMOVE_ALL, "real_client.cpp::removeAll", "RemoveAll") +PERF_KEY_DEF(RC_REMOVE_ALL_INTERNAL, "real_client.cpp::removeAll_internal", "RemoveAllInternal") +PERF_KEY_DEF(RC_BATCH_IS_EXIST, "real_client.cpp::batchIsExist", "BatchIsExist") +PERF_KEY_DEF(RC_BATCH_IS_EXIST_INTERNAL, "real_client.cpp::batchIsExist_internal", "BatchIsExistInternal") +PERF_KEY_DEF(RC_REGISTER_BUFFER, "real_client.cpp::register_buffer", "RegisterBuffer") +PERF_KEY_DEF(RC_REGISTER_BUFFER_INTERNAL, "real_client.cpp::register_buffer_internal","RegisterBufferInternal") +PERF_KEY_DEF(RC_BATCH_PUT_MULTI, "real_client.cpp::batch_put_from_multi_buffers", "BatchPutMultiBuf") +PERF_KEY_DEF(RC_BATCH_PUT_MULTI_INTERNAL, "real_client.cpp::batch_put_from_multi_buffers_internal","BatchPutMultiBufInternal") +PERF_KEY_DEF(RC_BATCH_GET_INTO_MULTI, "real_client.cpp::batch_get_into_multi_buffers", "BatchGetIntoMultiBuf") +PERF_KEY_DEF(RC_BATCH_GET_INTO_MULTI_INTERNAL, "real_client.cpp::batch_get_into_multi_buffers_internal","BatchGetIntoMultiBufInternal") +PERF_KEY_DEF(RC_BATCH_GET_REPLICA_DESC, "real_client.cpp::batch_get_replica_desc", "BatchGetReplicaDesc") + +// === Client 服务层(vllm 路径下沉,部分已有打点)=== +PERF_KEY_DEF(CLIENT_BATCH_QUERY, "client_service.cpp::BatchQuery", "BatchQuery") diff --git a/mooncake-integration/store/store_py.cpp b/mooncake-integration/store/store_py.cpp index 2970bbfd58..82f8e7f5d2 100644 --- a/mooncake-integration/store/store_py.cpp +++ b/mooncake-integration/store/store_py.cpp @@ -16,6 +16,7 @@ #include "memory_alloc.h" #include "ssd_register_client.h" #include "device/accelerator_registry.h" +#include "mooncake_logging.h" // MC_LOG #include // for atexit #include @@ -2084,6 +2085,9 @@ PYBIND11_MODULE(store, m) { const std::string &tenant_id = "default", bool enable_client_http_server = false, int client_http_port = DEFAULT_CLIENT_HTTP_PORT) { + SpDiag::PerfPoint pt(PerfKey::STORE_PY_SETUP, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); auto real_client = self.init_real_client(); std::shared_ptr transfer_engine = nullptr; @@ -2091,12 +2095,15 @@ PYBIND11_MODULE(store, m) { transfer_engine = engine.cast>(); } - return real_client->setup_real( + auto ret = real_client->setup_real( local_hostname, metadata_server, global_segment_size, local_buffer_size, protocol, rdma_devices, master_server_addr, transfer_engine, "", enable_ssd_offload, ssd_offload_path, tenant_id, enable_client_http_server, client_http_port); + pt.End(ret == 0 ? 0 : -1); + // MC_LOG 在下沉层 setup_real 输出(Q1b) + return ret; }, py::arg("local_hostname"), py::arg("metadata_server"), py::arg("global_segment_size"), py::arg("local_buffer_size"), @@ -2109,6 +2116,9 @@ PYBIND11_MODULE(store, m) { .def( "setup", [](MooncakeStorePyWrapper &self, const py::dict &config_dict) { + SpDiag::PerfPoint pt(PerfKey::STORE_PY_SETUP, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); auto real_client = self.init_real_client(); // Convert py::dict to ConfigDict (all values as strings) @@ -2120,8 +2130,11 @@ PYBIND11_MODULE(store, m) { } auto result = real_client->setup_internal(config); - return result.has_value() ? 0 - : static_cast(result.error()); + int ret = result.has_value() ? 0 + : static_cast(result.error()); + pt.End(ret == 0 ? 0 : -1); + // MC_LOG 在下沉层 setup_real 输出(Q1b) + return ret; }, py::arg("config"), "Setup the store with a configuration dictionary.\n" @@ -2230,8 +2243,16 @@ PYBIND11_MODULE(store, m) { .def( "remove_all", [](MooncakeStorePyWrapper &self, bool force) { + SpDiag::PerfPoint pt(PerfKey::STORE_PY_REMOVE_ALL, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); py::gil_scoped_release release; - return self.store_->removeAll(force); + auto ret = self.store_->removeAll(force); + pt.End(ret == 0 ? 0 : -1); + // 下沉 removeAll 不加 MC_LOG,入口层输出汇总(Q1b) + MC_LOG(INFO) << "[remove_all] elapsed_us=" << pt.ElapsedMicros() + << " success=" << (ret == 0 ? 1 : 0); + return ret; }, py::arg("force") = false, "Remove all objects from the store. If force=True, skip lease " @@ -2255,17 +2276,36 @@ PYBIND11_MODULE(store, m) { "batch_is_exist", [](MooncakeStorePyWrapper &self, const std::vector &keys) { + SpDiag::PerfPoint pt(PerfKey::STORE_PY_BATCH_IS_EXIST, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); py::gil_scoped_release release; - return self.store_->batchIsExist(keys); + auto ret = self.store_->batchIsExist(keys); + pt.End(0); + // MC_LOG 在下沉层 batchIsExist 输出 per-key(Q1b) + return ret; }, py::arg("keys"), "Check if multiple objects exist. Returns list of results: 1 if " "exists, 0 if not exists, -1 if error") .def("close", [](MooncakeStorePyWrapper &self) { - if (!self.store_) return 0; + SpDiag::PerfPoint pt(PerfKey::STORE_PY_CLOSE, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); + if (!self.store_) { + pt.End(0); + // 无下沉 MC_LOG,入口层输出汇总(Q1b) + MC_LOG(INFO) << "[close] elapsed_us=" << pt.ElapsedMicros() + << " success=1"; + return 0; + } int rc = self.store_->tearDownAll(); self.store_.reset(); + pt.End(rc == 0 ? 0 : -1); + // 下沉 tearDownAll 不加 MC_LOG,入口层输出汇总(Q1b) + MC_LOG(INFO) << "[close] elapsed_us=" << pt.ElapsedMicros() + << " success=" << (rc == 0 ? 1 : 0); return rc; }) .def("health_check", &MooncakeStorePyWrapper::health_check, @@ -2649,10 +2689,16 @@ PYBIND11_MODULE(store, m) { "register_buffer", [](MooncakeStorePyWrapper &self, uintptr_t buffer_ptr, size_t size) { + SpDiag::PerfPoint pt(PerfKey::STORE_PY_REGISTER_BUFFER, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); // Register memory buffer for RDMA operations void *buffer = reinterpret_cast(buffer_ptr); py::gil_scoped_release release; - return self.store_->register_buffer(buffer, size); + auto ret = self.store_->register_buffer(buffer, size); + pt.End(ret == 0 ? 0 : -1); + // MC_LOG 在下沉层 register_buffer 输出汇总(Q1b) + return ret; }, py::arg("buffer_ptr"), py::arg("size"), "Register a memory buffer for direct access operations") @@ -2897,13 +2943,20 @@ PYBIND11_MODULE(store, m) { const std::vector> &all_buffer_ptrs, const std::vector> &all_sizes, const ReplicateConfig &config = ReplicateConfig{}) { + SpDiag::PerfPoint pt(PerfKey::STORE_PY_BATCH_PUT_MULTI, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); if (!self.is_client_initialized()) { LOG(ERROR) << "Client is not initialized"; + pt.End(-1); return std::vector{}; } py::gil_scoped_release release; - return self.store_->batch_put_from_multi_buffers( + auto ret = self.store_->batch_put_from_multi_buffers( keys, CastAddrs2Ptrs(all_buffer_ptrs), all_sizes, config); + pt.End(ret.empty() ? -1 : 0); + // MC_LOG 在下沉层 *_internal 输出 汇总+per-key(Q1b) + return ret; }, py::arg("keys"), py::arg("all_buffer_ptrs"), py::arg("all_sizes"), py::arg("config") = ReplicateConfig{}, @@ -2917,10 +2970,16 @@ PYBIND11_MODULE(store, m) { const std::vector> &all_buffer_ptrs, const std::vector> &all_sizes, bool prefer_alloc_in_same_node = false) { + SpDiag::PerfPoint pt(PerfKey::STORE_PY_BATCH_GET_INTO_MULTI, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); py::gil_scoped_release release; - return self.store_->batch_get_into_multi_buffers( + auto ret = self.store_->batch_get_into_multi_buffers( keys, CastAddrs2Ptrs(all_buffer_ptrs), all_sizes, prefer_alloc_in_same_node); + pt.End(ret.empty() ? -1 : 0); + // MC_LOG 在下沉层 *_internal 输出 汇总+per-key(Q1b) + return ret; }, py::arg("keys"), py::arg("all_buffer_ptrs"), py::arg("all_sizes"), py::arg("prefer_alloc_in_same_node") = false, @@ -2938,8 +2997,14 @@ PYBIND11_MODULE(store, m) { "batch_get_replica_desc", [](MooncakeStorePyWrapper &self, const std::vector &keys) { + SpDiag::PerfPoint pt(PerfKey::STORE_PY_BATCH_GET_REPLICA_DESC, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); py::gil_scoped_release release; - return self.store_->batch_get_replica_desc(keys); + auto ret = self.store_->batch_get_replica_desc(keys); + pt.End(0); + // MC_LOG 在下沉层 batch_get_replica_desc 输出 per-key(Q1b) + return ret; }, py::arg("keys")) .def( diff --git a/mooncake-store/src/client_service.cpp b/mooncake-store/src/client_service.cpp index 501e5567fc..60d76a333b 100644 --- a/mooncake-store/src/client_service.cpp +++ b/mooncake-store/src/client_service.cpp @@ -1116,6 +1116,9 @@ std::vector> Client::BatchQuery( std::vector> Client::BatchQuery( const std::vector& object_keys, const std::string& tenant_id) { + SpDiag::PerfPoint pt(PerfKey::CLIENT_BATCH_QUERY, + SpDiag::PerfLevel::MODULE); + pt.Start(); std::chrono::steady_clock::time_point start_time = std::chrono::steady_clock::now(); auto response = master_client_.BatchGetReplicaList(object_keys, tenant_id); @@ -1130,6 +1133,7 @@ std::vector> Client::BatchQuery( for (size_t i = 0; i < object_keys.size(); ++i) { results.emplace_back(tl::unexpected(ErrorCode::RPC_FAIL)); } + pt.End(-1); return results; } std::vector> results; @@ -1145,6 +1149,7 @@ std::vector> Client::BatchQuery( results.emplace_back(tl::unexpected(response[i].error())); } } + pt.End(0); return results; } diff --git a/mooncake-store/src/real_client.cpp b/mooncake-store/src/real_client.cpp index c7eb5df178..7a184a3c52 100644 --- a/mooncake-store/src/real_client.cpp +++ b/mooncake-store/src/real_client.cpp @@ -1205,11 +1205,22 @@ int RealClient::setup_real( const std::string &ipc_socket_path, bool enable_ssd_offload, const std::string &ssd_offload_path, const std::string &tenant_id, bool enable_client_http_server, int client_http_port) { - return to_py_ret(setup_internal( + SpDiag::PerfPoint pt(PerfKey::RC_SETUP_REAL, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); + auto ret = to_py_ret(setup_internal( local_hostname, metadata_server, global_segment_size, local_buffer_size, protocol, rdma_devices, master_server_addr, transfer_engine, ipc_socket_path, 50052, enable_ssd_offload, true, ssd_offload_path, tenant_id, Environ::Get().GetOffloadRpcThreadNum(8), enable_client_http_server, client_http_port)); + pt.End(ret == 0 ? 0 : -1); + // S2 setup 下沉层 MC_LOG 汇总(Q1b:入口层只 PerfPoint) + MC_LOG(INFO) << "[setup] elapsed_us=" << pt.ElapsedMicros() + << " success=" << (ret == 0 ? 1 : 0) + << " hostname=" << local_hostname + << " metadata_server=" << metadata_server + << " protocol=" << protocol; + return ret; } namespace { @@ -2476,6 +2487,9 @@ int RealClient::isExist(const std::string &key) { std::vector RealClient::batchIsExist( const std::vector &keys) { + SpDiag::PerfPoint pt(PerfKey::RC_BATCH_IS_EXIST, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); auto internal_results = batchIsExist_internal(keys); std::vector results; results.reserve(internal_results.size()); @@ -2487,7 +2501,17 @@ std::vector RealClient::batchIsExist( results.push_back(toInt(result.error())); } } - + pt.End(0); + // S6 batch_is_exist 下沉层 MC_LOG 汇总+per-key(Q1b:入口层只 PerfPoint) + // 字段:success+key(无 size/replica/endpoint,缺失不输出) + std::ostringstream oss; + oss << "[batch_is_exist] elapsed_us=" << pt.ElapsedMicros() + << " num_keys=" << keys.size(); + for (size_t i = 0; i < keys.size(); ++i) { + int success = (i < results.size()) ? results[i] : 0; + oss << "\n success=" << success << " key=" << keys[i]; + } + MC_LOG(INFO) << oss.str(); return results; } @@ -3768,7 +3792,16 @@ tl::expected RealClient::register_buffer_internal( } int RealClient::register_buffer(void *buffer, size_t size) { - return to_py_ret(register_buffer_internal(buffer, size)); + SpDiag::PerfPoint pt(PerfKey::RC_REGISTER_BUFFER, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); + auto ret = to_py_ret(register_buffer_internal(buffer, size)); + pt.End(ret == 0 ? 0 : -1); + // S3 register_buffer 下沉层 MC_LOG 汇总(Q1b:入口层只 PerfPoint) + MC_LOG(INFO) << "[register_buffer] elapsed_us=" << pt.ElapsedMicros() + << " success=" << (ret == 0 ? 1 : 0) + << " base_addr=" << buffer << " size=" << size; + return ret; } tl::expected RealClient::unregister_buffer_internal( @@ -5854,19 +5887,26 @@ RealClient::batch_get_into_internal(const std::vector &keys, std::vector> RealClient::batchIsExist_internal( const std::vector &keys) { + SpDiag::PerfPoint pt(PerfKey::RC_BATCH_IS_EXIST_INTERNAL, + SpDiag::PerfLevel::MODULE); + pt.Start(); if (!client_) { LOG(ERROR) << "Client is not initialized"; + pt.End(-1); return std::vector>( keys.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); } if (keys.empty()) { LOG(WARNING) << "Empty keys vector provided to batchIsExist_internal"; + pt.End(0); return std::vector>(); } // Call client BatchIsExist and return the vector directly - return client_->BatchIsExist(keys); + auto ret = client_->BatchIsExist(keys); + pt.End(0); + return ret; } int RealClient::put_from_with_metadata(const std::string &key, void *buffer, @@ -5913,6 +5953,9 @@ std::vector RealClient::batch_put_from_multi_buffers( const std::vector> &all_buffers, const std::vector> &sizes, const ReplicateConfig &config) { + SpDiag::PerfPoint pt(PerfKey::RC_BATCH_PUT_MULTI, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); mooncake::logging::ScopedTraceId trace(mooncake::logging::NewTraceId()); auto internal_results = execute_timed_operation>>( @@ -5938,7 +5981,8 @@ std::vector RealClient::batch_put_from_multi_buffers( for (const auto &result : internal_results) { results.push_back(to_py_ret(result)); } - + pt.End(results.empty() ? -1 : 0); + // MC_LOG 在 *_internal 输出 汇总+per-key(Q1b) return results; } @@ -5948,8 +5992,13 @@ RealClient::batch_put_from_multi_buffers_internal( const std::vector> &all_buffers, const std::vector> &all_sizes, const ReplicateConfig &config) { + SpDiag::PerfPoint pt(PerfKey::RC_BATCH_PUT_MULTI_INTERNAL, + SpDiag::PerfLevel::MODULE); + pt.Start(); + auto t0 = std::chrono::steady_clock::now(); if (!client_) { LOG(ERROR) << "Client is not initialized"; + pt.End(-1); return std::vector>( keys.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); } @@ -5957,6 +6006,7 @@ RealClient::batch_put_from_multi_buffers_internal( if ((keys.size() != all_buffers.size()) || (all_buffers.size() != all_sizes.size())) { LOG(ERROR) << "Mismatched sizes for keys, buffers, and sizes"; + pt.End(-1); return std::vector>( keys.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); } @@ -5967,6 +6017,7 @@ RealClient::batch_put_from_multi_buffers_internal( const auto &sizes = all_sizes[i]; if (buffers.size() != sizes.size()) { LOG(ERROR) << "Mismatched buffers and sizes of key:" << keys[i]; + pt.End(-1); return std::vector>( keys.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); } @@ -5976,7 +6027,32 @@ RealClient::batch_put_from_multi_buffers_internal( } } // Call client BatchPut and return the vector directly - return client_->BatchPut(keys, batched_slices, config); + auto result = client_->BatchPut(keys, batched_slices, config); + auto t1 = std::chrono::steady_clock::now(); + auto total_us = std::chrono::duration_cast( + t1 - t0).count(); + pt.End(result.empty() ? -1 : 0); + // S4 batch_put_from_multi_buffers 下沉层 MC_LOG 汇总+per-key(Q1b) + // 字段:success+key+size(无 replica/endpoint,缺失不输出) + size_t total_bytes = 0; + for (const auto &sizes : all_sizes) + for (auto s : sizes) total_bytes += s; + std::ostringstream oss; + oss << "[batch_put_from_multi_buffers] elapsed_us=" << total_us + << " num_keys=" << keys.size() + << " total_bytes=" << total_bytes; + for (size_t i = 0; i < keys.size(); ++i) { + int success = (i < result.size() && result[i].has_value()) ? 1 : 0; + oss << "\n success=" << success << " key=" << keys[i]; + if (i < all_sizes.size()) { + size_t key_size = 0; + for (auto s : all_sizes[i]) key_size += s; + oss << " size=" << key_size; + } + // 无 replica/endpoint,不输出 + } + MC_LOG(INFO) << oss.str(); + return result; } std::vector RealClient::batch_get_into_multi_buffers( @@ -5984,6 +6060,9 @@ std::vector RealClient::batch_get_into_multi_buffers( const std::vector> &all_buffers, const std::vector> &all_sizes, bool prefer_alloc_in_same_node) { + SpDiag::PerfPoint pt(PerfKey::RC_BATCH_GET_INTO_MULTI, + SpDiag::PerfLevel::KEY_MODULE); + pt.Start(); auto internal_results = execute_timed_operation>>( [&]() { @@ -6008,6 +6087,8 @@ std::vector RealClient::batch_get_into_multi_buffers( for (const auto &result : internal_results) { results.push_back(to_py_ret(result)); } + pt.End(results.empty() ? -1 : 0); + // MC_LOG 在 *_internal 输出 汇总+per-key(Q1b) return results; } @@ -6017,9 +6098,14 @@ RealClient::batch_get_into_multi_buffers_internal( const std::vector> &all_buffers, const std::vector> &all_sizes, bool prefer_alloc_in_same_node) { + SpDiag::PerfPoint pt(PerfKey::RC_BATCH_GET_INTO_MULTI_INTERNAL, + SpDiag::PerfLevel::MODULE); + pt.Start(); + auto t0 = std::chrono::steady_clock::now(); // Validate preconditions if (!client_) { LOG(ERROR) << "Client is not initialized"; + pt.End(-1); return std::vector>( keys.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); } @@ -6028,6 +6114,7 @@ RealClient::batch_get_into_multi_buffers_internal( LOG(ERROR) << "Input vector sizes mismatch: keys=" << keys.size() << ", buffers=" << all_buffers.size() << ", sizes=" << all_sizes.size(); + pt.End(-1); return std::vector>( keys.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); } @@ -6035,7 +6122,11 @@ RealClient::batch_get_into_multi_buffers_internal( const size_t num_keys = keys.size(); std::vector> results; results.reserve(num_keys); + // Per-key replica/endpoint tracking for MC_LOG(Q1b:下沉层 per-key 全字段) + std::vector per_key_replica_type(num_keys); + std::vector per_key_endpoint(num_keys); if (num_keys == 0) { + pt.End(0); return results; } // Query metadata for all keys @@ -6094,6 +6185,21 @@ RealClient::batch_get_into_multi_buffers_internal( continue; } const auto replica = *best_replica; + // Capture replica info for MC_LOG per-key output + if (replica.is_memory_replica()) { + per_key_replica_type[i] = "memory"; + per_key_endpoint[i] = std::string(replica.get_memory_descriptor() + .buffer_descriptor + .transport_endpoint_); + } else if (replica.is_local_disk_replica()) { + per_key_replica_type[i] = "local_disk"; + per_key_endpoint[i] = + std::string(replica.get_local_disk_descriptor() + .transport_endpoint); + } else if (replica.is_disk_replica()) { + per_key_replica_type[i] = "disk"; + // DISK 副本无 transport_endpoint,缺失不输出 + } uint64_t total_size = calculate_total_size(replica); const auto &sizes = all_sizes[i]; uint64_t dst_total_size = 0; @@ -6149,6 +6255,7 @@ RealClient::batch_get_into_multi_buffers_internal( } // Early return if no valid operations if (valid_operations.empty() && valid_local_disk_ops.empty()) { + pt.End(0); return results; } @@ -6366,6 +6473,32 @@ RealClient::batch_get_into_multi_buffers_internal( } } + auto t1 = std::chrono::steady_clock::now(); + auto total_us = std::chrono::duration_cast( + t1 - t0).count(); + pt.End(results.empty() ? -1 : 0); + // S5 batch_get_into_multi_buffers 下沉层 MC_LOG 汇总+per-key(Q1b) + // 字段:success+key+size+replica+endpoint(缺失字段不输出) + std::ostringstream oss; + oss << "[batch_get_into_multi_buffers] elapsed_us=" << total_us + << " num_keys=" << keys.size(); + for (size_t i = 0; i < keys.size(); ++i) { + int success = + (i < results.size() && results[i].has_value()) ? 1 : 0; + oss << "\n success=" << success << " key=" << keys[i]; + if (i < all_sizes.size()) { + size_t key_size = 0; + for (auto s : all_sizes[i]) key_size += s; + oss << " size=" << key_size; + } + if (i < per_key_replica_type.size() && !per_key_replica_type[i].empty()) { + oss << " replica=" << per_key_replica_type[i]; + } + if (i < per_key_endpoint.size() && !per_key_endpoint[i].empty()) { + oss << " endpoint=" << per_key_endpoint[i]; + } + } + MC_LOG(INFO) << oss.str(); return results; } From 56abcdcf330f56ddcfc43fb63fc4047376d5a476 Mon Sep 17 00:00:00 2001 From: ZhiningPan Date: Thu, 27 Aug 2026 20:07:21 +0800 Subject: [PATCH 3/6] =?UTF-8?q?fix(spdiag):=20=E7=94=A8=20steady=5Fclock?= =?UTF-8?q?=20=E6=9B=BF=E4=BB=A3=20PerfPoint::ElapsedMicros?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit SpDiag::PerfPoint v1.0.0 仅暴露 Start()/End()/Abandon() 三个方法, 没有 ElapsedMicros()。服务器构建报错: real_client.cpp:1218: error: 'class SpDiag::PerfPoint' has no member named 'ElapsedMicros' 修复方式:在 setup_real / batchIsExist / register_buffer 以及 store_py.cpp 的 remove_all / close 入口层用 std::chrono::steady_clock 手动测量 elapsed_us,与下沉层 batch_get_into_multi_buffers_internal 已有的 t0/t1/total_us 风格保持一致。 --- mooncake-integration/store/store_py.cpp | 17 ++++++++++++++--- mooncake-store/src/real_client.cpp | 18 +++++++++++++++--- 2 files changed, 29 insertions(+), 6 deletions(-) diff --git a/mooncake-integration/store/store_py.cpp b/mooncake-integration/store/store_py.cpp index 82f8e7f5d2..21b8e53286 100644 --- a/mooncake-integration/store/store_py.cpp +++ b/mooncake-integration/store/store_py.cpp @@ -2246,11 +2246,15 @@ PYBIND11_MODULE(store, m) { SpDiag::PerfPoint pt(PerfKey::STORE_PY_REMOVE_ALL, SpDiag::PerfLevel::KEY_MODULE); pt.Start(); + auto t0 = std::chrono::steady_clock::now(); py::gil_scoped_release release; auto ret = self.store_->removeAll(force); + auto t1 = std::chrono::steady_clock::now(); + auto elapsed_us = std::chrono::duration_cast( + t1 - t0).count(); pt.End(ret == 0 ? 0 : -1); // 下沉 removeAll 不加 MC_LOG,入口层输出汇总(Q1b) - MC_LOG(INFO) << "[remove_all] elapsed_us=" << pt.ElapsedMicros() + MC_LOG(INFO) << "[remove_all] elapsed_us=" << elapsed_us << " success=" << (ret == 0 ? 1 : 0); return ret; }, @@ -2293,18 +2297,25 @@ PYBIND11_MODULE(store, m) { SpDiag::PerfPoint pt(PerfKey::STORE_PY_CLOSE, SpDiag::PerfLevel::KEY_MODULE); pt.Start(); + auto t0 = std::chrono::steady_clock::now(); if (!self.store_) { + auto t1 = std::chrono::steady_clock::now(); + auto elapsed_us = std::chrono::duration_cast( + t1 - t0).count(); pt.End(0); // 无下沉 MC_LOG,入口层输出汇总(Q1b) - MC_LOG(INFO) << "[close] elapsed_us=" << pt.ElapsedMicros() + MC_LOG(INFO) << "[close] elapsed_us=" << elapsed_us << " success=1"; return 0; } int rc = self.store_->tearDownAll(); self.store_.reset(); + auto t1 = std::chrono::steady_clock::now(); + auto elapsed_us = std::chrono::duration_cast( + t1 - t0).count(); pt.End(rc == 0 ? 0 : -1); // 下沉 tearDownAll 不加 MC_LOG,入口层输出汇总(Q1b) - MC_LOG(INFO) << "[close] elapsed_us=" << pt.ElapsedMicros() + MC_LOG(INFO) << "[close] elapsed_us=" << elapsed_us << " success=" << (rc == 0 ? 1 : 0); return rc; }) diff --git a/mooncake-store/src/real_client.cpp b/mooncake-store/src/real_client.cpp index 7a184a3c52..2d1ca23721 100644 --- a/mooncake-store/src/real_client.cpp +++ b/mooncake-store/src/real_client.cpp @@ -1208,14 +1208,18 @@ int RealClient::setup_real( SpDiag::PerfPoint pt(PerfKey::RC_SETUP_REAL, SpDiag::PerfLevel::KEY_MODULE); pt.Start(); + auto t0 = std::chrono::steady_clock::now(); auto ret = to_py_ret(setup_internal( local_hostname, metadata_server, global_segment_size, local_buffer_size, protocol, rdma_devices, master_server_addr, transfer_engine, ipc_socket_path, 50052, enable_ssd_offload, true, ssd_offload_path, tenant_id, Environ::Get().GetOffloadRpcThreadNum(8), enable_client_http_server, client_http_port)); + auto t1 = std::chrono::steady_clock::now(); + auto elapsed_us = std::chrono::duration_cast( + t1 - t0).count(); pt.End(ret == 0 ? 0 : -1); // S2 setup 下沉层 MC_LOG 汇总(Q1b:入口层只 PerfPoint) - MC_LOG(INFO) << "[setup] elapsed_us=" << pt.ElapsedMicros() + MC_LOG(INFO) << "[setup] elapsed_us=" << elapsed_us << " success=" << (ret == 0 ? 1 : 0) << " hostname=" << local_hostname << " metadata_server=" << metadata_server @@ -2490,7 +2494,11 @@ std::vector RealClient::batchIsExist( SpDiag::PerfPoint pt(PerfKey::RC_BATCH_IS_EXIST, SpDiag::PerfLevel::KEY_MODULE); pt.Start(); + auto t0 = std::chrono::steady_clock::now(); auto internal_results = batchIsExist_internal(keys); + auto t1 = std::chrono::steady_clock::now(); + auto elapsed_us = std::chrono::duration_cast( + t1 - t0).count(); std::vector results; results.reserve(internal_results.size()); @@ -2505,7 +2513,7 @@ std::vector RealClient::batchIsExist( // S6 batch_is_exist 下沉层 MC_LOG 汇总+per-key(Q1b:入口层只 PerfPoint) // 字段:success+key(无 size/replica/endpoint,缺失不输出) std::ostringstream oss; - oss << "[batch_is_exist] elapsed_us=" << pt.ElapsedMicros() + oss << "[batch_is_exist] elapsed_us=" << elapsed_us << " num_keys=" << keys.size(); for (size_t i = 0; i < keys.size(); ++i) { int success = (i < results.size()) ? results[i] : 0; @@ -3795,10 +3803,14 @@ int RealClient::register_buffer(void *buffer, size_t size) { SpDiag::PerfPoint pt(PerfKey::RC_REGISTER_BUFFER, SpDiag::PerfLevel::KEY_MODULE); pt.Start(); + auto t0 = std::chrono::steady_clock::now(); auto ret = to_py_ret(register_buffer_internal(buffer, size)); + auto t1 = std::chrono::steady_clock::now(); + auto elapsed_us = std::chrono::duration_cast( + t1 - t0).count(); pt.End(ret == 0 ? 0 : -1); // S3 register_buffer 下沉层 MC_LOG 汇总(Q1b:入口层只 PerfPoint) - MC_LOG(INFO) << "[register_buffer] elapsed_us=" << pt.ElapsedMicros() + MC_LOG(INFO) << "[register_buffer] elapsed_us=" << elapsed_us << " success=" << (ret == 0 ? 1 : 0) << " base_addr=" << buffer << " size=" << size; return ret; From 8347ca304e07ea295a288ba96ff1500dfd5fda20 Mon Sep 17 00:00:00 2001 From: ZhiningPan Date: Fri, 28 Aug 2026 10:40:16 +0800 Subject: [PATCH 4/6] =?UTF-8?q?bench:=20=E6=96=B0=E5=A2=9E=20store=5Fconne?= =?UTF-8?q?ctor=5Fbench=20=E4=B8=93=E6=B5=8B=20vLLM=20=E8=B7=AF=E5=BE=84?= =?UTF-8?q?=E6=89=93=E7=82=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 现有 stress_cluster_bench 只调用老接口 batch_get_into,无法触发 vllm 调用路径的新打点(batch_put_from_multi_buffers / batchIsExist / batch_get_into_multi_buffers)。 按方案文档第 6 章实现 store_connector_bench.cpp: - 4 个 scenario: write / is_exist / get / all - 模拟 vllm 多 layer KV cache 场景(每请求 num_layers 个 key + buffer) - 调用 RealClient::batch_put_from_multi_buffers / batchIsExist / batch_get_into_multi_buffers,触发 STORE_PY_* + RC_* + CLIENT_BATCH_QUERY SpDiag 打点 + MC_LOG 汇总日志 - 输出带宽 + 时延分位数(P50/P90/P99) - 同步追加 CMakeLists.txt 编译目标,链接库参考 stress_cluster_bench --- mooncake-store/benchmarks/CMakeLists.txt | 10 + .../benchmarks/store_connector_bench.cpp | 480 ++++++++++++++++++ 2 files changed, 490 insertions(+) create mode 100644 mooncake-store/benchmarks/store_connector_bench.cpp diff --git a/mooncake-store/benchmarks/CMakeLists.txt b/mooncake-store/benchmarks/CMakeLists.txt index 003f0b8724..817230b940 100644 --- a/mooncake-store/benchmarks/CMakeLists.txt +++ b/mooncake-store/benchmarks/CMakeLists.txt @@ -38,6 +38,16 @@ target_link_libraries( stress_cluster_bench PRIVATE mooncake_store transfer_engine asio_shared gflags::gflags glog::glog pthread) +# Benchmark for vLLM Store Connector path +# Triggers: batch_put_from_multi_buffers / batchIsExist / +# batch_get_into_multi_buffers (and setup/register_buffer/tearDownAll). +# Used to verify SpDiag perf points and MC_LOG logging. +add_executable(store_connector_bench store_connector_bench.cpp) +target_link_libraries( + store_connector_bench PRIVATE mooncake_store transfer_engine asio_shared + gflags::gflags glog::glog pthread) + + # Benchmark for RealClient::get_into_ranges with configurable value size, # fragments per key and keys per query. add_executable(stress_cluster_ranges_bench stress_cluster_ranges_bench.cpp) diff --git a/mooncake-store/benchmarks/store_connector_bench.cpp b/mooncake-store/benchmarks/store_connector_bench.cpp new file mode 100644 index 0000000000..039cb045a9 --- /dev/null +++ b/mooncake-store/benchmarks/store_connector_bench.cpp @@ -0,0 +1,480 @@ +// Store Connector Benchmark: 专测 vLLM 调用路径(batch_put_from_multi_buffers +// / batchIsExist / batch_get_into_multi_buffers),触发新增的 SpDiag 打点。 +// 参考方案文档 vllm_spdiag_logging_plan.md 第 6 章。 + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "gflags/gflags.h" +#include "glog/logging.h" +#include "mooncake_logging.h" +#include "real_client.h" + +namespace { +constexpr size_t KB = 1024; +constexpr size_t MB = 1024 * KB; +constexpr size_t GB = 1024 * MB; + +using Clock = std::chrono::steady_clock; +using Nanos = std::chrono::nanoseconds; + +inline int64_t ElapsedNanos(Clock::time_point t0, Clock::time_point t1) { + return std::chrono::duration_cast(t1 - t0).count(); +} +inline double NanosToUs(int64_t ns) { return static_cast(ns) / 1000.0; } +inline double NanosToSec(int64_t ns) { + return static_cast(ns) / 1e9; +} + +static std::string FormatBytes(size_t bytes) { + if (bytes == 0) return "0 B"; + const char* units[] = {"B", "KB", "MB", "GB", "TB"}; + int i = static_cast(std::floor(std::log2(bytes) / 10)); + if (i > 4) i = 4; + double val = static_cast(bytes) / std::pow(1024, i); + std::ostringstream oss; + oss << std::fixed << std::setprecision(2) << val << " " << units[i]; + return oss.str(); +} +} // namespace + +DEFINE_string(local_hostname, "localhost", "Local hostname"); +DEFINE_string( + metadata_server, "http://127.0.0.1:8080/metadata", + "Metadata server URL (e.g. http://127.0.0.1:8080/metadata or etcd://...)"); +DEFINE_string(master_server, "127.0.0.1:50051", "Master server address"); +DEFINE_string(protocol, "tcp", "Transport protocol: tcp, rdma, ub"); +DEFINE_string(device_name, "", "RDMA/UB device name (comma-separated)"); +DEFINE_uint64(global_segment_size, 4 * GB, "Global segment size in bytes"); +DEFINE_uint64(local_buffer_size, 512 * MB, "Local buffer size in bytes"); +DEFINE_bool(enable_ssd_offload, false, "Enable SSD offload on this client"); +DEFINE_string(ssd_offload_path, "", "SSD offload directory path"); + +DEFINE_string(scenario, "all", + "Benchmark scenario: write, is_exist, get, all"); +DEFINE_uint64(num_requests, 100, + "Number of requests (each has num_layers keys)"); +DEFINE_uint64(num_layers, 32, + "Layers per request (simulate vLLM transformer)"); +DEFINE_uint64(layer_size, 1 * MB, "Size of each layer buffer in bytes"); +DEFINE_uint64(num_threads, 1, "Number of concurrent threads"); +DEFINE_uint64(warmup_requests, 5, "Warmup requests (not counted)"); + +enum class Phase { WRITE, IS_EXIST, GET }; + +static std::string PhaseName(Phase p) { + switch (p) { + case Phase::WRITE: + return "WRITE [batch_put_from_multi_buffers]"; + case Phase::IS_EXIST: + return "IS_EXIST [batchIsExist]"; + case Phase::GET: + return "GET [batch_get_into_multi_buffers]"; + } + return "UNKNOWN"; +} + +struct ThreadResult { + std::vector latencies_ns; + size_t total_bytes = 0; + size_t total_keys = 0; + size_t total_queries = 0; + size_t failed_ops = 0; +}; + +class BenchmarkStats { + public: + void InitThreads(size_t n) { thread_results_.resize(n); } + ThreadResult& GetThreadResult(size_t tid) { return thread_results_[tid]; } + void StartTimer() { start_ = Clock::now(); } + void StopTimer() { end_ = Clock::now(); } + double WallSeconds() const { return NanosToSec(ElapsedNanos(start_, end_)); } + + void Finalize() { + merged_latencies_ns_.clear(); + total_bytes_ = total_keys_ = total_queries_ = total_failed_ = 0; + for (auto& tr : thread_results_) { + merged_latencies_ns_.insert(merged_latencies_ns_.end(), + tr.latencies_ns.begin(), + tr.latencies_ns.end()); + total_bytes_ += tr.total_bytes; + total_keys_ += tr.total_keys; + total_queries_ += tr.total_queries; + total_failed_ += tr.failed_ops; + } + std::sort(merged_latencies_ns_.begin(), merged_latencies_ns_.end()); + } + + double PercentileUs(double p) const { + if (merged_latencies_ns_.empty()) return 0.0; + double rank = (p / 100.0) * (merged_latencies_ns_.size() - 1); + size_t lo = static_cast(rank); + size_t hi = std::min(lo + 1, merged_latencies_ns_.size() - 1); + double frac = rank - lo; + int64_t ns_val = static_cast( + merged_latencies_ns_[lo] * (1.0 - frac) + + merged_latencies_ns_[hi] * frac); + return NanosToUs(ns_val); + } + + double MeanLatencyUs() const { + if (merged_latencies_ns_.empty()) return 0.0; + double sum = static_cast(std::accumulate( + merged_latencies_ns_.begin(), merged_latencies_ns_.end(), + int64_t(0))); + return NanosToUs(sum / static_cast(merged_latencies_ns_.size())); + } + + double ThroughputMBps() const { + double wall = WallSeconds(); + return (wall > 0) ? (static_cast(total_bytes_) / MB) / wall : 0; + } + + double KeysPerSec() const { + double wall = WallSeconds(); + return (wall > 0) ? static_cast(total_keys_) / wall : 0; + } + + void Print(const std::string& title, bool show_bandwidth) const { + std::cout << "\n========================================" + "========================================\n"; + std::cout << " " << title << "\n"; + std::cout << "========================================" + "========================================\n"; + std::cout << std::fixed << std::setprecision(2); + std::cout << " Wall time: " << WallSeconds() << " s\n"; + std::cout << " Total queries: " << total_queries_ + << " (failed: " << total_failed_ << ")\n"; + std::cout << " Total keys: " << total_keys_ << "\n"; + if (show_bandwidth) { + std::cout << " Total data: " << FormatBytes(total_bytes_) + << "\n"; + std::cout << " Throughput: " << ThroughputMBps() + << " MB/s"; + if (ThroughputMBps() > 1024) + std::cout << " (" << ThroughputMBps() / 1024 << " GB/s)"; + std::cout << "\n"; + } + std::cout << " Keys/sec: " << KeysPerSec() << "\n"; + + if (!merged_latencies_ns_.empty()) { + size_t n = merged_latencies_ns_.size(); + std::cout << "\n Latency (us) [n=" << n << ", per-query]\n"; + std::cout << " Min: " << std::setw(12) + << NanosToUs(merged_latencies_ns_.front()) << "\n"; + std::cout << " Avg: " << std::setw(12) << MeanLatencyUs() + << "\n"; + std::cout << " P50: " << std::setw(12) << PercentileUs(50) + << "\n"; + std::cout << " P90: " << std::setw(12) << PercentileUs(90) + << "\n"; + std::cout << " P99: " << std::setw(12) << PercentileUs(99) + << "\n"; + std::cout << " Max: " << std::setw(12) + << NanosToUs(merged_latencies_ns_.back()) << "\n"; + } + std::cout << "========================================" + "========================================\n\n"; + } + + private: + std::vector thread_results_; + std::vector merged_latencies_ns_; + size_t total_bytes_ = 0, total_keys_ = 0, total_queries_ = 0, + total_failed_ = 0; + Clock::time_point start_, end_; +}; + +class StoreConnectorBench { + public: + StoreConnectorBench() : client_(mooncake::RealClient::create()) {} + + ~StoreConnectorBench() { + for (auto& tb : thread_buffers_) { + if (tb.ptr) { + try { + client_->unregister_buffer(tb.ptr); + } catch (...) { + LOG(WARNING) + << "Failed to unregister thread buffer, ignoring"; + } + free(tb.ptr); + } + } + if (main_buffer_) { + try { + client_->unregister_buffer(main_buffer_); + } catch (...) { + LOG(WARNING) << "Failed to unregister main buffer, ignoring"; + } + free(main_buffer_); + } + } + + int Setup() { + int ret = client_->setup_real( + FLAGS_local_hostname, FLAGS_metadata_server, + FLAGS_global_segment_size, FLAGS_local_buffer_size, + FLAGS_protocol, FLAGS_device_name, FLAGS_master_server, nullptr, + "", FLAGS_enable_ssd_offload, FLAGS_ssd_offload_path); + if (ret != 0) { + LOG(ERROR) << "RealClient setup_real failed, ret=" << ret; + return ret; + } + LOG(INFO) << "RealClient setup succeeded"; + + // 主 buffer 用于 PrepareData 阶段 + main_buffer_size_ = FLAGS_num_layers * FLAGS_layer_size; + main_buffer_ = static_cast(malloc(main_buffer_size_)); + if (!main_buffer_) { + LOG(ERROR) << "Failed to allocate main buffer"; + return -1; + } + memset(main_buffer_, 0xAB, main_buffer_size_); // 填充测试数据 + ret = client_->register_buffer(main_buffer_, main_buffer_size_); + if (ret != 0) { + LOG(ERROR) << "register_buffer failed for main buffer"; + return ret; + } + + return AllocateThreadBuffers(FLAGS_num_threads); + } + + // is_exist/get 模式前先写入数据(不计入统计) + int PrepareData() { + LOG(INFO) << "Preparing data: writing " << FLAGS_num_requests + << " requests..."; + mooncake::ReplicateConfig config; + config.replica_num = 1; + + for (size_t r = 0; r < FLAGS_num_requests; ++r) { + auto keys = MakeRequestKeys(r); + auto all_buffers = MakeBufferList(main_buffer_); + auto all_sizes = MakeSizeList(); + auto ret = client_->batch_put_from_multi_buffers( + keys, all_buffers, all_sizes, config); + for (int v : ret) + if (v != 0) { + LOG(ERROR) << "PrepareData failed at request " << r; + return -1; + } + if ((r + 1) % 20 == 0) + LOG(INFO) << " Prepared " << (r + 1) << "/" + << FLAGS_num_requests; + } + LOG(INFO) << "Data preparation complete"; + return 0; + } + + int RunPhase(Phase phase, bool is_warmup) { + BenchmarkStats stats; + stats.InitThreads(FLAGS_num_threads); + stats.StartTimer(); + + std::latch start_latch(static_cast(FLAGS_num_threads)); + std::latch done_latch(static_cast(FLAGS_num_threads)); + + size_t total = FLAGS_num_requests; + std::vector threads; + for (size_t t = 0; t < FLAGS_num_threads; ++t) { + size_t my = total / FLAGS_num_threads + + (t < total % FLAGS_num_threads ? 1 : 0); + size_t offset = t * (total / FLAGS_num_threads) + + std::min(t, total % FLAGS_num_threads); + threads.emplace_back([&, t, my, offset]() { + PhaseWorker(t, my, offset, phase, stats, start_latch, + done_latch); + }); + } + done_latch.wait(); + stats.StopTimer(); + for (auto& th : threads) th.join(); + stats.Finalize(); + + if (!is_warmup) { + bool show_bw = (phase != Phase::IS_EXIST); + stats.Print("BENCHMARK " + PhaseName(phase), show_bw); + } + return 0; + } + + int Run() { + // Warmup(所有模式都 warmup get,确保连接建立) + if (FLAGS_warmup_requests > 0 && FLAGS_scenario != "write") { + LOG(INFO) << "Warmup: " << FLAGS_warmup_requests << " requests"; + size_t saved = FLAGS_num_requests; + FLAGS_num_requests = FLAGS_warmup_requests; + RunPhase(Phase::GET, true); + FLAGS_num_requests = saved; + } + + if (FLAGS_scenario == "all") { + RunPhase(Phase::WRITE, false); + RunPhase(Phase::IS_EXIST, false); + RunPhase(Phase::GET, false); + } else if (FLAGS_scenario == "write") { + RunPhase(Phase::WRITE, false); + } else if (FLAGS_scenario == "is_exist") { + if (PrepareData() != 0) return -1; + RunPhase(Phase::IS_EXIST, false); + } else if (FLAGS_scenario == "get") { + if (PrepareData() != 0) return -1; + RunPhase(Phase::GET, false); + } else { + LOG(ERROR) << "Unknown scenario: " << FLAGS_scenario; + return -1; + } + return 0; + } + + private: + static std::vector MakeRequestKeys(size_t req_id) { + std::vector keys; + keys.reserve(FLAGS_num_layers); + for (size_t l = 0; l < FLAGS_num_layers; ++l) + keys.push_back("layer." + std::to_string(l) + ".req_" + + std::to_string(req_id)); + return keys; + } + + // 每 key 对应 1 个 buffer(vLLM 场景),从大 buffer 切片 + std::vector> MakeBufferList(char* base) { + std::vector> all_buffers(FLAGS_num_layers); + for (size_t l = 0; l < FLAGS_num_layers; ++l) + all_buffers[l] = {base + l * FLAGS_layer_size}; + return all_buffers; + } + + std::vector> MakeSizeList() { + return std::vector>( + FLAGS_num_layers, {static_cast(FLAGS_layer_size)}); + } + + void PhaseWorker(size_t tid, size_t my_requests, size_t offset, + Phase phase, BenchmarkStats& stats, + std::latch& start_latch, std::latch& done_latch) { + ThreadResult& result = stats.GetThreadResult(tid); + result.latencies_ns.reserve(my_requests); + char* my_buf = thread_buffers_[tid].ptr; + + mooncake::ReplicateConfig config; + config.replica_num = 1; + + start_latch.arrive_and_wait(); + + size_t bytes_per_req = FLAGS_num_layers * FLAGS_layer_size; + + for (size_t i = 0; i < my_requests; ++i) { + size_t req_id = offset + i; + auto keys = MakeRequestKeys(req_id); + + auto t0 = Clock::now(); + std::vector ret; + + if (phase == Phase::WRITE) { + auto bufs = MakeBufferList(my_buf); + auto sizes = MakeSizeList(); + ret = client_->batch_put_from_multi_buffers(keys, bufs, sizes, + config); + } else if (phase == Phase::IS_EXIST) { + ret = client_->batchIsExist(keys); + } else { // GET + auto bufs = MakeBufferList(my_buf); + auto sizes = MakeSizeList(); + ret = client_->batch_get_into_multi_buffers(keys, bufs, sizes, + false); + } + auto t1 = Clock::now(); + result.latencies_ns.push_back(ElapsedNanos(t0, t1)); + + bool ok = true; + for (int v : ret) + if (v != 0) ok = false; + if (ok) { + result.total_keys += FLAGS_num_layers; + if (phase != Phase::IS_EXIST) + result.total_bytes += bytes_per_req; + } else { + result.failed_ops++; + } + result.total_queries++; + } + + done_latch.arrive_and_wait(); + } + + int AllocateThreadBuffers(size_t num_threads) { + thread_buffers_.resize(num_threads); + size_t per_buf = FLAGS_num_layers * FLAGS_layer_size; + for (size_t t = 0; t < num_threads; ++t) { + thread_buffers_[t].size = per_buf; + thread_buffers_[t].ptr = static_cast(malloc(per_buf)); + if (!thread_buffers_[t].ptr) { + LOG(ERROR) << "Failed to allocate buffer for thread " << t; + return -1; + } + memset(thread_buffers_[t].ptr, 0, per_buf); + int ret = + client_->register_buffer(thread_buffers_[t].ptr, per_buf); + if (ret != 0) { + LOG(ERROR) << "register_buffer failed for thread " << t; + return ret; + } + } + LOG(INFO) << "Allocated " << num_threads << " thread buffers, each " + << FormatBytes(per_buf); + return 0; + } + + std::shared_ptr client_; + char* main_buffer_ = nullptr; + size_t main_buffer_size_ = 0; + struct ThreadBuf { + char* ptr = nullptr; + size_t size = 0; + }; + std::vector thread_buffers_; +}; + +int main(int argc, char* argv[]) { + if (!google::IsGoogleLoggingInitialized()) { + google::InitGoogleLogging(argv[0]); + } + gflags::ParseCommandLineFlags(&argc, &argv, true); + + if (std::getenv("MC_LOG_DIR") == nullptr) { + FLAGS_logtostderr = true; + } + mooncake::logging::ApplyMooncakeLogEnableToGlog(); + + LOG(INFO) << "Mooncake Store Connector Benchmark (vLLM path)"; + LOG(INFO) << " Scenario: " << FLAGS_scenario; + LOG(INFO) << " Protocol: " << FLAGS_protocol; + LOG(INFO) << " Requests: " << FLAGS_num_requests; + LOG(INFO) << " Layers/req: " << FLAGS_num_layers; + LOG(INFO) << " Layer size: " << FormatBytes(FLAGS_layer_size); + LOG(INFO) << " Threads: " << FLAGS_num_threads; + size_t total_data = + FLAGS_num_requests * FLAGS_num_layers * FLAGS_layer_size; + LOG(INFO) << " Total data: " << FormatBytes(total_data); + + StoreConnectorBench bench; + int ret = bench.Setup(); + if (ret != 0) { + LOG(ERROR) << "Setup failed"; + return ret; + } + return bench.Run(); +} From 556c36bf5c3b96ad1a2e80e14df44a67229065fc Mon Sep 17 00:00:00 2001 From: ZhiningPan Date: Sat, 29 Aug 2026 16:22:59 +0800 Subject: [PATCH 5/6] =?UTF-8?q?fix(bench):=20=E6=94=B9=E7=94=A8=20numa=5Fa?= =?UTF-8?q?lloc=5Flocal=20=E5=88=86=E9=85=8D=20buffer?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit malloc 返回的内存不保证页对齐和 NUMA 本地性,UB/RDMA driver 会拒绝注册为 DMA 内存,导致 register_buffer 失败: Failed to register segment ... : Success [0] UbTransport: cannot register LocalMemory 参考 stress_cluster_bench 的做法,改用 numa_alloc_local (mmap + page-aligned + NUMA-aware),匹配 UB driver 的要求。 析构同步改用 numa_free。 --- mooncake-store/benchmarks/store_connector_bench.cpp | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/mooncake-store/benchmarks/store_connector_bench.cpp b/mooncake-store/benchmarks/store_connector_bench.cpp index 039cb045a9..f91f46e609 100644 --- a/mooncake-store/benchmarks/store_connector_bench.cpp +++ b/mooncake-store/benchmarks/store_connector_bench.cpp @@ -17,6 +17,9 @@ #include #include +#include +#include + #include "gflags/gflags.h" #include "glog/logging.h" #include "mooncake_logging.h" @@ -210,7 +213,7 @@ class StoreConnectorBench { LOG(WARNING) << "Failed to unregister thread buffer, ignoring"; } - free(tb.ptr); + numa_free(tb.ptr, tb.size); } } if (main_buffer_) { @@ -219,7 +222,7 @@ class StoreConnectorBench { } catch (...) { LOG(WARNING) << "Failed to unregister main buffer, ignoring"; } - free(main_buffer_); + numa_free(main_buffer_, main_buffer_size_); } } @@ -237,7 +240,8 @@ class StoreConnectorBench { // 主 buffer 用于 PrepareData 阶段 main_buffer_size_ = FLAGS_num_layers * FLAGS_layer_size; - main_buffer_ = static_cast(malloc(main_buffer_size_)); + main_buffer_ = reinterpret_cast( + numa_alloc_local(main_buffer_size_)); if (!main_buffer_) { LOG(ERROR) << "Failed to allocate main buffer"; return -1; @@ -420,7 +424,8 @@ class StoreConnectorBench { size_t per_buf = FLAGS_num_layers * FLAGS_layer_size; for (size_t t = 0; t < num_threads; ++t) { thread_buffers_[t].size = per_buf; - thread_buffers_[t].ptr = static_cast(malloc(per_buf)); + thread_buffers_[t].ptr = + reinterpret_cast(numa_alloc_local(per_buf)); if (!thread_buffers_[t].ptr) { LOG(ERROR) << "Failed to allocate buffer for thread " << t; return -1; From 6a1e4b092033127a3e758f47c28c5f2657799b3d Mon Sep 17 00:00:00 2001 From: ZhiningPan Date: Sat, 29 Aug 2026 16:35:22 +0800 Subject: [PATCH 6/6] =?UTF-8?q?fix(bench):=20=E4=BF=AE=E6=AD=A3=20IS=5FEXI?= =?UTF-8?q?ST/GET=20=E7=9A=84=E6=88=90=E5=8A=9F=E5=88=A4=E6=96=AD=E9=80=BB?= =?UTF-8?q?=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit batchIsExist 返回 1=存在(成功),0=不存在(失败) batch_get_into_multi_buffers 返回 >0=字节数(成功),<0=错误码(失败) 原代码统一用 v!=0 判断失败,导致 IS_EXIST 和 GET 全部误判为失败 修复:按 phase 分别判断 - IS_EXIST: v != 1 为失败 - GET: v <= 0 为失败 - WRITE: v != 0 为失败(保持不变) --- .../benchmarks/store_connector_bench.cpp | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/mooncake-store/benchmarks/store_connector_bench.cpp b/mooncake-store/benchmarks/store_connector_bench.cpp index f91f46e609..31a1b88b14 100644 --- a/mooncake-store/benchmarks/store_connector_bench.cpp +++ b/mooncake-store/benchmarks/store_connector_bench.cpp @@ -404,8 +404,18 @@ class StoreConnectorBench { result.latencies_ns.push_back(ElapsedNanos(t0, t1)); bool ok = true; - for (int v : ret) - if (v != 0) ok = false; + for (int v : ret) { + if (phase == Phase::IS_EXIST) { + // batchIsExist: 1=存在(成功), 0=不存在(失败) + if (v != 1) ok = false; + } else if (phase == Phase::GET) { + // batch_get: >0=字节数(成功), <0=错误码(失败) + if (v <= 0) ok = false; + } else { + // WRITE: 0=成功, 非0=失败 + if (v != 0) ok = false; + } + } if (ok) { result.total_keys += FLAGS_num_layers; if (phase != Phase::IS_EXIST)