NIXL push-mode KV 전송

NIXL push-mode KV 전송 (NIXL push-mode KV transfer)

기본 NIXL 커넥터는 pull 기반입니다. 즉 prefill이 끝난 뒤 decode(D) 인스턴스가 NIXL READ를 통해 prefill(P) 인스턴스로부터 KV 블록을 읽어 갑니다. NixlPushConnector는 push 기반 대안을 추가하는데, 이때는 P가 NIXL WRITE로 KV 블록을 D의 미리 할당된 메모리에 직접 써 넣습니다. 이 문서는 push 설계 특유의 스레딩·큐·스케줄링 상호작용을 설명합니다.

출처: 문서

본문

기본 NIXL 커넥터는 pull 기반입니다. 즉 prefill이 끝난 뒤 decode(D) 인스턴스가 NIXL READ로 prefill(P) 인스턴스에서 KV 블록을 읽어 갑니다. NixlPushConnector는 push 기반 대안을 추가하는데, P가 KV 블록을 NIXL WRITE로 D의 미리 할당된 메모리에 직접 써 넣는 방식입니다.

이 문서는 push 설계 특유의 스레딩·큐·스케줄링 상호작용을 설명합니다. pull-mode 설계는 변하지 않으며, push 커넥터는 가능한 한 동일한 핸드셰이크·NIXL 에이전트 설정·메타데이터 경로를 재사용합니다.

상위 레벨 흐름 (High-level flow)

sequenceDiagram
    autonumber
    participant Client
    participant Proxy
    participant DSched as D Scheduler
    participant DWorker as D Worker (main)
    participant DWriter as D Writer
    participant PWriter as P Writer
    participant PWorker as P Worker (main)
    participant PSched as P Scheduler

    Client->>Proxy: POST /v1/completions
    Proxy->>PSched: prefill leg (do_remote_decode=True, max_tokens=1)
    Proxy->>DSched: decode leg (do_remote_prefill=True, P coordinates)

    note over DSched,DWriter: D side - register blocks with P
    DSched->>DSched: update_state_after_alloc, stash registration, arm watchdog
    DSched->>DWorker: build_connector_meta -> meta.push_registrations
    DWorker->>DWriter: enqueue (req_id, reg_data) on _reg_send_inbox
    DWriter->>PWriter: NIXL send_notif PUSH_REG msgpack

    note over PSched,PWriter: P side - prefill, stage finished blocks
    PSched->>PSched: request_finished, stash blocks
    PSched->>PWorker: build_connector_meta -> meta.push_finished_blocks
    PWorker->>PWriter: enqueue (req_id, blocks) on _finished_blocks_inbox

    note over PWriter: P writer matches and WRITEs
    PWriter->>PWriter: get_new_notifs returns PUSH_REG, route via _handle_push_reg_notif
    alt PUSH_REG and finished blocks both present
        PWriter->>PWriter: pop matching pair, fire WRITE
    else only one side present
        PWriter->>PWriter: stash and wait, self-poll only when blocks unmatched
    end
    PWriter->>PWriter: _ensure_handshake to D (async; defer WRITE)
    PWriter->>PWriter: handshake callback re-queues on _deferred_push_inbox, wake
    PWriter->>DWriter: NIXL WRITE direct to D GPU + completion notif

    note over DWorker,DWriter: D side - completion accounting
    DWriter-->>DWorker: forward HB and completion notifs via _pending_completion_notifs
    DWorker->>DWorker: _get_new_notifs drains, HB extends lease, completion marks recv done
    DWorker->>DSched: update_connector_output(finished_recving)
    DSched->>DSched: clear watchdog deadline

    note over PWorker,PWriter: P side - reclaim
    PWorker->>PWorker: get_finished, drain _sending_transfers, queue eviction
    PWriter->>PWriter: drain _evict_finished_inbox, drop stale state
    PWorker->>PSched: update_connector_output(finished_sending)
    PSched->>PSched: free lease

    DWorker-->>Proxy: stream decode tokens
    Proxy-->>Client: response

스레드 (Threads)

NixlPushConnectorWorker는 워커당(즉 TP 랭크당) nixl-push-writer라는 전용 백그라운드 스레드 하나를 도입합니다. 각 스레드는 랭크에서 push 특유의 NIXL 연산을 소유합니다.

  • nixl_wrapper.get_new_notifs() — 알림 수신.
  • PUSH_REG:<msgpack>(D 측)과 WRITE별 완료 알림(P 측)을 위한 nixl_wrapper.send_notif(...).
  • WRITE 자체를 제출하는 nixl_wrapper.make_prepped_xfer(...) / transfer(...).

하트비트는 기존 base-worker의 _send_heartbeats 배관을 통해 엔진 메인 스레드에서 계속 나갑니다(start_load_kv 내부).

