diff --git a/components/src/dynamo/vllm/handlers.py b/components/src/dynamo/vllm/handlers.py index e548d9cadf34..6ef293976881 100644 --- a/components/src/dynamo/vllm/handlers.py +++ b/components/src/dynamo/vllm/handlers.py @@ -3010,28 +3010,22 @@ async def _encode_one(idx: int, prompt: Any): return # Default wire format: always base64 (the Rust frontend decodes back to - # float when the client's ``encoding_format`` is float or unset). - # 15x1024-float JSON arrays cost ~110 ms in Python json.dumps + Rust - # serde parse; base64 bytes are ~3x smaller and ~10x faster to - # (de)serialize. Client-visible wire format is preserved because Rust - # converts at the HTTP boundary. - embedding_objects: list[Dict[str, Any]] = [] - for idx, final_output in enumerate(outputs): - embedding = _pooling_output_to_list(final_output.outputs.data) - if dimensions is not None: - if dimensions > len(embedding): - raise ValueError( - f"dimensions={dimensions} exceeds model embedding " - f"dimension {len(embedding)}" - ) - embedding = embedding[:dimensions] - embedding_objects.append( - { - "object": "embedding", - "embedding": _encode_floats_to_base64(embedding), - "index": idx, - } - ) + # float when the client's ``encoding_format`` is float or unset). base64 + # bytes are ~3x smaller and far faster to (de)serialize than a JSON float + # array; the client-visible wire format is preserved because Rust converts + # at the HTTP boundary. The base64 is built straight from the pooling + # tensor (torch -> numpy.tobytes), skipping the per-embedding Python float + # list + struct.pack varargs. Bytes are identical on little-endian hosts. + embedding_objects: list[Dict[str, Any]] = [ + { + "object": "embedding", + "embedding": _pooling_output_to_base64( + final_output.outputs.data, dimensions + ), + "index": idx, + } + for idx, final_output in enumerate(outputs) + ] yield { "object": "list", @@ -3116,16 +3110,27 @@ def _classify_embedding_input(input_field: Any) -> list[Any]: ) +def _flatten_pooling_tensor(data: "torch.Tensor") -> "torch.Tensor": + """Flatten a vLLM ``PoolingOutput.data`` tensor to a 1-D float32 CPU tensor. + + Shared by :func:`_pooling_output_to_list` and + :func:`_pooling_output_to_base64` so the detach/cpu/flatten/cast step isn't + duplicated. ``float32`` matches the OpenAI base64 f32 wire format. + + vLLM's pooling pipeline can return a tensor with a singleton batch dim + (shape ``(1, hidden_dim)``) instead of a 1D vector; we flatten unconditionally. + """ + return data.detach().cpu().flatten().to(torch.float32) + + def _pooling_output_to_list(data: Any) -> list[float]: """Convert a vLLM PoolingOutput.data tensor (or list) to a flat list[float]. - vLLM's pooling pipeline can return a tensor with a singleton batch dim - (shape ``(1, hidden_dim)``) instead of a 1D vector (shape ``(hidden_dim,)``). The OpenAI ``/v1/embeddings`` response expects ``data[].embedding`` to be a - flat array of floats, so we flatten unconditionally. + flat array of floats. """ if isinstance(data, torch.Tensor): - return data.detach().cpu().flatten().tolist() + return _flatten_pooling_tensor(data).tolist() if isinstance(data, (list, tuple)): # Already a list — flatten one level if it's a list-of-lists. if data and isinstance(data[0], (list, tuple)): @@ -3150,6 +3155,36 @@ def _encode_floats_to_base64(floats: list[float]) -> str: return base64.b64encode(packed).decode("ascii") +def _pooling_output_to_base64(data: Any, dimensions: int | None = None) -> str: + """Serialize a vLLM ``PoolingOutput.data`` tensor straight to a base64 + float32 string, skipping the intermediate Python ``list[float]`` and the + ``struct.pack("<{N}f", *floats)`` varargs expansion. + + ``torch -> numpy.tobytes -> base64`` keeps the heavy work in C; output bytes + are identical to the ``struct``-based path on little-endian hosts. + """ + if isinstance(data, torch.Tensor): + vec = _flatten_pooling_tensor(data) + if dimensions is not None: + if dimensions > vec.numel(): + raise ValueError( + f"dimensions={dimensions} exceeds model embedding " + f"dimension {vec.numel()}" + ) + vec = vec[:dimensions] + return base64.b64encode(vec.contiguous().numpy().tobytes()).decode("ascii") + # Fallback for non-tensor pooling outputs (rare): reuse the list path. + floats = _pooling_output_to_list(data) + if dimensions is not None: + if dimensions > len(floats): + raise ValueError( + f"dimensions={dimensions} exceeds model embedding " + f"dimension {len(floats)}" + ) + floats = floats[:dimensions] + return _encode_floats_to_base64(floats) + + def _write_embeddings_to_shm(outputs: list[Any], dimensions: int | None) -> dict[str, Any]: """Stage all embedding vectors as one contiguous ``[count, dim]`` little-endian f32 buffer in a POSIX shared-memory segment; return the