CPU EC 커넥터 사용 가이드

CPU EC 커넥터 사용 가이드 (CPU EC Connector Usage Guide)

ECCPUConnector는 GPU 기반 인코더 캐시에 CPU 계층을 추가합니다. 인코더 출력(encoder_cache[mm_hash])을 공유 /dev/shm mmap 영역으로 오프로드해, 이후 단계와 이후 요청이 재계산 대신 그것을 재사용하게 합니다. GPU↔CPU 복사는 풀링된 CUDA 스트림에서 swap_blocks_batch로 모델 연산과 비동기적으로 실행됩니다.

출처: 문서

ec_connector_extra_configec_enable_nixl: true를 설정하면 피어 간(P2P) 전송도 활성화됩니다. consumer 인스턴스가 producer 인스턴스의 CPU 계층에서 인코딩을 NIXL을 통해 직접 당겨와 로컬 재계산을 하지 않습니다. E/PD 분리나 인코더/디코더 인스턴스 공유에 사용됩니다.

본문

사전 요구 사항 (Prerequisites)

  • ECCPUConnector는 V2 모델 러너가 필요합니다: VLLM_USE_V2_MODEL_RUNNER=1. 그렇지 않으면 생성 시 ValueError가 발생합니다.
  • 로컬 CPU 계층 오프로드(ec_enable_nixl 미설정 또는 false)는 추가 패키지가 필요 없습니다. 게이트가 꺼진 코드 경로(cpu/connector.py, cpu/scheduler/, cpu/worker/, cpu/common.py)는 nixl/zmq/msgspec을 import하지 않으며, 저장소 테스트(tests/v1/ec_connector/unit/test_no_nixl_imports.py)가 이를 강제합니다.
  • P2P NIXL 모드(ec_enable_nixl: true)는 nixl 패키지가 필요합니다: uv pip install nixl(requirements/kv_connectors.txt에서 nixl==1.3.2로 고정되며 NixlConnector와 공유). 플랫폼별 설치는 NIXL 저장소를 참고하세요. nixl을 import할 수 없으면 커넥터가 RuntimeError: ec_enable_nixl requires NIXL; install the nixl package or remove ec_enable_nixl from ec_connector_extra_config.을 발생시킵니다.

기본 사용법 (Basic Usage)

단일 엔진 인스턴스 안에서 로컬 CPU 계층 오프로드만 하는 경우:

vllm serve <model> --ec-transfer-config '{
  "ec_connector": "ECCPUConnector",
  "ec_role": "ec_both",
  "ec_connector_extra_config": {"ec_cpu_bytes": 1073741824}
}'
  • ec_role="ec_both": 같은 프로세스가 CPU 계층으로 오프로드하고 다시 로드합니다.
  • 계층은 인스턴스의 모든 TP/PCP 워커가 공유하는 하나의 mmap 영역(/dev/shm/vllm_ec_{instance_id}_dp{dp_rank}.mmap)입니다. 모든 랭크가 동일한 인코더 출력을 가지므로 저장 시에는 TP rank 0 / PCP rank 0만 씁니다.
  • 엔트리는 mm_hash로 키가 매겨집니다. EmbeddingCache는 공간이 필요할 때 ready+unpinned 엔트리를 FIFO로 퇴출합니다.
  • 각 배치 저장/로드는 풀링된 CUDA 스트림에서 실행되며, 전송의 종료 이벤트가 발생하면 완료가 스케줄러에 보고됩니다(ECCPUWorker.build_connector_worker_metaECCPUScheduler.update_connector_output). 이는 저장된 엔트리를 ready로 표시하고 로드된 엔트리를 unpin합니다.
  • 영역은 shutdown()에서 /dev/shm에서 연결 해제(unlink)됩니다.

P2P NIXL과 함께 사용하기 (Usage With P2P NIXL)

Producer — 자신의 CPU 계층으로 오프로드하고 consumer의 읽기를 서빙:

vllm serve <model> --ec-transfer-config '{
  "ec_connector": "ECCPUConnector",
  "ec_role": "ec_producer",
  "ec_connector_extra_config": {"ec_enable_nixl": true, "ec_cpu_bytes": 1073741824}
}'

ec_role="ec_producer"만으로도 멀티모달 구성에서 mm_encoder_only가 활성화되어 vllm_config.is_mm_encoder_only가 True가 됩니다(언어 모델·샘플러·풀러를 건너뜀). ec_transfer_config와 무관하게 인코더 전용 실행이 필요할 때만 --mm-encoder-only를 추가하세요.

Consumer — 로컬 인코딩으로 폴백하기 전에 요청의 ec_transfer_params에 이름이 있는 인코딩을 당겨옵니다:

vllm serve <model> --ec-transfer-config '{
  "ec_connector": "ECCPUConnector",
  "ec_role": "ec_consumer",
  "ec_connector_extra_config": {"ec_enable_nixl": true, "ec_cpu_bytes": 1073741824}
}'

오케스트레이션 흐름 (Orchestration flow)

  1. 요청이 producer에서 끝납니다. ECCPUConnector.request_finished()는 CPU 계층에 여전히 상주하는 각 mm_hash에 대해 다음을 반환합니다.
{mm_hash: {"metadata": {...}, "peer_host": str, "peer_port": int, "size_bytes": int}}

이는 호출자에게 ec_transfer_params(RequestOutput.ec_transfer_params / EngineCoreOutput.ec_transfer_params)로 드러납니다. metadata는 모델이 해당 모달리티에 대해 선언한 플레이스홀더 필드를 담으며, 미디어를 메타데이터 전용 참조로 다시 쓰는 오케스트레이터를 위한 것입니다. 나머지 키는 게시된 인코딩에 대한 커넥터 자신의 핸들입니다.

두 절반(metadata와 주소)은 함께 게시되거나 전혀 게시되지 않습니다. producer가 서빙할 수 없는 mm_hash(영역이 가득 차 저장된 적이 없거나, 그 이후 퇴출됨)는 빈 metadata와 주소 없이 보고되므로, 오케스트레이터는 미디어를 요청에 남겨두고 consumer가 로컬로 인코딩하게 합니다.

  1. 오케스트레이터는 같은 mm_hash로 follow-up 요청을 consumer 인스턴스에 발행하고, producer의 ec_transfer_paramsSamplingParams.extra_args["ec_transfer_params"]를 통해 전달합니다.

  2. consumer에서 ECCPUScheduler.ensure_cache_available()request.ec_transfer_params를 읽습니다. 로컬에 아직 캐시되지 않은 각 mm_hash에 대해 (peer_host, peer_port)로 ZMQ 세션을 열고 XferReq를 보낸 뒤, OK XferAck를 받으면 producer의 mmap에서 자신의 것으로 블록을 직접 당겨오는 NIXL READ를 발행합니다. 요청은 READ가 완료될 때까지 지연됩니다.

  3. NACK_NOT_READY는 producer가 인코딩을 공지했지만 GPU→mmap 저장이 아직 도착하지 않았음을 뜻합니다. consumer는 실패를 기록하지 않고 in-flight 엔트리를 해제하고 다음 스텝에 읽기를 다시 요청합니다. 몇 스텝 늦게 저장이 도착하면 재계산이 아니라 지연 비용만 듭니다.

  4. 다른 NACK(NACK_MISSING, NACK_INCOMPAT, NACK_VERSION, NACK_INTERNAL), ack 타임아웃, 읽기 타임아웃, 또는 피어 연결 끊김 시 consumer는 in-flight 엔트리를 버리고 그 mm_hash에 대해 로컬 인코딩으로 폴백합니다. P2P 실패가 요청을 무기한 차단하지는 않습니다.

프로토콜 (Protocol)

  • 컨트롤 플레인: ZMQ. producer는 VLLM_EC_SIDE_CHANNEL_HOST:VLLM_EC_SIDE_CHANNEL_PORTROUTER 소켓을 바인딩하고, 각 consumer는 producer 피어당 하나의 DEALER 연결을 열며 ZMQ 하트비트(2초 간격, 4초 타임아웃, 8초 TTL)로 죽은 피어를 감지합니다. XferReq/XferAckEC_CONNECTOR_VERSION(현재 1)으로 버전이 매겨진 msgspec msgpack 구조체이며, 버전 불일치는 NACK됩니다.
  • 호환성 검사: 모든 XferReq(vllm_version, model, dtype, block_size_bytes)에 대한 SHA-256 해시를 담습니다. producer는 해시가 다른 피어를 NACK(NACK_INCOMPAT)합니다.
  • Ack 상태: OK, NACK_MISSING(producer가 더 이상 인코딩을 보유하지 않음), NACK_NOT_READY(보유했지만 저장이 도착하지 않음), NACK_INCOMPAT, NACK_VERSION, NACK_INTERNAL. NACK_NOT_READY만 재시도 가능합니다. cpu/protocol.pyRETRYABLE_NACKS가 양쪽 끝이 참고하는 것이므로 분류는 각 호출 지점이 아니라 와이어 어휘와 함께 존재합니다.
  • 데이터 플레인: NIXL, UCX 백엔드, consumer가 시작한 READ — consumer가 producer의 등록된 mmap 영역에서 바이트를 직접 당겨오며, producer는 절대 push하지 않습니다.
  • Producer 재시작 복구: XferAck가 producer의 NIXL 에이전트 메타데이터를 담으므로, consumer는 재시작된 producer에 대해 새 handshake 왕복 없이 READ를 복구할 수 있습니다.
  • 타임아웃: consumer XferAck 대기 2초; NIXL 읽기 20초(그 후 60초간 중단 불가능한 DMA가 정리되도록 격리(quarantine) — 퇴출되지는 않음); producer는 30초 핀 임대 후 청구되지 않은 pinned grant를 해제합니다.