깨우기 모델 (Wake model)

라이터 스레드는 할 일이 없으면 _push_writer_wake(threading.Event)에서 블록합니다. 이벤트를 설정하는 호출자는 세 곳입니다.

  • start_load_kv(워커 메인 스레드, 스케줄러의 메타데이터와 함께 엔진 스텝마다 한 번 호출) — 스텝이 라이터에 새 작업을 실제로 넘길 때만, 즉 meta.push_registrations 또는 meta.push_finished_blocks가 비어 있지 않을 때만 wake를 설정합니다. 새 전송을 위한 wake입니다.
  • get_finished(워커 메인 스레드, 완료 보고를 위해 엔진 스텝마다 한 번 호출) — 항상 wake를 설정합니다. 라이터는 push에 대해 nixl_wrapper.get_new_notifs()의 유일한 소비자이므로, 처리할 새 메타데이터가 없어도 인바운드 알림(D의 하트비트, WRITE 후 완료 알림, 늦게 도착하는 PUSH_REG)을 소진할 기회를 얻습니다.
  • 핸드셰이크 완료 콜백(백그라운드 핸드셰이크 실행자 스레드) — 두 핸드셰이크 모두 실행자에서 돌며 라이터를 절대 블록하지 않습니다. 완료 콜백은 지연된 연산을 다시 큐에 넣고 wake를 설정합니다. send_notif와 NIXL WRITE 모두 라이터 스레드 밖에서는 실행될 수 없기 때문입니다.
    • D→P 핸드셰이크(PUSH_REG 보내기 전)는 등록을 _reg_send_inbox에 다시 큐에 넣습니다.
    • P→D 핸드셰이크(WRITE 전)는 매칭된 (req_id, blocks, reg_data)_deferred_push_inbox에 다시 큐에 넣습니다.

두 번째 패스에서 _ensure_handshakeNone을 반환합니다(에이전트가 이제 연결됨). 그래서 라이터는 PUSH_REG를 보내거나 WRITE를 직접 발행합니다. 핸드셰이크가 실패하면 콜백은 다시 큐에 넣는 대신 요청을 실패시키거나 버리므로, 재시도 루프가 없습니다(실패 처리 참조).

이벤트 기반 wake 외에도, P 측 finished 블록이 매칭되지 않은 PUSH_REG를 기다리는 동안 라이터는 _PUSH_WRITER_POLL_INTERVAL_MS = 1.0 ms 간격으로 스스로 폴링합니다.

요청이 P에서 완료되면(임대 만료 또는 WRITE 종료), get_finished은 요청 id를 _evict_finished_inbox에 큐에 넣고, 라이터는 이를 소진해 오래된 _push_finished_blocks / _pending_d_registrations를 버리고 자체 폴링을 멈춥니다.

라이터 로컬 매칭 테이블 (Writer-local matching tables)

Table Owner Holds
_pending_d_registrations writer 원격 D에서 받은 D 등록 — P의 블록을 기다림
_push_finished_blocks writer 스케줄러가 스테이징한 P 블록 — 원격 D 등록을 기다림

어느 쪽이 먼저 도착할 수 있습니다. 라이터는 양방향으로 매칭합니다. PUSH_REG가 도착하면 _push_finished_blocks를 조회하고, finished 블록이 도착하면 _pending_d_registrations를 조회합니다. 두 조회 모두 먼저 정확한 request_id 매칭을 시도한 뒤, 각 엔진별 무작위 접미사를 제거한 id를 비교하는 폴백(get_base_request_id)을 사용합니다. 이 폴백이 존재하는 이유는 프록시가 두 레그에 같은 X-Request-Id를 넘기므로, P와 D가 같은 cmpl-<uuid>-<index> 형태로 감싸되 input_processor.assign_request_id가 엔진별로 추가하는 8자리 hex 무작위 접미사만 다르기 때문입니다. 그 접미사만 잘라내면 완료 인덱스(completion index)는 보존하면서 양쪽을 같은 id로 정규화합니다(따라서 multi-prompt 서브요청은 구별됩니다). 이는 VLLM_DISABLE_REQUEST_ID_RANDOMIZATION 설정 여부와도 무관하게 동작하는데, 이 환경변수가 상류(upstream)에서 제거될 예정이기에 중요합니다.

와이어 포맷 (Wire format)

push 등록은 NIXL 알림으로 전송됩니다.

PUSH_REG:<msgpack-encoded dict>

딕셔너리의 필드:

