Spec-Zone.ru › PyTorch 2.14

Конвейерный параллелизм

Создано: 16 июн. 2025 | Последнее обновление: 24 июл. 2026

Примечание

torch.distributed.pipelining в настоящее время находится в альфа-версии и активно разрабатывается. Возможны изменения API. Пакет перенесён из проекта PiPPy.

Зачем нужен конвейерный параллелизм?

Конвейерный параллелизм — одна из базовых форм параллелизма в глубоком обучении. Он позволяет разделить выполнение модели таким образом, чтобы несколько микропакетов могли одновременно выполнять разные части кода модели. Конвейерный параллелизм может быть эффективным методом для следующих задач:

  • обучение крупномасштабных моделей
  • кластеры с ограниченной пропускной способностью
  • инференс больших моделей

Перечисленные выше сценарии объединяет то, что вычисления на одном устройстве не позволяют скрыть обмен данными, необходимый для традиционного параллелизма, например сбор всех весов в FSDP.

Что такое torch.distributed.pipelining?

Несмотря на перспективность конвейерного параллелизма для масштабирования, его часто сложно реализовать, поскольку необходимо разделить выполнение модели, а не только её веса. Разделение выполнения часто требует существенных изменений кода модели. Дополнительная сложность связана с планированием микропакетов в распределённой среде с учётом зависимостей потока данных.

Пакет pipelining предоставляет инструментарий, который выполняет все перечисленные действия автоматически, упрощая реализацию конвейерного параллелизма для произвольных моделей.

Он состоит из двух частей: интерфейса разделения и распределённой среды выполнения. Интерфейс разделения принимает код модели без изменений, разделяет его на «части модели» и фиксирует связи потока данных. Распределённая среда выполнения параллельно запускает этапы конвейера на разных устройствах и обрабатывает такие задачи, как разделение на микропакеты, планирование, обмен данными и распространение градиентов.

В целом пакет pipelining предоставляет следующие возможности:

  • Разделение кода модели на основе простой спецификации.
  • Широкая поддержка расписаний конвейера, включая GPipe, 1F1B, Interleaved 1F1B и Looped BFS, а также инфраструктура для создания пользовательских расписаний.
  • Полноценная поддержка конвейерного параллелизма между узлами, поскольку именно в таких условиях обычно используется PP (при более медленных межсоединениях).
  • Совместимость с другими методами параллелизма в PyTorch, например с параллелизмом данных (DDP, FSDP) или тензорным параллелизмом. В проекте TorchTitan демонстрируется применение «трёхмерного параллелизма» к модели Llama.

Шаг 1. Создание PipelineStage

Прежде чем использовать PipelineSchedule, необходимо создать объекты PipelineStage, оборачивающие часть модели, выполняемую на данном этапе. PipelineStage отвечает за выделение буферов обмена данными и создание операций отправки/получения для взаимодействия с соседними этапами. Он управляет промежуточными буферами, например для результатов прямого прохода, которые ещё не были использованы, и предоставляет функцию для выполнения обратного прохода модели этапа.

Объекту PipelineStage необходимо знать формы входных и выходных данных модели этапа, чтобы правильно выделить буферы обмена данными. Формы должны быть статическими: во время выполнения они не могут меняться от шага к шагу. Если формы во время выполнения не совпадут с ожидаемыми, будет вызвано исключение PipeliningShapeError. При совместном использовании с другими видами параллелизма или применении смешанной точности необходимо учитывать эти методы, чтобы PipelineStage знал правильную форму (и тип данных) выходных данных модуля этапа во время выполнения.

Пользователи могут создать экземпляр PipelineStage напрямую, передав nn.Module, представляющий ту часть модели, которая должна выполняться на этом этапе. Для этого может потребоваться изменить исходный код модели. См. пример в разделе Вариант 1. Разделение модели вручную.

Другой вариант: интерфейс разделения может автоматически разделить модель на последовательность объектов nn.Module с помощью разбиения графа. Для этого модель должна трассироваться с помощью torch.Export. Совместимость полученных объектов nn.Module с другими методами параллелизма пока носит экспериментальный характер и может потребовать обходных решений. Этот интерфейс может быть удобнее, если пользователь не может легко изменить код модели. Подробнее см. в разделе Вариант 2. Автоматическое разделение модели.

