Compare commits

...
7 Commits
5 changed files with 121 additions and 92 deletions
+11 -11
View File
@@ -1,5 +1,5 @@
ARG WORKER_CUDA_VERSION=11.8.0
FROM runpod/worker-vllm:base-0.2.0-cuda${WORKER_CUDA_VERSION} AS vllm-base
FROM runpod/worker-vllm:base-0.2.2-cuda${WORKER_CUDA_VERSION} AS vllm-base
RUN apt-get update -y \
&& apt-get install -y python3-pip
@@ -15,26 +15,26 @@ COPY src /src
# Setup for Option 2: Building the Image with the Model included
ARG MODEL_NAME=""
ARG MODEL_BASE_PATH="/runpod-volume"
ARG BASE_PATH="/runpod-volume"
ARG QUANTIZATION=""
ENV MODEL_BASE_PATH=$MODEL_BASE_PATH \
MODEL_NAME=$MODEL_NAME \
ENV MODEL_NAME=$MODEL_NAME \
BASE_PATH=$BASE_PATH \
QUANTIZATION=$QUANTIZATION \
HF_DATASETS_CACHE="${MODEL_BASE_PATH}/huggingface-cache/datasets" \
HUGGINGFACE_HUB_CACHE="${MODEL_BASE_PATH}/huggingface-cache/hub" \
HF_HOME="${MODEL_BASE_PATH}/huggingface-cache/hub" \
HF_DATASETS_CACHE="${BASE_PATH}/huggingface-cache/datasets" \
HUGGINGFACE_HUB_CACHE="${BASE_PATH}/huggingface-cache/hub" \
HF_HOME="${BASE_PATH}/huggingface-cache/hub" \
HF_TRANSFER=1
ENV PYTHONPATH="/:/vllm-installation"
RUN --mount=type=secret,id=HF_TOKEN,required=false \
if [ -f /run/secrets/HF_TOKEN ]; then \
export HF_TOKEN=$(cat /run/secrets/HF_TOKEN); \
fi && \
if [ -n "$MODEL_NAME" ]; then \
python3 /src/download_model.py --model $MODEL_NAME; \
python3 /src/download_model.py; \
fi
ENV PYTHONPATH="/:/vllm-installation"
# Start the handler
CMD ["python3", "/src/handler.py"]
CMD ["python3", "/src/handler.py"]
+10 -8
View File
@@ -42,7 +42,7 @@ We now offer a pre-built Docker Image for the vLLM Worker that you can configure
<div align="center">
Stable Image: ```runpod/worker-vllm:0.2.0```
Stable Image: ```runpod/worker-vllm:0.2.1```
Development Image: ```runpod/worker-vllm:dev```
@@ -57,9 +57,11 @@ Development Image: ```runpod/worker-vllm:dev```
- `MODEL_NAME`: Hugging Face Model Repository (e.g., `openchat/openchat-3.5-1210`).
**Optional**:
- Model Settings:
- LLM Settings:
- `TOKENIZER_NAME`: Tokenizer repository if you would like to use a different tokenizer than the one that comes with the model. (default: `None`)
- `CUSTOM_CHAT_TEMPLATE`: Custom chat jinja template, read more about Hugging Face chat templates [here](https://huggingface.co/docs/transformers/chat_templating). (default: `None`)
- `MAX_MODEL_LENGTH`: Maximum number of tokens for the engine to be able to handle. (default: maximum supported by the model)
- `MODEL_BASE_PATH`: Model storage directory (default: `/runpod-volume`).
- `BASE_PATH`: Storage directory where huggingface cache and model will be located. (default: `/runpod-volume`, which will utilize network storage if you attach it or create a local directory within the image if you don't)
- `LOAD_FORMAT`: Format to load model in (default: `auto`).
- `HF_TOKEN`: Hugging Face token for private and gated models (e.g., Llama, Falcon).
- `QUANTIZATION`: AWQ (`awq`), SqueezeLLM (`squeezellm`) or GPTQ (`gptq`) Quantization. The specified Model Repo must be of a quantized model. (default: `None`)
@@ -68,12 +70,12 @@ Development Image: ```runpod/worker-vllm:dev```
- Tensor Parallelism:
Note that the more GPUs you split a model's weights accross, the slower it will be due to inter-GPU communication overhead. If you can fit the model on a single GPU, it is recommended to do so.
- `USE_TENSOR_PARALLEL`: Enable (`1`) or disable (`0`) Tensor Parallelism. (default: `0`)
- `TENSOR_PARALLEL_SIZE`: Number of GPUs to shard the model across (default: `1`).
- If you are having issues loading your model with Tensor Parallelism, try decreasing `VLLM_CPU_FRACTION` (default: `1`).
- System Settings:
- `GPU_MEMORY_UTILIZATION`: GPU VRAM utilization (default: `0.98`).
- `MAX_PARALLEL_LOADING_WORKERS`: Maximum number of parallel workers for loading models (default: `number of available CPU cores`).
- `MAX_PARALLEL_LOADING_WORKERS`: Maximum number of parallel workers for loading models (default: `number of available CPU cores` if `TENSOR_PARALLEL_SIZE` is `1`, otherwise `None`).
- Serverless Settings:
@@ -94,7 +96,7 @@ To build an image with the model baked in, you must specify the following docker
- **Required**
- `MODEL_NAME`
- **Optional**
- `MODEL_BASE_PATH`: Defaults to `/runpod-volume` for network storage. Use `/models` or for local container storage.
- `BASE_PATH`: Storage directory where huggingface cache and model will be located. (default: `/runpod-volume`, which will utilize network storage if you attach it or create a local directory within the image if you don't. If your intention is to bake the model into the image, you should set this to something like `/models` to make sure there are no issues if you were to accidentally attach network storage.)
- `QUANTIZATION`
- `WORKER_CUDA_VERSION`: `11.8.0` or `12.1.0` (default: `11.8.0` due to a small amount of workers not having CUDA 12.1 support yet. `12.1.0` is recommended for optimal performance).
@@ -102,7 +104,7 @@ For the remaining settings, you may apply them as environment variables when run
#### Example: Building an image with OpenChat-3.5
```bash
sudo docker build -t username/image:tag --build-arg MODEL_NAME="openchat/openchat_3.5" --build-arg MODEL_BASE_PATH="/models" .
sudo docker build -t username/image:tag --build-arg MODEL_NAME="openchat/openchat_3.5" --build-arg BASE_PATH="/models" .
```
##### (Optional) Including Huggingface Token
@@ -151,7 +153,7 @@ You may either use a `prompt` or a list of `messages` as input. If you use `mess
|-----------------------|----------------------|--------------------|--------------------------------------------------------------------------------------------------------|
| `prompt` | str | | Prompt string to generate text based on. |
| `messages` | list[dict[str, str]] | | List of messages, which will automatically have the model's chat template applied. Overrides `prompt`. |
| `use_openai_format` | bool | False | Whether to return output in OpenAI format. `ALLOW_OPENAI_FORMAT` environment variable must be `1`, the input must be a `messages` list, and `stream` enabled. |
| `use_openai_format` | bool | False | Whether to return output in OpenAI format. `ALLOW_OPENAI_FORMAT` environment variable must be `1`, the input should preferably be a `messages` list, but `prompt` is accepted. |
| `apply_chat_template` | bool | False | Whether to apply the model's chat template to the `prompt`. |
| `sampling_params` | dict | {} | Sampling parameters to control the generation, like temperature, top_p, etc. |
| `stream` | bool | False | Whether to enable streaming of output. If True, responses are streamed as they are generated. |
+20 -17
View File
@@ -1,22 +1,25 @@
import argparse
import os
import logging
from vllm.model_executor.weight_utils import prepare_hf_model_weights
if __name__ == "__main__":
parser = argparse.ArgumentParser()
parser.add_argument("--model", type=str)
parser.add_argument(
"--download_dir", type=str, default=os.environ.get("MODEL_BASE_PATH")
)
args = parser.parse_args()
if not args.model or not args.download_dir:
raise ValueError("Must specify model and download_dir")
if not os.path.exists(args.download_dir):
os.makedirs(args.download_dir)
prepare_hf_model_weights(
model_name_or_path=args.model,
cache_dir=args.download_dir,
model = os.getenv("MODEL_NAME")
download_dir = os.getenv("HF_HOME")
if not model or not download_dir:
raise ValueError(f"Must specify model and download_dir. Model: {model}, download_dir: {download_dir}")
if not os.path.exists(download_dir):
os.makedirs(download_dir)
logging.info(f"Downloading model {model} to {download_dir}")
hf_folder, hf_weights_files, use_safetensors = prepare_hf_model_weights(
model_name_or_path=model,
cache_dir=download_dir,
)
logging.info(f"Finished downloading model {model} to {download_dir}")
# Wrie hf_folder to file
with open("/local_model_path.txt", "w") as f:
f.write(hf_folder)
+76 -55
View File
@@ -7,15 +7,18 @@ from vllm import AsyncLLMEngine, AsyncEngineArgs, SamplingParams
from vllm.entrypoints.openai.serving_chat import OpenAIServingChat
from vllm.entrypoints.openai.protocol import ChatCompletionRequest
from transformers import AutoTokenizer
from utils import count_physical_cores
from utils import count_physical_cores, DummyRequest
from constants import DEFAULT_MAX_CONCURRENCY
from dotenv import load_dotenv
class Tokenizer:
def __init__(self, model_name: str):
def __init__(self, model_name):
self.tokenizer = AutoTokenizer.from_pretrained(model_name)
self.has_chat_template = bool(self.tokenizer.chat_template)
self.custom_chat_template = os.getenv("CUSTOM_CHAT_TEMPLATE")
self.has_chat_template = bool(self.tokenizer.chat_template) or bool(self.custom_chat_template)
if self.custom_chat_template and isinstance(self.custom_chat_template, str):
self.tokenizer.chat_template = self.custom_chat_template
def apply_chat_template(self, input: Union[str, list[dict[str, str]]]) -> str:
if isinstance(input, list):
@@ -38,7 +41,7 @@ class vLLMEngine:
load_dotenv() # For local development
self.config = self._initialize_config()
logging.info("vLLM config: %s", self.config)
self.tokenizer = Tokenizer(self.config["model"])
self.tokenizer = Tokenizer(os.getenv("TOKENIZER_NAME", os.getenv("MODEL_NAME")))
self.llm = self._initialize_llm() if engine is None else engine
self.openai_engine = self._initialize_openai()
self.max_concurrency = int(os.getenv("MAX_CONCURRENCY", DEFAULT_MAX_CONCURRENCY))
@@ -106,70 +109,76 @@ class vLLMEngine:
async def generate_openai_chat(self, llm_input, validated_sampling_params, batch_size, stream, apply_chat_template, request_id: str) -> AsyncGenerator[dict, None]:
if not isinstance(llm_input, list):
raise ValueError("Input must be a list of messages")
if not stream:
raise ValueError("OpenAI Chat Completion Format only supports streaming")
if isinstance(llm_input, str):
llm_input = [{"role": "user", "content": llm_input}]
logging.warning("OpenAI Chat Completion format requires list input, converting to list and assigning 'user' role")
if not self.openai_engine:
raise ValueError("OpenAI Chat Completion format is disabled")
chat_completion_request = ChatCompletionRequest(
model=self.config["model"],
messages=llm_input,
stream=True,
stream=stream,
**validated_sampling_params,
)
response_generator = await self.openai_engine.create_chat_completion(chat_completion_request, None) # None for raw_request
batch_contents = {}
batch_latest_choices = {}
batch_token_counter = 0
last_chunk = {}
async for chunk_str in response_generator:
try:
chunk = json.loads(chunk_str.removeprefix("data: ").rstrip("\n\n"))
except:
continue
response_generator = await self.openai_engine.create_chat_completion(chat_completion_request, DummyRequest())
if not stream:
yield json.loads(response_generator.model_dump_json())
else:
batch_contents = {}
batch_latest_choices = {}
batch_token_counter = 0
last_chunk = {}
if "choices" in chunk:
for choice in chunk["choices"]:
choice_index = choice["index"]
if "delta" in choice and "content" in choice["delta"]:
batch_contents[choice_index] = batch_contents.get(choice_index, []) + [choice["delta"]["content"]]
batch_latest_choices[choice_index] = choice
batch_token_counter += 1
last_chunk = chunk
if batch_token_counter >= batch_size:
async for chunk_str in response_generator:
try:
chunk = json.loads(chunk_str.removeprefix("data: ").rstrip("\n\n"))
except:
continue
if "choices" in chunk:
for choice in chunk["choices"]:
choice_index = choice["index"]
if "delta" in choice and "content" in choice["delta"]:
batch_contents[choice_index] = batch_contents.get(choice_index, []) + [choice["delta"]["content"]]
batch_latest_choices[choice_index] = choice
batch_token_counter += 1
last_chunk = chunk
if batch_token_counter >= batch_size:
for choice_index in batch_latest_choices:
batch_latest_choices[choice_index]["delta"]["content"] = batch_contents[choice_index]
last_chunk["choices"] = list(batch_latest_choices.values())
yield last_chunk
batch_contents = {}
batch_latest_choices = {}
batch_token_counter = 0
if batch_contents:
for choice_index in batch_latest_choices:
batch_latest_choices[choice_index]["delta"]["content"] = batch_contents[choice_index]
last_chunk["choices"] = list(batch_latest_choices.values())
yield last_chunk
batch_contents = {}
batch_latest_choices = {}
batch_token_counter = 0
if batch_contents:
for choice_index in batch_latest_choices:
batch_latest_choices[choice_index]["delta"]["content"] = batch_contents[choice_index]
last_chunk["choices"] = list(batch_latest_choices.values())
yield last_chunk
def _initialize_config(self):
quantization = self._get_quantization()
dtype = "half" if quantization else "auto"
model, download_dir = self._get_model_name_and_path()
return {
"model": os.getenv("MODEL_NAME"),
"download_dir": os.getenv("MODEL_BASE_PATH", "/runpod-volume/"),
"model": model,
"download_dir": download_dir,
"quantization": quantization,
"load_format": os.getenv("LOAD_FORMAT", "auto"),
"dtype": dtype,
"dtype": "half" if quantization else "auto",
"tokenizer": os.getenv("TOKENIZER_NAME"),
"disable_log_stats": bool(int(os.getenv("DISABLE_LOG_STATS", 1))),
"disable_log_requests": bool(int(os.getenv("DISABLE_LOG_REQUESTS", 1))),
"trust_remote_code": bool(int(os.getenv("TRUST_REMOTE_CODE", 0))),
"gpu_memory_utilization": float(os.getenv("GPU_MEMORY_UTILIZATION", 0.95)),
"max_parallel_loading_workers": int(os.getenv("MAX_PARALLEL_LOADING_WORKERS", count_physical_cores())),
"max_parallel_loading_workers": self._get_max_parallel_loading_workers(),
"max_model_len": self._get_max_model_len(),
"tensor_parallel_size": self._get_num_gpu_shard(),
}
@@ -183,22 +192,34 @@ class vLLMEngine:
def _initialize_openai(self):
if bool(int(os.getenv("ALLOW_OPENAI_FORMAT", 1))) and self.tokenizer.has_chat_template:
return OpenAIServingChat(self.llm, self.config["model"], "assistant")
return OpenAIServingChat(self.llm, self.config["model"], "assistant", self.tokenizer.tokenizer.chat_template)
else:
return None
def _get_max_parallel_loading_workers(self):
if int(os.getenv("TENSOR_PARALLEL_SIZE", 1)) > 1:
return None
else:
return int(os.getenv("MAX_PARALLEL_LOADING_WORKERS", count_physical_cores()))
def _get_model_name_and_path(self):
if os.path.exists("/local_model_path.txt"):
model, download_dir = open("/local_model_path.txt", "r").read().strip(), None
logging.info("Using local model at %s", model)
else:
model, download_dir = os.getenv("MODEL_NAME"), os.getenv("HF_HOME")
return model, download_dir
def _get_num_gpu_shard(self):
final_num_gpu_shard = 1
if bool(int(os.getenv("USE_TENSOR_PARALLEL", 0))):
env_num_gpu_shard = int(os.getenv("TENSOR_PARALLEL_SIZE", 1))
num_gpu_shard = int(os.getenv("TENSOR_PARALLEL_SIZE", 1))
if num_gpu_shard > 1:
num_gpu_available = device_count()
final_num_gpu_shard = min(env_num_gpu_shard, num_gpu_available)
logging.info("Using %s GPU shards", final_num_gpu_shard)
return final_num_gpu_shard
num_gpu_shard = min(num_gpu_shard, num_gpu_available)
logging.info("Using %s GPU shards", num_gpu_shard)
return num_gpu_shard
def _get_max_model_len(self):
max_model_len = os.getenv("MAX_MODEL_LEN")
max_model_len = os.getenv("MAX_MODEL_LENGTH")
return int(max_model_len) if max_model_len is not None else None
def _get_n_current_jobs(self):
+4 -1
View File
@@ -46,4 +46,7 @@ class JobInput:
self.use_openai_format = job.get("use_openai_format", False)
self.validated_sampling_params = validate_sampling_params(job.get("sampling_params", {}))
self.request_id = random_uuid()
class DummyRequest:
async def is_disconnected(self):
return False