Skip to content

Commit b604581

Browse files
committed
[BugFix] Seperate prometheus multiproc dir for single-server multi-dp services
1 parent eb7ea99 commit b604581

4 files changed

Lines changed: 27 additions & 7 deletions

File tree

fastdeploy/engine/common_engine.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,7 @@
6767
)
6868
from fastdeploy.inter_communicator.fmq import FMQ
6969
from fastdeploy.metrics.metrics import main_process_metrics
70+
from fastdeploy.metrics.prometheus_multiprocess_setup import setup_dp_prometheus_dir
7071
from fastdeploy.model_executor.guided_decoding import schema_checker
7172
from fastdeploy.plugins.token_processor import load_token_processor_plugins
7273
from fastdeploy.spec_decode import SpecMethod
@@ -2570,6 +2571,7 @@ def launch_components(self):
25702571
self.launched_expert_service_signal.value[0] = 1
25712572
self.dp_processed = []
25722573
self.dp_engine_worker_queue_server = []
2574+
base_prom_dir = os.environ.get("PROMETHEUS_MULTIPROC_DIR")
25732575
for i in range(
25742576
1,
25752577
self.cfg.parallel_config.data_parallel_size // self.cfg.nnode,
@@ -2608,10 +2610,13 @@ def launch_components(self):
26082610
f"Engine is initialized successfully with {self.cfg.parallel_config.tensor_parallel_size}"
26092611
+ f" data parallel id {i}"
26102612
)
2613+
setup_dp_prometheus_dir(i, base_prom_dir)
26112614
self.dp_processed[-1].start()
26122615
while self.launched_expert_service_signal.value[i] == 0:
26132616
time.sleep(1)
26142617

2618+
setup_dp_prometheus_dir(0, base_prom_dir)
2619+
26152620
def check_worker_initialize_status(self):
26162621
"""
26172622
Check the initlialize status of workers by stdout logging

fastdeploy/engine/engine.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@
4545
from fastdeploy.engine.request import Request
4646
from fastdeploy.inter_communicator import EngineWorkerQueue, IPCSignal
4747
from fastdeploy.metrics.metrics import main_process_metrics
48+
from fastdeploy.metrics.prometheus_multiprocess_setup import setup_dp_prometheus_dir
4849
from fastdeploy.platforms import current_platform
4950
from fastdeploy.utils import EngineError, console_logger, envs, llm_logger
5051

@@ -805,6 +806,7 @@ def launch_components(self):
805806
self.launched_expert_service_signal.value[0] = 1
806807
self.dp_processed = []
807808
self.dp_engine_worker_queue_server = []
809+
base_prom_dir = os.environ.get("PROMETHEUS_MULTIPROC_DIR")
808810
for i in range(
809811
1,
810812
self.cfg.parallel_config.data_parallel_size // self.cfg.nnode,
@@ -842,8 +844,11 @@ def launch_components(self):
842844
f"Engine is initialized successfully with {self.cfg.parallel_config.tensor_parallel_size}"
843845
+ f" data parallel id {i}"
844846
)
847+
setup_dp_prometheus_dir(i, base_prom_dir)
845848
self.dp_processed[-1].start()
846849

850+
setup_dp_prometheus_dir(0, base_prom_dir)
851+
847852
for i in range(
848853
1,
849854
self.cfg.parallel_config.data_parallel_size // self.cfg.nnode,

fastdeploy/entrypoints/openai/multi_api_server.py

Lines changed: 2 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import sys
2121
import time
2222

23+
from fastdeploy.metrics.prometheus_multiprocess_setup import setup_dp_prometheus_dir
2324
from fastdeploy.platforms import current_platform
2425
from fastdeploy.utils import find_free_ports, get_logger, is_port_available
2526

@@ -108,13 +109,7 @@ def start_servers(
108109
env["FD_ENABLE_MULTI_API_SERVER"] = "1"
109110
env["FD_LOG_DIR"] = env.get("FD_LOG_DIR", "log") + f"/log_{i}"
110111
if "PROMETHEUS_MULTIPROC_DIR" in env:
111-
prom_dir = env.get("PROMETHEUS_MULTIPROC_DIR")
112-
prom_dir_i = os.path.join(os.path.dirname(prom_dir), os.path.basename(prom_dir) + f"_dp{i}")
113-
# Create the directory if it doesn't exist
114-
if not os.path.exists(prom_dir_i):
115-
os.makedirs(prom_dir_i, exist_ok=True)
116-
env["PROMETHEUS_MULTIPROC_DIR"] = prom_dir_i
117-
logger.info(f"Set PROMETHEUS_MULTIPROC_DIR for DP {i}: {prom_dir_i}")
112+
setup_dp_prometheus_dir(i, env["PROMETHEUS_MULTIPROC_DIR"], env)
118113

119114
cmd = [
120115
sys.executable,

fastdeploy/metrics/prometheus_multiprocess_setup.py

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,3 +44,18 @@ def setup_multiprocess_prometheus():
4444
"will properly handle cleanup."
4545
)
4646
return os.environ["PROMETHEUS_MULTIPROC_DIR"]
47+
48+
49+
def setup_dp_prometheus_dir(dp_id, base_dir, env_dict=None):
50+
"""Set up an isolated PROMETHEUS_MULTIPROC_DIR for a given data parallel rank.
51+
52+
Args:
53+
dp_id: The data parallel rank id.
54+
base_dir: The base PROMETHEUS_MULTIPROC_DIR to derive from.
55+
env_dict: If provided, write to this dict instead of os.environ.
56+
"""
57+
prom_dir_dp = os.path.join(base_dir, f"dp{dp_id}")
58+
os.makedirs(prom_dir_dp, exist_ok=True)
59+
target = env_dict if env_dict is not None else os.environ
60+
target["PROMETHEUS_MULTIPROC_DIR"] = prom_dir_dp
61+
llm_logger.info(f"Set PROMETHEUS_MULTIPROC_DIR for DP {dp_id}: {prom_dir_dp}")

0 commit comments

Comments
 (0)