Шаг 2. Использование PipelineSchedule для выполнения

Теперь можно подключить PipelineStage к расписанию конвейера и запустить его с входными данными. Ниже приведён пример для GPipe:

from torch.distributed.pipelining import ScheduleGPipe

# Create a schedule
schedule = ScheduleGPipe(stage, n_microbatches)

# Input data (whole batch)
x = torch.randn(batch_size, in_dim, device=device)

# Run the pipeline with input `x`
# `x` will be divided into microbatches automatically
if rank == 0:
    schedule.step(x)
else:
    output = schedule.step()

Если входной конвейер уже формирует микропакеты, передавайте их напрямую:

arg_mbs = [(x0,), (x1,)]
kwarg_mbs = [{"mask": mask0}, {"mask": mask1}]
target_mbs = [target0, target1]

schedule.step(
    arg_mbs=arg_mbs,
    kwarg_mbs=kwarg_mbs,
    target_mbs=target_mbs,
)

Обратите внимание, что приведённый выше код необходимо запускать на каждом рабочем процессе, поэтому для запуска нескольких процессов мы используем службу запуска:

torchrun --nproc_per_node=2 example.py

Варианты разделения модели

Вариант 1. Разделение модели вручную

При создании PipelineStage напрямую пользователь должен предоставить один экземпляр nn.Module, который содержит нужные nn.Parameters и nn.Buffers и определяет метод forward(), выполняющий операции, относящиеся к данному этапу. Например, сокращённая версия класса Transformer, определённого в Torchtitan, демонстрирует подход к созданию модели, которую легко разделить на части.

class Transformer(nn.Module):
    def __init__(self, model_args: ModelArgs):
        super().__init__()

        self.tok_embeddings = nn.Embedding(...)

        # Using a ModuleDict lets us delete layers without affecting names,
        # ensuring checkpoints will correctly save and load.
        self.layers = torch.nn.ModuleDict()
        for layer_id in range(model_args.n_layers):
            self.layers[str(layer_id)] = TransformerBlock(...)

        self.output = nn.Linear(...)

    def forward(self, tokens: torch.Tensor):
        # Handling layers being 'None' at runtime enables easy pipeline splitting
        h = self.tok_embeddings(tokens) if self.tok_embeddings else tokens

        for layer in self.layers.values():
            h = layer(h, self.freqs_cis)

        h = self.norm(h) if self.norm else h
        output = self.output(h).float() if self.output else h
        return output

Модель, определённую таким образом, можно легко настроить для каждого этапа: сначала инициализировать модель целиком (используя метаустройство, чтобы избежать ошибок нехватки памяти), удалить ненужные для этого этапа слои, а затем создать PipelineStage, оборачивающий модель. Например:

with torch.device("meta"):
    assert num_stages == 2, "This is a simple 2-stage example"

    # we construct the entire model, then delete the parts we do not need for this stage
    # in practice, this can be done using a helper function that automatically divides up layers across stages.
    model = Transformer()

    if stage_index == 0:
        # prepare the first stage model
        del model.layers["1"]
        model.norm = None
        model.output = None

    elif stage_index == 1:
        # prepare the second stage model
        model.tok_embeddings = None
        del model.layers["0"]

    from torch.distributed.pipelining import PipelineStage
    stage = PipelineStage(
        model,
        stage_index,
        num_stages,
        device,
    )

При совместном использовании с другими методами параллелизма данных или модели может также потребоваться output_args, если форма или тип данных выходных значений части модели изменятся.

Вариант 2. Автоматическое разделение модели

Если у вас есть полная модель и вы не хотите тратить время на её преобразование в последовательность «частей модели», вам поможет API pipeline. Вот краткий пример:

class Model(torch.nn.Module):
    def __init__(self) -> None:
        super().__init__()
        self.emb = torch.nn.Embedding(10, 3)
        self.layers = torch.nn.ModuleList(
            Layer() for _ in range(2)
        )
        self.lm = LMHead()

    def forward(self, x: torch.Tensor) -> torch.Tensor:
        x = self.emb(x)
        for layer in self.layers:
            x = layer(x)
        x = self.lm(x)
        return x