구성 (Configuration)

EC 전송은 --ec-transfer-config(CLI) 또는 VllmConfigec_transfer_config 필드(ECTransferConfig, vllm/config/ec_transfer.py)로 구성합니다.

| Field | Type | Default | Description | | ec_connector | str | None | None | Connector class name. Use "ECCPUConnector". | | ec_role | "ec_producer" | "ec_consumer" | "ec_both" | None | None | Required whenever ec_connector is set. ec_producer offloads GPU→CPU only, ec_consumer reloads CPU→GPU only, ec_both does both in the same process. | | ec_connector_extra_config | dict[str, Any] | {} | Connector-specific settings, including ec_enable_nixl — see ec_connector_extra_config Reference. | | engine_id | str | None | random UUID4 | Names the NIXL agent when ec_enable_nixl=True. | | ec_connector_module_path | str | None | None | Python module path to load an out-of-tree connector from, when ec_connector isn't in the built-in registry (ECExampleConnector, ECCPUConnector). |

ec_connector_extra_config 참조 (Reference)

| Key | Type | Required | Description | | ec_enable_nixl | bool | No (default false) | Enables NIXL P2P transfer in addition to local CPU offload. Omitted or false imports no NIXL/ZMQ. Extra config is not type coerced, so a string value is parsed: "true", "1", "yes" enable it, anything else does not. | | consumer_ack_timeout_s | float | No (default 2.0) | How long a consumer waits for an XferAck before giving up on a read. The producer answers XferReqs from its own scheduler step, so its reply latency scales with the encoder's --max-num-batched-tokens: a loaded encoder whose steps run longer than this makes consumers abandon reads the producer is about to grant. Raise it for large encoder batches. | | ec_cpu_bytes | int | Yes | Total size, in bytes, of the shared CPU mmap region. ECCPUConnector raises ValueError if unset. Block count = ec_cpu_bytes // block_size_bytes, where block_size_bytes = hidden_dim * dtype.element_size() (hidden_dim accounts for Qwen3-VL deepstack: out_hidden_size * (1 + num_deepstack_layers)). |

환경 변수 (Environment Variables)

| Variable | Default | Description | | VLLM_EC_SIDE_CHANNEL_HOST | localhost | Host the producer's ZMQ ROUTER socket binds to. Set to a routable address (e.g. the pod IP) for multi-instance/multi-node P2P — the default only works when producer and consumer share a host. | | VLLM_EC_SIDE_CHANNEL_PORT | 5601 | Port for the same ZMQ ROUTER socket. |

둘 다 producer(ec_role="ec_producer" 또는 "ec_both")에서 ec_enable_nixl=True일 때만 읽힙니다.

제한 사항 (Limitations)

  • 인코딩이 소비되기 전에 CPU 계층에서 퇴출됐을 때 오케스트레이터나 피어 인스턴스에 알릴 메커니즘이 없습니다. consumer는 자신의 XferReq가 NACK(NACK_MISSING)될 때만 미스를 발견하고 로컬 재계산으로 폴백합니다.
  • 현재 프로세스 종료 시 Mmap 정리는 best-effort입니다. 생성 프로세스가 ECSharedRegion.cleanup()이 실행되기 전에 SIGKILL되면 /dev/shm/vllm_ec_*.mmap 파일이 새어나와 수동으로 제거해야 합니다.
  • NixlDataTransportUCX 백엔드를 하드코딩합니다.
  • 재시도된 읽기는 자신의 목적지 블록을 해제하고 다음 스텝에 다시 할당하므로, 저장이 느린 producer는 consumer가 엔진 스텝마다 한 번씩 ready 엔트리를 퇴출해 같은 블록을 되찾게 만듭니다.
  • 로컬 인코딩으로 폴백하려면 미디어가 여전히 요청 위에 있어야 합니다. 서빙할 수 있을 때만 인코딩을 공지하면 오케스트레이터가 서빙할 수 없는 것으로 미디어를 다시 쓰는 것을 막을 수 있지만, 공지 이후에 인코딩이 유실된 경우(공지와 consumer의 읽기 사이에 퇴출) 다시 쓰여진 요청은 임베딩할 것이 없어져 워커의 sanity_check_mm_encoder_outputs에서 실패합니다. ensure_cache_available()은 요청을 지연만 시키고 실패시킬 수 없으므로, 이를 깨끗하게 보고하려면 스케줄러 측 실패 경로가 필요합니다.

더 알아보기 (Learn more)