Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 28 additions & 1 deletion xllm/core/scheduler/disagg_pd_scheduler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ limitations under the License.
#include <brpc/server.h>

#include <random>
#include <string>

#include "common/global_flags.h"
#include "common/macros.h"
Expand Down Expand Up @@ -453,14 +454,40 @@ void DisaggPDScheduler::dispatch_requests() {
}
for (size_t i = 0; i < requests.size(); ++i) {
if (resps.resps()[i].status_code() != 200) {
// push back to prefill_request_queue_
// D returned reject (e.g. 404 when try_allocate failed). Cap retries
// to avoid infinite retry.
const std::string& req_id = requests[i]->request_id();
int retry_count = 0;
{
std::lock_guard<std::mutex> lock(add_new_requests_retry_mutex_);
retry_count = ++add_new_requests_retry_count_[req_id];
}
if (retry_count >= kAddNewRequestsMaxRetryOnReject) {
{
std::lock_guard<std::mutex> lock(add_new_requests_retry_mutex_);
add_new_requests_retry_count_.erase(req_id);
}
response_processor_->process_failed_request(
requests[i],
{StatusCode::RESOURCE_EXHAUSTED,
"Decode instance rejected AddNewRequests (e.g. 404) after " +
std::to_string(kAddNewRequestsMaxRetryOnReject) +
" retries"});
continue;
}
Comment thread
magicheng0816 marked this conversation as resolved.
// push back to prefill_request_queue_ for retry
if (requests[i]->offline()) {
prefill_request_queue_offline_.enqueue(requests[i]);
} else {
prefill_request_queue_.enqueue(requests[i]);
}

} else {
// success: clear retry count so map does not grow
{
std::lock_guard<std::mutex> lock(add_new_requests_retry_mutex_);
add_new_requests_retry_count_.erase(requests[i]->request_id());
}
for (auto& sequence : requests[i]->sequences()) {
TransferKVInfo info;
info.request_id = requests[i]->request_id();
Expand Down
6 changes: 6 additions & 0 deletions xllm/core/scheduler/disagg_pd_scheduler.h
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,12 @@ class DisaggPDScheduler : public ContinuousScheduler {
moodycamel::BlockingConcurrentQueue<std::shared_ptr<Request>>
prefill_request_queue_offline_;

// Max retries when D returns 404 (e.g. try_allocate failed), avoid infinite
// retry.
static constexpr int kAddNewRequestsMaxRetryOnReject = 3;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

int ==> int32_t or int64_t

std::unordered_map<std::string, int> add_new_requests_retry_count_;
std::mutex add_new_requests_retry_mutex_;

// use threadpool to handle prefill-completed request
ThreadPool prefill_threadpool_;

Expand Down
Loading