Если вывести модель, можно увидеть несколько уровней иерархии, из-за которых вручную разделить её сложно:

Model(
  (emb): Embedding(10, 3)
  (layers): ModuleList(
    (0-1): 2 x Layer(
      (lin): Linear(in_features=3, out_features=3, bias=True)
    )
  )
  (lm): LMHead(
    (proj): Linear(in_features=3, out_features=3, bias=True)
  )
)

Рассмотрим, как работает API pipeline:

from torch.distributed.pipelining import pipeline, SplitPoint

# An example micro-batch input
x = torch.LongTensor([1, 2, 4, 5])

pipe = pipeline(
    module=mod,
    mb_args=(x,),
    split_spec={
        "layers.1": SplitPoint.BEGINNING,
    }
)

API pipeline разделяет модель согласно split_spec, где SplitPoint.BEGINNING обозначает добавление точки разделения перед выполнением определённого подмодуля в функции forward, а SplitPoint.END — точку разделения после его выполнения.

Если выполнить print(pipe), можно увидеть:

GraphModule(
  (submod_0): GraphModule(
    (emb): InterpreterModule()
    (layers): Module(
      (0): InterpreterModule(
        (lin): InterpreterModule()
      )
    )
  )
  (submod_1): GraphModule(
    (layers): Module(
      (1): InterpreterModule(
        (lin): InterpreterModule()
      )
    )
    (lm): InterpreterModule(
      (proj): InterpreterModule()
    )
  )
)

def forward(self, x):
    submod_0 = self.submod_0(x);  x = None
    submod_1 = self.submod_1(submod_0);  submod_0 = None
    return (submod_1,)

«Части модели» представлены подмодулями (submod_0, submod_1), каждый из которых воссоздаётся с исходными операциями модели, весами и иерархией. Кроме того, воссоздаётся функция forward «корневого уровня», фиксирующая поток данных между этими частями. Позже среда выполнения конвейера будет воспроизводить этот поток данных распределённым образом.

Объект Pipe предоставляет метод для получения «частей модели»:

stage_mod : nn.Module = pipe.get_stage_module(stage_idx)

Возвращаемый объект stage_mod — это nn.Module, с которым можно создать оптимизатор, сохранить или загрузить контрольные точки либо применить другие методы параллелизма.

Объект Pipe также позволяет создать среду выполнения распределённого этапа на устройстве, используя ProcessGroup:

stage = pipe.build_stage(stage_idx, device, group)

Если вы хотите создать среду выполнения этапа позднее, после изменения stage_mod, можно воспользоваться функциональным вариантом API build_stage. Например:

from torch.distributed.pipelining import build_stage
from torch.nn.parallel import DistributedDataParallel

dp_mod = DistributedDataParallel(stage_mod)
info = pipe.info()
stage = build_stage(dp_mod, stage_idx, info, device, group)

Примечание

Интерфейс pipeline использует трассировщик (torch.export), чтобы представить модель в виде единого графа. Если модель невозможно представить как полный граф, воспользуйтесь описанным ниже интерфейсом для ручного разделения.

Примеры Hugging Face

В репозитории PiPPy, где изначально был создан этот пакет, мы сохранили примеры на основе неизменённых моделей Hugging Face. См. каталог examples/huggingface.

Примеры включают:

  • GPT2
  • Llama

Подробный технический разбор

Как API pipeline разделяет модель?

Сначала API pipeline трассирует модель и преобразует её в ориентированный ациклический граф (DAG). Для трассировки используется torch.export — средство PyTorch 2 для захвата полного графа.

Затем он объединяет операции и параметры, необходимые этапу, в воссозданный подмодуль: submod_0, submod_1, …

В отличие от традиционных методов доступа к подмодулям, таких как Module.children(), API pipeline разделяет не только структуру модулей модели, но и её функцию прямого прохода.

Это необходимо, поскольку такая структура модели, как Module.children(), лишь фиксирует информацию во время Module.__init__() и не содержит сведений о Module.forward(). Иными словами, в Module.children() отсутствуют сведения о следующих важных для конвейерной обработки аспектах:

  • Порядок выполнения дочерних модулей в forward
  • Потоки активаций между дочерними модулями
  • Наличие функциональных операторов между дочерними модулями (например, операции relu или add не будут зафиксированы с помощью Module.children()).

