mirror of
https://github.com/NVIDIA/TensorRT-LLM.git
synced 2026-01-14 06:27:45 +08:00
1253 lines
118 KiB
HTML
1253 lines
118 KiB
HTML
|
||
|
||
<!DOCTYPE html>
|
||
|
||
|
||
<html lang="en" data-content_root="../../../" >
|
||
|
||
<head>
|
||
<meta charset="utf-8" />
|
||
<meta name="viewport" content="width=device-width, initial-scale=1.0" />
|
||
<title>tensorrt_llm.llmapi.mpi_session — TensorRT LLM</title>
|
||
|
||
|
||
|
||
<script data-cfasync="false">
|
||
document.documentElement.dataset.mode = localStorage.getItem("mode") || "";
|
||
document.documentElement.dataset.theme = localStorage.getItem("theme") || "";
|
||
</script>
|
||
<!--
|
||
this give us a css class that will be invisible only if js is disabled
|
||
-->
|
||
<noscript>
|
||
<style>
|
||
.pst-js-only { display: none !important; }
|
||
|
||
</style>
|
||
</noscript>
|
||
|
||
<!-- Loaded before other Sphinx assets -->
|
||
<link href="../../../_static/styles/theme.css?digest=8878045cc6db502f8baf" rel="stylesheet" />
|
||
<link href="../../../_static/styles/pydata-sphinx-theme.css?digest=8878045cc6db502f8baf" rel="stylesheet" />
|
||
|
||
<link rel="stylesheet" type="text/css" href="../../../_static/pygments.css?v=8f2a1f02" />
|
||
<link rel="stylesheet" type="text/css" href="../../../_static/styles/nvidia-sphinx-theme.css?v=933278ad" />
|
||
<link rel="stylesheet" type="text/css" href="../../../_static/copybutton.css?v=76b2166b" />
|
||
<link rel="stylesheet" type="text/css" href="../../../_static/autodoc_pydantic.css" />
|
||
<link rel="stylesheet" type="text/css" href="../../../_static/togglebutton.css?v=13237357" />
|
||
<link rel="stylesheet" type="text/css" href="../../../_static/custom.css?v=19d20f17" />
|
||
|
||
<!-- So that users can add custom icons -->
|
||
<script src="../../../_static/scripts/fontawesome.js?digest=8878045cc6db502f8baf"></script>
|
||
<!-- Pre-loaded scripts that we'll load fully later -->
|
||
<link rel="preload" as="script" href="../../../_static/scripts/bootstrap.js?digest=8878045cc6db502f8baf" />
|
||
<link rel="preload" as="script" href="../../../_static/scripts/pydata-sphinx-theme.js?digest=8878045cc6db502f8baf" />
|
||
|
||
|
||
|
||
<script src="../../../_static/documentation_options.js?v=5929fcd5"></script>
|
||
<script src="../../../_static/doctools.js?v=9a2dae69"></script>
|
||
<script src="../../../_static/sphinx_highlight.js?v=dc90522c"></script>
|
||
<script src="../../../_static/clipboard.min.js?v=a7894cd8"></script>
|
||
<script src="../../../_static/copybutton.js?v=65e89d2a"></script>
|
||
<script>let toggleHintShow = 'Click to show';</script>
|
||
<script>let toggleHintHide = 'Click to hide';</script>
|
||
<script>let toggleOpenOnPrint = 'true';</script>
|
||
<script src="../../../_static/togglebutton.js?v=4a39c7ea"></script>
|
||
<script>var togglebuttonSelector = '.toggle, .admonition.dropdown';</script>
|
||
<script>var togglebuttonSelector = '.toggle, .admonition.dropdown';</script>
|
||
<script>DOCUMENTATION_OPTIONS.pagename = '_modules/tensorrt_llm/llmapi/mpi_session';</script>
|
||
<script>
|
||
DOCUMENTATION_OPTIONS.theme_version = '0.16.1';
|
||
DOCUMENTATION_OPTIONS.theme_switcher_json_url = './_static/switcher.json';
|
||
DOCUMENTATION_OPTIONS.theme_switcher_version_match = '1.2.0rc4';
|
||
DOCUMENTATION_OPTIONS.show_version_warning_banner =
|
||
false;
|
||
</script>
|
||
|
||
<link rel="icon" href="../../../_static/favicon.png"/>
|
||
|
||
<link rel="index" title="Index" href="../../../genindex.html" />
|
||
<link rel="search" title="Search" href="../../../search.html" />
|
||
|
||
|
||
<meta name="viewport" content="width=device-width, initial-scale=1"/>
|
||
<meta name="docsearch:language" content="en"/>
|
||
<meta name="docsearch:version" content="1.2.0rc4" />
|
||
|
||
|
||
</head>
|
||
|
||
|
||
|
||
<body data-bs-spy="scroll" data-bs-target=".bd-toc-nav" data-offset="180" data-bs-root-margin="0px 0px -60%" data-default-mode="">
|
||
|
||
|
||
|
||
<div id="pst-skip-link" class="skip-link d-print-none"><a href="#main-content">Skip to main content</a></div>
|
||
|
||
|
||
|
||
<div id="pst-scroll-pixel-helper"></div>
|
||
|
||
<button type="button" class="btn rounded-pill" id="pst-back-to-top">
|
||
<i class="fa-solid fa-arrow-up"></i>Back to top</button>
|
||
|
||
|
||
<dialog id="pst-search-dialog">
|
||
|
||
<form class="bd-search d-flex align-items-center"
|
||
action="../../../search.html"
|
||
method="get">
|
||
<i class="fa-solid fa-magnifying-glass"></i>
|
||
<input type="search"
|
||
class="form-control"
|
||
name="q"
|
||
placeholder="Search the docs ..."
|
||
aria-label="Search the docs ..."
|
||
autocomplete="off"
|
||
autocorrect="off"
|
||
autocapitalize="off"
|
||
spellcheck="false"/>
|
||
<span class="search-button__kbd-shortcut"><kbd class="kbd-shortcut__modifier">Ctrl</kbd>+<kbd>K</kbd></span>
|
||
</form>
|
||
</dialog>
|
||
|
||
<div class="pst-async-banner-revealer d-none">
|
||
<aside id="bd-header-version-warning" class="d-none d-print-none" aria-label="Version warning"></aside>
|
||
</div>
|
||
|
||
|
||
<header class="bd-header navbar navbar-expand-lg bd-navbar d-print-none">
|
||
<div class="bd-header__inner bd-page-width">
|
||
<button class="pst-navbar-icon sidebar-toggle primary-toggle" aria-label="Site navigation">
|
||
<span class="fa-solid fa-bars"></span>
|
||
</button>
|
||
|
||
|
||
<div class="col-lg-3 navbar-header-items__start">
|
||
|
||
<div class="navbar-item">
|
||
|
||
|
||
|
||
|
||
|
||
<a class="navbar-brand logo" href="../../../index.html">
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
<img src="../../../_static/nvidia-logo-horiz-rgb-blk-for-screen.svg" class="logo__image only-light" alt="TensorRT LLM - Home"/>
|
||
<img src="../../../_static/nvidia-logo-horiz-rgb-wht-for-screen.svg" class="logo__image only-dark pst-js-only" alt="TensorRT LLM - Home"/>
|
||
|
||
|
||
<p class="title logo__title">TensorRT LLM</p>
|
||
|
||
</a></div>
|
||
|
||
</div>
|
||
|
||
<div class="col-lg-9 navbar-header-items">
|
||
|
||
<div class="me-auto navbar-header-items__center">
|
||
|
||
<div class="navbar-item">
|
||
|
||
|
||
<div class="version-switcher__container dropdown pst-js-only">
|
||
<button id="pst-version-switcher-button-2"
|
||
type="button"
|
||
class="version-switcher__button btn btn-sm dropdown-toggle"
|
||
data-bs-toggle="dropdown"
|
||
aria-haspopup="listbox"
|
||
aria-controls="pst-version-switcher-list-2"
|
||
aria-label="Version switcher list"
|
||
>
|
||
Choose version <!-- this text may get changed later by javascript -->
|
||
<span class="caret"></span>
|
||
</button>
|
||
<div id="pst-version-switcher-list-2"
|
||
class="version-switcher__menu dropdown-menu list-group-flush py-0"
|
||
role="listbox" aria-labelledby="pst-version-switcher-button-2">
|
||
<!-- dropdown will be populated by javascript on page load -->
|
||
</div>
|
||
</div></div>
|
||
|
||
</div>
|
||
|
||
|
||
<div class="navbar-header-items__end">
|
||
|
||
<div class="navbar-item navbar-persistent--container">
|
||
|
||
|
||
<button class="btn search-button-field search-button__button pst-js-only" title="Search" aria-label="Search" data-bs-placement="bottom" data-bs-toggle="tooltip">
|
||
<i class="fa-solid fa-magnifying-glass"></i>
|
||
<span class="search-button__default-text">Search</span>
|
||
<span class="search-button__kbd-shortcut"><kbd class="kbd-shortcut__modifier">Ctrl</kbd>+<kbd class="kbd-shortcut__modifier">K</kbd></span>
|
||
</button>
|
||
</div>
|
||
|
||
|
||
<div class="navbar-item">
|
||
|
||
<button class="btn btn-sm nav-link pst-navbar-icon theme-switch-button pst-js-only" aria-label="Color mode" data-bs-title="Color mode" data-bs-placement="bottom" data-bs-toggle="tooltip">
|
||
<i class="theme-switch fa-solid fa-sun fa-lg" data-mode="light" title="Light"></i>
|
||
<i class="theme-switch fa-solid fa-moon fa-lg" data-mode="dark" title="Dark"></i>
|
||
<i class="theme-switch fa-solid fa-circle-half-stroke fa-lg" data-mode="auto" title="System Settings"></i>
|
||
</button></div>
|
||
|
||
</div>
|
||
|
||
</div>
|
||
|
||
|
||
<div class="navbar-persistent--mobile">
|
||
|
||
<button class="btn search-button-field search-button__button pst-js-only" title="Search" aria-label="Search" data-bs-placement="bottom" data-bs-toggle="tooltip">
|
||
<i class="fa-solid fa-magnifying-glass"></i>
|
||
<span class="search-button__default-text">Search</span>
|
||
<span class="search-button__kbd-shortcut"><kbd class="kbd-shortcut__modifier">Ctrl</kbd>+<kbd class="kbd-shortcut__modifier">K</kbd></span>
|
||
</button>
|
||
</div>
|
||
|
||
|
||
|
||
</div>
|
||
|
||
</header>
|
||
|
||
|
||
<div class="bd-container">
|
||
<div class="bd-container__inner bd-page-width">
|
||
|
||
|
||
|
||
<dialog id="pst-primary-sidebar-modal"></dialog>
|
||
<div id="pst-primary-sidebar" class="bd-sidebar-primary bd-sidebar">
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
<a class="navbar-brand logo" href="../../../index.html">
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
<img src="../../../_static/nvidia-logo-horiz-rgb-blk-for-screen.svg" class="logo__image only-light" alt="TensorRT LLM - Home"/>
|
||
<img src="../../../_static/nvidia-logo-horiz-rgb-wht-for-screen.svg" class="logo__image only-dark pst-js-only" alt="TensorRT LLM - Home"/>
|
||
|
||
|
||
<p class="title logo__title">TensorRT LLM</p>
|
||
|
||
</a>
|
||
|
||
|
||
|
||
<div class="sidebar-header-items sidebar-primary__section">
|
||
|
||
|
||
<div class="sidebar-header-items__center">
|
||
|
||
|
||
|
||
<div class="navbar-item">
|
||
|
||
|
||
<div class="version-switcher__container dropdown pst-js-only">
|
||
<button id="pst-version-switcher-button-3"
|
||
type="button"
|
||
class="version-switcher__button btn btn-sm dropdown-toggle"
|
||
data-bs-toggle="dropdown"
|
||
aria-haspopup="listbox"
|
||
aria-controls="pst-version-switcher-list-3"
|
||
aria-label="Version switcher list"
|
||
>
|
||
Choose version <!-- this text may get changed later by javascript -->
|
||
<span class="caret"></span>
|
||
</button>
|
||
<div id="pst-version-switcher-list-3"
|
||
class="version-switcher__menu dropdown-menu list-group-flush py-0"
|
||
role="listbox" aria-labelledby="pst-version-switcher-button-3">
|
||
<!-- dropdown will be populated by javascript on page load -->
|
||
</div>
|
||
</div></div>
|
||
|
||
|
||
</div>
|
||
|
||
|
||
|
||
<div class="sidebar-header-items__end">
|
||
|
||
<div class="navbar-item">
|
||
|
||
<button class="btn btn-sm nav-link pst-navbar-icon theme-switch-button pst-js-only" aria-label="Color mode" data-bs-title="Color mode" data-bs-placement="bottom" data-bs-toggle="tooltip">
|
||
<i class="theme-switch fa-solid fa-sun fa-lg" data-mode="light" title="Light"></i>
|
||
<i class="theme-switch fa-solid fa-moon fa-lg" data-mode="dark" title="Dark"></i>
|
||
<i class="theme-switch fa-solid fa-circle-half-stroke fa-lg" data-mode="auto" title="System Settings"></i>
|
||
</button></div>
|
||
|
||
</div>
|
||
|
||
</div>
|
||
|
||
<div class="sidebar-primary-items__start sidebar-primary__section">
|
||
<div class="sidebar-primary-item">
|
||
|
||
|
||
|
||
<nav class="bd-docs-nav bd-links"
|
||
aria-label="Table of Contents">
|
||
<p class="bd-links__title" role="heading" aria-level="1">Table of Contents</p>
|
||
<div class="bd-toc-item navbar-nav"><p aria-level="2" class="caption" role="heading"><span class="caption-text">Getting Started</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../overview.html">Overview</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../quick-start-guide.html">Quick Start Guide</a></li>
|
||
<li class="toctree-l1 has-children"><a class="reference internal" href="../../../installation/index.html">Installation</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../installation/containers.html">Pre-built release container images on NGC</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../installation/linux.html">Installing on Linux via <code class="docutils literal notranslate"><span class="pre">pip</span></code></a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../installation/build-from-source-linux.html">Building from Source Code on Linux</a></li>
|
||
</ul>
|
||
</details></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">Deployment Guide</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1 has-children"><a class="reference internal" href="../../../examples/llm_api_examples.html">LLM Examples</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_inference.html">Generate text</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_inference_async.html">Generate text asynchronously</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_inference_async_streaming.html">Generate text in streaming</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_inference_distributed.html">Distributed LLM Generation</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_guided_decoding.html">Generate text with guided decoding</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_logits_processor.html">Control generated text using logits processor</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_multilora.html">Generate text with multiple LoRA adapters</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_sparse_attention.html">Sparse Attention</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_speculative_decoding.html">Speculative Decoding</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_kv_cache_connector.html">KV Cache Connector</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_kv_cache_offloading.html">KV Cache Offloading</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_runtime.html">Runtime Configuration Examples</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_sampling.html">Sampling Techniques Showcase</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_mgmn_llm_distributed.html">Run LLM-API with pytorch backend on Slurm</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_mgmn_trtllm_bench.html">Run trtllm-bench with pytorch backend on Slurm</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/llm_mgmn_trtllm_serve.html">Run trtllm-serve with pytorch backend on Slurm</a></li>
|
||
</ul>
|
||
</details></li>
|
||
<li class="toctree-l1 has-children"><a class="reference internal" href="../../../examples/trtllm_serve_examples.html">Online Serving Examples</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/curl_chat_client.html">Curl Chat Client</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/curl_chat_client_for_multimodal.html">Curl Chat Client For Multimodal</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/curl_completion_client.html">Curl Completion Client</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/deepseek_r1_reasoning_parser.html">Deepseek R1 Reasoning Parser</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/genai_perf_client.html">Genai Perf Client</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/genai_perf_client_for_multimodal.html">Genai Perf Client For Multimodal</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/openai_chat_client.html">OpenAI Chat Client</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/openai_chat_client_for_multimodal.html">OpenAI Chat Client for Multimodal</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/openai_completion_client.html">OpenAI Completion Client</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/openai_completion_client_for_lora.html">Openai Completion Client For Lora</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../examples/openai_completion_client_json_schema.html">OpenAI Completion Client with JSON Schema</a></li>
|
||
</ul>
|
||
</details></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../examples/dynamo_k8s_example.html">Dynamo K8s Example</a></li>
|
||
<li class="toctree-l1 has-children"><a class="reference internal" href="../../../deployment-guide/index.html">Model Recipes</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../deployment-guide/deployment-guide-for-deepseek-r1-on-trtllm.html">Deployment Guide for DeepSeek R1 on TensorRT LLM - Blackwell & Hopper Hardware</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../deployment-guide/deployment-guide-for-llama3.3-70b-on-trtllm.html">Deployment Guide for Llama3.3 70B on TensorRT LLM - Blackwell & Hopper Hardware</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../deployment-guide/deployment-guide-for-llama4-scout-on-trtllm.html">Deployment Guide for Llama4 Scout 17B on TensorRT LLM - Blackwell & Hopper Hardware</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../deployment-guide/deployment-guide-for-gpt-oss-on-trtllm.html">Deployment Guide for GPT-OSS on TensorRT-LLM - Blackwell Hardware</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../deployment-guide/deployment-guide-for-qwen3-next-on-trtllm.html">Deployment Guide for Qwen3 Next on TensorRT LLM - Blackwell & Hopper Hardware</a></li>
|
||
</ul>
|
||
</details></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">Models</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../models/supported-models.html">Supported Models</a></li>
|
||
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../models/adding-new-model.html">Adding a New Model</a></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">CLI Reference</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../commands/trtllm-bench.html">trtllm-bench</a></li>
|
||
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../commands/trtllm-eval.html">trtllm-eval</a></li>
|
||
<li class="toctree-l1 has-children"><a class="reference internal" href="../../../commands/trtllm-serve/index.html">trtllm-serve</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../commands/trtllm-serve/trtllm-serve.html">trtllm-serve</a></li>
|
||
<li class="toctree-l2"><a class="reference internal" href="../../../commands/trtllm-serve/run-benchmark-with-trtllm-serve.html">Run benchmarking with <code class="docutils literal notranslate"><span class="pre">trtllm-serve</span></code></a></li>
|
||
</ul>
|
||
</details></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">API Reference</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../llm-api/index.html">LLM API Introduction</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../llm-api/reference.html">API Reference</a></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">Features</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/feature-combination-matrix.html">Feature Combination Matrix</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/attention.html">Multi-Head, Multi-Query, and Group-Query Attention</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/disagg-serving.html">Disaggregated Serving</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/kvcache.html">KV Cache System</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/long-sequence.html">Long Sequences</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/lora.html">LoRA (Low-Rank Adaptation)</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/multi-modality.html">Multimodal Support in TensorRT LLM</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/overlap-scheduler.html">Overlap Scheduler</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/paged-attention-ifb-scheduler.html">Paged Attention, IFB, and Request Scheduling</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/parallel-strategy.html">Parallelism in TensorRT LLM</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/quantization.html">Quantization</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/sampling.html">Sampling</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/additional-outputs.html">Additional Outputs</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/speculative-decoding.html">Speculative Decoding</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/checkpoint-loading.html">Checkpoint Loading</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/auto_deploy/auto-deploy.html">AutoDeploy (Prototype)</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/ray-orchestrator.html">Ray Orchestrator (Prototype)</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../features/torch_compile_and_piecewise_cuda_graph.html">Torch Compile & Piecewise CUDA Graph</a></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">Developer Guide</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../developer-guide/overview.html">Architecture Overview</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../developer-guide/perf-analysis.html">Performance Analysis</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../developer-guide/perf-benchmarking.html">TensorRT LLM Benchmarking</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../developer-guide/ci-overview.html">Continuous Integration Overview</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../developer-guide/dev-containers.html">Using Dev Containers</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../developer-guide/api-change.html">LLM API Change Guide</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../developer-guide/kv-transfer.html">Introduction to KV Cache Transmission</a></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">Blogs</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog10_ADP_Balance_Strategy.html">ADP Balance Strategy</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog11_GPT_OSS_Eagle3.html">Running GPT-OSS-120B with Eagle3 Speculative Decoding on GB200/B200 (TensorRT LLM)</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog12_Combining_Guided_Decoding_and_Speculative_Decoding.html">Combining Guided Decoding and Speculative Decoding: Making CPU and GPU Cooperate Seamlessly</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog13_Inference_Time_Compute_Implementation_in_TensorRT-LLM.html">Inference Time Compute Implementation in TensorRT LLM</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog14_Scaling_Expert_Parallelism_in_TensorRT-LLM_part3.html">Scaling Expert Parallelism in TensorRT LLM (Part 3: Pushing the Performance Boundary)</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog1_Pushing_Latency_Boundaries_Optimizing_DeepSeek-R1_Performance_on_NVIDIA_B200_GPUs.html">Pushing Latency Boundaries: Optimizing DeepSeek-R1 Performance on NVIDIA B200 GPUs</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog2_DeepSeek_R1_MTP_Implementation_and_Optimization.html">DeepSeek R1 MTP Implementation and Optimization</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog3_Optimizing_DeepSeek_R1_Throughput_on_NVIDIA_Blackwell_GPUs.html">Optimizing DeepSeek R1 Throughput on NVIDIA Blackwell GPUs: A Deep Dive for Developers</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog4_Scaling_Expert_Parallelism_in_TensorRT-LLM.html">Scaling Expert Parallelism in TensorRT LLM (Part 1: Design and Implementation of Large-scale EP)</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog5_Disaggregated_Serving_in_TensorRT-LLM.html">Disaggregated Serving in TensorRT LLM</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog6_Llama4_maverick_eagle_guide.html">How to launch Llama4 Maverick + Eagle3 TensorRT LLM server</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog7_NGram_performance_Analysis_And_Auto_Enablement.html">N-Gram Speculative Decoding in TensorRT LLM</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog8_Scaling_Expert_Parallelism_in_TensorRT-LLM_part2.html">Scaling Expert Parallelism in TensorRT LLM (Part 2: Performance Status and Optimization)</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/tech_blog/blog9_Deploying_GPT_OSS_on_TRTLLM.html">Running a High Performance GPT-OSS-120B Inference Server with TensorRT LLM</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/Best_perf_practice_on_DeepSeek-R1_in_TensorRT-LLM.html">How to get best performance on DeepSeek-R1 in TensorRT LLM</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/H200launch.html">H200 achieves nearly 12,000 tokens/sec on Llama2-13B with TensorRT LLM</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/XQA-kernel.html">New XQA-kernel provides 2.4x more Llama-70B throughput within the same latency budget</a></li>
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../blogs/H100vsA100.html">H100 has 4.6x A100 Performance in TensorRT LLM, achieving 10,000 tok/s at 100ms to first token</a></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">Quick Links</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference external" href="https://github.com/NVIDIA/TensorRT-LLM/releases">Releases</a></li>
|
||
<li class="toctree-l1"><a class="reference external" href="https://github.com/NVIDIA/TensorRT-LLM">Github Code</a></li>
|
||
<li class="toctree-l1"><a class="reference external" href="https://github.com/NVIDIA/TensorRT-LLM/issues?q=is%3Aissue%20state%3Aopen%20label%3Aroadmap">Roadmap</a></li>
|
||
</ul>
|
||
<p aria-level="2" class="caption" role="heading"><span class="caption-text">Use TensorRT Engine</span></p>
|
||
<ul class="nav bd-sidenav">
|
||
<li class="toctree-l1"><a class="reference internal" href="../../../legacy/tensorrt_quickstart.html">LLM API with TensorRT Engine</a></li>
|
||
</ul>
|
||
</div>
|
||
</nav></div>
|
||
</div>
|
||
|
||
|
||
<div class="sidebar-primary-items__end sidebar-primary__section">
|
||
</div>
|
||
|
||
|
||
|
||
</div>
|
||
|
||
<main id="main-content" class="bd-main" role="main">
|
||
|
||
|
||
<div class="bd-content">
|
||
<div class="bd-article-container">
|
||
|
||
<div class="bd-header-article d-print-none">
|
||
<div class="header-article-items header-article__inner">
|
||
|
||
<div class="header-article-items__start">
|
||
|
||
<div class="header-article-item">
|
||
|
||
<nav aria-label="Breadcrumb" class="d-print-none">
|
||
<ul class="bd-breadcrumbs">
|
||
|
||
<li class="breadcrumb-item breadcrumb-home">
|
||
<a href="../../../index.html" class="nav-link" aria-label="Home">
|
||
<i class="fa-solid fa-home"></i>
|
||
</a>
|
||
</li>
|
||
|
||
<li class="breadcrumb-item"><a href="../../index.html" class="nav-link">Module code</a></li>
|
||
|
||
<li class="breadcrumb-item active" aria-current="page"><span class="ellipsis">tensorrt_llm.llmapi.mpi_session</span></li>
|
||
</ul>
|
||
</nav>
|
||
</div>
|
||
|
||
</div>
|
||
|
||
|
||
</div>
|
||
</div>
|
||
|
||
|
||
|
||
|
||
<div id="searchbox"></div>
|
||
<article class="bd-article">
|
||
|
||
<h1>Source code for tensorrt_llm.llmapi.mpi_session</h1><div class="highlight"><pre>
|
||
<span></span><span class="kn">import</span><span class="w"> </span><span class="nn">abc</span>
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">itertools</span>
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">os</span>
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">socket</span>
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">sys</span>
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">threading</span>
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">time</span>
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">traceback</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">collections.abc</span><span class="w"> </span><span class="kn">import</span> <span class="n">Callable</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">concurrent.futures</span><span class="w"> </span><span class="kn">import</span> <span class="n">Future</span><span class="p">,</span> <span class="n">ThreadPoolExecutor</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">typing</span><span class="w"> </span><span class="kn">import</span> <span class="n">Any</span><span class="p">,</span> <span class="n">Dict</span><span class="p">,</span> <span class="n">List</span><span class="p">,</span> <span class="n">NamedTuple</span><span class="p">,</span> <span class="n">Optional</span><span class="p">,</span> <span class="n">Tuple</span><span class="p">,</span> <span class="n">TypeVar</span>
|
||
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">zmq</span>
|
||
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">tensorrt_llm.bindings.BuildInfo</span><span class="w"> </span><span class="kn">import</span> <span class="n">ENABLE_MULTI_DEVICE</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">tensorrt_llm.logger</span><span class="w"> </span><span class="kn">import</span> <span class="n">logger</span>
|
||
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">.._utils</span><span class="w"> </span><span class="kn">import</span> <span class="n">global_mpi_rank</span><span class="p">,</span> <span class="n">mpi_barrier</span><span class="p">,</span> <span class="n">mpi_rank</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">.utils</span><span class="w"> </span><span class="kn">import</span> <span class="n">logger_debug</span><span class="p">,</span> <span class="n">print_colored</span>
|
||
|
||
<span class="k">if</span> <span class="n">ENABLE_MULTI_DEVICE</span><span class="p">:</span>
|
||
<span class="kn">import</span><span class="w"> </span><span class="nn">mpi4py</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">mpi4py.futures</span><span class="w"> </span><span class="kn">import</span> <span class="n">MPICommExecutor</span><span class="p">,</span> <span class="n">MPIPoolExecutor</span>
|
||
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">tensorrt_llm._utils</span><span class="w"> </span><span class="kn">import</span> <span class="n">global_mpi_size</span><span class="p">,</span> <span class="n">mpi_world_size</span>
|
||
|
||
<span class="n">T</span> <span class="o">=</span> <span class="n">TypeVar</span><span class="p">(</span><span class="s2">"T"</span><span class="p">)</span>
|
||
|
||
|
||
<span class="k">class</span><span class="w"> </span><span class="nc">MPINodeState</span><span class="p">:</span>
|
||
<span class="w"> </span><span class="sd">''' MPINodeState acts as a central global state shares between tasks on MPI node.</span>
|
||
|
||
<span class="sd"> An example:</span>
|
||
<span class="sd"> def task():</span>
|
||
<span class="sd"> if MPINodeState.state is None:</span>
|
||
<span class="sd"> MPINodeState.state = 0</span>
|
||
<span class="sd"> MPINodeState.state += 1</span>
|
||
<span class="sd"> return MPINodeState.state</span>
|
||
|
||
<span class="sd"> n_workers = 4</span>
|
||
<span class="sd"> with MPIPoolExecutor(max_workers=n_workers) as executor:</span>
|
||
<span class="sd"> for i in range(2):</span>
|
||
<span class="sd"> futures = [executor.submit(task) for i in range(n_workers)]</span>
|
||
|
||
<span class="sd"> This should produce the following output:</span>
|
||
<span class="sd"> - [1, 1, 1, 1]</span>
|
||
<span class="sd"> - [2, 2, 2, 2]</span>
|
||
<span class="sd"> '''</span>
|
||
|
||
<span class="n">state</span> <span class="o">=</span> <span class="kc">None</span>
|
||
<span class="c1"># Global MPICommExecutor instance to be reused across multiple MpiCommSession instances</span>
|
||
<span class="c1"># This is necessary because MPICommExecutor can only be created once per MPI process</span>
|
||
<span class="n">_global_comm_executor</span> <span class="o">=</span> <span class="kc">None</span>
|
||
<span class="n">_global_mpi_pool</span> <span class="o">=</span> <span class="kc">None</span>
|
||
|
||
<span class="nd">@staticmethod</span>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">is_initialized</span><span class="p">()</span> <span class="o">-></span> <span class="nb">bool</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="n">MPINodeState</span><span class="o">.</span><span class="n">state</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span>
|
||
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">external_mpi_comm_available</span><span class="p">(</span><span class="n">model_world_size</span><span class="p">:</span> <span class="nb">int</span><span class="p">)</span> <span class="o">-></span> <span class="nb">bool</span><span class="p">:</span>
|
||
<span class="w"> </span><span class="sd">''' Check if the current process is launched by mpirun and does not use MPIPoolExecutor to spawn processes.</span>
|
||
<span class="sd"> e.g. mpirun -np 4 python script.py</span>
|
||
<span class="sd"> '''</span>
|
||
<span class="k">if</span> <span class="n">ENABLE_MULTI_DEVICE</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="p">(</span><span class="n">get_mpi_world_size</span><span class="p">()</span> <span class="o">==</span> <span class="n">model_world_size</span>
|
||
<span class="ow">and</span> <span class="n">model_world_size</span> <span class="o">></span> <span class="mi">1</span><span class="p">)</span> <span class="ow">or</span> <span class="p">(</span><span class="n">global_mpi_size</span><span class="p">()</span>
|
||
<span class="o">></span> <span class="n">get_mpi_world_size</span><span class="p">())</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="kc">False</span>
|
||
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">need_spawn_mpi_workers</span><span class="p">(</span><span class="n">model_world_size</span><span class="p">:</span> <span class="nb">int</span><span class="p">)</span> <span class="o">-></span> <span class="nb">bool</span><span class="p">:</span>
|
||
<span class="w"> </span><span class="sd">''' Check if the current process needs to spawn MPI workers. '''</span>
|
||
<span class="k">if</span> <span class="n">ENABLE_MULTI_DEVICE</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="n">get_mpi_world_size</span><span class="p">()</span> <span class="o">==</span> <span class="mi">1</span> <span class="ow">and</span> <span class="n">model_world_size</span> <span class="o">></span> <span class="mi">1</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="kc">False</span>
|
||
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">set_mpi_session_cpp</span><span class="p">(</span><span class="n">comm</span><span class="p">):</span>
|
||
<span class="k">if</span> <span class="n">ENABLE_MULTI_DEVICE</span><span class="p">:</span>
|
||
<span class="n">comm_fortran</span> <span class="o">=</span> <span class="n">comm</span><span class="o">.</span><span class="n">py2f</span><span class="p">()</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">tensorrt_llm.bindings</span><span class="w"> </span><span class="kn">import</span> <span class="n">MpiComm</span>
|
||
<span class="n">MpiComm</span><span class="o">.</span><span class="n">set_raw_mpi_session_by_fortran_handle</span><span class="p">(</span><span class="n">comm_fortran</span><span class="p">)</span>
|
||
|
||
|
||
<span class="k">class</span><span class="w"> </span><span class="nc">MpiSession</span><span class="p">(</span><span class="n">abc</span><span class="o">.</span><span class="n">ABC</span><span class="p">):</span>
|
||
|
||
<span class="nd">@abc</span><span class="o">.</span><span class="n">abstractmethod</span>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">submit</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">],</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span>
|
||
<span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="n">List</span><span class="p">[</span><span class="n">Future</span><span class="p">[</span><span class="n">T</span><span class="p">]]:</span>
|
||
<span class="k">raise</span> <span class="ne">NotImplementedError</span><span class="p">()</span>
|
||
|
||
<span class="nd">@abc</span><span class="o">.</span><span class="n">abstractmethod</span>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">submit_sync</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">],</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="n">List</span><span class="p">[</span><span class="n">T</span><span class="p">]:</span>
|
||
<span class="k">raise</span> <span class="ne">NotImplementedError</span><span class="p">()</span>
|
||
|
||
<span class="nd">@abc</span><span class="o">.</span><span class="n">abstractmethod</span>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">shutdown</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">wait</span><span class="o">=</span><span class="kc">True</span><span class="p">):</span>
|
||
<span class="k">raise</span> <span class="ne">NotImplementedError</span><span class="p">()</span>
|
||
|
||
<span class="nd">@abc</span><span class="o">.</span><span class="n">abstractmethod</span>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">abort</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="k">raise</span> <span class="ne">NotImplementedError</span><span class="p">()</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">is_comm_session</span><span class="p">(</span><span class="bp">self</span><span class="p">)</span> <span class="o">-></span> <span class="nb">bool</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="nb">isinstance</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="p">(</span><span class="n">MpiCommSession</span><span class="p">,</span> <span class="n">RemoteMpiCommSessionClient</span><span class="p">))</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">_abort_on_timeout</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">fut</span><span class="p">:</span> <span class="n">Future</span><span class="p">,</span> <span class="n">timeout</span><span class="p">:</span> <span class="nb">float</span><span class="p">,</span> <span class="n">reason</span><span class="o">=</span><span class="kc">None</span><span class="p">):</span>
|
||
<span class="k">try</span><span class="p">:</span>
|
||
<span class="n">fut</span><span class="o">.</span><span class="n">result</span><span class="p">(</span><span class="n">timeout</span><span class="o">=</span><span class="n">timeout</span><span class="p">)</span>
|
||
<span class="k">except</span> <span class="ne">TimeoutError</span><span class="p">:</span>
|
||
<span class="n">logger</span><span class="o">.</span><span class="n">critical</span><span class="p">(</span><span class="s2">"MpiSession shutdown timeout, aborting..."</span><span class="p">)</span>
|
||
<span class="k">if</span> <span class="n">reason</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
|
||
<span class="n">logger</span><span class="o">.</span><span class="n">info</span><span class="p">(</span><span class="sa">f</span><span class="s2">"Reason to shutdown: </span><span class="si">{</span><span class="nb">repr</span><span class="p">(</span><span class="n">reason</span><span class="p">)</span><span class="si">}</span><span class="s2">"</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">abort</span><span class="p">()</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">shutdown_abort</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">grace</span><span class="p">:</span> <span class="nb">float</span> <span class="o">=</span> <span class="mi">60</span><span class="p">,</span> <span class="n">reason</span><span class="o">=</span><span class="kc">None</span><span class="p">):</span>
|
||
<span class="k">if</span> <span class="n">sys</span><span class="o">.</span><span class="n">is_finalizing</span><span class="p">():</span>
|
||
<span class="c1"># cannot start thread at interpreter shutdown</span>
|
||
<span class="c1"># simply don't wait to avoid hang</span>
|
||
<span class="k">return</span> <span class="bp">self</span><span class="o">.</span><span class="n">shutdown</span><span class="p">(</span><span class="n">wait</span><span class="o">=</span><span class="kc">False</span><span class="p">)</span>
|
||
|
||
<span class="n">fut</span> <span class="o">=</span> <span class="n">Future</span><span class="p">()</span>
|
||
<span class="n">killer</span> <span class="o">=</span> <span class="n">threading</span><span class="o">.</span><span class="n">Thread</span><span class="p">(</span><span class="n">group</span><span class="o">=</span><span class="kc">None</span><span class="p">,</span>
|
||
<span class="n">target</span><span class="o">=</span><span class="bp">self</span><span class="o">.</span><span class="n">_abort_on_timeout</span><span class="p">,</span>
|
||
<span class="n">name</span><span class="o">=</span><span class="s2">"MpiSessionTimeoutKiller"</span><span class="p">,</span>
|
||
<span class="n">args</span><span class="o">=</span><span class="p">(</span><span class="n">fut</span><span class="p">,</span> <span class="n">grace</span><span class="p">,</span> <span class="n">reason</span><span class="p">))</span>
|
||
<span class="n">killer</span><span class="o">.</span><span class="n">start</span><span class="p">()</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">shutdown</span><span class="p">()</span>
|
||
<span class="n">fut</span><span class="o">.</span><span class="n">set_result</span><span class="p">(</span><span class="kc">None</span><span class="p">)</span>
|
||
<span class="n">killer</span><span class="o">.</span><span class="n">join</span><span class="p">()</span>
|
||
|
||
|
||
<span class="k">class</span><span class="w"> </span><span class="nc">MpiPoolSession</span><span class="p">(</span><span class="n">MpiSession</span><span class="p">):</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="fm">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">n_workers</span><span class="p">:</span> <span class="nb">int</span><span class="p">):</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">n_workers</span> <span class="o">=</span> <span class="n">n_workers</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="p">:</span> <span class="n">Optional</span><span class="p">[</span><span class="n">MPIPoolExecutor</span><span class="p">]</span> <span class="o">=</span> <span class="kc">None</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">_start_mpi_pool</span><span class="p">()</span>
|
||
<span class="k">if</span> <span class="n">ENABLE_MULTI_DEVICE</span><span class="p">:</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">comm</span> <span class="o">=</span> <span class="n">mpi4py</span><span class="o">.</span><span class="n">MPI</span><span class="o">.</span><span class="n">COMM_WORLD</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">get_comm</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="k">return</span> <span class="bp">self</span><span class="o">.</span><span class="n">comm</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">submit</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">],</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span>
|
||
<span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="n">List</span><span class="p">[</span><span class="n">Future</span><span class="p">[</span><span class="n">T</span><span class="p">]]:</span>
|
||
<span class="k">return</span> <span class="p">[</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="o">.</span><span class="n">submit</span><span class="p">(</span><span class="n">task</span><span class="p">,</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span>
|
||
<span class="k">for</span> <span class="n">i</span> <span class="ow">in</span> <span class="nb">range</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">n_workers</span><span class="p">)</span>
|
||
<span class="p">]</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">submit_sync</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">],</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="n">List</span><span class="p">[</span><span class="n">T</span><span class="p">]:</span>
|
||
<span class="n">futures</span> <span class="o">=</span> <span class="p">[</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="o">.</span><span class="n">submit</span><span class="p">(</span><span class="n">task</span><span class="p">,</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span>
|
||
<span class="k">for</span> <span class="n">i</span> <span class="ow">in</span> <span class="nb">range</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">n_workers</span><span class="p">)</span>
|
||
<span class="p">]</span>
|
||
<span class="k">return</span> <span class="p">[</span><span class="n">future</span><span class="o">.</span><span class="n">result</span><span class="p">()</span> <span class="k">for</span> <span class="n">future</span> <span class="ow">in</span> <span class="n">futures</span><span class="p">]</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">shutdown</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">wait</span><span class="o">=</span><span class="kc">True</span><span class="p">):</span>
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="o">.</span><span class="n">shutdown</span><span class="p">(</span><span class="n">wait</span><span class="o">=</span><span class="n">wait</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span> <span class="o">=</span> <span class="kc">None</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">abort</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">get_comm</span><span class="p">()</span><span class="o">.</span><span class="n">Abort</span><span class="p">(</span><span class="mi">1</span><span class="p">)</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">_start_mpi_pool</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="k">assert</span> <span class="ow">not</span> <span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="p">,</span> <span class="s1">'MPI session already started'</span>
|
||
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span> <span class="o">=</span> <span class="n">MPIPoolExecutor</span><span class="p">(</span><span class="n">max_workers</span><span class="o">=</span><span class="bp">self</span><span class="o">.</span><span class="n">n_workers</span><span class="p">,</span>
|
||
<span class="n">path</span><span class="o">=</span><span class="n">sys</span><span class="o">.</span><span class="n">path</span><span class="p">)</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="fm">__del__</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">shutdown_abort</span><span class="p">()</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">__reduce__</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="k">raise</span> <span class="ne">TypeError</span><span class="p">(</span><span class="s1">'cannot pickle MPI session'</span><span class="p">)</span>
|
||
|
||
|
||
<div class="viewcode-block" id="MpiCommSession">
|
||
<a class="viewcode-back" href="../../../llm-api/reference.html#tensorrt_llm.llmapi.MpiCommSession">[docs]</a>
|
||
<span class="k">class</span><span class="w"> </span><span class="nc">MpiCommSession</span><span class="p">(</span><span class="n">MpiSession</span><span class="p">):</span>
|
||
|
||
<div class="viewcode-block" id="MpiCommSession.__init__">
|
||
<a class="viewcode-back" href="../../../llm-api/reference.html#tensorrt_llm.llmapi.MpiCommSession.__init__">[docs]</a>
|
||
<span class="k">def</span><span class="w"> </span><span class="fm">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">comm</span><span class="o">=</span><span class="kc">None</span><span class="p">,</span> <span class="n">n_workers</span><span class="p">:</span> <span class="nb">int</span> <span class="o">=</span> <span class="mi">1</span><span class="p">):</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">comm</span> <span class="o">=</span> <span class="n">comm</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">n_workers</span> <span class="o">=</span> <span class="n">n_workers</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">thread_pool</span><span class="p">:</span> <span class="n">Optional</span><span class="p">[</span><span class="n">ThreadPoolExecutor</span><span class="p">]</span> <span class="o">=</span> <span class="kc">None</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="p">:</span> <span class="n">Optional</span><span class="p">[</span><span class="n">MPIPoolExecutor</span><span class="p">]</span> <span class="o">=</span> <span class="kc">None</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">owns_mpi_pool</span> <span class="o">=</span> <span class="kc">False</span> <span class="c1"># Track if this instance owns the mpi_pool</span>
|
||
|
||
<span class="k">if</span> <span class="n">n_workers</span> <span class="o"><=</span> <span class="mi">0</span><span class="p">:</span>
|
||
<span class="k">raise</span> <span class="ne">ValueError</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s1">'n_workers must be non-negative, but got </span><span class="si">{</span><span class="n">n_workers</span><span class="si">}</span><span class="s1">'</span><span class="p">)</span>
|
||
|
||
<span class="k">if</span> <span class="n">ENABLE_MULTI_DEVICE</span><span class="p">:</span>
|
||
<span class="k">if</span> <span class="ow">not</span> <span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="p">:</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">comm</span> <span class="o">=</span> <span class="n">mpi4py</span><span class="o">.</span><span class="n">MPI</span><span class="o">.</span><span class="n">COMM_WORLD</span>
|
||
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="o">.</span><span class="n">Get_rank</span><span class="p">()</span> <span class="o">!=</span> <span class="mi">0</span><span class="p">:</span>
|
||
<span class="k">raise</span> <span class="ne">RuntimeError</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s1">'only rank 0 can start multi-node session, got </span><span class="si">{</span><span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="o">.</span><span class="n">Get_rank</span><span class="p">()</span><span class="si">}</span><span class="s1">'</span>
|
||
<span class="p">)</span>
|
||
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="o">.</span><span class="n">Get_size</span><span class="p">()</span> <span class="o">!=</span> <span class="n">n_workers</span><span class="p">:</span>
|
||
<span class="k">raise</span> <span class="ne">ValueError</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s1">'n_workers must be equal to the number of processes in MPI, got </span><span class="si">{</span><span class="n">n_workers</span><span class="si">}</span><span class="s1"> vs </span><span class="si">{</span><span class="n">get_mpi_world_size</span><span class="p">()</span><span class="si">}</span><span class="s1">'</span>
|
||
<span class="p">)</span>
|
||
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">_start_mpi_pool</span><span class="p">()</span></div>
|
||
|
||
|
||
<div class="viewcode-block" id="MpiCommSession.get_comm">
|
||
<a class="viewcode-back" href="../../../llm-api/reference.html#tensorrt_llm.llmapi.MpiCommSession.get_comm">[docs]</a>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">get_comm</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="k">return</span> <span class="bp">self</span><span class="o">.</span><span class="n">comm</span></div>
|
||
|
||
|
||
<div class="viewcode-block" id="MpiCommSession.submit">
|
||
<a class="viewcode-back" href="../../../llm-api/reference.html#tensorrt_llm.llmapi.MpiCommSession.submit">[docs]</a>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">submit</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">],</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span>
|
||
<span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="n">List</span><span class="p">[</span><span class="n">Future</span><span class="p">[</span><span class="n">T</span><span class="p">]]:</span>
|
||
<span class="w"> </span><span class="sd">''' Submit a task to MPI workers.</span>
|
||
|
||
<span class="sd"> Args:</span>
|
||
<span class="sd"> task: The task to be submitted.</span>
|
||
<span class="sd"> args: Positional arguments for the task.</span>
|
||
<span class="sd"> kwargs: Keyword arguments for the task.</span>
|
||
<span class="sd"> '''</span>
|
||
<span class="k">assert</span> <span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">,</span> <span class="s1">'MPI session not started'</span>
|
||
<span class="n">worker_futures</span> <span class="o">=</span> <span class="p">[</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="o">.</span><span class="n">submit</span><span class="p">(</span><span class="n">task</span><span class="p">,</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span>
|
||
<span class="k">for</span> <span class="n">i</span> <span class="ow">in</span> <span class="nb">range</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">n_workers</span> <span class="o">-</span> <span class="mi">1</span><span class="p">)</span>
|
||
<span class="p">]</span>
|
||
|
||
<span class="n">rank0_future</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">thread_pool</span><span class="o">.</span><span class="n">submit</span><span class="p">(</span><span class="n">task</span><span class="p">,</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span>
|
||
<span class="k">return</span> <span class="p">[</span><span class="n">rank0_future</span><span class="p">]</span> <span class="o">+</span> <span class="n">worker_futures</span></div>
|
||
|
||
|
||
<div class="viewcode-block" id="MpiCommSession.submit_sync">
|
||
<a class="viewcode-back" href="../../../llm-api/reference.html#tensorrt_llm.llmapi.MpiCommSession.submit_sync">[docs]</a>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">submit_sync</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">],</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="n">List</span><span class="p">[</span><span class="n">T</span><span class="p">]:</span>
|
||
<span class="n">futures</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">submit</span><span class="p">(</span><span class="n">task</span><span class="p">,</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span>
|
||
<span class="k">return</span> <span class="p">[</span><span class="n">future</span><span class="o">.</span><span class="n">result</span><span class="p">()</span> <span class="k">for</span> <span class="n">future</span> <span class="ow">in</span> <span class="n">futures</span><span class="p">]</span></div>
|
||
|
||
|
||
<div class="viewcode-block" id="MpiCommSession.shutdown">
|
||
<a class="viewcode-back" href="../../../llm-api/reference.html#tensorrt_llm.llmapi.MpiCommSession.shutdown">[docs]</a>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">shutdown</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">wait</span><span class="o">=</span><span class="kc">True</span><span class="p">):</span>
|
||
<span class="c1"># Only shutdown the mpi_pool if this instance created it</span>
|
||
<span class="c1"># For shared global mpi_pool, we don't shut it down</span>
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span> <span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">owns_mpi_pool</span><span class="p">:</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="o">.</span><span class="n">shutdown</span><span class="p">(</span><span class="n">wait</span><span class="o">=</span><span class="n">wait</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span> <span class="o">=</span> <span class="kc">None</span>
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">thread_pool</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">thread_pool</span><span class="o">.</span><span class="n">shutdown</span><span class="p">(</span><span class="n">wait</span><span class="o">=</span><span class="n">wait</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">thread_pool</span> <span class="o">=</span> <span class="kc">None</span></div>
|
||
|
||
|
||
<div class="viewcode-block" id="MpiCommSession.abort">
|
||
<a class="viewcode-back" href="../../../llm-api/reference.html#tensorrt_llm.llmapi.MpiCommSession.abort">[docs]</a>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">abort</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">get_comm</span><span class="p">()</span><span class="o">.</span><span class="n">Abort</span><span class="p">(</span><span class="mi">1</span><span class="p">)</span></div>
|
||
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">_start_mpi_pool</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="k">assert</span> <span class="ow">not</span> <span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span><span class="p">,</span> <span class="s1">'MPI session already started'</span>
|
||
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">thread_pool</span> <span class="o">=</span> <span class="n">ThreadPoolExecutor</span><span class="p">(</span><span class="n">max_workers</span><span class="o">=</span><span class="mi">2</span><span class="p">)</span>
|
||
|
||
<span class="c1"># Use global MPICommExecutor if using COMM_WORLD</span>
|
||
<span class="c1"># This is necessary because MPICommExecutor can only be created once per MPI process</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"_start_mpi_pool: ENABLE_MULTI_DEVICE=</span><span class="si">{</span><span class="n">ENABLE_MULTI_DEVICE</span><span class="si">}</span><span class="s2">, self.comm=</span><span class="si">{</span><span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"grey"</span><span class="p">)</span>
|
||
<span class="k">if</span> <span class="n">ENABLE_MULTI_DEVICE</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"_start_mpi_pool: Checking if self.comm == mpi4py.MPI.COMM_WORLD: </span><span class="si">{</span><span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="w"> </span><span class="o">==</span><span class="w"> </span><span class="n">mpi4py</span><span class="o">.</span><span class="n">MPI</span><span class="o">.</span><span class="n">COMM_WORLD</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"grey"</span><span class="p">)</span>
|
||
<span class="k">if</span> <span class="n">ENABLE_MULTI_DEVICE</span> <span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">comm</span> <span class="o">==</span> <span class="n">mpi4py</span><span class="o">.</span><span class="n">MPI</span><span class="o">.</span><span class="n">COMM_WORLD</span><span class="p">:</span>
|
||
<span class="k">if</span> <span class="n">MPINodeState</span><span class="o">.</span><span class="n">_global_comm_executor</span> <span class="ow">is</span> <span class="kc">None</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span><span class="s2">"Creating global MPICommExecutor for COMM_WORLD</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"yellow"</span><span class="p">)</span>
|
||
<span class="n">MPINodeState</span><span class="o">.</span><span class="n">_global_comm_executor</span> <span class="o">=</span> <span class="n">MPICommExecutor</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="p">)</span>
|
||
<span class="n">MPINodeState</span><span class="o">.</span><span class="n">_global_mpi_pool</span> <span class="o">=</span> <span class="n">MPINodeState</span><span class="o">.</span><span class="n">_global_comm_executor</span><span class="o">.</span><span class="fm">__enter__</span><span class="p">(</span>
|
||
<span class="p">)</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span><span class="s2">"Reusing global MPICommExecutor for COMM_WORLD</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"yellow"</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span> <span class="o">=</span> <span class="n">MPINodeState</span><span class="o">.</span><span class="n">_global_mpi_pool</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">owns_mpi_pool</span> <span class="o">=</span> <span class="kc">False</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"_start_mpi_pool: Creating new MPICommExecutor (not COMM_WORLD or ENABLE_MULTI_DEVICE=False)</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"grey"</span><span class="p">)</span>
|
||
<span class="c1"># For non-COMM_WORLD communicators, create a new executor</span>
|
||
<span class="n">comm_executor</span> <span class="o">=</span> <span class="n">MPICommExecutor</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">mpi_pool</span> <span class="o">=</span> <span class="n">comm_executor</span><span class="o">.</span><span class="fm">__enter__</span><span class="p">()</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">owns_mpi_pool</span> <span class="o">=</span> <span class="kc">True</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="fm">__del__</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">shutdown_abort</span><span class="p">()</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">__reduce__</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="k">raise</span> <span class="ne">TypeError</span><span class="p">(</span><span class="s1">'cannot pickle MPI session'</span><span class="p">)</span></div>
|
||
|
||
|
||
|
||
<span class="k">class</span><span class="w"> </span><span class="nc">RemoteTask</span><span class="p">(</span><span class="n">NamedTuple</span><span class="p">):</span>
|
||
<span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">]</span>
|
||
<span class="n">args</span><span class="p">:</span> <span class="n">Tuple</span><span class="p">[</span><span class="n">Any</span><span class="p">,</span> <span class="o">...</span><span class="p">]</span>
|
||
<span class="n">kwargs</span><span class="p">:</span> <span class="n">Dict</span><span class="p">[</span><span class="nb">str</span><span class="p">,</span> <span class="n">Any</span><span class="p">]</span>
|
||
<span class="n">sync</span><span class="p">:</span> <span class="nb">bool</span> <span class="o">=</span> <span class="kc">False</span> <span class="c1"># if True, the result will be sent back to the client</span>
|
||
|
||
|
||
<span class="k">class</span><span class="w"> </span><span class="nc">RemoteMpiCommSessionClient</span><span class="p">(</span><span class="n">MpiSession</span><span class="p">):</span>
|
||
<span class="w"> </span><span class="sd">'''</span>
|
||
<span class="sd"> RemoteMpiCommSessionClient is a variant of MpiCommSession that is used to connect to a remote MPI pool.</span>
|
||
|
||
<span class="sd"> Note: This class uses a global singleton pattern because ZeroMQ PAIR sockets only support</span>
|
||
<span class="sd"> one connection at a time. Multiple LLM instances will reuse the same client connection.</span>
|
||
<span class="sd"> '''</span>
|
||
<span class="n">_global_instance</span> <span class="o">=</span> <span class="kc">None</span>
|
||
<span class="n">_global_instance_lock</span> <span class="o">=</span> <span class="n">threading</span><span class="o">.</span><span class="n">Lock</span><span class="p">()</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="fm">__new__</span><span class="p">(</span><span class="bp">cls</span><span class="p">,</span> <span class="n">addr</span><span class="p">:</span> <span class="nb">str</span><span class="p">,</span> <span class="n">hmac_key</span><span class="p">:</span> <span class="n">Optional</span><span class="p">[</span><span class="nb">bytes</span><span class="p">]</span> <span class="o">=</span> <span class="kc">None</span><span class="p">):</span>
|
||
<span class="c1"># Implement singleton pattern to reuse the same client connection</span>
|
||
<span class="c1"># for multiple LLM instances, since PAIR sockets only support one connection</span>
|
||
<span class="k">with</span> <span class="bp">cls</span><span class="o">.</span><span class="n">_global_instance_lock</span><span class="p">:</span>
|
||
<span class="k">if</span> <span class="bp">cls</span><span class="o">.</span><span class="n">_global_instance</span> <span class="ow">is</span> <span class="kc">None</span> <span class="ow">or</span> <span class="bp">cls</span><span class="o">.</span><span class="n">_global_instance</span><span class="o">.</span><span class="n">addr</span> <span class="o">!=</span> <span class="n">addr</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"Creating new global RemoteMpiCommSessionClient for </span><span class="si">{</span><span class="n">addr</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"yellow"</span><span class="p">)</span>
|
||
<span class="n">instance</span> <span class="o">=</span> <span class="nb">super</span><span class="p">()</span><span class="o">.</span><span class="fm">__new__</span><span class="p">(</span><span class="bp">cls</span><span class="p">)</span>
|
||
<span class="bp">cls</span><span class="o">.</span><span class="n">_global_instance</span> <span class="o">=</span> <span class="n">instance</span>
|
||
<span class="n">instance</span><span class="o">.</span><span class="n">_initialized</span> <span class="o">=</span> <span class="kc">False</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"Reusing existing global RemoteMpiCommSessionClient for </span><span class="si">{</span><span class="n">addr</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"yellow"</span><span class="p">)</span>
|
||
<span class="k">return</span> <span class="bp">cls</span><span class="o">.</span><span class="n">_global_instance</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="fm">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">addr</span><span class="p">:</span> <span class="nb">str</span><span class="p">,</span> <span class="n">hmac_key</span><span class="p">:</span> <span class="n">Optional</span><span class="p">[</span><span class="nb">bytes</span><span class="p">]</span> <span class="o">=</span> <span class="kc">None</span><span class="p">):</span>
|
||
<span class="c1"># Only initialize once</span>
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">_initialized</span><span class="p">:</span>
|
||
<span class="k">return</span>
|
||
|
||
<span class="c1"># FIXME: this is a hack to avoid circular import, resolve later</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">tensorrt_llm.executor.ipc</span><span class="w"> </span><span class="kn">import</span> <span class="n">ZeroMqQueue</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">addr</span> <span class="o">=</span> <span class="n">addr</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span><span class="sa">f</span><span class="s2">"RemoteMpiCommSessionClient connecting to </span><span class="si">{</span><span class="n">addr</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"yellow"</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">queue</span> <span class="o">=</span> <span class="n">ZeroMqQueue</span><span class="p">((</span><span class="n">addr</span><span class="p">,</span> <span class="n">hmac_key</span><span class="p">),</span>
|
||
<span class="n">is_server</span><span class="o">=</span><span class="kc">False</span><span class="p">,</span>
|
||
<span class="n">socket_type</span><span class="o">=</span><span class="n">zmq</span><span class="o">.</span><span class="n">PAIR</span><span class="p">,</span>
|
||
<span class="n">use_hmac_encryption</span><span class="o">=</span><span class="nb">bool</span><span class="p">(</span><span class="n">hmac_key</span><span class="p">))</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">_is_shutdown</span> <span class="o">=</span> <span class="kc">False</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">_initialized</span> <span class="o">=</span> <span class="kc">True</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">submit</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span>
|
||
<span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">],</span>
|
||
<span class="o">*</span><span class="n">args</span><span class="p">,</span>
|
||
<span class="n">sync</span><span class="p">:</span> <span class="nb">bool</span> <span class="o">=</span> <span class="kc">False</span><span class="p">,</span>
|
||
<span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="nb">list</span><span class="p">:</span>
|
||
<span class="w"> </span><span class="sd">''' Submit a task to the remote MPI pool. '''</span>
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">_is_shutdown</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span><span class="s2">"RemoteMpiCommSessionClient is already shut down</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"yellow"</span><span class="p">)</span>
|
||
<span class="k">return</span> <span class="p">[]</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"RemoteMpiCommSessionClient [rank</span><span class="si">{</span><span class="n">global_mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2">] sending task </span><span class="si">{</span><span class="n">task</span><span class="si">}</span><span class="s2"> to </span><span class="si">{</span><span class="bp">self</span><span class="o">.</span><span class="n">addr</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"yellow"</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">queue</span><span class="o">.</span><span class="n">put</span><span class="p">(</span><span class="n">RemoteTask</span><span class="p">(</span><span class="n">task</span><span class="p">,</span> <span class="n">args</span><span class="p">,</span> <span class="n">kwargs</span><span class="p">,</span> <span class="n">sync</span><span class="o">=</span><span class="n">sync</span><span class="p">))</span>
|
||
<span class="k">return</span> <span class="p">[]</span>
|
||
|
||
<span class="n">SYNC_IDLE_INTERVAL</span> <span class="o">=</span> <span class="mi">8</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">submit_sync</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">task</span><span class="p">,</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="n">List</span><span class="p">[</span><span class="n">T</span><span class="p">]:</span>
|
||
<span class="w"> </span><span class="sd">''' Submit a task to the remote MPI pool and wait for task completion. '''</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">submit</span><span class="p">(</span><span class="n">task</span><span class="p">,</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="n">sync</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span>
|
||
|
||
<span class="k">while</span> <span class="ow">not</span> <span class="p">((</span><span class="n">res</span> <span class="o">:=</span> <span class="bp">self</span><span class="o">.</span><span class="n">poll</span><span class="p">())</span> <span class="ow">or</span> <span class="bp">self</span><span class="o">.</span><span class="n">_is_shutdown</span><span class="p">):</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span><span class="sa">f</span><span class="s2">"Waiting for task completion... </span><span class="si">{</span><span class="n">res</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span> <span class="s2">"grey"</span><span class="p">)</span>
|
||
<span class="n">time</span><span class="o">.</span><span class="n">sleep</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">SYNC_IDLE_INTERVAL</span><span class="p">)</span>
|
||
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"rank</span><span class="si">{</span><span class="n">global_mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2"> RemoteMpiCommSessionClient.send_sync received results: </span><span class="si">{</span><span class="n">res</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"green"</span><span class="p">)</span>
|
||
|
||
<span class="k">if</span> <span class="ow">not</span> <span class="n">res</span><span class="p">:</span>
|
||
<span class="k">raise</span> <span class="ne">RuntimeError</span><span class="p">(</span>
|
||
<span class="s2">"RemoteMpiCommSessionClient received unexpected response"</span><span class="p">)</span>
|
||
<span class="k">return</span> <span class="n">res</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">poll</span><span class="p">(</span><span class="bp">self</span><span class="p">)</span> <span class="o">-></span> <span class="nb">bool</span><span class="p">:</span>
|
||
<span class="w"> </span><span class="sd">''' Poll the queue for a response.</span>
|
||
<span class="sd"> Returns:</span>
|
||
<span class="sd"> True if a response is received, False otherwise.</span>
|
||
<span class="sd"> '''</span>
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">_is_shutdown</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="kc">False</span>
|
||
<span class="n">response</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">queue</span><span class="o">.</span><span class="n">poll</span><span class="p">(</span><span class="mf">0.1</span><span class="p">)</span>
|
||
<span class="k">if</span> <span class="n">response</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="bp">self</span><span class="o">.</span><span class="n">queue</span><span class="o">.</span><span class="n">get</span><span class="p">()</span> <span class="c1"># should get a True if success</span>
|
||
<span class="k">return</span> <span class="kc">False</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">abort</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">shutdown</span><span class="p">()</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">shutdown</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">wait</span><span class="o">=</span><span class="kc">True</span><span class="p">):</span>
|
||
<span class="c1"># NOTE: We do NOT close the queue or mark as shutdown for the singleton instance.</span>
|
||
<span class="c1"># The RemoteMpiCommSessionClient is a global singleton that's reused across multiple</span>
|
||
<span class="c1"># LLM instances. Marking it as shutdown would prevent subsequent LLM instances from</span>
|
||
<span class="c1"># using it. The connection stays open for the entire lifetime of the mgmn setup.</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"RemoteMpiCommSessionClient.shutdown() called (no-op for singleton)</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"grey"</span><span class="p">)</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">shutdown_abort</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">grace</span><span class="p">:</span> <span class="nb">float</span> <span class="o">=</span> <span class="mi">60</span><span class="p">,</span> <span class="n">reason</span><span class="o">=</span><span class="kc">None</span><span class="p">):</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">shutdown</span><span class="p">()</span>
|
||
|
||
|
||
<span class="k">class</span><span class="w"> </span><span class="nc">RemoteMpiCommSessionServer</span><span class="p">():</span>
|
||
<span class="w"> </span><span class="sd">'''</span>
|
||
<span class="sd"> RemoteMpiCommSessionServer is a variant of MpiCommSession that is used to create a remote MPI pool.</span>
|
||
<span class="sd"> '''</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="fm">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span>
|
||
<span class="n">n_workers</span><span class="p">:</span> <span class="nb">int</span> <span class="o">=</span> <span class="mi">0</span><span class="p">,</span>
|
||
<span class="n">addr</span><span class="p">:</span> <span class="nb">str</span> <span class="o">=</span> <span class="sa">f</span><span class="s1">'tcp://127.0.0.1:*'</span><span class="p">,</span>
|
||
<span class="n">hmac_key</span><span class="p">:</span> <span class="n">Optional</span><span class="p">[</span><span class="nb">bytes</span><span class="p">]</span> <span class="o">=</span> <span class="kc">None</span><span class="p">,</span>
|
||
<span class="n">comm</span><span class="o">=</span><span class="kc">None</span><span class="p">,</span>
|
||
<span class="n">is_comm</span><span class="p">:</span> <span class="nb">bool</span> <span class="o">=</span> <span class="kc">False</span><span class="p">):</span>
|
||
<span class="c1"># FIXME: this is a hack to avoid circular import, resolve later</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">tensorrt_llm.executor.ipc</span><span class="w"> </span><span class="kn">import</span> <span class="n">ZeroMqQueue</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">addr</span> <span class="o">=</span> <span class="n">addr</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">queue</span> <span class="o">=</span> <span class="n">ZeroMqQueue</span><span class="p">((</span><span class="n">addr</span><span class="p">,</span> <span class="n">hmac_key</span><span class="p">),</span>
|
||
<span class="n">is_server</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span>
|
||
<span class="n">socket_type</span><span class="o">=</span><span class="n">zmq</span><span class="o">.</span><span class="n">PAIR</span><span class="p">,</span>
|
||
<span class="n">use_hmac_encryption</span><span class="o">=</span><span class="nb">bool</span><span class="p">(</span><span class="n">hmac_key</span><span class="p">))</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">comm</span> <span class="o">=</span> <span class="n">comm</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">results</span> <span class="o">=</span> <span class="p">[]</span> <span class="c1"># the results may arrive in any order</span>
|
||
|
||
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">comm</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">session</span> <span class="o">=</span> <span class="n">MpiCommSession</span><span class="p">(</span><span class="n">n_workers</span><span class="o">=</span><span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="o">.</span><span class="n">Get_size</span><span class="p">(),</span>
|
||
<span class="n">comm</span><span class="o">=</span><span class="bp">self</span><span class="o">.</span><span class="n">comm</span><span class="p">)</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">session</span> <span class="o">=</span> <span class="n">MpiCommSession</span><span class="p">(</span>
|
||
<span class="n">n_workers</span><span class="o">=</span><span class="n">n_workers</span><span class="p">)</span> <span class="k">if</span> <span class="n">is_comm</span> <span class="k">else</span> <span class="n">MpiPoolSession</span><span class="p">(</span>
|
||
<span class="n">n_workers</span><span class="o">=</span><span class="n">n_workers</span><span class="p">)</span>
|
||
|
||
<span class="nd">@staticmethod</span>
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">task_wrapper</span><span class="p">(</span><span class="n">task</span><span class="p">:</span> <span class="n">Callable</span><span class="p">[</span><span class="o">...</span><span class="p">,</span> <span class="n">T</span><span class="p">],</span> <span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span> <span class="o">-></span> <span class="n">T</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"MpiCommSession rank</span><span class="si">{</span><span class="n">mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2"> with world_size </span><span class="si">{</span><span class="n">mpi_world_size</span><span class="p">()</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"green"</span><span class="p">)</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"MpiCommSession rank</span><span class="si">{</span><span class="n">mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2"> start task [</span><span class="si">{</span><span class="n">task</span><span class="si">}</span><span class="s2">] with args: </span><span class="si">{</span><span class="n">args</span><span class="si">}</span><span class="s2"> and kwargs: </span><span class="si">{</span><span class="n">kwargs</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"green"</span><span class="p">)</span>
|
||
|
||
<span class="c1"># wait for all ranks to start the task</span>
|
||
<span class="n">mpi_barrier</span><span class="p">()</span>
|
||
|
||
<span class="k">try</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="n">task</span><span class="p">(</span><span class="o">*</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">kwargs</span><span class="p">)</span>
|
||
<span class="k">except</span> <span class="ne">Exception</span> <span class="k">as</span> <span class="n">e</span><span class="p">:</span>
|
||
<span class="n">print_colored</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"MpiCommSession rank</span><span class="si">{</span><span class="n">mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2"> task [</span><span class="si">{</span><span class="n">task</span><span class="si">}</span><span class="s2">] failed with exception: </span><span class="si">{</span><span class="n">e</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"red"</span><span class="p">)</span>
|
||
<span class="n">traceback</span><span class="o">.</span><span class="n">print_exc</span><span class="p">()</span>
|
||
<span class="k">raise</span> <span class="n">e</span>
|
||
<span class="k">finally</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"MpiCommSession rank</span><span class="si">{</span><span class="n">mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2"> task [</span><span class="si">{</span><span class="n">task</span><span class="si">}</span><span class="s2">] finished</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"green"</span><span class="p">)</span>
|
||
<span class="n">mpi_barrier</span><span class="p">()</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">serve</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span><span class="sa">f</span><span class="s2">"RemoteMpiCommSessionServer listening on </span><span class="si">{</span><span class="bp">self</span><span class="o">.</span><span class="n">addr</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"yellow"</span><span class="p">)</span>
|
||
<span class="n">pending_futures</span> <span class="o">=</span> <span class="p">[]</span>
|
||
<span class="k">while</span> <span class="kc">True</span><span class="p">:</span>
|
||
<span class="c1"># Wait for any pending futures from previous tasks to complete</span>
|
||
<span class="c1"># This ensures all ranks are ready before accepting the next task</span>
|
||
<span class="k">if</span> <span class="n">pending_futures</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"RemoteMpiCommSessionServer waiting for </span><span class="si">{</span><span class="nb">len</span><span class="p">(</span><span class="n">pending_futures</span><span class="p">)</span><span class="si">}</span><span class="s2"> pending futures to complete</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"grey"</span><span class="p">)</span>
|
||
<span class="k">for</span> <span class="n">future</span> <span class="ow">in</span> <span class="n">pending_futures</span><span class="p">:</span>
|
||
<span class="k">try</span><span class="p">:</span>
|
||
<span class="n">future</span><span class="o">.</span><span class="n">result</span><span class="p">()</span> <span class="c1"># Wait for completion</span>
|
||
<span class="k">except</span> <span class="ne">Exception</span> <span class="k">as</span> <span class="n">e</span><span class="p">:</span>
|
||
<span class="n">print_colored</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"RemoteMpiCommSessionServer future failed with exception: </span><span class="si">{</span><span class="n">e</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"red"</span><span class="p">)</span>
|
||
<span class="n">pending_futures</span><span class="o">.</span><span class="n">clear</span><span class="p">()</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="s2">"RemoteMpiCommSessionServer all pending futures completed</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"grey"</span><span class="p">)</span>
|
||
|
||
<span class="n">message</span><span class="p">:</span> <span class="n">Optional</span><span class="p">[</span><span class="n">RemoteTask</span><span class="p">]</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">queue</span><span class="o">.</span><span class="n">get</span><span class="p">()</span>
|
||
<span class="k">if</span> <span class="n">message</span> <span class="ow">is</span> <span class="kc">None</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"RemoteMpiCommSessionServer [rank</span><span class="si">{</span><span class="n">global_mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2">] received shutdown signal</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"green"</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">session</span><span class="o">.</span><span class="n">shutdown_abort</span><span class="p">()</span>
|
||
<span class="k">break</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"RemoteMpiCommSessionServer [rank</span><span class="si">{</span><span class="n">global_mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2">] received task [</span><span class="si">{</span><span class="n">message</span><span class="o">.</span><span class="n">task</span><span class="si">}</span><span class="s2">] from </span><span class="si">{</span><span class="bp">self</span><span class="o">.</span><span class="n">addr</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"green"</span><span class="p">)</span>
|
||
<span class="n">futures</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">session</span><span class="o">.</span><span class="n">submit</span><span class="p">(</span>
|
||
<span class="n">RemoteMpiCommSessionServer</span><span class="o">.</span><span class="n">task_wrapper</span><span class="p">,</span> <span class="n">message</span><span class="o">.</span><span class="n">task</span><span class="p">,</span>
|
||
<span class="o">*</span><span class="n">message</span><span class="o">.</span><span class="n">args</span><span class="p">,</span> <span class="o">**</span><span class="n">message</span><span class="o">.</span><span class="n">kwargs</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">num_results</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">session</span><span class="o">.</span><span class="n">n_workers</span>
|
||
<span class="k">assert</span> <span class="nb">len</span><span class="p">(</span><span class="n">futures</span><span class="p">)</span> <span class="o">==</span> <span class="bp">self</span><span class="o">.</span><span class="n">num_results</span> <span class="o">==</span> <span class="n">mpi_world_size</span><span class="p">()</span>
|
||
<span class="c1"># Store futures to wait for them before the next task</span>
|
||
<span class="n">pending_futures</span> <span class="o">=</span> <span class="nb">list</span><span class="p">(</span><span class="n">futures</span><span class="p">)</span>
|
||
<span class="k">if</span> <span class="n">message</span><span class="o">.</span><span class="n">sync</span><span class="p">:</span>
|
||
<span class="k">for</span> <span class="n">future</span> <span class="ow">in</span> <span class="n">futures</span><span class="p">:</span>
|
||
<span class="n">future</span><span class="o">.</span><span class="n">add_done_callback</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">mpi_future_callback</span><span class="p">)</span>
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">mpi_future_callback</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">future</span><span class="p">):</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span><span class="sa">f</span><span class="s2">"rank</span><span class="si">{</span><span class="n">global_mpi_rank</span><span class="p">()</span><span class="si">}</span><span class="s2"> got future: </span><span class="si">{</span><span class="n">future</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span> <span class="s2">"red"</span><span class="p">)</span>
|
||
<span class="k">if</span> <span class="n">future</span><span class="o">.</span><span class="n">exception</span><span class="p">()</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"mpi_future got exception: </span><span class="si">{</span><span class="n">future</span><span class="o">.</span><span class="n">exception</span><span class="p">()</span><span class="si">}</span><span class="s2">, quitting</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"red"</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">queue</span><span class="o">.</span><span class="n">put</span><span class="p">(</span><span class="n">future</span><span class="o">.</span><span class="n">exception</span><span class="p">())</span>
|
||
<span class="k">return</span>
|
||
|
||
<span class="n">result</span> <span class="o">=</span> <span class="n">future</span><span class="o">.</span><span class="n">result</span><span class="p">()</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">results</span><span class="o">.</span><span class="n">append</span><span class="p">(</span><span class="n">result</span><span class="p">)</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"RemoteMpiCommSessionServer working status: </span><span class="si">{</span><span class="nb">len</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">results</span><span class="p">)</span><span class="si">}</span><span class="s2">/</span><span class="si">{</span><span class="bp">self</span><span class="o">.</span><span class="n">num_results</span><span class="si">}</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"grey"</span><span class="p">)</span>
|
||
<span class="k">if</span> <span class="nb">len</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">results</span><span class="p">)</span> <span class="o">==</span> <span class="bp">self</span><span class="o">.</span><span class="n">num_results</span><span class="p">:</span>
|
||
<span class="n">logger_debug</span><span class="p">(</span>
|
||
<span class="sa">f</span><span class="s2">"RemoteMpiCommSessionServer received all results, sending to client</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"green"</span><span class="p">)</span>
|
||
<span class="k">try</span><span class="p">:</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">queue</span><span class="o">.</span><span class="n">put_noblock</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">results</span><span class="p">,</span> <span class="n">retry</span><span class="o">=</span><span class="mi">2</span><span class="p">)</span>
|
||
<span class="k">except</span> <span class="n">zmq</span><span class="o">.</span><span class="n">ZMQError</span> <span class="k">as</span> <span class="n">e</span><span class="p">:</span>
|
||
<span class="c1"># The client could be shutdown first.</span>
|
||
<span class="k">if</span> <span class="n">e</span><span class="o">.</span><span class="n">errno</span> <span class="o">==</span> <span class="n">zmq</span><span class="o">.</span><span class="n">EAGAIN</span><span class="p">:</span>
|
||
<span class="k">pass</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="k">raise</span> <span class="n">e</span>
|
||
|
||
<span class="n">logger_debug</span><span class="p">(</span><span class="sa">f</span><span class="s2">"RemoteMpiCommSessionServer sent results to client</span><span class="se">\n</span><span class="s2">"</span><span class="p">,</span>
|
||
<span class="s2">"green"</span><span class="p">)</span>
|
||
<span class="bp">self</span><span class="o">.</span><span class="n">results</span><span class="o">.</span><span class="n">clear</span><span class="p">()</span>
|
||
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">find_free_port</span><span class="p">()</span> <span class="o">-></span> <span class="nb">int</span><span class="p">:</span>
|
||
<span class="k">with</span> <span class="n">socket</span><span class="o">.</span><span class="n">socket</span><span class="p">(</span><span class="n">socket</span><span class="o">.</span><span class="n">AF_INET</span><span class="p">,</span> <span class="n">socket</span><span class="o">.</span><span class="n">SOCK_STREAM</span><span class="p">)</span> <span class="k">as</span> <span class="n">s</span><span class="p">:</span>
|
||
<span class="n">s</span><span class="o">.</span><span class="n">bind</span><span class="p">((</span><span class="s1">''</span><span class="p">,</span> <span class="mi">0</span><span class="p">))</span>
|
||
<span class="k">return</span> <span class="n">s</span><span class="o">.</span><span class="n">getsockname</span><span class="p">()[</span><span class="mi">1</span><span class="p">]</span>
|
||
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">get_mpi_world_size</span><span class="p">()</span> <span class="o">-></span> <span class="nb">int</span><span class="p">:</span>
|
||
<span class="c1"># avoid cyclic import</span>
|
||
<span class="kn">from</span><span class="w"> </span><span class="nn">..executor.utils</span><span class="w"> </span><span class="kn">import</span> <span class="n">get_spawn_proxy_process_env</span>
|
||
|
||
<span class="c1"># If the proxy process is spawned, the MPI-related env will be cleaned in the proxy process, thus we made another env for the mpi_world_size</span>
|
||
<span class="k">if</span> <span class="n">get_spawn_proxy_process_env</span><span class="p">():</span>
|
||
<span class="k">return</span> <span class="nb">int</span><span class="p">(</span><span class="n">os</span><span class="o">.</span><span class="n">getenv</span><span class="p">(</span><span class="s2">"tllm_mpi_size"</span><span class="p">)</span> <span class="ow">or</span> <span class="mi">1</span><span class="p">)</span>
|
||
<span class="k">else</span><span class="p">:</span>
|
||
<span class="k">return</span> <span class="n">mpi_world_size</span><span class="p">()</span>
|
||
|
||
|
||
<span class="k">def</span><span class="w"> </span><span class="nf">split_mpi_env</span><span class="p">(</span><span class="n">mpi_env_keys</span><span class="p">:</span> <span class="n">List</span><span class="p">[</span><span class="nb">str</span><span class="p">]</span> <span class="o">|</span> <span class="kc">None</span> <span class="o">=</span> <span class="kc">None</span><span class="p">)</span> <span class="o">-></span> <span class="n">Tuple</span><span class="p">[</span><span class="nb">dict</span><span class="p">,</span> <span class="nb">dict</span><span class="p">]:</span>
|
||
<span class="w"> </span><span class="sd">'''</span>
|
||
<span class="sd"> Splits the environment variables into MPI-related and non-MPI-related dictionaries.</span>
|
||
|
||
<span class="sd"> Args:</span>
|
||
<span class="sd"> mpi_env_keys: Additional environment variables to be considered as MPI-related.</span>
|
||
|
||
<span class="sd"> Returns:</span>
|
||
<span class="sd"> Tuple[dict, dict]: (non_mpi_env, mpi_env)</span>
|
||
<span class="sd"> - non_mpi_env: Environment dictionary without MPI-related variables</span>
|
||
<span class="sd"> - mpi_env: Environment dictionary containing only MPI-related variables</span>
|
||
<span class="sd"> '''</span>
|
||
<span class="n">current_env</span> <span class="o">=</span> <span class="n">os</span><span class="o">.</span><span class="n">environ</span><span class="o">.</span><span class="n">copy</span><span class="p">()</span>
|
||
|
||
<span class="c1"># Identify MPI-related variables</span>
|
||
<span class="n">mpi_vars</span> <span class="o">=</span> <span class="nb">set</span><span class="p">(</span>
|
||
<span class="n">itertools</span><span class="o">.</span><span class="n">chain</span><span class="p">([</span>
|
||
<span class="n">var</span> <span class="k">for</span> <span class="n">var</span> <span class="ow">in</span> <span class="n">current_env</span> <span class="k">if</span> <span class="n">var</span><span class="o">.</span><span class="n">startswith</span><span class="p">((</span>
|
||
<span class="s1">'MPI_'</span><span class="p">,</span>
|
||
<span class="s1">'OMPI_'</span><span class="p">,</span>
|
||
<span class="s1">'PMIX_'</span><span class="p">,</span>
|
||
<span class="s1">'PMI_'</span><span class="p">,</span>
|
||
<span class="s1">'OMPI_'</span><span class="p">,</span>
|
||
<span class="s1">'PMIX_'</span><span class="p">,</span>
|
||
<span class="s1">'PMI_'</span><span class="p">,</span>
|
||
<span class="s1">'SLURM_'</span><span class="p">,</span>
|
||
<span class="s1">'MPI_'</span><span class="p">,</span>
|
||
<span class="s1">'UCX_'</span><span class="p">,</span>
|
||
<span class="s1">'I_MPI_'</span><span class="p">,</span>
|
||
<span class="s1">'HYDRA_'</span><span class="p">,</span>
|
||
<span class="s1">'KMP_'</span><span class="p">,</span>
|
||
<span class="s1">'MPICH_'</span><span class="p">,</span>
|
||
<span class="s1">'MV2_'</span><span class="p">,</span>
|
||
<span class="s1">'CRAY_'</span><span class="p">,</span>
|
||
<span class="p">))</span>
|
||
<span class="p">],</span> <span class="n">mpi_env_keys</span> <span class="ow">or</span> <span class="p">[]))</span>
|
||
|
||
<span class="c1"># Split into two dictionaries</span>
|
||
<span class="n">non_mpi_env</span> <span class="o">=</span> <span class="p">{</span><span class="n">k</span><span class="p">:</span> <span class="n">v</span> <span class="k">for</span> <span class="n">k</span><span class="p">,</span> <span class="n">v</span> <span class="ow">in</span> <span class="n">current_env</span><span class="o">.</span><span class="n">items</span><span class="p">()</span> <span class="k">if</span> <span class="n">k</span> <span class="ow">not</span> <span class="ow">in</span> <span class="n">mpi_vars</span><span class="p">}</span>
|
||
<span class="n">mpi_env</span> <span class="o">=</span> <span class="p">{</span><span class="n">k</span><span class="p">:</span> <span class="n">v</span> <span class="k">for</span> <span class="n">k</span><span class="p">,</span> <span class="n">v</span> <span class="ow">in</span> <span class="n">current_env</span><span class="o">.</span><span class="n">items</span><span class="p">()</span> <span class="k">if</span> <span class="n">k</span> <span class="ow">in</span> <span class="n">mpi_vars</span><span class="p">}</span>
|
||
|
||
<span class="k">return</span> <span class="n">non_mpi_env</span><span class="p">,</span> <span class="n">mpi_env</span>
|
||
</pre></div>
|
||
|
||
</article>
|
||
|
||
|
||
|
||
|
||
|
||
<footer class="prev-next-footer d-print-none">
|
||
|
||
<div class="prev-next-area">
|
||
</div>
|
||
</footer>
|
||
|
||
</div>
|
||
|
||
|
||
|
||
<div class="bd-sidebar-secondary"></div>
|
||
|
||
|
||
|
||
|
||
|
||
</div>
|
||
<footer class="bd-footer-content">
|
||
|
||
</footer>
|
||
|
||
</main>
|
||
</div>
|
||
</div>
|
||
|
||
|
||
<!-- Scripts loaded after <body> so the DOM is not blocked -->
|
||
<script defer src="../../../_static/scripts/bootstrap.js?digest=8878045cc6db502f8baf"></script>
|
||
<script defer src="../../../_static/scripts/pydata-sphinx-theme.js?digest=8878045cc6db502f8baf"></script>
|
||
|
||
|
||
<footer class="bd-footer">
|
||
<div class="bd-footer__inner bd-page-width">
|
||
|
||
<div class="footer-items__start">
|
||
|
||
<div class="footer-item">
|
||
<a class="footer-brand logo" href="https://www.nvidia.com">
|
||
<img src="../../../_static/nvidia-logo-horiz-rgb-1c-blk-for-screen.svg" class="logo__image only-light" alt="NVIDIA"/>
|
||
<img src="../../../_static/nvidia-logo-horiz-rgb-1c-wht-for-screen.svg" class="logo__image only-dark" alt="NVIDIA"/>
|
||
</a></div>
|
||
|
||
<div class="footer-item">
|
||
|
||
<div class="footer-links">
|
||
|
||
|
||
<a class="external" href="https://www.nvidia.com/en-us/about-nvidia/privacy-policy/">Privacy Policy</a>
|
||
|
|
||
|
||
|
||
|
||
<a class="external" href="https://www.nvidia.com/en-us/about-nvidia/privacy-center/">Your Privacy Choices</a>
|
||
|
|
||
|
||
|
||
|
||
<a class="external" href="https://www.nvidia.com/en-us/about-nvidia/terms-of-service/">Terms of Service</a>
|
||
|
|
||
|
||
|
||
|
||
<a class="external" href="https://www.nvidia.com/en-us/about-nvidia/accessibility/">Accessibility</a>
|
||
|
|
||
|
||
|
||
|
||
<a class="external" href="https://www.nvidia.com/en-us/about-nvidia/company-policies/">Corporate Policies</a>
|
||
|
|
||
|
||
|
||
|
||
<a class="external" href="https://www.nvidia.com/en-us/product-security/">Product Security</a>
|
||
|
|
||
|
||
|
||
|
||
<a class="external" href="https://www.nvidia.com/en-us/contact/">Contact</a>
|
||
|
||
|
||
|
||
</div>
|
||
</div>
|
||
|
||
<div class="footer-item">
|
||
|
||
|
||
|
||
|
||
<p class="copyright">
|
||
|
||
Copyright © 2025, NVidia.
|
||
<br/>
|
||
|
||
</p>
|
||
</div>
|
||
|
||
<div class="footer-item">
|
||
<div class="extra_footer">
|
||
|
||
<p>Last updated on November 23, 2025.</p>
|
||
|
||
<p>This page is generated by TensorRT-LLM commit <a href="https://github.com/NVIDIA/TensorRT-LLM/tree/a761585">a761585</a>.</p>
|
||
|
||
</div></div>
|
||
|
||
</div>
|
||
|
||
|
||
|
||
</div>
|
||
|
||
</footer>
|
||
</body>
|
||
</html> |