评价此页

使用 Join 上下文管理器进行不均匀输入的分布式训练#

创建日期:2021 年 8 月 4 日 | 最后更新:2025 年 9 月 3 日 | 最后验证:2024 年 11 月 5 日

作者: Andrew Gu

注意

editgithub 上查看并编辑本教程。

注意

Join 是 PyTorch 1.10 中引入的原型功能。此 API 可能会发生变化。

在本教程中,你将看到

  • Join 上下文管理器的概述。

  • 将此上下文管理器与 DistributedDataParallel 一起使用的示例。

  • 同时将此上下文管理器与 DistributedDataParallelZeroRedundancyOptimizer 一起使用的示例。

  • 向上下文管理器传递关键字参数的示例。

  • 深入了解 Join 上下文管理器的工作原理。

  • 展示如何使自定义类兼容该上下文管理器的示例。

要求#

什么是 Join#

分布式数据并行入门 - 基本用例 中,你看到了使用 DistributedDataParallel 进行数据并行训练的一般框架。它在每次反向传播中隐式调度 all-reduce 操作,以在各进程(rank)间同步梯度。此类 集合通信 需要进程组中的所有进程参与,因此如果某个进程的输入较少,其他进程就会挂起或报错(取决于后端)。更一般地说,对于任何在每次迭代中执行同步集合通信的类,都存在此问题。

Join 是一个上下文管理器,用于包裹各进程的训练循环,以辅助处理不均匀输入的情况。该上下文管理器允许那些提前耗尽输入的进程(即提前“加入”)来覆盖那些尚未加入的进程所执行的集合通信。通信的覆盖方式由钩子(hook)指定。

JoinDistributedDataParallel 配合使用#

PyTorch 的 DistributedDataParallel 可以直接与 Join 上下文管理器开箱即用。以下是使用示例

import os
import torch
import torch.distributed as dist
import torch.multiprocessing as mp
from torch.distributed.algorithms.join import Join
from torch.nn.parallel import DistributedDataParallel as DDP

BACKEND = "nccl"
WORLD_SIZE = 2
NUM_INPUTS = 5

def worker(rank):
    os.environ['MASTER_ADDR'] = 'localhost'
    os.environ['MASTER_PORT'] = '29500'
    dist.init_process_group(BACKEND, rank=rank, world_size=WORLD_SIZE)

    model = DDP(torch.nn.Linear(1, 1).to(rank), device_ids=[rank])
    # Rank 1 gets one more input than rank 0
    inputs = [torch.tensor([1]).float() for _ in range(NUM_INPUTS + rank)]

    num_inputs = 0
    with Join([model]):
        for input in inputs:
            num_inputs += 1
            loss = model(input).sum()
            loss.backward()

    print(f"Rank {rank} has exhausted all {num_inputs} of its inputs!")

def main():
    mp.spawn(worker, nprocs=WORLD_SIZE, join=True)

if __name__ == "__main__":
    main()

这将产生以下输出(其中来自进程 0 和进程 1 的 print() 顺序可能是任意的)

Rank 0 has exhausted all 5 of its inputs!
Rank 1 has exhausted all 6 of its inputs!

注意

DistributedDataParallel 在引入此通用 Join 上下文管理器之前,提供了自己的 join() 上下文管理器。在上面的示例中,使用 with Join([model]): 等同于使用 with model.join():。现有的 DistributedDataParallel.join() 的一个局限性在于它不允许有多个参与类,例如同时使用 DistributedDataParallelZeroRedundancyOptimizer

JoinDistributedDataParallelZeroRedundancyOptimizer 配合使用#

Join 上下文管理器不仅可以与单个类配合使用,还可以与多个类共同使用。PyTorch 的 ZeroRedundancyOptimizer 也与此上下文管理器兼容,因此在这里,我们研究如何修改之前的示例以同时使用 DistributedDataParallelZeroRedundancyOptimizer

from torch.distributed.optim import ZeroRedundancyOptimizer as ZeRO
from torch.optim import Adam

def worker(rank):
    os.environ['MASTER_ADDR'] = 'localhost'
    os.environ['MASTER_PORT'] = '29500'
    dist.init_process_group(BACKEND, rank=rank, world_size=WORLD_SIZE)

    model = DDP(torch.nn.Linear(1, 1).to(rank), device_ids=[rank])
    optim = ZeRO(model.parameters(), Adam, lr=0.01)
    # Rank 1 gets one more input than rank 0
    inputs = [torch.tensor([1]).float() for _ in range(NUM_INPUTS + rank)]

    num_inputs = 0
    # Pass both `model` and `optim` into `Join()`
    with Join([model, optim]):
        for input in inputs:
            num_inputs += 1
            loss = model(input).sum()
            loss.backward()
            optim.step()

    print(f"Rank {rank} has exhausted all {num_inputs} of its inputs!")

这将产生与之前相同的输出。显着的变化是将 ZeroRedundancyOptimizer 实例也传入到了 Join() 中。

传递关键字参数#