Напротив, API pipeline гарантирует точное сохранение поведения forward. Он также фиксирует потоки активаций между частями, помогая распределённой среде выполнения правильно вызывать операции отправки и получения без участия пользователя.

Ещё одно преимущество API pipeline — возможность задавать точки разделения на произвольных уровнях иерархии модели. В разделённых частях исходная иерархия модели, относящаяся к каждой части, будет восстановлена без дополнительных усилий с вашей стороны. Поэтому полные квалифицированные имена (FQN), указывающие на подмодуль или параметр, останутся действительными, а службы, использующие FQN (например, FSDP, TP или контрольные точки), смогут работать с разделёнными модулями практически без изменений кода.

Создание собственного расписания

Собственное расписание конвейера можно реализовать, расширив один из следующих двух классов:

  • PipelineScheduleSingle
  • PipelineScheduleMulti

PipelineScheduleSingle предназначен для расписаний, назначающих каждому рангу только один этап. PipelineScheduleMulti предназначен для расписаний, назначающих каждому рангу несколько этапов.

Например, ScheduleGPipe и Schedule1F1B являются подклассами PipelineScheduleSingle. А ScheduleInterleaved1F1B, ScheduleLoopedBFS, ScheduleInterleavedZeroBubble и ScheduleZBVZeroBubble являются подклассами PipelineScheduleMulti.

Ведение журнала

Дополнительное журналирование можно включить с помощью переменной среды TORCH_LOGS из модуля torch._logging:

  • TORCH_LOGS=+pp будет отображать сообщения уровня logging.DEBUG и всех более высоких уровней.
  • TORCH_LOGS=pp будет отображать сообщения уровня logging.INFO и выше.
  • TORCH_LOGS=-pp будет отображать сообщения уровня logging.WARNING и выше.

Справочник по API

API разделения модели

Следующий набор API преобразует модель в представление конвейера.

class torch.distributed.pipelining.SplitPoint(value) [исходный код]

Перечисление, задающее точки, в которых можно разделить выполнение подмодуля. :ivar BEGINNING: Добавление точки разделения перед выполнением определённого подмодуля в функции forward. :ivar END: Добавление точки разделения после выполнения определённого подмодуля в функции forward.

torch.distributed.pipelining.pipeline(module, mb_args, mb_kwargs=None, split_spec=None, split_policy=None) [исходный код]

Разделяет модуль на основе спецификации.

Подробнее см. в Pipe.

Параметры:
  • module (Module) – Модуль, который нужно разделить.
  • mb_args (tuple[Any, ...]) – Пример позиционных входных данных в виде микропакета.
  • mb_kwargs (dict[str, Any] | None) – Пример входных данных, передаваемых по ключевым словам, в виде микропакета. (по умолчанию: None)
  • split_spec (dict[str, SplitPoint] | None) – Словарь, в котором имена подмодулей используются в качестве маркеров разделения. (по умолчанию: None)
  • split_policy (Callable[[GraphModule], GraphModule] | None) – Политика разделения модуля. (по умолчанию: None)
Тип возвращаемого значения:

Представление конвейера класса Pipe.

class torch.distributed.pipelining.Pipe(split_gm, num_stages, has_loss_and_backward, loss_spec) [исходный код]
torch.distributed.pipelining.pipe_split() [исходный код]

pipe_split — специальный оператор, который отмечает границу между этапами в модуле. Он используется для разделения модуля на этапы. При непосредственном выполнении модуля с такими аннотациями оператор ничего не делает.

Пример

>>> def forward(self, x):
>>>     x = torch.mm(x, self.mm_param)
>>>     x = torch.relu(x)
>>>     pipe_split()
>>>     x = self.lin(x)
>>>     return x

Пример выше будет разделён на два этапа.

Вспомогательные средства для микропакетов

class torch.distributed.pipelining.microbatch.TensorChunkSpec(split_dim) [исходный код]

Класс для задания разбиения входных данных на части