Field Set by Meaning
request_id D D 자신의 vLLM 요청 id. P의 매칭 키이며 완료 알림에 그대로 반향됩니다.
decode_engine_id D D의 엔진 id(P가 역방향 핸드셰이크에 사용)
decode_host D D의 NIXL 사이드채널 호스트
decode_port D D의 NIXL 사이드채널 포트
decode_tp_size D D의 텐서 병렬 크기
local_block_ids D D의 논리 블록 id들(미리 할당된)의 그룹별 목록
remote_engine_id D P의 엔진 id(기존 P 측 핸드셰이크용)
remote_host D P의 NIXL 사이드채널 호스트
remote_port D P의 NIXL 사이드채널 포트
remote_tp_size D P의 텐서 병렬 크기

D는 논리 블록 id를 보내고, P는 WRITE 제출 시점에 NIXL 핸드셰이크에서 배운 비율(remote_physical_blocks_per_logical)로 이를 물리 블록 id로 확장합니다. 이는 pull-mode 계약과 일치합니다 — 스케줄러는 논리 id를 보내고 워커는 제출 시 물리 id로 확장합니다.

WRITE 후 P에서 D로 보내는 완료 알림은 pull 모드에서 쓰는 기존 <request_id>:<tp_size> 형식입니다(여기서 request_id는 등록에서 가져온 D 자신의 요청 id). 그래서 D 측 회계 코드는 변하지 않습니다.

파이프라인 병렬화와 하이브리드 KV 캐시 (Pipeline parallelism and hybrid KV caches)

pull 커넥터는 로컬·원격 워커가 일치하는 KV 영역 목록을 노출해야 합니다 — 여기의 영역 i는 저기의 영역 i에 대응합니다. 이 가정은 파이프라인 병렬(PP)과 하이브리드(HMA) KV 레이아웃이 결합되면 깨집니다.

  • PP 샤딩된 프리필러(P)는 모델 레이어의 일부만 보유하는 반면, PP=1 디코더(D)는 전부 보유하므로 영역 수가 다릅니다.
  • HMA에서는 여러 레이어 이름이 하나의 영역으로 풀링되는데, 풀링된 영역을 대표하는 레이어가 P와 D에서 다를 수 있습니다.

NixlPushConnector는 PP 샤딩된 프로듀서에 대해 영역 인덱스 대신 레이어 이름(member) 정체성으로 라우팅하여 처리합니다. 각 워커는 핸드셰이크 메타데이터(NixlAgentMetadata.region_members)에서 각 NIXL 영역을 백업하는 레이어 이름을 광고합니다. member 라우팅이 필요한 프로듀서는 KV 캐시를 등록할 때 member-major 레이아웃을 한 번 파생하고, 발행하는 모든 전송이 그 순서를 사용합니다. 그러면 add_remote_agent가 이 스테이지가 소유한 정확한 원격 영역을 선택해 일치하도록 재정렬하므로, 각 원격 랭크가 메타데이터를 우연히 어떻게 정렬하든 양쪽이 쌍을 유지합니다.

주소·블록 길이·스트라이드·영역별 용량도 같은 member 순서를 따릅니다. 디스크립터 오프셋은 각 member의 영역 용량을 사용하므로, P와 D가 같은 수의 블록을 할당할 필요가 없습니다. 물리 할당은 여러 member가 공유해도(여러 member가 공유할 때조차) 한 번만 등록됩니다.

원격 영역이 정렬될 때 강제되는 불변식:

  • 로컬 소유의 모든 member는 원격이 정확히 한 번 광고해야 합니다. member가 빠지면 그 레이어의 KV를 조용히 오래된 채로 두는 대신 핸드셰이크가 실패합니다. 원격 전용 member(다른 PP 스테이지 소유)는 무시됩니다.
  • 로컬 레이아웃이 member 라우팅을 요구하는데 원격이 member 메타데이터를 생략하면, 영역 인덱스 라우팅으로 폴백하는 대신 핸드셰이크가 실패합니다.
  • member 순서는 로컬 레이아웃 단독의 속성이므로, 같은 로컬 소스 디스크립터가 모든 원격 엔진·TP 랭크에 서비스됩니다. 원격 디스크립터 목록만 랭크별로 재구성됩니다.

완료가 소비자 랭크별로 집계되기 때문에 Decode 측 PP는 지원되지 않습니다. Mamba/SSM 하이브리드는 PP 아래에서 지원되지 않습니다.

스케줄러 측 책임 (Scheduler-side responsibilities)

