From 66a813111d4064543af70982e77e90adb8f5db91 Mon Sep 17 00:00:00 2001 From: Saturday Date: Tue, 23 Jun 2026 11:12:16 +0800 Subject: [PATCH 1/4] =?UTF-8?q?feat(offload):=20=E6=96=B0=E5=A2=9E?= =?UTF-8?q?=E6=8E=A8=E9=80=81=E6=A8=A1=E5=BC=8F=E5=8D=B8=E8=BD=BD=E4=BB=A5?= =?UTF-8?q?=E4=BC=98=E5=8C=96=E8=BF=9C=E7=A8=8B=E6=95=B0=E6=8D=AE=E8=8E=B7?= =?UTF-8?q?=E5=8F=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 引入推送模式卸载(push-mode offload),数据持有者直接通过单边 WRITE 将数据写入请求方内存,取代原有的拉取模式(获取指针 + READ + 释放缓冲区)的三步序列。新增 RPC 方法 `batch_get_offload_object_push`、相关数据结构及传输任务,并通过环境变量 `MC_OFFLOAD_PUSH` 控制启用。此优化减少了 RPC 往返次数,提升了卸载性能。 --- mooncake-store/include/client_service.h | 15 +++ mooncake-store/include/pyclient.h | 12 +++ mooncake-store/include/real_client.h | 14 +++ mooncake-store/include/rpc_types.h | 38 +++++++ mooncake-store/include/transfer_task.h | 11 ++ mooncake-store/src/client_service.cpp | 20 ++++ mooncake-store/src/real_client.cpp | 133 +++++++++++++++++++++++- mooncake-store/src/real_client_main.cpp | 2 + mooncake-store/src/transfer_task.cpp | 40 +++++++ 9 files changed, 284 insertions(+), 1 deletion(-) diff --git a/mooncake-store/include/client_service.h b/mooncake-store/include/client_service.h index 0094f4f14d..3005d0943b 100644 --- a/mooncake-store/include/client_service.h +++ b/mooncake-store/include/client_service.h @@ -497,6 +497,21 @@ class Client { const std::unordered_map>& batch_slices); + /** + * @brief Push counterpart of BatchGetOffloadObject, run on the data owner. + * WRITEs each key's staged ClientBuffer blob (src_pointers[i]) directly + * into the requester's destination slices over the transfer engine. + * @param requester_te_addr The requester's Transfer Engine endpoint. + * @param keys List of keys (one per src pointer / dst slice group). + * @param src_pointers Owner-local ClientBuffer addresses, one per key. + * @param dst_slices Per-key destination slices on the requester. + */ + tl::expected BatchPushOffloadObject( + const std::string& requester_te_addr, + const std::vector& keys, + const std::vector& src_pointers, + const std::vector>& dst_slices); + /** * @brief Notifies the master that offloading of specified objects has * succeeded. diff --git a/mooncake-store/include/pyclient.h b/mooncake-store/include/pyclient.h index c2e2201c5d..8e39d45de2 100644 --- a/mooncake-store/include/pyclient.h +++ b/mooncake-store/include/pyclient.h @@ -153,6 +153,18 @@ class ClientRequester { const std::vector &keys, const std::vector &sizes); + /** + * @brief Push-mode offload: invoke the owner's push handler, which WRITEs + * the data straight into our destination slices. A single RPC replaces the + * pull path's get-pointers + READ + release_offload_buffer sequence. + * @param client_addr Network address of the remote FileStorage owner. + * @param req keys/sizes plus our TE endpoint and destination slices. + */ + tl::expected + batch_get_offload_object_push( + const std::string &client_addr, + const BatchGetOffloadObjectPushRequest &req); + /** * @brief Notifies remote FileStorage to release buffer after transfer * completion. This is a fire-and-forget call - errors are logged but not diff --git a/mooncake-store/include/real_client.h b/mooncake-store/include/real_client.h index ed7be07634..373e672c6f 100644 --- a/mooncake-store/include/real_client.h +++ b/mooncake-store/include/real_client.h @@ -682,6 +682,20 @@ class RealClient : public PyClient { batch_get_offload_object(const std::vector &keys, const std::vector &sizes); + /** + * @brief Push-mode offload handler, run on the data owner. Reads the + * requested keys from SSD into the local ClientBuffer, WRITEs them + * directly into the requester's destination slices over the transfer + * engine, then releases the ClientBuffer locally. The requester therefore + * needs no follow-up READ nor a separate release_offload_buffer RPC. + * @param req keys/sizes plus the requester's TE endpoint and dst slices. + * @return per-batch result; data already landed in requester memory. + */ + async_simple::coro::Lazy< + tl::expected> + batch_get_offload_object_push( + const BatchGetOffloadObjectPushRequest &req); + /** * @brief Releases buffer associated with a specific batch_id. * Called by remote client after transfer completion. diff --git a/mooncake-store/include/rpc_types.h b/mooncake-store/include/rpc_types.h index fbf4fe3779..e1f3091240 100644 --- a/mooncake-store/include/rpc_types.h +++ b/mooncake-store/include/rpc_types.h @@ -273,4 +273,42 @@ struct BatchGetOffloadObjectResponse { YLT_REFL(BatchGetOffloadObjectResponse, batch_id, pointers, transfer_engine_addr, gc_ttl_ms); +// One destination slice on the requester side. In push mode the data owner +// writes (one-sided) the on-disk blob of a key directly into these addresses. +struct OffloadDstSlice { + uint64_t addr; + uint64_t size; + + OffloadDstSlice() : addr(0), size(0) {} + OffloadDstSlice(uint64_t addr_param, uint64_t size_param) + : addr(addr_param), size(size_param) {} +}; +YLT_REFL(OffloadDstSlice, addr, size); + +// Push-mode offload request. Unlike the pull path (where the requester gets +// back ClientBuffer pointers and issues the RDMA READ itself), here the +// requester hands the owner its own transfer engine endpoint and destination +// slice addresses. The owner reads SSD into its registered ClientBuffer and +// then WRITEs straight into the requester's memory, so the requester needs no +// follow-up READ nor a separate release_offload_buffer RPC. +struct BatchGetOffloadObjectPushRequest { + std::vector keys; // tenant-scoped storage keys + std::vector sizes; // total bytes per key (for FileStorage) + std::string requester_te_addr; // requester's transfer engine endpoint + std::vector> dst_slices; // per-key dst slices + + BatchGetOffloadObjectPushRequest() = default; +}; +YLT_REFL(BatchGetOffloadObjectPushRequest, keys, sizes, requester_te_addr, + dst_slices); + +struct BatchGetOffloadObjectPushResponse { + ErrorCode error_code; // overall result; data is already in requester memory + + BatchGetOffloadObjectPushResponse() : error_code(ErrorCode::OK) {} + explicit BatchGetOffloadObjectPushResponse(ErrorCode error_code_param) + : error_code(error_code_param) {} +}; +YLT_REFL(BatchGetOffloadObjectPushResponse, error_code); + } // namespace mooncake diff --git a/mooncake-store/include/transfer_task.h b/mooncake-store/include/transfer_task.h index 5fcb60cfa1..52f44e412f 100644 --- a/mooncake-store/include/transfer_task.h +++ b/mooncake-store/include/transfer_task.h @@ -18,6 +18,7 @@ #include "transfer_engine.h" #include "types.h" #include "replica.h" +#include "rpc_types.h" #include "storage_backend.h" #include "client_metric.h" #ifdef USE_NOF @@ -585,6 +586,16 @@ class TransferSubmitter { const std::unordered_map>& batched_slices); + // Push counterpart of submit_batch_get_offload_object, run on the data + // owner. WRITEs each key's on-disk blob (staged in the owner's local + // ClientBuffer at src_pointers[i]) directly into the requester's + // destination slices, opening the requester's segment once. + std::optional submit_batch_push_offload_object( + const std::string& requester_te_addr, + const std::vector& keys, + const std::vector& src_pointers, + const std::vector>& dst_slices); + /** * @brief Pure comparison helper: returns true iff both endpoints are * non-empty and identical. Exposed for unit testing of the locality diff --git a/mooncake-store/src/client_service.cpp b/mooncake-store/src/client_service.cpp index a826ebf151..fcca53752c 100644 --- a/mooncake-store/src/client_service.cpp +++ b/mooncake-store/src/client_service.cpp @@ -3223,6 +3223,26 @@ tl::expected Client::BatchGetOffloadObject( return {}; } +tl::expected Client::BatchPushOffloadObject( + const std::string& requester_te_addr, + const std::vector& keys, + const std::vector& src_pointers, + const std::vector>& dst_slices) { + auto future = transfer_submitter_->submit_batch_push_offload_object( + requester_te_addr, keys, src_pointers, dst_slices); + if (!future) { + MC_LOG(ERROR) << "Failed to submit push transfer operation"; + return tl::make_unexpected(ErrorCode::TRANSFER_FAIL); + } + MC_VLOG(1) << "Using transfer strategy: " << future->strategy(); + auto result = future->get(); + if (result != ErrorCode::OK) { + MC_LOG(ERROR) << "Push transfer failed, error code is " << result; + return tl::make_unexpected(result); + } + return {}; +} + tl::expected Client::NotifyOffloadSuccess( const std::vector& keys, const std::vector& metadatas) { diff --git a/mooncake-store/src/real_client.cpp b/mooncake-store/src/real_client.cpp index 4e67ecc78f..35a3085a63 100644 --- a/mooncake-store/src/real_client.cpp +++ b/mooncake-store/src/real_client.cpp @@ -927,6 +927,8 @@ tl::expected RealClient::setup_internal( std::make_unique(1, 0, "0.0.0.0"); offload_rpc_server_ ->register_handler<&RealClient::batch_get_offload_object>(this); + offload_rpc_server_ + ->register_handler<&RealClient::batch_get_offload_object_push>(this); offload_rpc_server_ ->register_handler<&RealClient::release_offload_buffer>(this); offload_rpc_server_->async_start(); @@ -6235,6 +6237,65 @@ RealClient::batch_get_offload_object(const std::vector &keys, file_storage_->config_.client_buffer_gc_ttl_ms); } +async_simple::coro::Lazy< + tl::expected> +RealClient::batch_get_offload_object_push( + const BatchGetOffloadObjectPushRequest &req) { + if (!file_storage_) { + LOG(ERROR) + << "batch_get_offload_object_push called but file_storage_ is null"; + co_return tl::make_unexpected(ErrorCode::INVALID_PARAMS); + } + if (req.keys.size() != req.sizes.size() || + req.keys.size() != req.dst_slices.size()) { + LOG(ERROR) << "batch_get_offload_object_push size mismatch: keys=" + << req.keys.size() << ", sizes=" << req.sizes.size() + << ", dst_slices=" << req.dst_slices.size(); + co_return tl::make_unexpected(ErrorCode::INVALID_PARAMS); + } + // Do the SSD read, the one-sided WRITE and the buffer release on the + // dedicated thread pool so the coro_rpc IO thread stays free. Same + // heap-owned state pattern as batch_get_offload_object: the lambda captures + // only a raw pointer. + struct CallState { + BatchGetOffloadObjectPushRequest req; + std::shared_ptr file_storage; + Client *client; + }; + auto state = std::make_unique(); + state->req = req; + state->file_storage = file_storage_; + state->client = client_.get(); + auto *s = state.get(); + auto try_result = + co_await coro_io::post([s]() -> tl::expected { + // Stage the on-disk blobs into the local ClientBuffer. + auto result = s->file_storage->BatchGet(s->req.keys, s->req.sizes); + if (!result) { + LOG(ERROR) << "Push offload BatchGet failed, err_code = " + << result.error(); + return tl::make_unexpected(result.error()); + } + const uint64_t batch_id = result.value().batch_id; + // WRITE the staged blobs straight into the requester's memory. + auto write_result = s->client->BatchPushOffloadObject( + s->req.requester_te_addr, s->req.keys, result.value().pointers, + s->req.dst_slices); + // The WRITE has completed (BatchPushOffloadObject blocks on the + // transfer future), so the ClientBuffer can be reclaimed + // immediately instead of waiting for the GC lease. + s->file_storage->ReleaseBuffer(batch_id); + return write_result; + }); + auto pushed = try_result.value(); + if (!pushed) { + LOG(ERROR) << "Push offload transfer failed, err_code = " + << pushed.error(); + co_return tl::make_unexpected(pushed.error()); + } + co_return BatchGetOffloadObjectPushResponse(ErrorCode::OK); +} + bool RealClient::release_offload_buffer(uint64_t batch_id) { if (!file_storage_) { LOG(WARNING) @@ -6253,14 +6314,69 @@ RealClient::batch_get_into_offload_object_internal( std::vector keys; std::vector storage_keys; std::vector sizes; + // Per-key destination slices, kept in lockstep with storage_keys so the + // push path can hand the owner positionally-aligned (key, dst) pairs. + std::vector> dst_slices; for (const auto &object_it : objects) { keys.emplace_back(object_it.first); storage_keys.emplace_back( MakeTenantScopedStorageKey(client_->tenant_id(), object_it.first)); int64_t total = 0; - for (const auto &s : object_it.second) total += s.size; + std::vector key_dst; + key_dst.reserve(object_it.second.size()); + for (const auto &s : object_it.second) { + total += s.size; + key_dst.emplace_back(reinterpret_cast(s.ptr), s.size); + } sizes.emplace_back(total); + dst_slices.emplace_back(std::move(key_dst)); + } + + // Push mode (MC_OFFLOAD_PUSH=true): the owner WRITEs the SSD data straight + // into our destination slices, so a single RPC replaces the pull path's + // get-pointers + READ + release_offload_buffer sequence. Falls back to the + // pull path otherwise (and for transports without one-sided WRITE). + static const bool kOffloadPush = []() { + const char *v = std::getenv("MC_OFFLOAD_PUSH"); + return v && (std::string_view(v) == "true" || + std::string_view(v) == "1"); + }(); + if (kOffloadPush) { + BatchGetOffloadObjectPushRequest push_req; + push_req.keys = storage_keys; + push_req.sizes = sizes; + push_req.requester_te_addr = client_->GetSegmentEndpoint(); + push_req.dst_slices = std::move(dst_slices); + + UbDiag::PerfPoint pt_push(PerfKey::GET_SSD_OFFLOAD_RPC, + UbDiag::PerfLevel::MODULE); + pt_push.Start(); + auto pushResp = client_requester_->batch_get_offload_object_push( + target_rpc_service_addr, push_req); + pt_push.End(pushResp ? 0 : -1); + if (!pushResp) { + LOG(ERROR) << "Push offload failed with error: " + << pushResp.error(); + return tl::make_unexpected(pushResp.error()); + } + if (pushResp->error_code != ErrorCode::OK) { + LOG(ERROR) << "Push offload owner-side error: " + << pushResp->error_code; + return tl::make_unexpected(pushResp->error_code); + } + auto end_time = std::chrono::steady_clock::now(); + auto elapsed_time = static_cast( + std::chrono::duration_cast(end_time - + start_time) + .count()); + LOG(INFO) << "Time taken for batch_get_into_offload_object_internal " + "(push): " + << elapsed_time << "ms, with target_rpc_service_addr: " + << target_rpc_service_addr + << ", key size: " << objects.size(); + return {}; } + UbDiag::PerfPoint pt_rpc(PerfKey::GET_SSD_OFFLOAD_RPC, UbDiag::PerfLevel::MODULE); pt_rpc.Start(); @@ -6355,6 +6471,21 @@ ClientRequester::batch_get_offload_object(const std::string &client_addr, return result; } +tl::expected +ClientRequester::batch_get_offload_object_push( + const std::string &client_addr, + const BatchGetOffloadObjectPushRequest &req) { + auto result = invoke_rpc<&RealClient::batch_get_offload_object_push, + BatchGetOffloadObjectPushResponse>(client_addr, + req); + if (!result) { + LOG(ERROR) + << "Failed to invoke batch_get_offload_object_push, client_addr = " + << client_addr << ", error is: " << result.error(); + } + return result; +} + void ClientRequester::release_offload_buffer(const std::string &client_addr, uint64_t batch_id) { // Fire-and-forget: attempt to release buffer, log errors but don't block diff --git a/mooncake-store/src/real_client_main.cpp b/mooncake-store/src/real_client_main.cpp index 8f0a4ba80c..10cd3a71c7 100644 --- a/mooncake-store/src/real_client_main.cpp +++ b/mooncake-store/src/real_client_main.cpp @@ -86,6 +86,8 @@ void RegisterClientRpcService(coro_rpc::coro_rpc_server &server, server.register_handler<&RealClient::query_task>(&real_client); server.register_handler<&RealClient::batch_get_offload_object>( &real_client); + server.register_handler<&RealClient::batch_get_offload_object_push>( + &real_client); server.register_handler<&RealClient::release_offload_buffer>(&real_client); } } // namespace mooncake diff --git a/mooncake-store/src/transfer_task.cpp b/mooncake-store/src/transfer_task.cpp index ce58d2325b..1dd440f0b0 100644 --- a/mooncake-store/src/transfer_task.cpp +++ b/mooncake-store/src/transfer_task.cpp @@ -1114,6 +1114,46 @@ TransferSubmitter::submit_batch_get_offload_object( return submitTransfer(requests); } +std::optional +TransferSubmitter::submit_batch_push_offload_object( + const std::string& requester_te_addr, + const std::vector& keys, + const std::vector& src_pointers, + const std::vector>& dst_slices) { + if (keys.size() != src_pointers.size() || + keys.size() != dst_slices.size()) { + MC_LOG(ERROR) << "submit_batch_push_offload_object size mismatch: keys=" + << keys.size() << ", src_pointers=" << src_pointers.size() + << ", dst_slices=" << dst_slices.size(); + return std::nullopt; + } + std::vector requests; + // Open the requester's segment once — all keys share requester_te_addr. + SegmentHandle seg = engine_.openSegment(requester_te_addr); + if (seg == static_cast(ERR_INVALID_ARGUMENT)) { + MC_LOG(ERROR) << "Failed to open segment " << requester_te_addr; + return std::nullopt; + } + for (size_t i = 0; i < keys.size(); ++i) { + // src is the owner's local ClientBuffer blob for this key (contiguous, + // registered via FileStorage::RegisterLocalMemory). It is split across + // the requester's possibly non-contiguous destination slices. + const uint64_t src = src_pointers[i]; + uint64_t offset = 0; + for (const auto& dst : dst_slices[i]) { + TransferRequest request; + request.opcode = TransferRequest::WRITE; + request.source = reinterpret_cast(src + offset); + request.target_id = seg; + request.target_offset = dst.addr; + request.length = dst.size; + requests.emplace_back(request); + offset += dst.size; + } + } + return submitTransfer(requests); +} + std::optional TransferSubmitter::submitMemcpyOperation( const AllocatedBuffer::Descriptor& handle, const std::vector& slices, const TransferRequest::OpCode op_code, uint64_t src_offset) { From 9cfdfbad268ff24c022fd703cc23219d1a3e304a Mon Sep 17 00:00:00 2001 From: Saturday Date: Tue, 23 Jun 2026 12:04:05 +0800 Subject: [PATCH 2/4] =?UTF-8?q?feat(perf):=20=E4=B8=BA=E6=89=80=E6=9C=89?= =?UTF-8?q?=E8=80=85=E7=AB=AFSSD=E5=8D=B8=E8=BD=BD=E6=93=8D=E4=BD=9C?= =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E6=80=A7=E8=83=BD=E7=9B=91=E6=8E=A7=E7=82=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 在mooncake_perf_points.def中定义三个新的性能监控点: - GET_SSD_OWNER_READ:所有者端SSD读取操作 - GET_SSD_OWNER_PUSH_WRITE:所有者端推送写入操作 - GET_SSD_OWNER_RELEASE:所有者端缓冲区释放操作 在real_client.cpp中相应位置添加性能监控代码,覆盖以下场景: 1. 拉取路径的批量获取(batch_get_offload_object) 2. 推送路径的批量获取和写入(batch_get_offload_object_push) 3. 请求端触发的缓冲区释放(release_offload_buffer) 这些监控点将帮助分析所有者端在SSD卸载操作中的性能表现。 --- .../store/mooncake_perf_points.def | 5 +++ mooncake-store/src/real_client.cpp | 33 ++++++++++++++++--- 2 files changed, 34 insertions(+), 4 deletions(-) diff --git a/mooncake-integration/store/mooncake_perf_points.def b/mooncake-integration/store/mooncake_perf_points.def index abf3a228c5..efec3b828e 100644 --- a/mooncake-integration/store/mooncake_perf_points.def +++ b/mooncake-integration/store/mooncake_perf_points.def @@ -108,6 +108,11 @@ PERF_KEY_DEF(GET_SSD_OFFLOAD_RPC, "real_client.cpp::batch_get_into_of PERF_KEY_DEF(GET_SSD_TRANSFER_DATA, "real_client.cpp::batch_get_into_offload_object_internal", "TransferData") PERF_KEY_DEF(GET_SSD_RELEASE_BUFFER, "real_client.cpp::batch_get_into_offload_object_internal", "ReleaseBuffer") +// === Owner-side offload sub-steps (pull & push) === +PERF_KEY_DEF(GET_SSD_OWNER_READ, "real_client.cpp::offload_owner", "OwnerSsdRead") +PERF_KEY_DEF(GET_SSD_OWNER_PUSH_WRITE, "real_client.cpp::batch_get_offload_object_push", "OwnerPushWrite") +PERF_KEY_DEF(GET_SSD_OWNER_RELEASE, "real_client.cpp::offload_owner", "OwnerReleaseBuffer") + // === UB / URMA endpoint first connection path === PERF_KEY_DEF(UB_HANDSHAKE_ENCODE, "transfer_metadata.cpp::TransferHandshakeUtil::encode", "Encode") PERF_KEY_DEF(UB_HANDSHAKE_DECODE, "transfer_metadata.cpp::TransferHandshakeUtil::decode", "Decode") diff --git a/mooncake-store/src/real_client.cpp b/mooncake-store/src/real_client.cpp index 35a3085a63..27a8116543 100644 --- a/mooncake-store/src/real_client.cpp +++ b/mooncake-store/src/real_client.cpp @@ -6223,8 +6223,15 @@ RealClient::batch_get_offload_object(const std::vector &keys, state->sizes = sizes; state->file_storage = file_storage_; auto *s = state.get(); - auto try_result = co_await coro_io::post( - [s]() { return s->file_storage->BatchGet(s->keys, s->sizes); }); + auto try_result = co_await coro_io::post([s]() { + // Owner-side SSD -> ClientBuffer read (pull path). + UbDiag::PerfPoint pt_read(PerfKey::GET_SSD_OWNER_READ, + UbDiag::PerfLevel::MODULE); + pt_read.Start(); + auto r = s->file_storage->BatchGet(s->keys, s->sizes); + pt_read.End(r ? 0 : -1); + return r; + }); auto result = try_result.value(); if (!result) { LOG(ERROR) << "Batch get offload object failed,err_code = " @@ -6269,8 +6276,12 @@ RealClient::batch_get_offload_object_push( auto *s = state.get(); auto try_result = co_await coro_io::post([s]() -> tl::expected { - // Stage the on-disk blobs into the local ClientBuffer. + // Owner-side SSD -> ClientBuffer read (push path). + UbDiag::PerfPoint pt_read(PerfKey::GET_SSD_OWNER_READ, + UbDiag::PerfLevel::MODULE); + pt_read.Start(); auto result = s->file_storage->BatchGet(s->req.keys, s->req.sizes); + pt_read.End(result ? 0 : -1); if (!result) { LOG(ERROR) << "Push offload BatchGet failed, err_code = " << result.error(); @@ -6278,13 +6289,21 @@ RealClient::batch_get_offload_object_push( } const uint64_t batch_id = result.value().batch_id; // WRITE the staged blobs straight into the requester's memory. + UbDiag::PerfPoint pt_write(PerfKey::GET_SSD_OWNER_PUSH_WRITE, + UbDiag::PerfLevel::MODULE); + pt_write.Start(); auto write_result = s->client->BatchPushOffloadObject( s->req.requester_te_addr, s->req.keys, result.value().pointers, s->req.dst_slices); + pt_write.End(write_result ? 0 : -1); // The WRITE has completed (BatchPushOffloadObject blocks on the // transfer future), so the ClientBuffer can be reclaimed // immediately instead of waiting for the GC lease. + UbDiag::PerfPoint pt_release(PerfKey::GET_SSD_OWNER_RELEASE, + UbDiag::PerfLevel::MODULE); + pt_release.Start(); s->file_storage->ReleaseBuffer(batch_id); + pt_release.End(0); return write_result; }); auto pushed = try_result.value(); @@ -6302,7 +6321,13 @@ bool RealClient::release_offload_buffer(uint64_t batch_id) { << "release_offload_buffer called but file_storage_ is null"; return false; } - return file_storage_->ReleaseBuffer(batch_id); + // Owner-side buffer release triggered by the requester's RPC (pull path). + UbDiag::PerfPoint pt_release(PerfKey::GET_SSD_OWNER_RELEASE, + UbDiag::PerfLevel::MODULE); + pt_release.Start(); + bool released = file_storage_->ReleaseBuffer(batch_id); + pt_release.End(released ? 0 : -1); + return released; } tl::expected From aaa4456ca85f7632aab0ef79384c32c6a086123f Mon Sep 17 00:00:00 2001 From: Saturday Date: Tue, 23 Jun 2026 16:12:12 +0800 Subject: [PATCH 3/4] =?UTF-8?q?perf(offload):=20=E5=B0=86=E6=80=A7?= =?UTF-8?q?=E8=83=BD=E7=82=B9=E7=A7=BB=E8=87=B3=20FileStorage=20=E4=BB=A5?= =?UTF-8?q?=E7=B2=BE=E7=A1=AE=E6=B5=8B=E9=87=8F=20SSD=20=E8=AF=BB=E5=8F=96?= =?UTF-8?q?=E9=98=B6=E6=AE=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 将 GET_SSD_OWNER_READ 和 GET_SSD_OWNER_RELEASE 性能点的测量位置从 real_client.cpp 移至 file_storage.cpp 的实现内部。同时,在 FileStorage::BatchGet 中新增 GET_SSD_OWNER_ALLOC 和 GET_SSD_OWNER_LOAD 性能点,以细分 SSD 读取路径中的缓冲区分配和磁盘加载阶段。 这使得性能测量更精确地反映实际耗时操作,避免了在协程调度器中测量可能引入的误差,并确保了 push 和 pull 路径测量的一致性。 --- .../store/mooncake_perf_points.def | 6 ++-- mooncake-store/src/file_storage.cpp | 31 +++++++++++++++++ mooncake-store/src/real_client.cpp | 33 +++++-------------- 3 files changed, 43 insertions(+), 27 deletions(-) diff --git a/mooncake-integration/store/mooncake_perf_points.def b/mooncake-integration/store/mooncake_perf_points.def index efec3b828e..85983deb05 100644 --- a/mooncake-integration/store/mooncake_perf_points.def +++ b/mooncake-integration/store/mooncake_perf_points.def @@ -109,9 +109,11 @@ PERF_KEY_DEF(GET_SSD_TRANSFER_DATA, "real_client.cpp::batch_get_into_of PERF_KEY_DEF(GET_SSD_RELEASE_BUFFER, "real_client.cpp::batch_get_into_offload_object_internal", "ReleaseBuffer") // === Owner-side offload sub-steps (pull & push) === -PERF_KEY_DEF(GET_SSD_OWNER_READ, "real_client.cpp::offload_owner", "OwnerSsdRead") +PERF_KEY_DEF(GET_SSD_OWNER_READ, "file_storage.cpp::BatchGet", "OwnerSsdRead") +PERF_KEY_DEF(GET_SSD_OWNER_ALLOC, "file_storage.cpp::BatchGet", "OwnerAllocBuffer") +PERF_KEY_DEF(GET_SSD_OWNER_LOAD, "file_storage.cpp::BatchGet", "OwnerDiskLoad") PERF_KEY_DEF(GET_SSD_OWNER_PUSH_WRITE, "real_client.cpp::batch_get_offload_object_push", "OwnerPushWrite") -PERF_KEY_DEF(GET_SSD_OWNER_RELEASE, "real_client.cpp::offload_owner", "OwnerReleaseBuffer") +PERF_KEY_DEF(GET_SSD_OWNER_RELEASE, "file_storage.cpp::ReleaseBuffer", "OwnerReleaseBuffer") // === UB / URMA endpoint first connection path === PERF_KEY_DEF(UB_HANDSHAKE_ENCODE, "transfer_metadata.cpp::TransferHandshakeUtil::encode", "Encode") diff --git a/mooncake-store/src/file_storage.cpp b/mooncake-store/src/file_storage.cpp index 4830c13eb3..754fbfda17 100644 --- a/mooncake-store/src/file_storage.cpp +++ b/mooncake-store/src/file_storage.cpp @@ -12,6 +12,14 @@ #include "file_interface.h" #endif +// ubdiag perf points for the offload owner-side SSD path. UBDIAG_PROGRAM_NAME +// must match the other TUs (store_py/real_client/client_service); each TU's +// static initializer writes the shared detail::AutoProgramName(), and leaving +// it nullptr here could clobber the program name depending on init order. +#define UBDIAG_PERF_DEF_FILE "mooncake_perf_points.def" +#define UBDIAG_PROGRAM_NAME "mooncake_store" +#include "ubdiag/auto_perf.h" + namespace mooncake { using gpu_staging::CopyDeviceToHost; @@ -323,13 +331,28 @@ tl::expected FileStorage::Init() { tl::expected FileStorage::BatchGet( const std::vector& keys, const std::vector& sizes) { auto start_time = std::chrono::steady_clock::now(); + // Owner-side SSD read total (OwnerSsdRead); broken down into buffer + // allocation (OwnerAllocBuffer) and disk load (OwnerDiskLoad). End() is on + // the success path only — the error early-returns let the dtor Abandon the + // unfinished total sample. + UbDiag::PerfPoint pt_read(PerfKey::GET_SSD_OWNER_READ, + UbDiag::PerfLevel::MODULE); + pt_read.Start(); + UbDiag::PerfPoint pt_alloc(PerfKey::GET_SSD_OWNER_ALLOC, + UbDiag::PerfLevel::MODULE); + pt_alloc.Start(); auto allocate_res = AllocateBatch(keys, sizes); + pt_alloc.End(allocate_res ? 0 : -1); if (!allocate_res) { LOG(ERROR) << "Failed to allocate batch objects"; return tl::make_unexpected(allocate_res.error()); } auto allocated_batch = allocate_res.value(); + UbDiag::PerfPoint pt_load(PerfKey::GET_SSD_OWNER_LOAD, + UbDiag::PerfLevel::MODULE); + pt_load.Start(); auto result = BatchLoad(allocated_batch->slices); + pt_load.End(result ? 0 : -1); if (!result) { LOG(ERROR) << "Batch load object failed,err_code = " << result.error(); return tl::make_unexpected(result.error()); @@ -358,6 +381,7 @@ tl::expected FileStorage::BatchGet( .count(); VLOG(1) << "Time taken for FileStorage::BatchGet: " << elapsed_time << "us, key size: " << keys.size() << ", batch_id: " << batch_id; + pt_read.End(0); return batch_result; } @@ -1010,16 +1034,23 @@ void FileStorage::ClientBufferGCThreadFunc() { } bool FileStorage::ReleaseBuffer(uint64_t batch_id) { + // Owner-side ClientBuffer release (OwnerReleaseBuffer); shared by the pull + // path (release_offload_buffer RPC) and the push path (inline after WRITE). + UbDiag::PerfPoint pt_release(PerfKey::GET_SSD_OWNER_RELEASE, + UbDiag::PerfLevel::MODULE); + pt_release.Start(); MutexLocker locker(&client_buffer_mutex_); auto it = client_buffer_allocated_batches_.find(batch_id); if (it != client_buffer_allocated_batches_.end()) { VLOG(1) << "Releasing buffer for batch_id: " << batch_id << " (transfer completed)"; client_buffer_allocated_batches_.erase(it); + pt_release.End(0); return true; } VLOG(1) << "batch_id " << batch_id << " not found (may have been GC'd already)"; + pt_release.End(-1); return false; } diff --git a/mooncake-store/src/real_client.cpp b/mooncake-store/src/real_client.cpp index 27a8116543..0ff8c48c82 100644 --- a/mooncake-store/src/real_client.cpp +++ b/mooncake-store/src/real_client.cpp @@ -6223,15 +6223,8 @@ RealClient::batch_get_offload_object(const std::vector &keys, state->sizes = sizes; state->file_storage = file_storage_; auto *s = state.get(); - auto try_result = co_await coro_io::post([s]() { - // Owner-side SSD -> ClientBuffer read (pull path). - UbDiag::PerfPoint pt_read(PerfKey::GET_SSD_OWNER_READ, - UbDiag::PerfLevel::MODULE); - pt_read.Start(); - auto r = s->file_storage->BatchGet(s->keys, s->sizes); - pt_read.End(r ? 0 : -1); - return r; - }); + auto try_result = co_await coro_io::post( + [s]() { return s->file_storage->BatchGet(s->keys, s->sizes); }); auto result = try_result.value(); if (!result) { LOG(ERROR) << "Batch get offload object failed,err_code = " @@ -6276,12 +6269,9 @@ RealClient::batch_get_offload_object_push( auto *s = state.get(); auto try_result = co_await coro_io::post([s]() -> tl::expected { - // Owner-side SSD -> ClientBuffer read (push path). - UbDiag::PerfPoint pt_read(PerfKey::GET_SSD_OWNER_READ, - UbDiag::PerfLevel::MODULE); - pt_read.Start(); + // Stage the on-disk blobs into the local ClientBuffer + // (FileStorage::BatchGet carries the OwnerSsdRead breakdown). auto result = s->file_storage->BatchGet(s->req.keys, s->req.sizes); - pt_read.End(result ? 0 : -1); if (!result) { LOG(ERROR) << "Push offload BatchGet failed, err_code = " << result.error(); @@ -6298,12 +6288,9 @@ RealClient::batch_get_offload_object_push( pt_write.End(write_result ? 0 : -1); // The WRITE has completed (BatchPushOffloadObject blocks on the // transfer future), so the ClientBuffer can be reclaimed - // immediately instead of waiting for the GC lease. - UbDiag::PerfPoint pt_release(PerfKey::GET_SSD_OWNER_RELEASE, - UbDiag::PerfLevel::MODULE); - pt_release.Start(); + // immediately instead of waiting for the GC lease + // (FileStorage::ReleaseBuffer carries OwnerReleaseBuffer). s->file_storage->ReleaseBuffer(batch_id); - pt_release.End(0); return write_result; }); auto pushed = try_result.value(); @@ -6322,12 +6309,8 @@ bool RealClient::release_offload_buffer(uint64_t batch_id) { return false; } // Owner-side buffer release triggered by the requester's RPC (pull path). - UbDiag::PerfPoint pt_release(PerfKey::GET_SSD_OWNER_RELEASE, - UbDiag::PerfLevel::MODULE); - pt_release.Start(); - bool released = file_storage_->ReleaseBuffer(batch_id); - pt_release.End(released ? 0 : -1); - return released; + // FileStorage::ReleaseBuffer carries the OwnerReleaseBuffer perf point. + return file_storage_->ReleaseBuffer(batch_id); } tl::expected From 46c3cc43ec69e3766041badbd1883eb5c5b9b2d7 Mon Sep 17 00:00:00 2001 From: Saturday Date: Wed, 24 Jun 2026 15:03:24 +0800 Subject: [PATCH 4/4] =?UTF-8?q?feat(offload):=20=E6=96=B0=E5=A2=9E=20SSD?= =?UTF-8?q?=20offload=20push=20=E6=A8=A1=E5=BC=8F=E4=BB=A5=E4=BC=98?= =?UTF-8?q?=E5=8C=96=E8=AF=BB=E5=8F=96=E6=80=A7=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 在 mooncake_perf_points.def 中添加三个性能监控点用于跟踪 owner 端磁盘加载各阶段耗时 - 在 storage_backend.cpp 中实现性能监控,分别测量读取计划构建、io_uring 读取和 POSIX 读取时间 - 新增详细设计文档 offload-push-mode.md,描述 push 模式架构、实现细节和性能优势 - push 模式通过环境变量 MC_OFFLOAD_PUSH 控制,可将 offload 读取从 3 次 RPC 减少为 1 次 - 对端 owner 主动执行 URMA write 将数据推送到请求方内存,减少网络往返延迟 --- docs/offload-push-mode.md | 374 ++++++++++++++++++ .../store/mooncake_perf_points.def | 3 + mooncake-store/src/storage_backend.cpp | 22 +- 3 files changed, 398 insertions(+), 1 deletion(-) create mode 100644 docs/offload-push-mode.md diff --git a/docs/offload-push-mode.md b/docs/offload-push-mode.md new file mode 100644 index 0000000000..2118501bf2 --- /dev/null +++ b/docs/offload-push-mode.md @@ -0,0 +1,374 @@ +# Mooncake SSD Offload —— Push 读取模式(URMA write 为例) + +> 本文档描述 `feature/offloadrpc` 分支引入的 **Push 读取模式**,作为 [offload-mechanism.md](offload-mechanism.md) 第 2.2 节「Load(SSD → 请求方)」的扩展。 +> 传输层以 **UB / URMA** 为例(`URMA_OPC_WRITE` 单边写);RDMA 等其它 transport 走同一套抽象,不单独展开。 +> 关联开关:`MC_OFFLOAD_PUSH`。 + +--- + +## 1. 背景与动机 + +当一个对象只在远端节点的 SSD 上时,请求方需要把它从对端磁盘读回本地内存。原有实现(下称 **Pull 模式**)由**请求方主动发起**,一次 Load 要 **3 次 RPC / 单边操作**;**Push 模式**把单边传输的方向对调,由**数据持有方(owner)主动 URMA write**,一次 Load 只需 **1 次 RPC 往返**。 + +### 1.1 Pull vs Push 流程图 + +**Pull 模式(请求方主动拉,3 次往返)** + +```mermaid +sequenceDiagram + autonumber + participant R as 请求方 + participant O as 对端 owner + R->>O: ① RPC batch_get_offload_object(keys) + Note over O: BatchGet:SSD → ClientBuffer + O-->>R: pointers + gc_ttl + R->>O: ② URMA READ:从对端 ClientBuffer 拉数据 + Note over O: (被动,数据被读走) + R->>O: ③ RPC release_offload_buffer(batch_id) + Note over O: ReleaseBuffer +``` + +**Push 模式(持有方主动推,1 次往返)** + +```mermaid +sequenceDiagram + autonumber + participant R as 请求方 + participant O as 对端 owner + R->>O: ① RPC batch_get_offload_object_push
(keys + 自身 TE 端点 + dst_slices) + Note over O: BatchGet:SSD → ClientBuffer + O->>R: ② URMA WRITE:ClientBuffer → 请求方内存 + Note over O: ReleaseBuffer(本地,写完即放) + O-->>R: error_code(数据已落在请求方内存) +``` + +**收益**:消除请求方侧的 URMA READ 往返,以及单独的 `release_offload_buffer` RPC —— 由对端写完后就地释放 buffer。 + +**不变的部分**:对端仍必须先把 SSD 数据读进**已注册的 ClientBuffer**(`FileStorage::BatchGet`)。URMA 单边写的源必须是注册过的内存段(`urma_register_seg` 得到的 `urma_target_seg_t`),SSD 上的数据无法绕过这块中转直接走网卡。Push 省的是后续步骤,不是这次 SSD→内存的拷贝。 + +--- + +## 2. 设计要点 + +### 2.1 方向对调(Pull READ ↔ Push WRITE) + +两条传输函数互为镜像,只差三个字段;最终都落到 URMA 的同一套 `urma_jfs_wr_t`,仅 `opcode` 与 SGE 方向不同: + +| | Pull `submit_batch_get_offload_object` | Push `submit_batch_push_offload_object` | +|---|---|---| +| `TransferRequest::opcode` | `READ` | `WRITE` | +| `openSegment` 的对象 | 对端(owner)的 segment | **请求方**的 segment | +| `source`(本地) | 请求方目标分片 `slice.ptr` | **对端 ClientBuffer** `src_pointer` | +| `target_offset`(远端) | 对端 ClientBuffer 地址 | **请求方目标分片地址** `dst.addr` | +| URMA `wr.opcode` | `URMA_OPC_READ` | `URMA_OPC_WRITE` | +| 发起方 | 请求方 | 对端 | + +### 2.2 成立前提:两端内存都注册为 URMA segment + +- **对端 ClientBuffer**:`FileStorage::RegisterLocalMemory()` 注册到对端 transfer engine,底层经 `urma_register_seg` 得到本地段 `l_seg`,可作为 WRITE 的源(local SGE)。 +- **请求方目标内存**:应用 GET 时传入的 `objects` 分片本就注册过;对端 `openSegment(requester_te_addr)` 把它**导入**为远端段 `r_seg`(remote SGE)+ 远端 `tjetty`,WRITE 才能寻址过去。 + +### 2.3 非连续目标分片 + +对端某个 key 的数据是**一整块连续** ClientBuffer;请求方接收内存可能是**多段非连续**分片(GPU 显存常见)。Push 按 `dst.size` 累加 `offset`,把连续源切到各目标分片,逐段生成一条 `TransferRequest` → 一条 `urma_jfs_wr_t`。 + +### 2.4 ub / urma 适配 + +Push 改动全部位于 `TransferSubmitter::submit_*` 层,只改 `opcode/source/target_offset`,**未触碰任何具体 transport**。`MultiTransport` 按端点自动选到 `UbTransport`;`urma_endpoint.cpp` 已支持 `URMA_OPC_WRITE` 且 SGE 方向自动反转,注册时 access 已开 `READ|WRITE|ATOMIC`。Push 无需为 ub/urma 写第二份逻辑。 + +--- + +## 3. RPC 数据结构(`mooncake-store/include/rpc_types.h`) + +```cpp +// 请求方一个目标分片:对端 URMA write 时写入的目的地址 + 长度 +struct OffloadDstSlice { + uint64_t addr; // 请求方目标内存虚拟地址 + uint64_t size; // 该分片字节数 +}; +YLT_REFL(OffloadDstSlice, addr, size); + +// Push 请求体 +struct BatchGetOffloadObjectPushRequest { + std::vector keys; // 租户作用域 storage key + std::vector sizes; // 每个 key 总字节数 + std::string requester_te_addr; // 请求方 transfer engine 端点 + std::vector> dst_slices; // 每个 key 一组目标分片 +}; +YLT_REFL(BatchGetOffloadObjectPushRequest, keys, sizes, requester_te_addr, dst_slices); + +// Push 响应体(数据已落在请求方内存,仅回状态码) +struct BatchGetOffloadObjectPushResponse { + ErrorCode error_code; +}; +YLT_REFL(BatchGetOffloadObjectPushResponse, error_code); +``` + +- `YLT_REFL` 是序列化反射宏;不加它,struct_pack 无法在 RPC 中编解码这些结构。 +- **不变量**:`keys`、`sizes`、`dst_slices` 是三个**平行数组**,长度必须相等,下标 `i` 描述同一个 key。handler 在进入任何按下标循环前先校验等长,坏请求直接 `INVALID_PARAMS` 拒绝(fail-fast,防越界/错位)。 + +--- + +## 4. 改动清单(逐文件) + +| 文件 | 改动 | +|---|---| +| `include/rpc_types.h` | 新增 `OffloadDstSlice` / `BatchGetOffloadObjectPushRequest` / `…PushResponse` | +| `include/transfer_task.h`、`src/transfer_task.cpp` | 新增 `submit_batch_push_offload_object`(WRITE 版传输);header 增加 `#include "rpc_types.h"` | +| `include/client_service.h`、`src/client_service.cpp` | 新增 `Client::BatchPushOffloadObject`(提交传输 + 等 future 完成) | +| `include/real_client.h`、`src/real_client.cpp` | 新增对端 handler `batch_get_offload_object_push`;请求方 `batch_get_into_offload_object_internal` 增加 `MC_OFFLOAD_PUSH` 分支 | +| `include/pyclient.h`、`src/real_client.cpp` | 新增 `ClientRequester::batch_get_offload_object_push`(invoke_rpc 封装) | +| `src/real_client.cpp`、`src/real_client_main.cpp` | 两处 server 各 `register_handler` 新 handler | + +> ⚠️ **handler 必须两处都注册**:内嵌 server(`real_client.cpp` 的 `offload_rpc_server_`)和独立进程 server(`real_client_main.cpp`)。漏一处,对应部署形态下 push 会因「RPC 未注册」失败。 + +--- + +## 5. 调用链总览 + +```mermaid +flowchart TB + subgraph REQ[请求方 本端] + A["batch_get_into_offload_object_internal"] + B["ClientRequester::batch_get_offload_object_push"] + C["invoke_rpc 模板 (&RealClient::batch_get_offload_object_push)"] + D["coro_rpc_client.send_request
struct_pack 序列化"] + A --> B --> C --> D + end + subgraph OWN[对端 owner] + E["RealClient::batch_get_offload_object_push"] + F["co_await coro_io::post(lambda)"] + G["FileStorage::BatchGet
SSD → ClientBuffer"] + G2["BucketStorageBackend::BatchLoad
preadv / io_uring"] + H["Client::BatchPushOffloadObject
URMA write"] + I["TransferSubmitter::submit_batch_push_offload_object
openSegment + submitTransfer"] + J["UbTransport::submitTransferTask"] + K["UrmaEndpoint::submitPostSend"] + L["urma_post_jetty_send_wr ★ URMA_OPC_WRITE"] + M["FileStorage::ReleaseBuffer
写完即释放"] + E --> F + F --> G --> G2 + F --> H --> I --> J --> K --> L + F --> M + end + D -- "网络 RPC" --> E + E -. "co_return error_code" .-> D +``` + +--- + +## 6. 关键代码解读 + +### 6.1 请求方入口:`batch_get_into_offload_object_internal`(`src/real_client.cpp`) + +收集目标地址,按 `MC_OFFLOAD_PUSH` 分流到 Push 或保留 Pull: + +```cpp +std::vector> dst_slices; +for (const auto &object_it : objects) { + storage_keys.emplace_back(MakeTenantScopedStorageKey(...)); + int64_t total = 0; + std::vector key_dst; + for (const auto &s : object_it.second) { + total += s.size; + key_dst.emplace_back(reinterpret_cast(s.ptr), s.size); // 应用目标内存地址+长度 + } + sizes.emplace_back(total); + dst_slices.emplace_back(std::move(key_dst)); // 与 storage_keys 对齐 +} + +static const bool kOffloadPush = []() { // 只读一次环境变量并缓存 + const char *v = std::getenv("MC_OFFLOAD_PUSH"); + return v && (std::string_view(v) == "true" || std::string_view(v) == "1"); +}(); +if (kOffloadPush) { + BatchGetOffloadObjectPushRequest push_req; + push_req.keys = storage_keys; + push_req.sizes = sizes; + push_req.requester_te_addr = client_->GetSegmentEndpoint(); // 自身 TE 端点 + push_req.dst_slices = std::move(dst_slices); + auto pushResp = client_requester_->batch_get_offload_object_push(target_rpc_service_addr, push_req); + if (!pushResp) { return tl::make_unexpected(pushResp.error()); } // RPC 层失败 + if (pushResp->error_code != ErrorCode::OK) { return tl::make_unexpected(pushResp->error_code); } + return {}; // ★ Push 到此结束:无 READ、无 release +} +// 否则走下方原有 Pull 链路(完全保留) +``` + +`s.ptr` 是应用 GET 时给定的目标内存(已注册段),转成 `uint64_t` 即 `OffloadDstSlice::addr`,与 URMA `r_sge.addr` 语义一致。 + +### 6.2 对端 handler:`RealClient::batch_get_offload_object_push`(`src/real_client.cpp`) + +读盘 + URMA write + 释放,在一个线程池任务里完成: + +```cpp +async_simple::coro::Lazy> +RealClient::batch_get_offload_object_push(const BatchGetOffloadObjectPushRequest &req) { + if (!file_storage_) { co_return tl::make_unexpected(ErrorCode::INVALID_PARAMS); } + if (req.keys.size() != req.sizes.size() || + req.keys.size() != req.dst_slices.size()) { co_return ... INVALID_PARAMS; } // 平行数组校验 + + struct CallState { req; file_storage; client; }; // 堆上打包,lambda 只捕获裸指针 + auto state = std::make_unique(); ... + auto *s = state.get(); + + auto try_result = co_await coro_io::post([s]() -> tl::expected { + auto result = s->file_storage->BatchGet(s->req.keys, s->req.sizes); // ① SSD → ClientBuffer + if (!result) { return tl::make_unexpected(result.error()); } + const uint64_t batch_id = result.value().batch_id; + auto write_result = s->client->BatchPushOffloadObject( // ② URMA write → 请求方内存 + s->req.requester_te_addr, s->req.keys, + result.value().pointers, // 对端 ClientBuffer 每个 key 的地址 = URMA write 的源 + s->req.dst_slices); + s->file_storage->ReleaseBuffer(batch_id); // ③ 写完即释放 + return write_result; + }); + + auto pushed = try_result.value(); + if (!pushed) { co_return tl::make_unexpected(pushed.error()); } + co_return BatchGetOffloadObjectPushResponse(ErrorCode::OK); +} +``` + +`coro_io::post` 把「读盘 + 等 URMA write 完成 + 释放」这些会阻塞的慢操作提交到阻塞线程池,`co_await` 让出 coro_rpc 的 IO 线程,使其继续处理 ping 等其它 RPC。 + +### 6.3 Client 封装:`Client::BatchPushOffloadObject`(`src/client_service.cpp`) + +```cpp +auto future = transfer_submitter_->submit_batch_push_offload_object(...); +if (!future) { return tl::make_unexpected(ErrorCode::TRANSFER_FAIL); } +auto result = future->get(); // ★ 阻塞到 URMA write 完成(jfc 收到完成事件) +if (result != ErrorCode::OK) { return tl::make_unexpected(result); } +return {}; +``` + +`future->get()` 阻塞到 URMA write 完成,这是对端能安全释放 ClientBuffer 的前提(否则源 buffer 在传输中途被回收会损坏数据)。 + +### 6.4 传输层:`submit_batch_push_offload_object`(`src/transfer_task.cpp`) + +把「一块连续源 → 多个目标分片」翻译成 transfer engine 的 WRITE 请求: + +```cpp +SegmentHandle seg = engine_.openSegment(requester_te_addr); // 打开/导入“请求方”的 segment(只开一次) +for (size_t i = 0; i < keys.size(); ++i) { + const uint64_t src = src_pointers[i]; // 这个 key 在对端 ClientBuffer 里的连续起始地址 + uint64_t offset = 0; + for (const auto& dst : dst_slices[i]) { + TransferRequest request; + request.opcode = TransferRequest::WRITE; // ★ WRITE + request.source = reinterpret_cast(src + offset); // 源 = 对端本地 buffer + request.target_id = seg; // 目标 = 请求方 segment + request.target_offset = dst.addr; // 目标地址 = 请求方分片地址 + request.length = dst.size; + requests.emplace_back(request); + offset += dst.size; + } +} +return submitTransfer(requests); +``` + +### 6.5 落到 URMA:`TransferRequest` → `urma_jfs_wr_t` + +`submitTransfer` → `UbTransport::submitTransferTask` 把每个 `TransferRequest` 切成 `Slice`(`ub_transport.cpp`): + +```cpp +slice->opcode = request.opcode; // WRITE 透传 +slice->source_addr = request.source; // 本地源(对端 ClientBuffer) +slice->ub.dest_addr = request.target_offset + offset; // 远端目标(请求方分片地址) +``` + +再由 `UrmaEndpoint::submitPostSend`(`urma_endpoint.cpp`)组装成 URMA 工作请求并提交: + +```cpp +// 本地 SGE(源):对端 ClientBuffer +l_sge.addr = (uint64_t)slice->source_addr; +l_sge.tseg = slice->ub.l_seg; // 本地注册段(urma_register_seg) +// 远端 SGE(目标):请求方内存 +r_sge.addr = slice->ub.dest_addr; +r_sge.tseg = slice->ub.r_seg; // openSegment 导入的远端段 + +wr.opcode = (slice->opcode == READ) ? URMA_OPC_READ : URMA_OPC_WRITE; // ★ 本路径 = URMA_OPC_WRITE +wr.rw.src.sge = (READ) ? &r_sge : &l_sge; // WRITE:源 = 本地 l_sge +wr.rw.dst.sge = (READ) ? &l_sge : &r_sge; // WRITE:目标 = 远端 r_sge +wr.tjetty = imported_jetty_map_[jetty]; // 导入的远端 jetty + +urma_post_jetty_send_wr(jetty_list_[jetty_index], wr_list, &bad_wr); // 提交,完成经 jfc 通知 +``` + +> URMA 术语对照:`jetty` ≈ 收发队列(类 QP),`urma_target_seg_t` ≈ 注册内存段(类 MR),`jfc` ≈ 完成队列(类 CQ)。Push 路径就是构造 `URMA_OPC_WRITE` 的 `jfs_wr`,源 SGE 指向对端 ClientBuffer,目标 SGE 指向请求方导入段。 + +### 6.6 请求方 RPC 封装:`ClientRequester::batch_get_offload_object_push` + +```cpp +auto result = invoke_rpc<&RealClient::batch_get_offload_object_push, + BatchGetOffloadObjectPushResponse>(client_addr, req); +``` + +`invoke_rpc(addr, 参数...)` 用成员函数指针 `&RealClient::batch_get_offload_object_push` 作为「调对端哪个 handler」的编译期 ID;对端 `register_handler<同一个指针>` 把该 ID 映射回 handler。 + +--- + +## 7. 时序图(含线程模型) + +```mermaid +sequenceDiagram + autonumber + participant R as 请求方线程
(syncAwait 阻塞) + participant IO as 对端 IO 线程
(server 仅 1 条) + participant W as 对端 worker 线程
(coro_io 阻塞池) + participant U as URMA / 网卡 + R->>IO: send_request(请求) + Note over IO: 反序列化 req + 平行数组校验 + IO->>W: co_await coro_io::post + Note over IO: IO 线程让出,去收别的 RPC + Note over W: FileStorage::BatchGet
preadv 读盘 → ClientBuffer + W->>U: submitPostSend / urma_post_jetty_send_wr
(URMA_OPC_WRITE) + Note over U: ClientBuffer → 请求方内存 + U-->>W: jfc 完成事件(future->get 返回) + Note over W: FileStorage::ReleaseBuffer + W-->>IO: lambda 完成,协程恢复 + IO-->>R: 响应(error_code) +``` + +- **IO 线程**:coro_rpc server 事件循环,只做收包/反序列化/分发/回包等快操作,**绝不阻塞**;内嵌 offload server 仅 1 条(`coro_rpc_server(1, 0, ...)`)。 +- **worker 线程**:coro_io 共享阻塞线程池,跑 `coro_io::post` 提交的慢活(读盘、等 URMA write、释放)。 +- coro_rpc **默认不会**自动给每个请求分配工作线程 —— handler 默认就在 IO 线程跑。任务挪到 worker 线程,是 handler **主动 `coro_io::post`** 的结果。这也是 offload server 只配 1 条 IO 线程也不会被读盘/传输拖死的原因。 + +--- + +## 8. 开关:`MC_OFFLOAD_PUSH` + +| 取值 | 行为 | +|---|---| +| 未设置 / 其它 | **Pull 模式**(默认),原 3 步链路完全保留 | +| `true` 或 `1` | **Push 模式** | + +- 在 `batch_get_into_offload_object_internal` 中通过 `static const bool` + lambda 读取一次并缓存,之后分支零开销。 +- Pull / Push **互斥**:一次运行只走其中一条路径。 +- 没有单边 WRITE 能力的 transport 应保持 Pull(回退路径已保留)。 + +--- + +## 9. 关键代码索引 + +| 角色 | 位置 | +|---|---| +| RPC 数据结构 | `mooncake-store/include/rpc_types.h` | +| 请求方入口/分支 | `RealClient::batch_get_into_offload_object_internal`(`src/real_client.cpp`) | +| 请求方 RPC 封装 | `ClientRequester::batch_get_offload_object_push` / `invoke_rpc`(`src/real_client.cpp`) | +| 对端 handler | `RealClient::batch_get_offload_object_push`(`src/real_client.cpp`) | +| Client 封装 | `Client::BatchPushOffloadObject`(`src/client_service.cpp`) | +| Push 传输(WRITE) | `TransferSubmitter::submit_batch_push_offload_object`(`src/transfer_task.cpp`) | +| Slice 构建 | `UbTransport::submitTransferTask`(`mooncake-transfer-engine/.../ub_transport.cpp`) | +| URMA 提交 | `UrmaEndpoint::submitPostSend` → `urma_post_jetty_send_wr`(`.../urma/urma_endpoint.cpp`) | +| handler 注册 | `RealClient::setup_internal` 内嵌 server / `RegisterClientRpcService`(`src/real_client_main.cpp`) | + +--- + +## 10. 限制与待办 + +1. **响应粒度**:`BatchGetOffloadObjectPushResponse` 目前只回整体 `error_code`,未支持逐 key 部分失败。如需,可扩成 `std::vector status`。 +2. **gc_ttl 语义**:Pull 路径中「`elapsed >= gc_ttl → OBJECT_HAS_LEASE`」的判定在 Push 下不再需要(对端写完即释放),已省略。 +3. **传输回退**:Push 仅在确认 transport 支持单边 WRITE 时启用;TCP/共享内存等应继续走 Pull。当前由 `MC_OFFLOAD_PUSH` 手动控制,尚未做 transport 能力自动探测。 +4. **构建/测试**:改动尚未在 Linux 目标上编译验证与端到端压测;Windows 开发机无法构建 Mooncake(依赖 RDMA/etcd 等)。 +``` diff --git a/mooncake-integration/store/mooncake_perf_points.def b/mooncake-integration/store/mooncake_perf_points.def index 85983deb05..251a02bfcb 100644 --- a/mooncake-integration/store/mooncake_perf_points.def +++ b/mooncake-integration/store/mooncake_perf_points.def @@ -112,6 +112,9 @@ PERF_KEY_DEF(GET_SSD_RELEASE_BUFFER, "real_client.cpp::batch_get_into_of PERF_KEY_DEF(GET_SSD_OWNER_READ, "file_storage.cpp::BatchGet", "OwnerSsdRead") PERF_KEY_DEF(GET_SSD_OWNER_ALLOC, "file_storage.cpp::BatchGet", "OwnerAllocBuffer") PERF_KEY_DEF(GET_SSD_OWNER_LOAD, "file_storage.cpp::BatchGet", "OwnerDiskLoad") +PERF_KEY_DEF(GET_SSD_OWNER_LOAD_PLAN, "storage_backend.cpp::BatchLoad", "OwnerLoadPlan") +PERF_KEY_DEF(GET_SSD_OWNER_LOAD_URING, "storage_backend.cpp::BatchLoad", "OwnerLoadUring") +PERF_KEY_DEF(GET_SSD_OWNER_LOAD_POSIX, "storage_backend.cpp::BatchLoad", "OwnerLoadPosix") PERF_KEY_DEF(GET_SSD_OWNER_PUSH_WRITE, "real_client.cpp::batch_get_offload_object_push", "OwnerPushWrite") PERF_KEY_DEF(GET_SSD_OWNER_RELEASE, "file_storage.cpp::ReleaseBuffer", "OwnerReleaseBuffer") diff --git a/mooncake-store/src/storage_backend.cpp b/mooncake-store/src/storage_backend.cpp index 7a6792351e..93c5150cf8 100644 --- a/mooncake-store/src/storage_backend.cpp +++ b/mooncake-store/src/storage_backend.cpp @@ -23,6 +23,11 @@ #include #include "storage/distributed/distributed_storage_backend.h" +// ubdiag perf points for the offload owner-side disk load breakdown. +#define UBDIAG_PERF_DEF_FILE "mooncake_perf_points.def" +#define UBDIAG_PROGRAM_NAME "mooncake_store" +#include "ubdiag/auto_perf.h" + namespace mooncake { bool FilePerKeyConfig::Validate() const { @@ -1393,6 +1398,12 @@ tl::expected BucketStorageBackend::BatchLoad( std::unordered_map> bucket_read_plans; std::vector bucket_guards; // RAII guards for all buckets + // Phase 1: build read plan under lock (OwnerLoadPlan). Pure in-memory + // metadata lookups; the error early-returns let the dtor Abandon the + // unfinished sample. + UbDiag::PerfPoint pt_plan(PerfKey::GET_SSD_OWNER_LOAD_PLAN, + UbDiag::PerfLevel::MODULE); + pt_plan.Start(); { SharedMutexLocker lock(&mutex_, shared_lock); for (const auto& [key, dest_slice] : batch_object) { @@ -1443,6 +1454,7 @@ tl::expected BucketStorageBackend::BatchLoad( metadata.key_size, metadata.data_size, dest_slice}); } } + pt_plan.End(0); // Lock released here - bucket files protected by BucketReadGuards // which remain alive until this function returns (~line bucket_guards // destructor), keeping inflight_reads_ > 0 throughout the I/O phase. @@ -1488,8 +1500,12 @@ tl::expected BucketStorageBackend::BatchLoad( // Zero-copy path: read directly into the slice buffer. // dest_slice.ptr is 4096-aligned and oversized (from // AllocateBatch) to accommodate the full aligned read range. + UbDiag::PerfPoint pt_uring(PerfKey::GET_SSD_OWNER_LOAD_URING, + UbDiag::PerfLevel::MODULE); + pt_uring.Start(); read_res = uring_file->read_aligned( plan.dest_slice.ptr, aligned_size, aligned_offset); + pt_uring.End(read_res ? 0 : -1); if (read_res) { // Adjust ptr to point to actual data start (no memcpy) @@ -1501,9 +1517,13 @@ tl::expected BucketStorageBackend::BatchLoad( } else #endif { - // Fallback to vector_read for non-UringFile + // Fallback to vector_read for non-UringFile (preadv). iovec iov{plan.dest_slice.ptr, plan.dest_slice.size}; + UbDiag::PerfPoint pt_posix(PerfKey::GET_SSD_OWNER_LOAD_POSIX, + UbDiag::PerfLevel::MODULE); + pt_posix.Start(); read_res = file->vector_read(&iov, 1, actual_offset); + pt_posix.End(read_res ? 0 : -1); } if (!read_res) {