torch.distributed.pipelining.microbatch.split_args_kwargs_into_chunks(args, kwargs, chunks, args_chunk_spec=None, kwargs_chunk_spec=None) [исходный код]

Принимает последовательность args и kwargs и разделяет их на указанное количество частей согласно соответствующим спецификациям разбиения.

Параметры:
  • args (tuple[Any, ...]) – Кортеж аргументов
  • kwargs (dict[str, Any] | None) – Словарь аргументов, передаваемых по ключевым словам
  • chunks (int) – Количество частей, на которые нужно разделить args и kwargs
  • args_chunk_spec (tuple[TensorChunkSpec, ...] | None) – Спецификации разбиения для args, соответствующие форме args
  • kwargs_chunk_spec (dict[str, TensorChunkSpec] | None) – Спецификации разбиения для kwargs, соответствующие форме kwargs
Возвращает:

Список разделённых args kwargs_split: список разделённых kwargs

Тип возвращаемого значения:

args_split

torch.distributed.pipelining.microbatch.merge_chunks(chunks, chunk_spec) [исходный код]

Принимает список частей и объединяет их в одно значение согласно спецификации разбиения.

Параметры:
  • chunks (list[Any]) – список частей
  • chunk_spec – Спецификация разбиения для частей
Возвращает:

Объединённое значение

Тип возвращаемого значения:

value

Этапы конвейера

class torch.distributed.pipelining.stage.PipelineStage(submodule, stage_index, num_stages, device, input_args=None, output_args=None, output_grads=None, input_grads=None, group=None, dw_builder=None, get_mesh=None) [исходный код]

Этап конвейера для конвейерного параллелизма с последовательным разделением модели.

Поддерживается вывод метаданных в статическом и динамическом режимах:

Статический режим:

Значения input_args, output_args (а также input_grads/output_grads при наличии DTensor) передаются во время создания объекта.

Динамический режим:

Метаданные определяются по первому микропакету во время выполнения; любые статически переданные аргументы используются только для проверки.

Параметры:
  • submodule (Module) – nn.Module, оборачиваемый этим этапом.
  • stage_index (int) – Идентификатор этапа, отсчитываемый от нуля.
  • num_stages (int) – Общее количество этапов в конвейере.
  • device (device) – Устройство, на котором выполняется этот этап.
  • input_args (Tensor | tuple[Tensor, ...] | None) – Примеры входных тензоров (один тензор или кортеж). Необязательный параметр.
  • output_args (Tensor | tuple[Tensor, ...] | None) – Примеры выходных тензоров. Необязательный параметр.
  • output_grads (Tensor | tuple[Tensor | None, ...] | None) – Примеры градиентов выходных данных (получаемых от следующего этапа). Необязательный параметр.
  • input_grads (Tensor | tuple[Tensor | None, ...] | None) – Примеры градиентов входных данных (отправляемых предыдущему этапу). Необязательный параметр.
  • group (ProcessGroup | None) – Группа процессов для обмена данными P2P. По умолчанию используется глобальная группа процессов.
  • dw_builder (Callable[[], Callable[[...], None]] | None) – Функция создания отложенных обработчиков обновления весов, используемых в расписаниях с нулевым пузырьком (F/I/W).
  • get_mesh (GetMeshCallback | None) – GetMeshCallback, используемый при динамическом выводе DTensor. Игнорируется в полностью статическом режиме DTensor.
torch.distributed.pipelining.stage.build_stage(stage_module, stage_index, pipe_info, device, group=None) [исходный код]

Создаёт этап конвейера на основе stage_module, который будет обёрнут этим этапом, и информации о конвейере.

Параметры:
  • stage_module (torch.nn.Module) – модуль, который будет обёрнут этим этапом
  • stage_index (int) – индекс этого этапа в конвейере
  • pipe_info (PipeInfo) – информация о конвейере, которую можно получить с помощью pipe.info()
  • device (torch.device) – устройство, используемое этим этапом
  • group (Optional[dist.ProcessGroup]) – группа процессов, используемая этим этапом
Возвращает:

Этап конвейера, который можно запустить с помощью PipelineSchedules.

Тип возвращаемого значения:

_PipelineStage