NixlPushConnectorScheduler는 기본 스케줄러를 확장합니다.

  • D 측 — update_state_after_alloc은 등록 데이터를 _push_pending_registrations에 쌓아 두고 소프트 워치독(_push_registration_deadlines)을 작동시킵니다. build_connector_meta는 쌓인 것을 meta.push_registrations로 소진하며, 만료된 항목은 경고와 함께 버립니다.
  • P 측 — request_finished는 블록 id를 _finished_request_blocks(임대 및 has_pending_push_work용)와 _newly_finished_push_blocks(다음 워커 스텝용, meta.push_finished_blocks를 통해)에 쌓아 둡니다.
  • 양쪽 — has_pending_push_work는 진행 중인 push 상태가 있는 동안 엔진 메인 루프가 계속 스테핑하도록 해, 라이터가 스텝마다 최소 한 번은 wake를 받게 합니다.

update_connector_output:

  • finished_sending(P 측)은 임대 항목을 지웁니다.
  • finished_recving(D 측)은 워치독 마감 시한을 지웁니다.

타임아웃과 워치독 (Timeouts and watchdogs)

스케줄러에 두 개의 요청별 타이머가 작동합니다.

  • D 측 등록 워치독_push_registration_deadlines. 등록된 요청이 push_registration_timeout 초(기본 decoder_kv_blocks_ttl) 내에 push 완료를 보지 못하면, build_connector_meta는 오래된 등록과 pending 항목을 버리고 경고를 로깅한 뒤 등록 재전송을 중단합니다. 해당 요청은 _reqs_need_recv에 남아 있고, 결국 요청을 실패시키는 것은 엔진의 요청 레벨 중단 경로(또는 사용자/프록시가 HTTP 호출을 타임아웃시키는 것)입니다.
  • P 측 블록 임대 — pull 모드와 같은 _kv_lease_duration. request_finished는 만료를 _reqs_need_send에 설정하고, update_connector_output(finished_sending=...)은 WRITE 성공 시 이를 지웁니다. 오래된 임대는 base worker의 get_finished가 정리(reap)하고, 이어서 eviction을 _evict_finished_inbox에 큐에 넣어 라이터도 자체 폴링을 멈춥니다.

실패 처리 (Failure handling)

  • D 측 핸드셰이크 실패(PUSH_REG를 보내기 전 P→D 핸드셰이크) — future의 done-콜백이 _handle_failed_transfer(rid, None)을 호출해 D의 미리 할당된 블록을 무효로 표시하고 _failed_recv_reqs에 큐에 넣어, 다음 get_finished이 요청을 실패한 recv로 보고하게 합니다. pull 모드와 같은 recv 측 회계입니다.
  • PUSH_REG를 P로 보낼 때 D 측 send_notif 실패 — 동일하게 처리합니다. _handle_failed_transfer가 recv를 실패로 표시합니다.
  • P 측 핸드셰이크 실패(WRITE 전 P→D 핸드셰이크) — done-콜백은 push_handshake_failed를 로깅하고 재큐잉 없이 요청을 버립니다. 의도적으로 _handle_failed_transfer를 호출하지 않습니다(아래 WRITE 제출 실패와 같은 이유로 프로듀서 측에 무효화할 _recving_metadata 항목이 없습니다). P의 블록은 _kv_lease_duration 임대로 회수되고 D의 오래된 등록은 워치독이 처리합니다.
  • P 측 WRITE 제출 실패 — WRITE 핸들(있다면)을 해제하고 xfer_stats.record_failed_transfer()로 실패 카운터를 올립니다. 의도적으로 _handle_failed_transfer를 호출하지 않습니다. P 측의 req_id_recving_metadata에 항목이 없고(P는 수신자가 아님), 헬퍼가 P 로컬 요청 id를 _failed_recv_reqs에 넣어 base worker의 get_finished에서 assertion을 터뜨리기 때문입니다. 나가는 WRITE는 바닥에 버려지고, D의 임대 워치독이 누락된 완료를 처리합니다.

요약 (Summary)

push 설계는 기존 NIXL 커넥터 위의 작고 잘 격리된 확장입니다.

  • 새 커넥터 클래스 하나, 새 스케줄러 클래스 하나, 새 워커 클래스 하나 — 모두 기존 base 클래스의 하위 클래스.
  • 워커당 전용 백그라운드 스레드 하나.
  • 소비자가 하나(라이터)인 크로스 스레드 큐 몇 개. 대부분은 프로듀서가 하나지만, 엔진 메인 스레드와 핸드셰이크 완료 콜백이 모두 공급하는 재생(replay) 큐 두 개가 있습니다: _reg_send_inbox(D→P 핸드셰이크 후 재생되는 등록)와 _deferred_push_inbox(P→D 핸드셰이크 후 재생되는 매칭 push).
  • 새 알림 타입 하나(PUSH_REG:<msgpack>).

엔진 메인 스레드의 동작은 그 외에 변하지 않습니다. 라이터 스레드는 이벤트 기반이며 push 작업이 없으면 유휴 상태입니다.

더 알아보기 (Learn more)