This commit is contained in:
Jorg Doku
2023-08-30 00:08:11 -05:00
parent 476498b9bc
commit 5e4865e708
3 changed files with 63 additions and 23 deletions
+43
View File
@@ -0,0 +1,43 @@
#!/bin/bash
# Number of concurrent requests
num_instances=40
# Start time
start_time=$(date +%s)
# Function to run the command
run_command() {
curl -v --request POST \
--url https://api.runpod.ai/v2/u4f0txz9dssz6q/runsync \
--header 'accept: application/json' \
--header 'authorization: UGXLPIBFYJXUCKOBNR8CCJ3KP8GWJNI97F91QT6Y' \
--header 'content-type: application/json' \
--data '
{
"input": {
"prompt": "Who is the president of the United States?",
"sampling_params": {
"max_tokens": "300",
"ignore_eos": true
}
}
}
'
}
# Run instances in the background
for ((i = 1; i <= num_instances; i++)); do
run_command &
done
# Wait for all instances to finish
wait
# End time
end_time=$(date +%s)
# Calculate total time
total_time=$((end_time - start_time))
echo "Total time taken for $num_instances instances (1k tokens each): $total_time seconds"
+15 -9
View File
@@ -1,15 +1,12 @@
import concurrent.futures
import requests
import json
import time
import os
RUNPOD_ENDPOINT = os.environ('RUNPOD_ENDPOINT')
RUNPOD_API_KEY = os.environ('RUNPOD_API_KEY')
url = "https://api.runpod.ai/v2/{RUNPOD_ENDPOINT}/run"
url = "https://api.runpod.ai/v2/4hlrhh430u5tz7/runsync"
headers = {
"Authorization": RUNPOD_API_KEY,
"Authorization":"UGXLPIBFYJXUCKOBNR8CCJ3KP8GWJNI97F91QT6Y",
"Content-Type": "application/json"
}
@@ -31,18 +28,21 @@ payload = {
}
}
import concurrent.futures
import time
def make_request(url, headers, payload):
response = requests.post(url, headers=headers, json=payload)
return response
url = "https://api.runpod.ai/v2/4hlrhh430u5tz7/run"
while True:
# Number of concurrent requests to make per second.
# Number of concurrent requests to make
num_requests = 100
with concurrent.futures.ThreadPoolExecutor(max_workers=num_requests) as executor:
futures = [executor.submit(make_request, url, headers, payload)
for _ in range(num_requests)]
futures = [executor.submit(make_request, url, headers, payload) for _ in range(num_requests)]
# Wait for all requests to complete
for future in concurrent.futures.as_completed(futures):
@@ -52,3 +52,9 @@ while True:
# Sleep for 1 second before starting the next iteration
time.sleep(1)
# response_json = json.loads(response.text)
# get_status = requests.get(status_url, headers=headers)
# print(get_status.text)
+5 -14
View File
@@ -34,9 +34,9 @@ engine_args = AsyncEngineArgs(
tensor_parallel_size=NUM_GPU_SHARD,
dtype="auto",
seed=0,
#max_num_batched_tokens=8192,
#max_num_seqs=4096,
disable_log_stats=False
max_num_batched_tokens=8192,
disable_log_stats=False,
#max_num_seqs=256,
)
# Create the vLLM asynchronous engine
@@ -50,15 +50,13 @@ def concurrency_controller() -> bool:
total_pending_sequences = len(llm.engine.scheduler.waiting) + len(llm.engine.scheduler.swapped)
print("vLLM has {} pending sequences in its internal queue.".format(total_pending_sequences))
# If we have over 30 pending sequences, then we'll start auto-scaling.
return total_pending_sequences > 30
# If we have pending sequences, then we'll start auto-scaling.
return total_pending_sequences > 0
def prepare_metrics() -> dict:
# The vLLM metrics are updated every 5 seconds, see metrics.py for the _LOGGING_INTERVAL_SEC field.
if hasattr(llm.engine, 'metrics'):
print("Showing metrics")
print(llm.engine.metrics)
return llm.engine.metrics
else:
return {}
@@ -154,10 +152,8 @@ async def handler_streaming(job: dict) -> Generator[dict[str, list], None, None]
results_generator = llm.generate(prompt, sampling_params, request_id)
# Streaming case
print("Phase B")
positions = None
async for request_output in results_generator:
print("Phase C")
prompt = request_output.prompt
text_outputs = []
@@ -167,7 +163,6 @@ async def handler_streaming(job: dict) -> Generator[dict[str, list], None, None]
'token_pos': 0
}] * len(request_output.outputs)
print("Phase D")
for idx, output in enumerate(request_output.outputs):
# Extract the chunk position
text_pos = positions[idx]['text_pos']
@@ -177,7 +172,6 @@ async def handler_streaming(job: dict) -> Generator[dict[str, list], None, None]
text_chunk = " ".join(output.text.split(" ")[text_pos:])
text_outputs.append(text_chunk)
print("Phase E")
# Metrics for the vLLM serverless worker
runpod_metrics = prepare_metrics()
metrics = {}
@@ -188,7 +182,6 @@ async def handler_streaming(job: dict) -> Generator[dict[str, list], None, None]
# The input tokens is the prompt. For each 'num_seqs' we'll have that many of them.
metrics['input_tokens'] = len(request_output.prompt_token_ids)
print("Phase F")
metrics['output_tokens'] = []
for output in request_output.outputs:
token_pos = positions[idx]['token_pos']
@@ -196,7 +189,6 @@ async def handler_streaming(job: dict) -> Generator[dict[str, list], None, None]
metrics['output_tokens'].append(num_output_tokens)
print("Phase G")
# Update positions
for idx, output in enumerate(request_output.outputs):
positions[idx] = {
@@ -204,7 +196,6 @@ async def handler_streaming(job: dict) -> Generator[dict[str, list], None, None]
'token_pos': len(output.token_ids)
}
print("Phase H")
ret = {
"text": text_outputs,
"metrics": metrics,