Расписания конвейера

class torch.distributed.pipelining.schedules.ScheduleGPipe(stage, n_microbatches, loss_fn=None, args_chunk_spec=None, kwargs_chunk_spec=None, output_merge_spec=None, scale_grads=True) [исходный код]

Расписание GPipe. Последовательно обрабатывает все микропакеты в режиме заполнения и опустошения.

class torch.distributed.pipelining.schedules.Schedule1F1B(stage, n_microbatches, loss_fn=None, args_chunk_spec=None, kwargs_chunk_spec=None, output_merge_spec=None, scale_grads=True) [исходный код]

Расписание 1F1B. В установившемся режиме выполняет для микропакетов один прямой и один обратный проход.

class torch.distributed.pipelining.schedules.ScheduleInterleaved1F1B(stages, n_microbatches, loss_fn=None, args_chunk_spec=None, kwargs_chunk_spec=None, output_merge_spec=None, scale_grads=True, backward_requires_autograd=True, defer_pp_recv=False) [исходный код]

Чередующееся расписание 1F1B. Подробности см. в https://arxiv.org/pdf/2104.04473. В установившемся режиме выполняет для микропакетов один прямой и один обратный проход и поддерживает несколько этапов на ранг. Когда микропакеты готовы для нескольких локальных этапов, чередующееся расписание 1F1B отдает приоритет более раннему микропакету (этот подход также называют «сначала в глубину»).

