Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@
# - NPU: HCSP operator timeout detects the dead peer.
# - Deadline (45s): slowest fallback (30s) + margin.
CPU_DISTRIBUTED_TIMEOUT_S = 30
FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS = 15
FT_COMMUNICATION_ABORT_TIMEOUT_S = 15
FAULT_DETECTION_DEADLINE_S = 45


Expand Down Expand Up @@ -129,7 +129,7 @@ def _ft_server_args() -> list[str]:
"--fault-tolerance-config",
'{"engine_recovery_timeout_sec": 120}',
"--additional-config",
f'{{"ft_communication_ops_abort_timeout_ms": {FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS}}}',
f'{{"ft_communication_abort_timeout": {FT_COMMUNICATION_ABORT_TIMEOUT_S}}}',
]


Expand Down
22 changes: 10 additions & 12 deletions vllm_ascend/ascend_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -200,20 +200,18 @@ def __init__(self, vllm_config: "VllmConfig"):
"VLLM_ASCEND_FUSION_OP_TRANSPOSE_KV_CACHE_BY_BLOCK",
ascend_envs.VLLM_ASCEND_FUSION_OP_TRANSPOSE_KV_CACHE_BY_BLOCK,
)
self.ft_communication_ops_abort_timeout_ms = self._get_config_value(
additional_config,
"ft_communication_ops_abort_timeout_ms",
"FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS",
ascend_envs.FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS,
)
if (
not isinstance(self.ft_communication_ops_abort_timeout_ms, int)
or self.ft_communication_ops_abort_timeout_ms < 0
):
# Fault-tolerance communication abort timeout (seconds); when fault
# tolerance is enabled and > 0, a hung NPU comm op is aborted after
# this many seconds. Drives HCCL_EVENT_TIMEOUT / HCCL_EXEC_TIMEOUT
# (= timeout - 1) and set_op_timeout_ms(timeout * 1000). 0 disables.
abort_timeout = additional_config.get("ft_communication_abort_timeout", 0)
if not isinstance(abort_timeout, int) or not (abort_timeout == 0 or abort_timeout >= 2):
raise ValueError(
"ft_communication_ops_abort_timeout_ms must be a non-negative "
f"integer, got {self.ft_communication_ops_abort_timeout_ms}"
"ft_communication_abort_timeout must be 0 (disabled) or an "
f"integer of at least 2 seconds, "
f"got {abort_timeout}"
)
self.ft_communication_abort_timeout = abort_timeout

self.pd_tp_ratio = 1
self.pd_head_ratio = 1
Expand Down
6 changes: 0 additions & 6 deletions vllm_ascend/envs.py
Original file line number Diff line number Diff line change
Expand Up @@ -100,12 +100,6 @@
# Control the aclrtMemcpyBatchAsync compile path for KV cache offloading.
# "1": force enable, "0": force disable, None: auto-detect from CANN headers.
"VLLM_ASCEND_ENABLE_BATCH_MEMCPY": lambda: os.getenv("VLLM_ASCEND_ENABLE_BATCH_MEMCPY", None),
# Fault-tolerance timeout (in milliseconds) after which a hung NPU operator
# (e.g. a communication op) is aborted with a timeout exception, so fault
# tolerance can detect it and trigger recovery. 0 (default) disables the
# timeout. Applied via torch_npu.npu.set_op_timeout_ms when fault tolerance
# is enabled.
"FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS": lambda: int(os.getenv("FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS", "0")),
}

# end-env-vars-definition
Expand Down
17 changes: 16 additions & 1 deletion vllm_ascend/worker/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import copy
import gc
import logging
import os
from types import NoneType
from typing import Any

Expand Down Expand Up @@ -370,7 +371,21 @@ def _init_device(self):
if self.parallel_config.enable_fault_tolerance:
import torch_npu

torch_npu.npu.set_op_timeout_ms(get_ascend_config().ft_communication_ops_abort_timeout_ms)
abort_timeout = get_ascend_config().ft_communication_abort_timeout
if abort_timeout > 0:
# User-provided HCCL timeouts win; otherwise derive them

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

不用加这么多注释

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

已修改

# from the config value. HCCL_EVENT_TIMEOUT must be
# greater than HCCL_EXEC_TIMEOUT, hence EXEC defaults to
# abort_timeout - 1.
os.environ.setdefault("HCCL_EVENT_TIMEOUT", str(abort_timeout))
os.environ.setdefault("HCCL_EXEC_TIMEOUT", str(abort_timeout - 1))
if int(os.environ["HCCL_EVENT_TIMEOUT"]) <= int(os.environ["HCCL_EXEC_TIMEOUT"]):
raise ValueError(
f"HCCL_EVENT_TIMEOUT ({os.environ['HCCL_EVENT_TIMEOUT']}) "
"must be greater than HCCL_EXEC_TIMEOUT "
f"({os.environ['HCCL_EXEC_TIMEOUT']})"
)
torch_npu.npu.set_op_timeout_ms(abort_timeout * 1000)
if (
parallel_config.distributed_executor_backend not in ("ray", "external_launcher")
and parallel_config.data_parallel_backend != "ray"
Expand Down
Loading