类可以提供关键字参数,以在运行时修改它们在上下文管理器中的行为。例如,DistributedDataParallel 提供了一个参数 divide_by_initial_world_size,它决定了梯度是按初始世界大小还是按有效世界大小(即未加入进程的数量)进行除法。此类关键字参数可以直接传递到上下文管理器中。

with Join([model, optim], divide_by_initial_world_size=False):
    for input in inputs:
        ...

警告

传入上下文管理器的关键字参数在所有参与的类之间共享。这不应成为限制,因为我们不期望出现多个 Joinable 需要对同一参数进行不同设置的情况。尽管如此,请务必记住这一点。

Join 的工作原理是什么?#

既然我们已经了解了如何使用 Join 上下文管理器的初步示例,让我们深入探讨一下它的工作原理。这将让你更深入地了解它提供的全部功能,并帮助你构建自己的兼容类。在这里,我们将介绍 Join 类以及支撑类 JoinableJoinHook

Joinable#

首先,与 Join 上下文管理器兼容的类必须继承自抽象基类 Joinable。特别是,Joinable 必须实现

  • join_hook(self, **kwargs) -> JoinHook

这会返回该 JoinableJoinHook 实例,决定了已加入的进程应如何覆盖该 Joinable 执行的每次迭代的集合通信。

  • join_device(self) -> torch.device

这会返回一个由 Join 上下文管理器用于执行集合通信的设备,例如 torch.device("cuda:0")torch.device("cpu")

  • join_process_group(self) -> ProcessGroup

这会返回一个由 Join 上下文管理器用于执行集合通信的进程组。

特别是,join_devicejoin_process_group 是必需的属性,以确保上下文管理器可以在已加入和未加入的进程之间调度集合通信。一种用法是使用 all-reduce 在每次迭代时计算未加入进程的数量。另一种用法是实现 throw_on_early_termination=True 所需的机制,我们将在下文解释。

DistributedDataParallelZeroRedundancyOptimizer 已经继承了 Joinable 并实现了上述方法,这就是为什么我们可以在之前的示例中直接使用它们。

Joinable 类应确保调用 Joinable 的构造函数,因为它会初始化一个 JoinConfig 实例,供上下文管理器内部使用以确保正确性。这将作为名为 _join_config 的字段保存在每个 Joinable 中。

JoinHook#

接下来,让我们拆解 JoinHook 类。JoinHook 为上下文管理器提供了两个入口点:

  • main_hook(self) -> None

此钩子在尚有进程未加入时,由每个已加入的进程重复调用。其目的是覆盖 Joinable 在每次训练迭代(例如在一次前向传播、反向传播和优化器步骤中)执行的集合通信。

  • post_hook(self, is_last_joiner: bool) -> None

此钩子在所有进程均已加入后调用一次。它被传入一个额外的 bool 参数 is_last_joiner,指示该进程是否为最后加入的进程之一。该参数对于同步可能很有用。

为了提供这些钩子的具体示例,所提供的 ZeroRedundancyOptimizer 主钩子会照常执行一次优化器步骤,因为已加入的进程仍需负责更新和同步其分片的参数;而所提供的 DistributedDataParallel 后钩子会广播最后更新的模型,以确保所有进程的模型保持一致。

Join#

最后,让我们看看这些是如何融入 Join 类本身的。

  • __init__(self, joinables: List[Joinable], enable: bool = True, throw_on_early_termination: bool = False)

如我们在之前的示例中所见,构造函数接收一个参与训练循环的 Joinable 列表。这些应该是每次迭代中执行集合通信的类。

enable 是一个 bool 值,如果你确定不会出现不均匀输入,可以将其设置为 False,在这种情况下,上下文管理器将变为类似于 contextlib.nullcontext() 的空操作。这也可能禁用参与的 Joinable 类中与加入相关的计算。

throw_on_early_termination 是一个 bool 值,可以设置为 True,以便在检测到不均匀输入时让每个进程引发异常。这对于不符合上下文管理器要求的情况非常有用,最常见的情况是当存在来自不同类的集合通信且可能被任意交错时,例如在使用带有 SyncBatchNorm 层的模型配合 DistributedDataParallel 时。在这种情况下,应将此参数设置为 True,以便应用程序逻辑可以捕获异常并决定如何继续。

  • 核心逻辑发生在 __exit__() 方法中,它在仍有未加入进程时循环,调用每个 Joinable 的主钩子;待所有进程加入后,再调用它们的后钩子。主钩子和后钩子均按照 Joinable 传入的顺序进行迭代。

  • 上下文管理器需要来自未加入进程的心跳信号。因此,每个 Joinable 类都应在其每次迭代的集合通信之前调用 Join.notify_join_context()。上下文管理器将确保只有传入的第一个 Joinable 会实际发送心跳。

警告

如上所述,关于 throw_on_early_terminationJoin 上下文管理器与某些类组合方式不兼容。JoinableJoinHook 必须是可序列化的,因为每个钩子在继续下一个钩子之前都会完全执行。换句话说,两个钩子不能重叠。此外,目前主钩子和后钩子都是以相同的确定性顺序进行迭代的。如果这成为主要限制,我们可能会修改 API 以允许自定义顺序。