Это расписание в основном повторяет оригинальную статью. Отличие состоит в ослаблении требования num_microbatch % pp_size == 0. При использовании расписания flex_pp получаем num_rounds = max(1, n_microbatches // pp_group_size); оно работает, если n_microbatches % num_rounds равно 0. Например, поддерживаются следующие случаи:

  1. pp_group_size = 4, n_microbatches = 10. Получаем num_rounds = 2, и n_microbatches % 2 равно 0.
  2. pp_group_size = 4, n_microbatches = 3. Получаем num_rounds = 1, и n_microbatches % 1 равно 0.
class torch.distributed.pipelining.schedules.ScheduleLoopedBFS(stages, n_microbatches, loss_fn=None, output_merge_spec=None, scale_grads=True, backward_requires_autograd=True, defer_pp_recv=False) [исходный код]

Параллелизм конвейера с обходом в ширину. Подробности см. в https://arxiv.org/abs/2211.05953. Как и чередующееся расписание 1F1B, Looped BFS поддерживает несколько этапов на ранг. Отличие в том, что, когда микропакеты готовы для нескольких локальных этапов, Looped BFS отдает приоритет более раннему этапу, запуская одновременно все доступные микропакеты.

class torch.distributed.pipelining.schedules.ScheduleInterleavedZeroBubble(stages, n_microbatches, loss_fn=None, args_chunk_spec=None, kwargs_chunk_spec=None, output_merge_spec=None, scale_grads=True, backward_requires_autograd=True, defer_pp_recv=False) [исходный код]

Чередующееся расписание с нулевым пузырем. Подробности см. в https://arxiv.org/pdf/2401.10241. В установившемся режиме выполняет для входных данных микропакетов один прямой и один обратный проход и поддерживает несколько этапов на ранг. Для заполнения пузыря конвейера используется обратный проход для весов.

В частности, здесь реализовано расписание ZB1P из статьи.

class torch.distributed.pipelining.schedules.ScheduleZBVZeroBubble(stages, n_microbatches, loss_fn=None, args_chunk_spec=None, kwargs_chunk_spec=None, output_merge_spec=None, scale_grads=True, backward_requires_autograd=True, defer_pp_recv=False) [исходный код]

Расписание с нулевым пузырем (вариант ZBV). Подробности см. в разделе 6 статьи https://arxiv.org/pdf/2401.10241.

Для этого расписания требуется ровно два этапа на ранг.

В установившемся режиме расписание выполняет для входных данных микропакетов один прямой и один обратный проход и поддерживает несколько этапов на ранг. Для заполнения пузыря конвейера используется обратный проход относительно весов.

Расписание ZB-V обладает свойством «нулевого пузыря» только при условии time forward == time backward input == time backward weights. На практике для реальных моделей это маловероятно, поэтому в качестве альтернативы можно реализовать жадный планировщик для неодинакового/несбалансированного времени выполнения.

class torch.distributed.pipelining.schedules.ScheduleDualPipeV(stages, n_microbatches, loss_fn=None, args_chunk_spec=None, kwargs_chunk_spec=None, output_merge_spec=None, scale_grads=True, backward_requires_autograd=True, defer_pp_recv=False) [исходный код]

Расписание DualPipeV. Более эффективный вариант расписания на основе DualPipe, представленного DeepSeek в https://arxiv.org/pdf/2412.19437.

Основано на опубликованном в открытом доступе коде deepseek-ai/DualPipe.

class torch.distributed.pipelining.schedules.PipelineScheduleSingle(stage, n_microbatches, loss_fn=None, args_chunk_spec=None, kwargs_chunk_spec=None, output_merge_spec=None, scale_grads=True) [исходный код]

Базовый класс расписаний для одного этапа. Реализует метод step. Производные классы должны реализовать _step_microbatches.

Градиенты масштабируются на num_microbatches в зависимости от аргумента scale_grads, который по умолчанию равен True. Эта настройка должна соответствовать конфигурации вашего loss_fn, который может усреднять потери (scale_grads=True) или суммировать их (scale_grads=False).

step(*args, target=None, losses=None, return_outputs=True, loss_kwargs=None, arg_mbs=None, kwarg_mbs=None, target_mbs=None, **kwargs) [исходный код]

Выполняет одну итерацию обучения расписания конвейера с одним этапом.

По умолчанию args, kwargs и target содержат данные всего пакета, который расписание делит на микропакеты. Если вызывающий код уже разделил входные данные на микропакеты, передавайте их вместо этого через arg_mbs, kwarg_mbs и target_mbs.

Параметры:
  • *args (Any) – Позиционные входные данные всего пакета для первого этапа конвейера. Не передавайте позиционные входные данные вместе с предварительно разделенными входными данными.
  • target (Any, optional) – Целевые значения всего пакета для вычисления потерь. При передаче предварительно разделенных входных данных передавайте целевые значения через target_mbs. Значение по умолчанию: None.
  • losses (list, optional) – Изменяемый список, в который записываются потери для каждого микропакета, если это расписание отвечает за последний этап и настроен loss_fn. Значение по умолчанию: None.
  • return_outputs (bool, optional) – Нужно ли объединить и вернуть выходные фрагменты на последнем этапе. Значение по умолчанию: True.
  • loss_kwargs (dict, optional) – Дополнительные именованные аргументы, передаваемые настроенному loss_fn. Значение по умолчанию: None.
  • arg_mbs (list[tuple], optional) – Предварительно разделенные позиционные входные данные: один кортеж на микропакет. Значение по умолчанию: None.
  • kwarg_mbs (list[dict], optional) – Предварительно разделенные именованные входные данные: один словарь на микропакет. Значение по умолчанию: None.
  • target_mbs (list, optional) – Предварительно разделенные целевые значения: одна запись на микропакет. Значение по умолчанию: None.
  • **kwargs (Any) – Именованные входные данные всего пакета для первого этапа конвейера. Не передавайте именованные входные данные вместе с предварительно разделенными входными данными.
Возвращает:
Объединенный выход последнего этапа, если этот ранг

отвечает за последний этап и return_outputs=True; в противном случае — None.

Тип возвращаемого значения:

Any or None

Вызывает исключения:
  • RuntimeError – Если включены обратные вычисления, а step вызывается в контексте torch.no_grad(). Для выполнения только прямого прохода используйте eval().
  • TypeError – Если контейнер предварительно разделенных микропакетов имеет неверный тип.
  • ValueError – Если смешаны входные данные всего пакета и предварительно разделенные входные данные либо если контейнер предварительно разделенных данных не содержит по одной записи на микропакет.

Примеры:

>>> output = schedule.step(x, mask=mask, target=target)
>>> arg_mbs = [(x0,), (x1,)]
>>> kwarg_mbs = [{"mask": mask0}, {"mask": mask1}]
>>> target_mbs = [target0, target1]
>>> output = schedule.step(
...     arg_mbs=arg_mbs,
...     kwarg_mbs=kwarg_mbs,
...     target_mbs=target_mbs,
... )
class torch.distributed.pipelining.schedules.PipelineScheduleMulti(stages, n_microbatches, loss_fn=None, args_chunk_spec=None, kwargs_chunk_spec=None, output_merge_spec=None, use_full_backward=None, scale_grads=True, backward_requires_autograd=True) [исходный код]

Базовый класс расписаний для нескольких этапов. Реализует метод step.

Градиенты масштабируются на num_microbatches в зависимости от аргумента scale_grads, который по умолчанию равен True. Эта настройка должна соответствовать конфигурации вашего loss_fn, который может усреднять потери (scale_grads=True) или суммировать их (scale_grads=False).

step(*args, target=None, losses=None, return_outputs=True, loss_kwargs=None, arg_mbs=None, kwarg_mbs=None, target_mbs=None, **kwargs) [исходный код]

Выполняет одну итерацию обучения расписания конвейера с несколькими этапами.

По умолчанию args, kwargs и target содержат данные всего пакета, который расписание делит на микропакеты. Если вызывающий код уже разделил входные данные на микропакеты, передавайте их вместо этого через arg_mbs, kwarg_mbs и target_mbs.

Параметры:
  • *args (Any) – Позиционные корневые входные данные всего пакета, если этот ранг отвечает за первый этап конвейера. Не передавайте позиционные входные данные вместе с предварительно разделенными входными данными.
  • target (Any, optional) – Целевые значения всего пакета для вычисления потерь. При передаче предварительно разделенных входных данных передавайте целевые значения через target_mbs. Значение по умолчанию: None.
  • losses (list, optional) – Изменяемый список, в который записываются потери для каждого микропакета, если это расписание отвечает за последний этап и настроен loss_fn. Значение по умолчанию: None.
  • return_outputs (bool, optional) – Нужно ли объединить и вернуть выходные фрагменты на последнем этапе. Значение по умолчанию: True.
  • loss_kwargs (dict, optional) – Дополнительные именованные аргументы, передаваемые настроенному loss_fn. Значение по умолчанию: None.
  • arg_mbs (list[tuple], optional) – Предварительно разделенные позиционные входные данные: один кортеж на микропакет. Значение по умолчанию: None.
  • kwarg_mbs (list[dict], optional) – Предварительно разделенные именованные входные данные: один словарь на микропакет. Значение по умолчанию: None.
  • target_mbs (list, optional) – Предварительно разделенные целевые значения: одна запись на микропакет. Значение по умолчанию: None.
  • **kwargs (Any) – Именованные корневые входные данные всего пакета, если этот ранг отвечает за первый этап конвейера. Не передавайте именованные входные данные вместе с предварительно разделенными входными данными.
Возвращает:
Объединенный выход последнего этапа, если этот ранг

отвечает за последний этап и return_outputs=True; в противном случае — None.

Тип возвращаемого значения:

Any or None

Вызывает исключения:
  • RuntimeError – Если включены обратные вычисления, а step вызывается в контексте torch.no_grad(). Для выполнения только прямого прохода используйте eval().
  • TypeError – Если контейнер предварительно разделенных микропакетов имеет неверный тип.
  • ValueError – Если смешаны входные данные всего пакета и предварительно разделенные входные данные либо если контейнер предварительно разделенных данных не содержит по одной записи на микропакет.

Примеры:

>>> output = schedule.step(x, mask=mask, target=target)
>>> arg_mbs = [(x0,), (x1,)]
>>> kwarg_mbs = [{"mask": mask0}, {"mask": mask1}]
>>> target_mbs = [target0, target1]
>>> output = schedule.step(
...     arg_mbs=arg_mbs,
...     kwarg_mbs=kwarg_mbs,
...     target_mbs=target_mbs,
... )
torch.distributed.pipelining.schedules.get_schedule_class(schedule_name) [исходный код]

Сопоставляет название расписания (без учета регистра) с соответствующим объектом класса.

Параметры:

schedule_name (str) – Название расписания.

© 2026, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://docs.pytorch.org/docs/2.14/distributed.pipelining.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API