mirror of
https://github.com/PaddlePaddle/FastDeploy.git
synced 2026-04-23 00:17:25 +08:00
[PD Disaggregation] support DP via v1 router and decouple DP and EP (#5197)
* [fix] support DP via v1 router and decouple DP and EP * [fix] fix scripts * [fix] reset model path * [fix] dp use get_output_ep, fix router port type, update scripts * [merge] merge with latest code * [chore] remove some debug log * [fix] fix code style check * [fix] fix test_multi_api_server for log_dir name * [chore] reduce logs * Apply suggestions from code review Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
@@ -0,0 +1,141 @@
|
||||
#!/bin/bash
|
||||
set -e
|
||||
|
||||
# Test splitwise deployment
|
||||
# There are two methods for splitwise deployment:
|
||||
# v0: using splitwise_scheduler or dp_scheduler
|
||||
# v1: using local_scheduler + router
|
||||
|
||||
MODEL_NAME="PaddlePaddle/ERNIE-4.5-0.3B-Paddle"
|
||||
DATA_PARALLEL_SIZE=2
|
||||
TENSOR_PARALLEL_SIZE=1
|
||||
NUM_GPUS=$(($DATA_PARALLEL_SIZE * $TENSOR_PARALLEL_SIZE))
|
||||
LOG_DATE=$(date +%Y%m%d_%H%M%S)
|
||||
|
||||
export FD_DEBUG=1
|
||||
export ENABLE_V1_KVCACHE_SCHEDULER=1
|
||||
export KVCACHE_GDRCOPY_FLUSH_ENABLE=1
|
||||
export FD_ENABLE_MULTI_API_SERVER=1
|
||||
|
||||
SCRIPT_PATH=$(readlink -f "$0")
|
||||
SCRIPT_DIR=$(dirname "$SCRIPT_PATH")
|
||||
export $(bash ${SCRIPT_DIR}/../../scripts/get_rdma_nics.sh gpu)
|
||||
echo "KVCACHE_RDMA_NICS:${KVCACHE_RDMA_NICS}"
|
||||
if [ -z "${KVCACHE_RDMA_NICS}" ]; then
|
||||
echo "KVCACHE_RDMA_NICS is empty, please check the output of get_rdma_nics.sh"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
unset http_proxy && unset https_proxy
|
||||
source ${SCRIPT_DIR}/utils.sh
|
||||
|
||||
# start router
|
||||
ROUTER_PORT=$(get_free_ports 1)
|
||||
echo "---------------------------"
|
||||
echo ROUTER_PORT: $ROUTER_PORT
|
||||
|
||||
export FD_LOG_DIR="log/$LOG_DATE/router"
|
||||
rm -rf $FD_LOG_DIR
|
||||
mkdir -p ${FD_LOG_DIR}
|
||||
|
||||
nohup python -m fastdeploy.router.launch \
|
||||
--port ${ROUTER_PORT} \
|
||||
--splitwise \
|
||||
2>&1 >${FD_LOG_DIR}/nohup &
|
||||
sleep 1
|
||||
|
||||
|
||||
# start prefill
|
||||
P_SERVER_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
P_METRICS_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
P_ENGINE_WORKER_QUEUE_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
P_CACHE_QUEUE_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
P_RDMA_COMM_PORTS=$(get_free_ports $NUM_GPUS)
|
||||
P_PD_COMM_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
echo "---------------------------"
|
||||
echo P_SERVER_PORTS: $P_SERVER_PORTS
|
||||
echo P_METRICS_PORTS: $P_METRICS_PORTS
|
||||
echo P_ENGINE_WORKER_QUEUE_PORTS: $P_ENGINE_WORKER_QUEUE_PORTS
|
||||
echo P_CACHE_QUEUE_PORTS: $P_CACHE_QUEUE_PORTS
|
||||
echo P_RDMA_COMM_PORTS: $P_RDMA_COMM_PORTS
|
||||
echo P_PD_COMM_PORTS: $P_PD_COMM_PORTS
|
||||
|
||||
export CUDA_VISIBLE_DEVICES="0,1"
|
||||
export FD_LOG_DIR="log/$LOG_DATE/prefill"
|
||||
rm -rf $FD_LOG_DIR
|
||||
mkdir -p ${FD_LOG_DIR}
|
||||
|
||||
nohup python -m fastdeploy.entrypoints.openai.multi_api_server \
|
||||
--num-servers ${DATA_PARALLEL_SIZE}\
|
||||
--ports ${P_SERVER_PORTS} \
|
||||
--metrics-port ${P_METRICS_PORTS} \
|
||||
--args --model ${MODEL_NAME} \
|
||||
--engine-worker-queue-port ${P_ENGINE_WORKER_QUEUE_PORTS} \
|
||||
--cache-queue-port ${P_CACHE_QUEUE_PORTS} \
|
||||
--max-model-len 32768 \
|
||||
--data-parallel-size ${DATA_PARALLEL_SIZE} \
|
||||
--tensor-parallel-size ${TENSOR_PARALLEL_SIZE} \
|
||||
--splitwise-role "prefill" \
|
||||
--cache-transfer-protocol "rdma" \
|
||||
--rdma-comm-ports ${P_RDMA_COMM_PORTS} \
|
||||
--pd-comm-port ${P_PD_COMM_PORTS} \
|
||||
--router "0.0.0.0:${ROUTER_PORT}" \
|
||||
2>&1 >${FD_LOG_DIR}/nohup &
|
||||
|
||||
echo "--- Health Check Status ---"
|
||||
wait_for_health ${P_SERVER_PORTS}
|
||||
|
||||
|
||||
# start decode
|
||||
D_SERVER_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
D_ENGINE_WORKER_QUEUE_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
D_CACHE_QUEUE_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
D_METRICS_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
D_RDMA_COMM_PORTS=$(get_free_ports $NUM_GPUS)
|
||||
D_PD_COMM_PORTS=$(get_free_ports $DATA_PARALLEL_SIZE)
|
||||
echo "---------------------------"
|
||||
echo D_SERVER_PORTS: $D_SERVER_PORTS
|
||||
echo D_ENGINE_WORKER_QUEUE_PORTS: $D_ENGINE_WORKER_QUEUE_PORTS
|
||||
echo D_CACHE_QUEUE_PORTS: $D_CACHE_QUEUE_PORTS
|
||||
echo D_METRICS_PORTS: $D_METRICS_PORTS
|
||||
echo D_RDMA_COMM_PORTS: $D_RDMA_COMM_PORTS
|
||||
echo D_PD_COMM_PORTS: $D_PD_COMM_PORTS
|
||||
|
||||
export CUDA_VISIBLE_DEVICES="2,3"
|
||||
export FD_LOG_DIR="log/$LOG_DATE/decode"
|
||||
rm -rf $FD_LOG_DIR
|
||||
mkdir -p ${FD_LOG_DIR}
|
||||
|
||||
nohup python -m fastdeploy.entrypoints.openai.multi_api_server \
|
||||
--num-servers ${DATA_PARALLEL_SIZE}\
|
||||
--ports ${D_SERVER_PORTS} \
|
||||
--metrics-port ${D_METRICS_PORTS} \
|
||||
--args --model ${MODEL_NAME} \
|
||||
--engine-worker-queue-port ${D_ENGINE_WORKER_QUEUE_PORTS} \
|
||||
--cache-queue-port ${D_CACHE_QUEUE_PORTS} \
|
||||
--max-model-len 32768 \
|
||||
--data-parallel-size ${DATA_PARALLEL_SIZE} \
|
||||
--tensor-parallel-size ${TENSOR_PARALLEL_SIZE} \
|
||||
--splitwise-role "decode" \
|
||||
--cache-transfer-protocol "rdma" \
|
||||
--rdma-comm-ports ${D_RDMA_COMM_PORTS} \
|
||||
--pd-comm-port ${D_PD_COMM_PORTS} \
|
||||
--router "0.0.0.0:${ROUTER_PORT}" \
|
||||
2>&1 >${FD_LOG_DIR}/nohup &
|
||||
|
||||
echo "--- Health Check Status ---"
|
||||
wait_for_health ${D_SERVER_PORTS}
|
||||
|
||||
|
||||
# send request
|
||||
echo "------ Request Check ------"
|
||||
sleep 10 # make sure server is registered to router
|
||||
curl -X POST "http://0.0.0.0:${ROUTER_PORT}/v1/chat/completions" \
|
||||
-H "Content-Type: application/json" \
|
||||
-d '{
|
||||
"messages": [
|
||||
{"role": "user", "content": "hello"}
|
||||
],
|
||||
"max_tokens": 100,
|
||||
"stream": false
|
||||
}'
|
||||
@@ -1,8 +1,16 @@
|
||||
#!/bin/bash
|
||||
|
||||
is_port_free() {
|
||||
local port=$1
|
||||
if ss -ltn | awk '{print $4}' | grep -q ":${port}$"; then
|
||||
return 1 # Port is occupied
|
||||
fi
|
||||
return 0 # Port is free
|
||||
}
|
||||
|
||||
check_ports() {
|
||||
for port in "$@"; do
|
||||
if ss -tuln | grep -q ":$port "; then
|
||||
if ! is_port_free $port; then
|
||||
echo "❌ Port $port is already in use"
|
||||
return 1
|
||||
fi
|
||||
@@ -11,14 +19,79 @@ check_ports() {
|
||||
}
|
||||
|
||||
wait_for_health() {
|
||||
local server_port=$1
|
||||
IFS=',' read -r -a server_ports <<< "$1"
|
||||
local num_ports=${#server_ports[@]}
|
||||
local total_lines=$((num_ports + 1))
|
||||
local first_run=true
|
||||
local GREEN='\033[0;32m'
|
||||
local RED='\033[0;31m'
|
||||
local NC='\033[0m' # No Color
|
||||
local start_time=$(date +%s)
|
||||
|
||||
while true; do
|
||||
status_code=$(curl -s -o /dev/null -w "%{http_code}" "http://0.0.0.0:${server_port}/health" || echo "000")
|
||||
if [ "$status_code" -eq 200 ]; then
|
||||
local all_ready=true
|
||||
for port in "${server_ports[@]}"; do
|
||||
status_code=$(curl -s --max-time 1 -o /dev/null -w "%{http_code}" "http://0.0.0.0:${port}/health" || echo "000")
|
||||
if [ "$status_code" -eq 200 ]; then
|
||||
printf "Port %s: ${GREEN}[OK] 200${NC}\033[K\n" "$port"
|
||||
else
|
||||
all_ready=false
|
||||
printf "Port %s: ${RED}[WAIT] %s${NC}\033[K\n" "$port" "$status_code"
|
||||
fi
|
||||
done
|
||||
cur_time=$(date +%s)
|
||||
if [ "$all_ready" = "true" ]; then
|
||||
echo "All services are ready! [$((cur_time-start_time))s]"
|
||||
break
|
||||
else
|
||||
echo "Service not ready. Retrying in 4s..."
|
||||
sleep 4
|
||||
fi
|
||||
else
|
||||
echo "Waiting for services... [$((cur_time-start_time))s]"
|
||||
printf "\033[%dA" "$total_lines" # roll back cursor
|
||||
sleep 1
|
||||
fi
|
||||
done
|
||||
}
|
||||
|
||||
get_free_ports() {
|
||||
free_ports_num=${1:-1}
|
||||
start_port=${2:-8000}
|
||||
end_port=${3:-9000}
|
||||
|
||||
free_ports=()
|
||||
if [[ ! -n ${free_ports_num} || "${free_ports_num}" -le 0 ]]; then
|
||||
log_warn "param can't be empty, and should > 0"
|
||||
echo ${free_ports[@]}
|
||||
return 1
|
||||
fi
|
||||
|
||||
used_ports1=$(netstat -an | grep -E "(0.0.0.0|127.0.0.1|${POD_IP}|tcp6)" | awk '{n=split($4,a,":"); if(a[n]~/^[0-9]+$/) print a[n];}' | sort -u)
|
||||
used_ports2=$(netstat -an | grep -E "(0.0.0.0|127.0.0.1|${POD_IP}|tcp6)" | awk '{n=split($5,a,":"); if(a[n]~/^[0-9]+$/) print a[n];}' | sort -u)
|
||||
all_used_ports=$(printf "%s\n" "${used_ports1}" "${used_ports2}" | sort -u)
|
||||
|
||||
# Generate random number between 0 and 32767
|
||||
random_num=$(( RANDOM ))
|
||||
port=$(( random_num % (end_port - start_port + 1) + start_port ))
|
||||
|
||||
while true; do
|
||||
(( port++ ))
|
||||
if [[ ${port} -ge ${end_port} ]]; then
|
||||
port=${start_port}
|
||||
fi
|
||||
|
||||
if [[ "${all_used_ports[@]}" =~ "${port}" ]]; then
|
||||
continue
|
||||
fi
|
||||
|
||||
if is_port_free ${port}; then
|
||||
free_ports+=("${port}")
|
||||
(( free_ports_num-- ))
|
||||
if [[ ${free_ports_num} = 0 ]]; then
|
||||
break
|
||||
fi
|
||||
fi
|
||||
|
||||
done
|
||||
|
||||
# echo ${free_ports[@]}
|
||||
IFS=',' && echo "${free_ports[*]}"
|
||||
return 0
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user