使自定义类与 Join 协同工作#

由于上一节引入了几个概念,让我们通过一个示例来实践一下。在这里,我们将实现一个类,它统计在自身进程加入之前,所有进程观察到的输入数量。这应该能为你提供关于如何使你自己的类与 Join 上下文管理器兼容的基本思路。

具体来说,以下代码使每个进程打印出 (1) 在其加入之前观察到的跨所有进程的输入数量,以及 (2) 跨所有进程的输入总量。

import os
import torch
import torch.distributed as dist
import torch.multiprocessing as mp
from torch.distributed.algorithms.join import Join, Joinable, JoinHook

BACKEND = "nccl"
WORLD_SIZE = 2
NUM_INPUTS = 5

class CounterJoinHook(JoinHook):
    r"""
    Join hook for :class:`Counter`.

    Arguments:
        counter (Counter): the :class:`Counter` object using this hook.
        sync_max_count (bool): whether to sync the max count once all ranks
            join.
    """
    def __init__(
        self,
        counter,
        sync_max_count
    ):
        self.counter = counter
        self.sync_max_count = sync_max_count

    def main_hook(self):
        r"""
        Shadows the counter's all-reduce by all-reducing a dim-1 zero tensor.
        """
        t = torch.zeros(1, device=self.counter.device)
        dist.all_reduce(t)

    def post_hook(self, is_last_joiner: bool):
        r"""
        Synchronizes the max count across all :class:`Counter` s if
        ``sync_max_count=True``.
        """
        if not self.sync_max_count:
            return
        rank = dist.get_rank(self.counter.process_group)
        common_rank = self.counter.find_common_rank(rank, is_last_joiner)
        if rank == common_rank:
            self.counter.max_count = self.counter.count.detach().clone()
        dist.broadcast(self.counter.max_count, src=common_rank)

class Counter(Joinable):
    r"""
    Example :class:`Joinable` that counts the number of training iterations
    that it participates in.
    """
    def __init__(self, device, process_group):
        super(Counter, self).__init__()
        self.device = device
        self.process_group = process_group
        self.count = torch.tensor([0], device=device).float()
        self.max_count = torch.tensor([0], device=device).float()

    def __call__(self):
        r"""
        Counts the number of inputs processed on this iteration by all ranks
        by all-reducing a dim-1 one tensor; increments its own internal count.
        """
        Join.notify_join_context(self)
        t = torch.ones(1, device=self.device).float()
        dist.all_reduce(t)
        self.count += t

    def join_hook(self, **kwargs) -> JoinHook:
        r"""
        Return a join hook that shadows the all-reduce in :meth:`__call__`.

        This join hook supports the following keyword arguments:
            sync_max_count (bool, optional): whether to synchronize the maximum
                count across all ranks once all ranks join; default is ``False``.
        """
        sync_max_count = kwargs.get("sync_max_count", False)
        return CounterJoinHook(self, sync_max_count)

    @property
    def join_device(self) -> torch.device:
        return self.device

    @property
    def join_process_group(self):
        return self.process_group

    def find_common_rank(self, rank, to_consider):
        r"""
        Returns the max rank of the ones to consider over the process group.
        """
        common_rank = torch.tensor([rank if to_consider else -1], device=self.device)
        dist.all_reduce(common_rank, op=dist.ReduceOp.MAX, group=self.process_group)
        common_rank = common_rank.item()
        return common_rank

def worker(rank):
    assert torch.cuda.device_count() >= WORLD_SIZE
    os.environ['MASTER_ADDR'] = 'localhost'
    os.environ['MASTER_PORT'] = '29500'
    dist.init_process_group(BACKEND, rank=rank, world_size=WORLD_SIZE)

    counter = Counter(torch.device(f"cuda:{rank}"), dist.group.WORLD)
    inputs = [torch.tensor([1]).float() for _ in range(NUM_INPUTS + rank)]

    with Join([counter], sync_max_count=True):
        for _ in inputs:
            counter()

    print(f"{int(counter.count.item())} inputs processed before rank {rank} joined!")
    print(f"{int(counter.max_count.item())} inputs processed across all ranks!")

def main():
    mp.spawn(worker, nprocs=WORLD_SIZE, join=True)

if __name__ == "__main__":
    main()

由于进程 0 看到 5 个输入,进程 1 看到 6 个,因此输出如下

10 inputs processed before rank 0 joined!
11 inputs processed across all ranks!
11 inputs processed before rank 1 joined!
11 inputs processed across all ranks!

一些需要强调的关键点

  • 一个 Counter 实例在每次迭代中执行一次 all-reduce,因此主钩子也执行一次 all-reduce 来覆盖它。

  • Counter 类在其 __call__() 方法的开头调用 Join.notify_join_context(),因为这是其每次迭代集合通信(即 all-reduce)之前的位置。

  • is_last_joiner 参数用于确定后钩子中的广播源。

  • 我们将 sync_max_count 关键字参数传递给上下文管理器,该参数随后被转发到 Counter 的 join 钩子中。