1、重点关注问题:

        1.1、nvidia-smi topo -m 是怎么获取topo结构的?调用了什么api?

        1.2、以下接口有什么用,怎么实现的?

                nvmlDeveiceGetNvLinkVersion()
                nvmlDeveiceGetNvLinkCapability()
                nvmlDeveiceGetNvLinkState()
                nvmlDeveiceGetNvLinkErrorCounter()
                nvmlDeveiceResetNvLinkErrorCounters()
                nvmlDeveiceGetNvLinkRemotePciInfo()
                nvmlDeveiceGetP2PStatus()
                nvml接口是专用于查看GPU硬件信息的接口,如温度/风扇/显存/频率/性能计数/nvlink信息/版本等。

        1.3、P2P:单机单进程P2P场景使用一下接口,具体有什么用,怎么实现的?

                cudaDeviceCanAccessPeer()
                cudaDeviceEnablePeerAccess()
                cudaMemcpyPeer()
                以上接口是单进程进行GPU之间传输数据的api接口。

        1.4、CUDA IPC:以下是单机多进程P2P,具体有什么用,怎么实现的?

                cudaIpcGetMemHandle()
                cudaIpcOpenMemHandle()
                cudaIpcCloseMemHandle()
                cudaIpcGetEventHandle()
                cudaIpcOpenEventHandle()
                以上接口是进程或主机间传递mem句柄的和事件的接口,不通进程或不通主机使用句柄对同一内存进行读写。

        1.5、统一内存(Unified Memory):

                cudaMallocManaged();
                cudaMemAdvise();
                cudaMemPrefetchAsync();
                cudaMemRangeGetAttribute();
                cudaMemAdviseSetPreferredLocation();
                cudaMemAdviseSetAccessedBy();
                cudaStreamAttachMemAsync();
                concurrentManagedAccess();
                cuMemAllocManaged();
                cuMemPrefetchAsync();
                cuMemAdvise();

  • “单指针”数据模型:你只需要分配一次内存(例如使用 cudaMallocManaged()),返回的指针可以在CPU代码和GPU内核代码中直接使用,无需区分主机指针还是设备指针。

  • 按需页面迁移:这是其智能化的核心。当CPU或GPU访问统一内存中的数据时,如果数据不在当前处理器的物理内存中,系统会产生一个页面错误(Page Fault),并自动将所需的数据页(通常为4KB大小)从当前所在位置迁移到访问者的物理内存中。这个过程对开发者是透明的,数据总是“流向”正在使用它的处理器。

  • 硬件与软件协同:这种自动迁移能力在Pascal及更新架构的GPU上通过硬件页面迁移引擎获得了更佳性能。在最新的CPU-GPU一体化架构(如Grace Hopper)上,甚至可以通过硬件实现CPU和GPU间缓存一致性(Hardware Coherency),使得CPU访问常驻GPU的数据也无需触发昂贵的页面迁移。

        1.6、NCCL是怎么获取topo结构的?

        1.7、NCCL通过nvml获取本地拓扑,多机拓扑是怎么获取的呢?

        1.8、pytorch如何调用DDP?

        1.9、DDP如何调用NCCP?

        1.0、NCCL如何调用RDMA?

        1.11、怎么理解RING和TREE?

2、实验环境

        在之前的一篇博文《在4060TI8g GPU上使用Yolov5》中,训练过程是在windows上进行的,但是windows不支持nccl,所以这次迁移到ubuntu环境进行测试,系统采用双机双卡RDMA连接,系统环境如下:

  • 主机0硬件环境:i7 + 4060ti16g + MCX354A-QCBT
  • 主机1硬件环境:i7 + 3090ti24g + MCX354A-QCBT
  • 操作系统:ubuntu22.04
  • 主机0 IP地址:192.168.1.101
  • 主机1 IP地址:192.168.1.102

3、pytorch

        这里与mpirun不一样的是,mpirun会直接使用免密ssh登录到其他主机,所以只需要在rank0主机上执行命令即可,pytorch则不会使用ssh,需要在每台主机上执行命令。

        主机0上运行以下脚本:

torchrun \
  --nnodes=2 \
  --nproc_per_node=1 \
  --node_rank=0 \
  --master_addr=192.168.1.101 \
  --master_port=29500 \
  train.py \
  --data data.yaml \
  --weights yolov5s.pt \
  --img 640 \
  --batch 32 \
  --epochs 100 \
  --device 0 \
  --workers 8

        主机1上运行以下脚本:

torchrun \
  --nnodes=2 \
  --nproc_per_node=1 \
  --node_rank=1 \
  --master_addr=192.168.1.101 \
  --master_port=29500 \
  train.py \
  --data data.yaml \
  --weights yolov5s.pt \
  --img 640 \
  --batch 32 \
  --epochs 100 \
  --device 0 \
  --workers 8
  1. YOLOv5 代码本身是PyTorch写的,不能直接拿去TensorFlow、MindSpore、PaddlePaddle里跑;
  2. 单机单卡训练命令:python train.py --data data.yaml --weights yolov5s.pt --img 640 --batch 16 --epochs 100,脚本在PyTorch框架下运行;
  3. YOLOv5 双机双卡分布式训练基于PyTorch框架,在分布式训练中主要使用了PyTorch的三个模块:torch.distributed(GPU通信)torch.nn.parallel(DDP)DistributedSampler(数据切分),下面分别介绍:
torch
 ├── distributed/	# torch.distributed(GPU通信)
					# NCCL属于该模块,如dist.all_reduce,dist.init_process_group
					# 底层会调用ProcessGroupNCCL->ncclAllReduce
 ├── nn/parallel/	# torch.nn.parallel(DDP/DistributedDataParallel)
					# 训练同步协调器,负责hook梯度/bucket梯度/决定何时同步/overlap通信与计算/管理reducer
					# 实际通信在torch.distributed模块
 └── utils/data/	# DistributedSampler(数据切分)
                    # 属于dataloader

        4、torch.distributed(GPU通信)/torch.nn.parallel(DDP)/DistributedSampler(数据切分)在train.py脚本中怎么被使用的:

  • torch.distributed(GPU通信):dist.init_process_group(),dist.all_reduce();
  • torch.nn.parallel(DDP):smart_DDP(),smart_DDP()会处理了一些兼容性问题,核心是调用DistributedDataParallel(model),这样让模型在backward(scaler.scale(loss).backward())时触发梯度同步,也就是执行dist.all_reduce();
  • DistributedSampler(数据切分):create_dataloader()会调用DistributedSampler()函数。

        以下是yolov5 train.py的核心代码,代码已标注注释:

# Ultralytics 🚀 AGPL-3.0 License - https://ultralytics.com/license
"""
Train a YOLOv5 model on a custom dataset. Models and datasets download automatically from the latest YOLOv5 release.

Usage - Single-GPU training:
    $ python train.py --data coco128.yaml --weights yolov5s.pt --img 640  # from pretrained (recommended)
    $ python train.py --data coco128.yaml --weights '' --cfg yolov5s.yaml --img 640  # from scratch

Usage - Multi-GPU DDP training:
    $ python -m torch.distributed.run --nproc_per_node 4 --master_port 1 train.py --data coco128.yaml --weights yolov5s.pt --img 640 --device 0,1,2,3

Models:     https://github.com/ultralytics/yolov5/tree/master/models
Datasets:   https://github.com/ultralytics/yolov5/tree/master/data
Tutorial:   https://docs.ultralytics.com/yolov5/tutorials/train_custom_data
"""

import argparse
import math
import os
import random
import subprocess
import sys
import time
from copy import deepcopy
from datetime import datetime, timedelta
from pathlib import Path

try:
    import comet_ml  # must be imported before torch (if installed)
except ImportError:
    comet_ml = None

import numpy as np
# PyTorch 深度学习框架核心库
import torch
# PyTorch 分布式训练模块(DDP/NCCL/Gloo)
import torch.distributed as dist
import torch.nn as nn
import yaml
from torch.optim import lr_scheduler
from tqdm import tqdm

FILE = Path(__file__).resolve()
ROOT = FILE.parents[0]  # YOLOv5 root directory
if str(ROOT) not in sys.path:
    sys.path.append(str(ROOT))  # add ROOT to PATH
ROOT = Path(os.path.relpath(ROOT, Path.cwd()))  # relative

from ultralytics.utils.patches import torch_load

import val as validate  # for end-of-epoch mAP
from models.experimental import attempt_load
from models.yolo import Model
from utils.autoanchor import check_anchors
from utils.autobatch import check_train_batch_size
from utils.callbacks import Callbacks
from utils.dataloaders import create_dataloader
from utils.downloads import attempt_download, is_url
from utils.general import (
    LOGGER,
    TQDM_BAR_FORMAT,
    check_amp,
    check_dataset,
    check_file,
    check_git_info,
    check_git_status,
    check_img_size,
    check_requirements,
    check_suffix,
    check_yaml,
    colorstr,
    get_latest_run,
    increment_path,
    init_seeds,
    intersect_dicts,
    labels_to_class_weights,
    labels_to_image_weights,
    methods,
    one_cycle,
    print_args,
    print_mutation,
    strip_optimizer,
    yaml_save,
)
from utils.loggers import LOGGERS, Loggers
from utils.loggers.comet.comet_utils import check_comet_resume
from utils.loss import ComputeLoss
from utils.metrics import fitness
from utils.plots import plot_evolve
from utils.torch_utils import (
    EarlyStopping,
    ModelEMA,
    de_parallel,
    select_device,
    smart_DDP,
    smart_optimizer,
    smart_resume,
    torch_distributed_zero_first,
)

# 当前进程绑定的 GPU 编号(DDP 模式下由 torchrun 自动注入)
LOCAL_RANK = int(os.getenv("LOCAL_RANK", -1))  # https://pytorch.org/docs/stable/elastic/run.html
RANK = int(os.getenv("RANK", -1))
WORLD_SIZE = int(os.getenv("WORLD_SIZE", 1))
GIT_INFO = check_git_info()


# YOLOv5 主训练函数
# hyp: 超参数配置
# opt: 命令行参数
# device: 训练设备
# callbacks: 回调函数

def train(hyp, opt, device, callbacks):
    """Train a YOLOv5 model on a custom dataset using specified hyperparameters, options, and device, managing datasets,
    model architecture, loss computation, and optimizer steps.

    Args:
        hyp (str | dict): Path to the hyperparameters YAML file or a dictionary of hyperparameters.
        opt (argparse.Namespace): Parsed command-line arguments containing training options.
        device (torch.device): Device on which training occurs, e.g., 'cuda' or 'cpu'.
        callbacks (Callbacks): Callback functions for various training events.

    Returns:
        None

    Examples:
        Single-GPU training:
        ```bash
        $ python train.py --data coco128.yaml --weights yolov5s.pt --img 640  # from pretrained (recommended)
        $ python train.py --data coco128.yaml --weights '' --cfg yolov5s.yaml --img 640  # from scratch
        ```

        Multi-GPU DDP training:
        ```bash
        $ python -m torch.distributed.run --nproc_per_node 4 --master_port 1 train.py --data coco128.yaml --weights
        yolov5s.pt --img 640 --device 0,1,2,3
        ```

        For more usage details, refer to:
        - Models: https://github.com/ultralytics/yolov5/tree/master/models
        - Datasets: https://github.com/ultralytics/yolov5/tree/master/data
        - Tutorial: https://docs.ultralytics.com/yolov5/tutorials/train_custom_data

    Notes:
        Models and datasets download automatically from the latest YOLOv5 release.
    """
    save_dir, epochs, batch_size, weights, single_cls, evolve, data, cfg, resume, noval, nosave, workers, freeze = (
        Path(opt.save_dir),
        opt.epochs,
        opt.batch_size,
        opt.weights,
        opt.single_cls,
        opt.evolve,
        opt.data,
        opt.cfg,
        opt.resume,
        opt.noval,
        opt.nosave,
        opt.workers,
        opt.freeze,
    )
    callbacks.run("on_pretrain_routine_start")

    # =========================
    # 创建权重保存目录
    # =========================
    w = save_dir / "weights"  # weights dir
    (w.parent if evolve else w).mkdir(parents=True, exist_ok=True)  # make dir
    last, best = w / "last.pt", w / "best.pt"

    # =========================
    # 加载超参数
    # =========================
    if isinstance(hyp, str):
        with open(hyp, errors="ignore") as f:
            hyp = yaml.safe_load(f)  # load hyps dict
    LOGGER.info(colorstr("hyperparameters: ") + ", ".join(f"{k}={v}" for k, v in hyp.items()))
    opt.hyp = hyp.copy()  # for saving hyps to checkpoints

    # Save run settings
    if not evolve:
        yaml_save(save_dir / "hyp.yaml", hyp)
        yaml_save(save_dir / "opt.yaml", vars(opt))

    # Loggers
    data_dict = None
    if RANK in {-1, 0}:
        include_loggers = list(LOGGERS)
        if getattr(opt, "ndjson_console", False):
            include_loggers.append("ndjson_console")
        if getattr(opt, "ndjson_file", False):
            include_loggers.append("ndjson_file")

        loggers = Loggers(
            save_dir=save_dir,
            weights=weights,
            opt=opt,
            hyp=hyp,
            logger=LOGGER,
            include=tuple(include_loggers),
        )

        # Register actions
        for k in methods(loggers):
            callbacks.register_action(k, callback=getattr(loggers, k))

        # Process custom dataset artifact link
        data_dict = loggers.remote_dataset
        if resume:  # If resuming runs from remote artifact
            weights, epochs, hyp, batch_size = opt.weights, opt.epochs, opt.hyp, opt.batch_size

    # Config
    plots = not evolve and not opt.noplots  # create plots
    cuda = device.type != "cpu"
    init_seeds(opt.seed + 1 + RANK, deterministic=True)
    with torch_distributed_zero_first(LOCAL_RANK):
        data_dict = data_dict or check_dataset(data)  # check if None
    train_path, val_path = data_dict["train"], data_dict["val"]
    nc = 1 if single_cls else int(data_dict["nc"])  # number of classes
    names = {0: "item"} if single_cls and len(data_dict["names"]) != 1 else data_dict["names"]  # class names
    is_coco = isinstance(val_path, str) and val_path.endswith("coco/val2017.txt")  # COCO dataset

    # =========================
    # 构建模型
    # 支持:
    # 1. 加载预训练模型
    # 2. 从 cfg 创建新模型
    # =========================
    check_suffix(weights, ".pt")  # check weights
    pretrained = weights.endswith(".pt")
    if pretrained:
        with torch_distributed_zero_first(LOCAL_RANK):
            weights = attempt_download(weights)  # download if not found locally
        ckpt = torch_load(weights, map_location="cpu")  # load checkpoint to CPU to avoid CUDA memory leak
        model = Model(cfg or ckpt["model"].yaml, ch=3, nc=nc, anchors=hyp.get("anchors")).to(device)  # create
        exclude = ["anchor"] if (cfg or hyp.get("anchors")) and not resume else []  # exclude keys
        csd = ckpt["model"].float().state_dict()  # checkpoint state_dict as FP32
        csd = intersect_dicts(csd, model.state_dict(), exclude=exclude)  # intersect
        model.load_state_dict(csd, strict=False)  # load
        LOGGER.info(f"Transferred {len(csd)}/{len(model.state_dict())} items from {weights}")  # report
    else:
        model = Model(cfg, ch=3, nc=nc, anchors=hyp.get("anchors")).to(device)  # create
    amp = check_amp(model)  # check AMP

    # =========================
    # 冻结指定层
    # 常用于迁移学习
    # =========================
    freeze = [f"model.{x}." for x in (freeze if len(freeze) > 1 else range(freeze[0]))]  # layers to freeze
    for k, v in model.named_parameters():
        v.requires_grad = True  # train all layers
        # v.register_hook(lambda x: torch.nan_to_num(x))  # NaN to 0 (commented for erratic training results)
        if any(x in k for x in freeze):
            LOGGER.info(f"freezing {k}")
            v.requires_grad = False

    # Image size
    gs = max(int(model.stride.max()), 32)  # grid size (max stride)
    imgsz = check_img_size(opt.imgsz, gs, floor=gs * 2)  # verify imgsz is gs-multiple

    # Batch size
    if RANK == -1 and batch_size == -1:  # single-GPU only, estimate best batch size
        batch_size = check_train_batch_size(model, imgsz, amp)
        loggers.on_params_update({"batch_size": batch_size})

    # =========================
    # 创建优化器(SGD/Adam/AdamW)
    # =========================
    nbs = 64  # nominal batch size
    accumulate = max(round(nbs / batch_size), 1)  # accumulate loss before optimizing
    hyp["weight_decay"] *= batch_size * accumulate / nbs  # scale weight_decay
    optimizer = smart_optimizer(model, opt.optimizer, hyp["lr0"], hyp["momentum"], hyp["weight_decay"])

    # =========================
    # 学习率调度器
    # 支持余弦退火 / 线性衰减
    # =========================
    if opt.cos_lr:
        lf = one_cycle(1, hyp["lrf"], epochs)  # cosine 1->hyp['lrf']
    else:

        def lf(x):
            """Linear learning rate scheduler function with decay calculated by epoch proportion."""
            return (1 - x / epochs) * (1.0 - hyp["lrf"]) + hyp["lrf"]  # linear

    scheduler = lr_scheduler.LambdaLR(optimizer, lr_lambda=lf)  # plot_lr_scheduler(optimizer, scheduler, epochs)

    # =========================
    # EMA(指数滑动平均模型)
    # 推理效果通常更稳定
    # =========================
    ema = ModelEMA(model) if RANK in {-1, 0} else None

    # Resume
    best_fitness, start_epoch = 0.0, 0
    if pretrained:
        if resume:
            best_fitness, start_epoch, epochs = smart_resume(ckpt, optimizer, ema, weights, epochs, resume)
        del ckpt, csd

    # DP mode
    if cuda and RANK == -1 and torch.cuda.device_count() > 1:
        LOGGER.warning(
            "WARNING ⚠️ DP not recommended, use torch.distributed.run for best DDP Multi-GPU results.\n"
            "See Multi-GPU Tutorial at https://docs.ultralytics.com/yolov5/tutorials/multi_gpu_training to get started."
        )
        model = torch.nn.DataParallel(model)

    # SyncBatchNorm
    if opt.sync_bn and cuda and RANK != -1:
        model = torch.nn.SyncBatchNorm.convert_sync_batchnorm(model).to(device)
        LOGGER.info("Using SyncBatchNorm()")

    # =========================
    # 创建训练集 DataLoader
    # 内部会自动创建 DistributedSampler
    # DDP 多卡时每张卡只读取部分数据
    # =========================
    train_loader, dataset = create_dataloader(
        train_path,
        imgsz,
        batch_size // WORLD_SIZE,
        gs,
        single_cls,
        hyp=hyp,
        augment=True,
        cache=None if opt.cache == "val" else opt.cache,
        rect=opt.rect,
        rank=LOCAL_RANK,
        workers=workers,
        image_weights=opt.image_weights,
        quad=opt.quad,
        prefix=colorstr("train: "),
        shuffle=True,
        seed=opt.seed,
    )
    labels = np.concatenate(dataset.labels, 0)
    mlc = int(labels[:, 0].max())  # max label class
    assert mlc < nc, f"Label class {mlc} exceeds nc={nc} in {data}. Possible class labels are 0-{nc - 1}"

    # Process 0
    if RANK in {-1, 0}:
        val_loader = create_dataloader(
            val_path,
            imgsz,
            batch_size // WORLD_SIZE * 2,
            gs,
            single_cls,
            hyp=hyp,
            cache=None if noval else opt.cache,
            rect=True,
            rank=-1,
            workers=workers * 2,
            pad=0.5,
            prefix=colorstr("val: "),
        )[0]

        if not resume:
            if not opt.noautoanchor:
                check_anchors(dataset, model=model, thr=hyp["anchor_t"], imgsz=imgsz)  # run AutoAnchor
            model.half().float()  # pre-reduce anchor precision

        callbacks.run("on_pretrain_routine_end", labels, names)

    # =========================
    # DDP 分布式训练模式
    # DDP 会在 backward 后自动同步梯度
    # =========================
    if cuda and RANK != -1:
        model = smart_DDP(model)

    # =========================
    # 构建模型
    # 支持:
    # 1. 加载预训练模型
    # 2. 从 cfg 创建新模型
    # ========================= attributes
    nl = de_parallel(model).model[-1].nl  # number of detection layers (to scale hyps)
    hyp["box"] *= 3 / nl  # scale to layers
    hyp["cls"] *= nc / 80 * 3 / nl  # scale to classes and layers
    hyp["obj"] *= (imgsz / 640) ** 2 * 3 / nl  # scale to image size and layers
    hyp["label_smoothing"] = opt.label_smoothing
    model.nc = nc  # attach number of classes to model
    model.hyp = hyp  # attach hyperparameters to model
    model.class_weights = labels_to_class_weights(dataset.labels, nc).to(device) * nc  # attach class weights
    model.names = names

    # =========================
    # 开始正式训练循环
    # =========================
    t0 = time.time()
    nb = len(train_loader)  # number of batches
    nw = max(round(hyp["warmup_epochs"] * nb), 100)  # number of warmup iterations, max(3 epochs, 100 iterations)
    # nw = min(nw, (epochs - start_epoch) / 2 * nb)  # limit warmup to < 1/2 of training
    last_opt_step = -1
    maps = np.zeros(nc)  # mAP per class
    results = (0, 0, 0, 0, 0, 0, 0)  # P, R, mAP@.5, mAP@.5-.95, val_loss(box, obj, cls)
    scheduler.last_epoch = start_epoch - 1  # do not move
    scaler = torch.cuda.amp.GradScaler(enabled=amp)
    stopper, stop = EarlyStopping(patience=opt.patience), False
    compute_loss = ComputeLoss(model)  # init loss class
    callbacks.run("on_train_start")
    LOGGER.info(
        f"Image sizes {imgsz} train, {imgsz} val\n"
        f"Using {train_loader.num_workers * WORLD_SIZE} dataloader workers\n"
        f"Logging results to {colorstr('bold', save_dir)}\n"
        f"Starting training for {epochs} epochs..."
    )
    for epoch in range(start_epoch, epochs):  # epoch ------------------------------------------------------------------
        callbacks.run("on_train_epoch_start")
        model.train()

        # Update image weights (optional, single-GPU only)
        if opt.image_weights:
            cw = model.class_weights.cpu().numpy() * (1 - maps) ** 2 / nc  # class weights
            iw = labels_to_image_weights(dataset.labels, nc=nc, class_weights=cw)  # image weights
            dataset.indices = random.choices(range(dataset.n), weights=iw, k=dataset.n)  # rand weighted idx

        # Update mosaic border (optional)
        # b = int(random.uniform(0.25 * imgsz, 0.75 * imgsz + gs) // gs * gs)
        # dataset.mosaic_border = [b - imgsz, -b]  # height, width borders

        mloss = torch.zeros(3, device=device)  # mean losses
        if RANK != -1:
            train_loader.sampler.set_epoch(epoch)
        pbar = enumerate(train_loader)
        LOGGER.info(("\n" + "%11s" * 7) % ("Epoch", "GPU_mem", "box_loss", "obj_loss", "cls_loss", "Instances", "Size"))
        if RANK in {-1, 0}:
            pbar = tqdm(pbar, total=nb, bar_format=TQDM_BAR_FORMAT)  # progress bar
        optimizer.zero_grad()
        for i, (imgs, targets, paths, _) in pbar:  # batch -------------------------------------------------------------
            callbacks.run("on_train_batch_start")
            ni = i + nb * epoch  # number integrated batches (since train start)
            imgs = imgs.to(device, non_blocking=True).float() / 255  # uint8 to float32, 0-255 to 0.0-1.0

            # =========================
            # Warmup 学习率预热
            # 训练初期逐渐增大学习率
            # 避免梯度震荡
            # =========================
            if ni <= nw:
                xi = [0, nw]  # x interp
                # compute_loss.gr = np.interp(ni, xi, [0.0, 1.0])  # iou loss ratio (obj_loss = 1.0 or iou)
                accumulate = max(1, np.interp(ni, xi, [1, nbs / batch_size]).round())
                for j, x in enumerate(optimizer.param_groups):
                    # bias lr falls from 0.1 to lr0, all other lrs rise from 0.0 to lr0
                    x["lr"] = np.interp(ni, xi, [hyp["warmup_bias_lr"] if j == 0 else 0.0, x["initial_lr"] * lf(epoch)])
                    if "momentum" in x:
                        x["momentum"] = np.interp(ni, xi, [hyp["warmup_momentum"], hyp["momentum"]])

            # Multi-scale
            if opt.multi_scale:
                sz = random.randrange(int(imgsz * 0.5), int(imgsz * 1.5) + gs) // gs * gs  # size
                sf = sz / max(imgs.shape[2:])  # scale factor
                if sf != 1:
                    ns = [math.ceil(x * sf / gs) * gs for x in imgs.shape[2:]]  # new shape (stretched to gs-multiple)
                    imgs = nn.functional.interpolate(imgs, size=ns, mode="bilinear", align_corners=False)

            # =========================
            # 前向传播
            # =========================
            with torch.cuda.amp.autocast(amp):
                pred = model(imgs)  # forward
                loss, loss_items = compute_loss(pred, targets.to(device))  # loss scaled by batch_size
                if RANK != -1:
                    loss *= WORLD_SIZE  # gradient averaged between devices in DDP mode
                if opt.quad:
                    loss *= 4.0

            # =========================
            # 反向传播
            # DDP 会在这里自动 all_reduce 同步梯度
            # =========================
            scaler.scale(loss).backward()

            # =========================
            # 优化器更新参数
            # 包含 AMP 混合精度训练
            # =========================
            if ni - last_opt_step >= accumulate:
                scaler.unscale_(optimizer)  # unscale gradients
                torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=10.0)  # clip gradients
                scaler.step(optimizer)  # optimizer.step
                scaler.update()
                optimizer.zero_grad()
                if ema:
                    ema.update(model)
                last_opt_step = ni

            # Log
            if RANK in {-1, 0}:
                mloss = (mloss * i + loss_items) / (i + 1)  # update mean losses
                mem = f"{torch.cuda.memory_reserved() / 1e9 if torch.cuda.is_available() else 0:.3g}G"  # (GB)
                pbar.set_description(
                    ("%11s" * 2 + "%11.4g" * 5)
                    % (f"{epoch}/{epochs - 1}", mem, *mloss, targets.shape[0], imgs.shape[-1])
                )
                callbacks.run("on_train_batch_end", model, ni, imgs, targets, paths, list(mloss))
                if callbacks.stop_training:
                    return
            # end batch ------------------------------------------------------------------------------------------------

        # =========================
        # 学习率调度器
        # 支持余弦退火 / 线性衰减
        # =========================
        lr = [x["lr"] for x in optimizer.param_groups]  # for loggers
        scheduler.step()

        if RANK in {-1, 0}:
            # mAP
            callbacks.run("on_train_epoch_end", epoch=epoch)
            ema.update_attr(model, include=["yaml", "nc", "hyp", "names", "stride", "class_weights"])
            final_epoch = (epoch + 1 == epochs) or stopper.possible_stop
            if not noval or final_epoch:  # Calculate mAP
                results, maps, _ = validate.run(
                    data_dict,
                    batch_size=batch_size // WORLD_SIZE * 2,
                    imgsz=imgsz,
                    half=amp,
                    model=ema.ema,
                    single_cls=single_cls,
                    dataloader=val_loader,
                    save_dir=save_dir,
                    plots=False,
                    callbacks=callbacks,
                    compute_loss=compute_loss,
                )

            # Update best mAP
            fi = fitness(np.array(results).reshape(1, -1))  # weighted combination of [P, R, mAP@.5, mAP@.5-.95]
            stop = stopper(epoch=epoch, fitness=fi)  # early stop check
            if fi > best_fitness:
                best_fitness = fi
            log_vals = list(mloss) + list(results) + lr
            callbacks.run("on_fit_epoch_end", log_vals, epoch, best_fitness, fi)

            # Save model
            if (not nosave) or (final_epoch and not evolve):  # if save
                ckpt = {
                    "epoch": epoch,
                    "best_fitness": best_fitness,
                    "model": deepcopy(de_parallel(model)).half(),
                    "ema": deepcopy(ema.ema).half(),
                    "updates": ema.updates,
                    "optimizer": optimizer.state_dict(),
                    "opt": vars(opt),
                    "git": GIT_INFO,  # {remote, branch, commit} if a git repo
                    "date": datetime.now().isoformat(),
                }

                # Save last, best and delete
                torch.save(ckpt, last)
                if best_fitness == fi:
                    torch.save(ckpt, best)
                if opt.save_period > 0 and epoch % opt.save_period == 0:
                    torch.save(ckpt, w / f"epoch{epoch}.pt")
                del ckpt
                callbacks.run("on_model_save", last, epoch, final_epoch, best_fitness, fi)

        # EarlyStopping
        if RANK != -1:  # if DDP training
            broadcast_list = [stop if RANK == 0 else None]
            dist.broadcast_object_list(broadcast_list, 0)  # broadcast 'stop' to all ranks
            if RANK != 0:
                stop = broadcast_list[0]
        if stop:
            break  # must break all DDP ranks

        # end epoch ----------------------------------------------------------------------------------------------------
    # end training -----------------------------------------------------------------------------------------------------
    if RANK in {-1, 0}:
        LOGGER.info(f"\n{epoch - start_epoch + 1} epochs completed in {(time.time() - t0) / 3600:.3f} hours.")
        for f in last, best:
            if f.exists():
                strip_optimizer(f)  # strip optimizers
                if f is best:
                    LOGGER.info(f"\nValidating {f}...")
                    results, _, _ = validate.run(
                        data_dict,
                        batch_size=batch_size // WORLD_SIZE * 2,
                        imgsz=imgsz,
                        model=attempt_load(f, device).half(),
                        iou_thres=0.65 if is_coco else 0.60,  # best pycocotools at iou 0.65
                        single_cls=single_cls,
                        dataloader=val_loader,
                        save_dir=save_dir,
                        save_json=is_coco,
                        verbose=True,
                        plots=plots,
                        callbacks=callbacks,
                        compute_loss=compute_loss,
                    )  # val best model with plots
                    if is_coco:
                        callbacks.run("on_fit_epoch_end", list(mloss) + list(results) + lr, epoch, best_fitness, fi)

        callbacks.run("on_train_end", last, best, epoch, results)

    torch.cuda.empty_cache()
    return results


# 解析命令行参数
# 例如:
# python train.py --weights yolov5s.pt --data coco128.yaml

def parse_opt(known=False):
    """Parse command-line arguments for YOLOv5 training, validation, and testing.

    Args:
        known (bool, optional): If True, parses known arguments, ignoring the unknown. Defaults to False.

    Returns:
        (argparse.Namespace): Parsed command-line arguments containing options for YOLOv5 execution.

    Examples:
        ```python
        from ultralytics.yolo import parse_opt
        opt = parse_opt()
        print(opt)
        ```

    Links:
        - Models: https://github.com/ultralytics/yolov5/tree/master/models
        - Datasets: https://github.com/ultralytics/yolov5/tree/master/data
        - Tutorial: https://docs.ultralytics.com/yolov5/tutorials/train_custom_data
    """
    parser = argparse.ArgumentParser()
    parser.add_argument("--weights", type=str, default=ROOT / "yolov5s.pt", help="initial weights path")
    parser.add_argument("--cfg", type=str, default="", help="model.yaml path")
    parser.add_argument("--data", type=str, default=ROOT / "data/coco128.yaml", help="dataset.yaml path")
    parser.add_argument("--hyp", type=str, default=ROOT / "data/hyps/hyp.scratch-low.yaml", help="hyperparameters path")
    parser.add_argument("--epochs", type=int, default=100, help="total training epochs")
    parser.add_argument("--batch-size", type=int, default=16, help="total batch size for all GPUs, -1 for autobatch")
    parser.add_argument("--imgsz", "--img", "--img-size", type=int, default=640, help="train, val image size (pixels)")
    parser.add_argument("--rect", action="store_true", help="rectangular training")
    parser.add_argument("--resume", nargs="?", const=True, default=False, help="resume most recent training")
    parser.add_argument("--nosave", action="store_true", help="only save final checkpoint")
    parser.add_argument("--noval", action="store_true", help="only validate final epoch")
    parser.add_argument("--noautoanchor", action="store_true", help="disable AutoAnchor")
    parser.add_argument("--noplots", action="store_true", help="save no plot files")
    parser.add_argument("--evolve", type=int, nargs="?", const=300, help="evolve hyperparameters for x generations")
    parser.add_argument(
        "--evolve_population", type=str, default=ROOT / "data/hyps", help="location for loading population"
    )
    parser.add_argument("--resume_evolve", type=str, default=None, help="resume evolve from last generation")
    parser.add_argument("--bucket", type=str, default="", help="gsutil bucket")
    parser.add_argument("--cache", type=str, nargs="?", const="ram", help="image --cache ram/disk")
    parser.add_argument("--image-weights", action="store_true", help="use weighted image selection for training")
    parser.add_argument("--device", default="", help="cuda device, i.e. 0 or 0,1,2,3 or cpu")
    parser.add_argument("--multi-scale", action="store_true", help="vary img-size +/- 50%%")
    parser.add_argument("--single-cls", action="store_true", help="train multi-class data as single-class")
    parser.add_argument("--optimizer", type=str, choices=["SGD", "Adam", "AdamW"], default="SGD", help="optimizer")
    parser.add_argument("--sync-bn", action="store_true", help="use SyncBatchNorm, only available in DDP mode")
    parser.add_argument("--workers", type=int, default=8, help="max dataloader workers (per RANK in DDP mode)")
    parser.add_argument("--project", default=ROOT / "runs/train", help="save to project/name")
    parser.add_argument("--name", default="exp", help="save to project/name")
    parser.add_argument("--exist-ok", action="store_true", help="existing project/name ok, do not increment")
    parser.add_argument("--quad", action="store_true", help="quad dataloader")
    parser.add_argument("--cos-lr", action="store_true", help="cosine LR scheduler")
    parser.add_argument("--label-smoothing", type=float, default=0.0, help="Label smoothing epsilon")
    parser.add_argument("--patience", type=int, default=100, help="EarlyStopping patience (epochs without improvement)")
    parser.add_argument("--freeze", nargs="+", type=int, default=[0], help="Freeze layers: backbone=10, first3=0 1 2")
    parser.add_argument("--save-period", type=int, default=-1, help="Save checkpoint every x epochs (disabled if < 1)")
    parser.add_argument("--seed", type=int, default=0, help="Global training seed")
    parser.add_argument("--local_rank", type=int, default=-1, help="Automatic DDP Multi-GPU argument, do not modify")

    # Logger arguments
    parser.add_argument("--entity", default=None, help="Entity")
    parser.add_argument("--upload_dataset", nargs="?", const=True, default=False, help='Upload data, "val" option')
    parser.add_argument("--bbox_interval", type=int, default=-1, help="Set bounding-box image logging interval")
    parser.add_argument("--artifact_alias", type=str, default="latest", help="Version of dataset artifact to use")

    # NDJSON logging
    parser.add_argument("--ndjson-console", action="store_true", help="Log ndjson to console")
    parser.add_argument("--ndjson-file", action="store_true", help="Log ndjson to file")

    return parser.parse_known_args()[0] if known else parser.parse_args()


# 主入口函数
# 包括:
# 1. 初始化 DDP
# 2. 检查配置
# 3. 调用 train()

def main(opt, callbacks=Callbacks()):
    """Runs the main entry point for training or hyperparameter evolution with specified options and optional callbacks.

    Args:
        opt (argparse.Namespace): The command-line arguments parsed for YOLOv5 training and evolution.
        callbacks (ultralytics.utils.callbacks.Callbacks, optional): Callback functions for various training stages.
            Defaults to Callbacks().

    Returns:
        None

    Notes:
        For detailed usage, refer to:
        https://github.com/ultralytics/yolov5/tree/master/models
    """
    if RANK in {-1, 0}:
        print_args(vars(opt))
        check_git_status()
        check_requirements(ROOT / "requirements.txt")

    # Resume (from specified or most recent last.pt)
    if opt.resume and not check_comet_resume(opt) and not opt.evolve:
        last = Path(check_file(opt.resume) if isinstance(opt.resume, str) else get_latest_run())
        opt_yaml = last.parent.parent / "opt.yaml"  # train options yaml
        opt_data = opt.data  # original dataset
        if opt_yaml.is_file():
            with open(opt_yaml, errors="ignore") as f:
                d = yaml.safe_load(f)
        else:
            d = torch_load(last, map_location="cpu")["opt"]
        opt = argparse.Namespace(**d)  # replace
        opt.cfg, opt.weights, opt.resume = "", str(last), True  # reinstate
        if is_url(opt_data):
            opt.data = check_file(opt_data)  # avoid HUB resume auth timeout
    else:
        opt.data, opt.cfg, opt.hyp, opt.weights, opt.project = (
            check_file(opt.data),
            check_yaml(opt.cfg),
            check_yaml(opt.hyp),
            str(opt.weights),
            str(opt.project),
        )  # checks
        assert len(opt.cfg) or len(opt.weights), "either --cfg or --weights must be specified"
        if opt.evolve:
            if opt.project == str(ROOT / "runs/train"):  # if default project name, rename to runs/evolve
                opt.project = str(ROOT / "runs/evolve")
            opt.exist_ok, opt.resume = opt.resume, False  # pass resume to exist_ok and disable resume
        if opt.name == "cfg":
            opt.name = Path(opt.cfg).stem  # use model.yaml as name
        opt.save_dir = str(increment_path(Path(opt.project) / opt.name, exist_ok=opt.exist_ok))

    # =========================
    # DDP 分布式训练模式
    # DDP 会在 backward 后自动同步梯度
    # =========================
    device = select_device(opt.device, batch_size=opt.batch_size)
    if LOCAL_RANK != -1:
        msg = "is not compatible with YOLOv5 Multi-GPU DDP training"
        assert not opt.image_weights, f"--image-weights {msg}"
        assert not opt.evolve, f"--evolve {msg}"
        assert opt.batch_size != -1, f"AutoBatch with --batch-size -1 {msg}, please pass a valid --batch-size"
        assert opt.batch_size % WORLD_SIZE == 0, f"--batch-size {opt.batch_size} must be multiple of WORLD_SIZE"
        assert torch.cuda.device_count() > LOCAL_RANK, "insufficient CUDA devices for DDP command"
        torch.cuda.set_device(LOCAL_RANK)
        device = torch.device("cuda", LOCAL_RANK)
        # 初始化 PyTorch 分布式通信组
        # Linux 下一般使用 NCCL
        # Windows 下一般只能使用 Gloo
        dist.init_process_group(
            backend="nccl" if dist.is_nccl_available() else "gloo", timeout=timedelta(seconds=10800)
        )

    # Train
    train(opt.hyp, opt, device, callbacks)


if __name__ == "__main__":
    opt = parse_opt()
    main(opt)

        3.1、重点代码分析

        3.1.1、如果是分布式训练,则初始化pyTorch分布式通信组,主要执行init_process_group()函数,该函数的底层就会调用ncclCommInitRank()函数。

if LOCAL_RANK != -1:
    msg = "is not compatible with YOLOv5 Multi-GPU DDP training"
    assert not opt.image_weights, f"--image-weights {msg}"
    assert not opt.evolve, f"--evolve {msg}"
    assert opt.batch_size != -1, f"AutoBatch with --batch-size -1 {msg}, please pass a valid --batch-size"
    assert opt.batch_size % WORLD_SIZE == 0, f"--batch-size {opt.batch_size} must be multiple of WORLD_SIZE"
    assert torch.cuda.device_count() > LOCAL_RANK, "insufficient CUDA devices for DDP command"
    torch.cuda.set_device(LOCAL_RANK)
    device = torch.device("cuda", LOCAL_RANK)
    # 初始化 PyTorch 分布式通信组
    # Linux 下一般使用 NCCL
    # Windows 下一般只能使用 Gloo
    dist.init_process_group(
        backend="nccl" if dist.is_nccl_available() else "gloo", timeout=timedelta(seconds=10800)
    )

        3.1.2、把普通 BatchNorm 转成多卡同步版 BatchNorm,让所有GPU一起统计BN均值和方差,而不是每张GPU各算各的

# SyncBatchNorm
if opt.sync_bn and cuda and RANK != -1:
    model = torch.nn.SyncBatchNorm.convert_sync_batchnorm(model).to(device)
    LOGGER.info("Using SyncBatchNorm()")

        3.1.3、将普通模型包装成分布式模型,主要操作是DistributedDataParallel(model),smart_DDP处理了一些兼容性问题,核心是调用DistributedDataParallel(model),这样让模型在backward时触发梯度同步all_reduce(梯度),DDP 自动 hook 了 backward:loss.backward(),这样在scaler.scale(loss).backward()执行时就会触发all_reduce(梯度)。

if cuda and RANK != -1:
    model = smart_DDP(model)


# =========================
# 反向传播
# DDP 会在这里自动 all_reduce 同步梯度
# =========================
scaler.scale(loss).backward()

        3.1.4、告诉 DistributedSampler: 现在进入新的 epoch 了,请重新 shuffle 数据

if RANK != -1:
    train_loader.sampler.set_epoch(epoch)

        3.1.5、修正 DDP “梯度平均”行为,AllReduce时做了梯度平均,这里补回来

if RANK != -1:
    loss *= WORLD_SIZE  # gradient averaged between devices in DDP mode

        3.1.6、在 DDP 多卡训练中,同步是否停止训练,让所有GPU一起停止,避免有的GPU停了有的GPU还在继续训练

# EarlyStopping
if RANK != -1:  # if DDP training
    broadcast_list = [stop if RANK == 0 else None]
    dist.broadcast_object_list(broadcast_list, 0)  # broadcast 'stop' to all ranks
    if RANK != 0:
        stop = broadcast_list[0]

        3.1.7、yolov5是数据并行模型,每张GPU有完整模型,将数据进行切片,每张GPU训练自己的切片后的数据集,如果整个数据集是10240张图片,两张GPU(WORLD_SIZE),那么切片后就是每个GPU读取各自的5120张图片,rank=LOCAL_RANK表示本地GPU读取数据的位置,数据读取后存在数据集dataset里,即这里的数据集为5120张图片;
        create_dataloader函数内部会调用DistributedSampler进行数据切片;
        如果batch_size=16,则表示每次训练,所有的GPU总共训练16张图片,如果是两张GPU,则每张GPU训练8张图片;
        训练有两个循环,epoch表示对全部数据集训练的轮数,batch表示GPU每次训练,一次batch训练16张图片,需要640次batch才能完成一次全数据集的训练:
                for epoch in range(start_epoch, epochs):  # epoch
                        for i, (imgs, targets, paths, _) in pbar:  # batch

# =========================
# 创建训练集 DataLoader
# 内部会自动创建 DistributedSampler
# DDP 多卡时每张卡只读取部分数据
# =========================
train_loader, dataset = create_dataloader(
    train_path,
    imgsz,
    batch_size // WORLD_SIZE,
    gs,
    single_cls,
    hyp=hyp,
    augment=True,
    cache=None if opt.cache == "val" else opt.cache,
    rect=opt.rect,
    rank=LOCAL_RANK,
    workers=workers,
    image_weights=opt.image_weights,
    quad=opt.quad,
    prefix=colorstr("train: "),
    shuffle=True,
    seed=opt.seed,
)

        3.1.8、创建模型,可以加载预训练模型(即加载训练过的.pt文件继续训练,以加快收敛速度),也可以通过cfg从0创建模型

# =========================
# 构建模型
# 支持:
# 1. 加载预训练模型
# 2. 从 cfg 创建新模型
# =========================
check_suffix(weights, ".pt")  # check weights
pretrained = weights.endswith(".pt")
if pretrained:
    with torch_distributed_zero_first(LOCAL_RANK):
        weights = attempt_download(weights)  # download if not found locally
    ckpt = torch_load(weights, map_location="cpu")  # load checkpoint to CPU to avoid CUDA memory leak
    model = Model(cfg or ckpt["model"].yaml, ch=3, nc=nc, anchors=hyp.get("anchors")).to(device)  # create
    exclude = ["anchor"] if (cfg or hyp.get("anchors")) and not resume else []  # exclude keys
    csd = ckpt["model"].float().state_dict()  # checkpoint state_dict as FP32
    csd = intersect_dicts(csd, model.state_dict(), exclude=exclude)  # intersect
    model.load_state_dict(csd, strict=False)  # load
    LOGGER.info(f"Transferred {len(csd)}/{len(model.state_dict())} items from {weights}")  # report
else:
    model = Model(cfg, ch=3, nc=nc, anchors=hyp.get("anchors")).to(device)  # create
amp = check_amp(model)  # check AMP

        3.1.9、单线程多GPU模式训练,这里的DataParallel()跟分布式没关系,不是指DDP

# DP mode
if cuda and RANK == -1 and torch.cuda.device_count() > 1:
    LOGGER.warning(
        "WARNING ⚠️ DP not recommended, use torch.distributed.run for best DDP Multi-GPU results.\n"
        "See Multi-GPU Tutorial at https://docs.ultralytics.com/yolov5/tutorials/multi_gpu_training to get started."
    )
    model = torch.nn.DataParallel(model)

4、DDP(DistributedDataParallel)

        4.1、DDP在代码里是什么地方进行梯度同步的?DDP是怎么调用NCCL的?
        详见3.1.3章节。
        4.2、梯度同步的依据是什么,每个rank怎么知道自己的切分数据?
        torchrun时带入了总的rank数和当前rank值,torchrun会设置LOCAL_RANK/RANK/WORLD_SIZE这三个环境变量,DDP从这些环境变量获取的。

5、DistributedSampler

        5.1、DistributedSampler在代码里哪里进行切分数据的?
        详见3.1.7章节。
        5.2、DistributedSampler切分数据的依据是什么?每个rank怎么知道取那部分数据的?切分数据需要传输吗?用什么传输的(普通网卡还是RDMA)?
        切分数据时会带入总的rank数和当前rank值;切分数据不用传输,每台主机都有全部的样本,nccl也只会传递梯度;用什么传输是在ncclCommInitRank()时计算出来的GPU与GPU之间的最快传输硬件,ncclAllReduce()时会在GPU与GPU已经建立硬件通路情况下寻找系统级的最优路线。

6、MPI

        MPI会使用免密SSH登录其他主机,用MPI启动程序只需在一台主机上执行启动命令即可,MPI会自动登录其他主机执行启动命令。

7、NCCL

        下面是一个4机4GPU,每机1GPU环境下使用nccl做AllReduce的例程:
        4机4卡nccl allreduce例程

        7.1、编译代码:
make NCCL_HOME=$HOME/Desktop/NCCL_Source/nccl/build  # 路径为自己编译nccl库的路径
        7.2、确保4台机器之间SSH免密
        7.3、配置环境变量:
export LD_LIBRARY_PATH=/usr/local/cuda/lib64:/usr/local/lib:$LD_LIBRARY_PATH
export LD_LIBRARY_PATH=$HOME/Desktop/NCCL_Source/nccl/build/lib:/usr/local/cuda/lib64:$LD_LIBRARY_PATH  # 如果为自己编译nncl库
        7.4、运行程序(其中host0为对应的主机IP,使用rdma网络):
mpirun -np 4 \
  -H host0,host1,host2,host3 \
  -x LD_LIBRARY_PATH \
  -x NCCL_ALGO=Ring \
  -x NCCL_IB_DISABLE=0 \
  -x NCCL_IB_HCA=mlx5_0 \
  -x NCCL_IB_GID_INDEX=3 \
  -x NCCL_SOCKET_IFNAME=eth0 \
  -x NCCL_NET_GDR_LEVEL=2 \
  -x NCCL_DEBUG=INFO \
  -x NCCL_DEBUG_SUBSYS=INIT,NET,GRAPH \
  ./nccl_ring_allreduce 1048576
        7.5、结果检查(初始值:rank0/1/2/3=1/2/3/4,结果应该是每个gpu都得到sum=10):

                log显示:Rank 0 PASS: count=1048576 expected=10.0 first=10.0 last=10.0

        7.6、4卡执行All-Reduce数据流分析
  • 4卡执行AllReduce时,使用ring拓扑结构,nccl识别出如下的ring拓扑结构,数据按固定方向传输:

  • AllReduce分为Reduce-Scatter和All-Gather两个步骤,4卡AllReduce,每卡数据被分成4个trunk,同一时间每张卡传输一个trunk的数据,流程如下,共需3步Reduce-Scatter和3步all-gather:

        7.7、ncclAllReduce函数说明,从接口可以看出,每个GPU拥有完整的梯度,即sendbuff,recvbuff为完成同步的梯度:
ncclResult_t ncclAllReduce(
    const void* sendbuff,    // 发送缓冲区,可以与接收缓冲区同一地址
    void* recvbuff,        // 接收缓冲区,可以与发送缓冲区同一地址
    size_t count,        // 数据个数
    ncclDataType_t datatype,    // 数据类型:ncclFloat/ncclHalf/ncclDouble/ncclInt/ncclBfloat16
    ncclRedOp_t op,        // 规约操作:ncclSum/ncclProd/ncclMax/ncclMin/ncclAvg
    ncclComm_t comm,        // 通信器:由rank0 ncclCommInitRank()创建后广播给其他rank
    cudaStream_t stream);

        

        重点:熟悉该例程代码流程、熟悉nccl API(包括详细参数)、简要熟悉nccl源码

        7.8、nccl源码分析

        7.8.1、nccl源码文件架构:

nccl-master/
  ├── .github/        # GitHub CI、issue/PR 配置,和源码学习关系不大
  ├── bindings/       # 语言绑定,比如 Python/Cython 接口,把 NCCL C API 包装给其他语言用
  ├── cmake/          # CMake 构建脚本
  ├── contrib/        # 贡献/扩展代码,现在包含 nccl_ep 等实验或扩展功能
  ├── docs/           # 文档和示例,适合先跑 demo
  ├── makefiles/      # Makefile 构建规则,比如 CUDA 架构、编译选项
  ├── pkg/            # 打包相关,生成 deb/rpm/tar 包
  ├── plugins/        # NCCL 插件接口/插件实现,如网络插件、tuner 等
  ├── src/            # NCCL 核心源码,最重要
      ├── init.cc       # ncclCommInitRank 等初始化
      ├── group.cc      # ncclGroupStart / ncclGroupEnd
      ├── enqueue.cc    # collective 任务入队
      ├── collectives/  # allreduce、broadcast 等集合通信入口
      ├── device/       # GPU kernel/device 侧代码
      ├── transport/    # P2P、SHM、NET 等传输层
      ├── graph/        # 拓扑分析、ring/tree/channel 构造
      └── proxy.cc      # CPU proxy 线程,处理网络/异步通信

        7.8.2、核心源文件:

include/API 思路
↓
src/init.cc
↓
src/group.cc
↓
src/enqueue.cc
↓
src/collectives/
↓
src/transport/
↓
src/device/

        7.8.3、使用nccl进行AllReduce操作,需要执行ncclCommInitRank()和ncclAllReduce()两个函数。

        ncclCommInitRank():

ncclCommInitRank()
文件:src/init.cc
作用:用户入口 API,检查参数,创建 communicator 初始化任务
  ↓
  ncclCommInitRankDev()
  文件:src/init.cc
  作用:确认当前 CUDA device,处理 device/rank 相关初始化
    ↓
    ncclAsyncInit()
    文件:src/group.cc / src/include/group.h
    作用:把初始化过程包装成异步任务,支持 group 语义
      ↓
      ncclCommInitRankFunc()
      文件:src/init.cc
      作用:真正执行 communicator 初始化的核心函数
        ↓
        commAlloc()
        文件:src/init.cc
        作用:分配 ncclComm 结构体,记录 rank、nranks、cudaDev、busId 等基础信息
          ↓
          bootstrapInit()
          文件:src/bootstrap.cc
          作用:建立 bootstrap 控制通道,让不同 rank 能互相发现和交换信息
            ↓
            initTransportsRank()
            文件:src/init.cc
            作用:初始化 rank 间通信资源,是整个初始化最核心阶段之一
              ↓
              bootstrapAllGather()
              文件:src/bootstrap.cc
              作用:所有 rank 交换 peerInfo、rank、busId、hostHash、pidHash、CUDA 计算能力等信息
                ↓
                ncclTopoGetSystem()
                文件:src/graph/topo.cc
                作用:获取本机 GPU、CPU、PCIe、NVLink、NIC 等硬件拓扑
                  ↓
                  ncclTopoCompute() ──────────────────────────────────────────
                  文件:src/graph/search.cc                                   │
                  作用:根据拓扑搜索通信路径,决定 ring、tree、channel 等通信结构     ↓
                    ↓                                                        selectTransport()
                    bootstrapAllGather()                                     文件:src/transport.cc
                    文件:src/bootstrap.cc                                    作用:为每条 peer 连接选择具体 transport(选择GPU之间的硬件连接方式,P2P/NVLink/PCIeP2P/SHM/NET/RDMA)。
                    作用:再次交换 graphInfo、channel 数量、拓扑 rank 映射等信息
                      ↓
                      ncclTransportP2pSetup()
                      文件:src/transport/p2p.cc
                      作用:建立 GPU-GPU P2P 通信连接,例如 NVLink 或 PCIe P2P
                        ↓
                        ncclTransportShmSetup()
                        文件:src/transport/shm.cc
                        作用:建立同机 GPU 间共享内存通信路径
                          ↓
                          ncclTransportNetSetup()
                          文件:src/transport/net.cc
                          作用:建立跨机器网络通信路径,例如 socket、IB、RoCE
                            ↓
                            ncclProxyCreate()
                            文件:src/proxy.cc
                            作用:创建 proxy 线程,用于辅助网络通信、异步进度推进
                              ↓
                              ncclProxyStart()
                              文件:src/proxy.cc
                              作用:启动 proxy 服务线程
                                ↓
                                返回 ncclComm_t
                                作用:communicator 创建完成,后续 ncclAllReduce、ncclBroadcast 等 collective 可以使用它

        7.8.4、ncclAllReduce():

ncclAllReduce()
文件:src/collectives.cc
作用:AllReduce 用户接口入口,创建 collective 通信任务。
  ↓
  ncclEnqueueCheck()
  文件:src/enqueue.cc
  作用:检查参数并把 AllReduce 封装为 task。
    ↓
    ncclGroupStartInternal()
    文件:src/group.cc
    作用:进入 NCCL group 管理模式。
      ↓
      taskAppend()
      文件:src/enqueue.cc
      作用:把 AllReduce task 加入 communicator 队列。
        ↓
        ncclSaveKernel()
        文件:src/enqueue.cc
        作用:生成 collective kernel 的工作描述信息。
          ↓
          ncclGroupEndInternal()
          文件:src/group.cc
          作用:结束 group 并开始提交 collective task。
            ↓
            ncclPrepareTasks()
            文件:src/enqueue.cc
            作用:切分 chunk、分配 channel、准备 work queue。
              ↓
              computeColl()
              文件:src/enqueue.cc
              作用:选择 collective 算法与通信协议(这两步搜索最优路径)。
                ↓
                ncclTopoGetAlgoTime()
                文件:src/graph/tuning.cc
                作用:评估 Ring/Tree/协议 的预计通信耗时(这两步搜索最优路径)。
                  ↓
                  scheduleCollTasksToPlan()
                  文件:src/enqueue.cc
                  作用:把 collective task 安排到 launch plan。
                    ↓
                    ncclLaunchPrepare()
                    文件:src/launch.cc
                    作用:准备 CUDA kernel launch 参数。
                      ↓
                      ncclProxySaveColl()
                      文件:src/proxy.cc
                      作用:为网络通信创建 proxy 通信任务。
                        ↓
                        ncclLaunchKernel()
                        文件:src/launch.cc
                        作用:向 CUDA stream 提交 NCCL GPU kernel。
                          ↓
                          ├─ 同机通信(优先 P2P/NVLink/PCIeP2P,必要时 SHM/NET,主要 GPU 运行,CPU 可能并行协作):
                          │     ↓
                          │     ncclKernel_Main()
                          │     文件:src/device/kernel.cc
                          │     作用:GPU collective kernel 统一入口。
                          │       ↓
                          │       runRing() / runTreeUpDown()
                          │       文件:src/device/all_reduce.h
                          │       作用:执行 Ring 或 Tree AllReduce 算法流程。
                          │         ↓
                          │         prims.send/recvReduceSend
                          │         文件:src/device/prims_*.h
                          │         作用:GPU 侧执行发送、接收、规约,底层 transport 可能是 P2P/SHM/NET。
                          │
                          └─ 跨机 RDMA/NET 通信(CPU 与 GPU 并行协作):
                              ↓
                              ├─ CPU 运行:
                              │    ↓
                              │    proxyService / proxyProgress
                              │    文件:src/proxy.cc
                              │    作用:CPU proxy 推进网络通信。
                              │      ↓
                              │      netSendProxy / netRecvProxy
                              │      文件:src/transport/net.cc
                              │      作用:网络 transport 发送/接收。
                              │        ↓
                              │        ncclNetIsend / ncclNetIrecv
                              │        文件:src/net.cc / src/include/net.h
                              │        作用:调用网络后端接口。
                              │          ↓
                              │          ibIsend / ibIrecv
                              │          文件:src/transport/net_ib.cc
                              │          作用:IB/RoCE 后端发起异步传输。
                              │            ↓
                              │            ibv_post_send()
                              │            文件:libibverbs
                              │            作用:真正提交 RDMA 请求给网卡。
                              │
                              └─ GPU 运行:
                                  ↓
                                  ncclKernel_Main()
                                  文件:src/device/kernel.cc
                                  作用:GPU collective kernel 统一入口。
                                    ↓
                                    runRing() / runTreeUpDown()
                                    文件:src/device/all_reduce.h
                                    作用:执行 Ring 或 Tree AllReduce 算法流程。
                                      ↓
                                      prims.send / prims.recvReduceSend / prims.directRecvReduceCopySend
                                      文件:src/device/prims_simple.h
                                           src/device/prims_ll.h
                                           src/device/prims_ll128.h
                                      作用:执行 GPU 侧发送、接收、规约、拷贝。
                                        ↓
                                        如果是跨机通信:
                                        GPU prims 通过 step/FIFO 与 proxy 协作

        7.8.5、GPU与GPU之间使用哪种硬件通路,是在构建拓扑时确定的,即ncclCommInitRank()阶段,详细实现在src/transport.cc、src/transport/下,策略如下所示:

ncclCommInitRank 阶段:
  决定 GPU 之间用什么路径/transport:
  NVLink、PCIe P2P、SHM、NET/RDMA
  选择硬件路径的策略如下:
    ↓
    ├─ 同机 GPU-GPU:优先 P2P / NVLink / PCIe P2P
    │    ↓
    │    同机但 P2P 不可用:可能走 SHM
    │  
    └─ 跨机器 GPU-GPU:走 NET transport
         ↓
         如果是 IB/RoCE:src/transport/net_ib.cc,最终走 RDMA

        7.8.6、拓扑阶段确定了GPU与GPU之间使用的硬件通路,选择最优路径则在ncclAllreduce()阶段实现:

ncclAllReduce 阶段:
  决定本次 collective 用什么算法/协议:
  Ring、Tree、CollNet、NVLS
  Simple、LL、LL128
  即寻找最优路径

nccl源码分析(一)——数据接收和发送的处理
nccl分析(二)——RDMA带外建链过程
nccl分析(三)——GPU-Initiated Networking(gin)数据发送过程分析
NCCL源码详解1:NCCL官网使用/调用案例 Example : One Device per Process or Thread包含视频教程
NCCL源码详解2:通信初始化如何获取唯一ID UniqueId,ncclGetUniqueId()中ncclInit()、bootstrapGetUniqueId()包含视频教程
NCCL源码详解3:通信器初始化ncclCommInitRank() 含视频教程
NCCL源码解读3.1:double binary tree双二叉树构建算法,相比ring环算法的优势
NCCL源码详解4:bootstrapInit()引导网络bootstrap网络连接建立 视频教程
NCCL源码解读5:拓扑识别感知整体思路总览
NCCL源码详解6:通信拓扑识别感知构建 物理拓扑xml文件 ncclTopoGetSystem() 视频教程
NCCL源码解读1:官网案例详解 单进程单设备使用/调用案例
NCCL源码解析:建图过程
NCCL源码解析:路径计算
NCCL与RDMA和MPI基本框架源码分析

8、RDMA驱动框架

1、驱动框架:

网卡相关的模块包括以下3个内容:

  • 用户态驱动:libmlx5.so
  • 内核态驱动:mlx5_ib.ko、mlx5_core.ko
2、数据流
3、GPU-Direct RDMA

        非GPU-Direct RDMA时,代码如下:

buf = malloc(size);
ibv_reg_mr(pd, buf, size, ...);
ibv_post_send(...)

        当RDMA传输数据在host内存时,应用层使用malloc()函数申请现存,内核mlx5_ib.ko模块处理不连续性问题(如果连续不用处理),将虚拟地址转换为NIC可直接读取的物理地址。

        GPU-Direct RDMA时,代码如下:

cudaMalloc(&buf, size);
ibv_reg_mr(pd, buf, size, ...);
ibv_post_send(...)

        当RDMA传输数据在显存的时候,就实现了GPU-Direct RDMA,NIC直接读取显存数据进行发送,应用层使用cudaMalloc()函数申请现存,内核mlx5_ib.ko模块识别到这个地址为现存时,转换为NIC可直接读取的物理地址。

4、mlx5源码分析
  • libmlx5.so

        核心源码:rdma-core-master/providers/mlx5/mlx5.c

        libmlx5.so作为用户层驱动,是一个provider(或者叫HAL),核心操作就是使用PROVIDER_DRIVER(mlx5, mlx5_dev_ops)向libibverbs.so注册provider,其中提供的ops包含4个重要函数,如下所示:

        1、注册provider:

static const struct verbs_device_ops mlx5_dev_ops = {
	.name = "mlx5",
	.match_min_abi_version = MLX5_UVERBS_MIN_ABI_VERSION,
	.match_max_abi_version = MLX5_UVERBS_MAX_ABI_VERSION,
	.match_table = mlx5_hca_table,
	.alloc_device = mlx5_device_alloc,
	.uninit_device = mlx5_uninit_device,
	.alloc_context = mlx5_alloc_context,
	.import_context = mlx5_import_context,
};

PROVIDER_DRIVER(mlx5, mlx5_dev_ops);

        2、mlx5_device_alloc()函数:使用PROVIDER_DRIVER(mlx5, mlx5_dev_ops)向libibverbs.so注册provider时,libibverbs.so会调用ibv_get_device_list()函数扫描/dev/infiniband/uverbsX进行device枚举,流程如下,最终mlx5_device_alloc()函数会被调用,创建ibv_device。

PROVIDER_DRIVER(mlx5, mlx5_dev_ops);
libibverbs.so设备枚举执行:
ibv_get_device_list()
    ↓
    device->ops->device_alloc()		// libibverbs.so
        ↓
		mlx5_device_alloc()			// libmlx5.so
		分配并初始化 mlx5_device

        3、mlx5_alloc_context()函数:当用户层执行ibv_open_device()时,会调用到该函数,libibverbs.so会将打开的/dev/infiniband/uverbs0的fd传到该函数,该函数创建verbs_contex、mlx5_context,将verbs_contex和fd保存到mlx5_context,将verbs_context_ops挂载到verbs_contex、将mlx5_dv_context_ops挂载到mlx5_context,返回verbs_contex给libibverbs.so(verbs_init_and_alloc_context()函数将其赋值给ibv_device),这样libibverbs.so就获得了verbs_context_ops,libmlx5.so也获得了fd(ioctl/mmap用),实际上新版代码中libibverbs.so将ioctl进行了封装,提供封装ioctl的ibv_cmd_*接口。

ibv_open_device()					// 用户层
    ↓
	fd=open(/dev/infiniband/uverbsX)
    device->ops->alloc_context()	// libibverbs.so
        ↓
		mlx5_alloc_context()		// libmlx5.so
        1.mmap UAR/doorbell/capabilities/初始化DV/DevX
		2.创建verbs_context、mlx5_context
		3.mlx5_context保存verbs_context和fd(ioctl)
		4.将verbs_context_ops挂在到verbs_context
		5.将mlx5_dv_context_ops挂在到mlx5_context

        mlx5_alloc_context()函数使用verbs_init_and_alloc_context()申请verbs_context,将其初始化,然后通过ibv_cmd_get_context()函数使用ioctl的IB_USER_VERBS_CMD_GET_CONTEXT将verbs_context传递给内核,让内核创建内核mlx5_ib_ucontext/ib_ucontext,并获取内核返回参数,这样libmlx5.so才与内核产生了联系。

调用关系大概是:
ibv_open_device()
  ↓
mlx5_alloc_context()
  ↓
verbs_init_and_alloc_context()
  ↓
ibv_cmd_get_context()
  ↓ ioctl 到内核
ib_uverbs
  ↓
mlx5_ib_alloc_ucontext()

主要做几件事:
1. 填充 GET_CONTEXT 命令
2. 通过 context->cmd_fd 发送 ioctl
3. 让内核创建 ucontext
4. 接收内核返回的 response
5. provider 根据 response 保存能力信息、mmap 偏移、UAR/BF 参数等

        4、ibv_open_device_from_fd()函数:

ibv_open_device_from_fd()			// 用户层
    ↓
    device->ops->import_context()	// libibverbs.so
		↓
		mlx5_import_context()		// libmlx5.so
		复用已有context

        5、ibv_close_device()函数:

ibv_close_device()					// 用户层
    ↓
	device->ops->uninit_context()	// libibverbs.so
		↓
		mlx5_uninit_device()		// libmlx5.so

        到此,用户调用libibverbs.so的ibv_*函数(ibv_create_cq()、ibv_reg_mr()、ibv_alloc_pd())的访问就会直接通过调用verbs_context_ops到达provider libmlx5.so,厂家自定义的高性能接口mlx5dv_*也能直接到达libmlx5.so。

        但是libmlx5.so又怎么与内核交互呢,分两种情况,一种是linux内核提供的rdma标准ABI,另外一种是厂商提供的高性能接口。

        1、rdma标准ABI:使用libibverbs.so传给libmlx5.so的fd(/dev/infiniband/uverbs0)的ioctl/mmap,其中ioctl用于控制类的慢路径,mmap用于数据类的快路径,慢路径ioctl会经过内核ib_uverbs.ko模块,快路径则会绕过ib_uverbs.ko,用mmap直接操作硬件缓冲区或寄存器,如下表所示:

        在新版代码里libibverbs.so将ioctl封装成了ibv_cmd_*接口函数,ibv_cmd_*函数完成参数填充,再调用ioctl(),ibv_cmd_*与ibv_*接口及ioctl枚举值的对应关系如上述表格。ibv_cmd_*的调用过程如下:

大致如下:
mlx5_create_qp()
  ↓
ibv_cmd_create_qp()
  ↓
execute_ioctl()
  ↓
ioctl(context->cmd_fd, RDMA_VERBS_IOCTL, command_buffer)
  ↓
ib_uverbs.ko

        2、厂商提供的高性能接口:不使用或部分使用libibverbs.so传给libmlx5.so的fd(/dev/infiniband/uverbs0),更多是用厂家内核驱动创建的高性能虚拟文件接口实现。

        rdma标准ABI厂商提供的高性能接口关系如下:

应用
 ├─ ibv_create_qp()
 │    ↓
 │  libibverbs
 │    ↓
 │  verbs_context_ops.create_qp
 │    ↓
 │  mlx5_create_qp()
 │
 └─ mlx5dv_devx_obj_create()
      ↓
    mlx5dv API
      ↓
    mlx5_dv_context_ops.devx_obj_create
      ↓
    mlx5_devx_obj_create()

        libmlx5.so中两个重要的ops:verbs_context_opsmlx5_dv_context_ops,前者作为ibv_*标准接口的承接者会传递给libibverbs.so,后者不属于标准数据,仅在libmlx5.so内部,两个ops定义如下。

        verbs_context_ops:承接rdma标准ibv_*接口,会传递到libibverbs.so,使用/dev/infiniband/uverbs0的ioctl/mmap

static const struct verbs_context_ops mlx5_ctx_common_ops = {
	.query_port    = mlx5_query_port,
	.alloc_pd      = mlx5_alloc_pd,
	.async_event   = mlx5_async_event,
	.dealloc_pd    = mlx5_free_pd,
	.reg_mr	       = mlx5_reg_mr,
	.reg_dmabuf_mr = mlx5_reg_dmabuf_mr,
	.rereg_mr      = mlx5_rereg_mr,
	.dereg_mr      = mlx5_dereg_mr,
	.alloc_mw      = mlx5_alloc_mw,
	.dealloc_mw    = mlx5_dealloc_mw,
	.bind_mw       = mlx5_bind_mw,
	.create_cq     = mlx5_create_cq,
	.poll_cq       = mlx5_poll_cq,
	.req_notify_cq = mlx5_arm_cq,
	.cq_event      = mlx5_cq_event,
	.resize_cq     = mlx5_resize_cq,
	.destroy_cq    = mlx5_destroy_cq,
	.create_srq    = mlx5_create_srq,
	.modify_srq    = mlx5_modify_srq,
	.query_srq     = mlx5_query_srq,
	.destroy_srq   = mlx5_destroy_srq,
	.post_srq_recv = mlx5_post_srq_recv,
	.create_qp     = mlx5_create_qp,
	.query_qp      = mlx5_query_qp,
	.modify_qp     = mlx5_modify_qp,
	.destroy_qp    = mlx5_destroy_qp,
	.post_send     = mlx5_post_send,
	.post_recv     = mlx5_post_recv,
	.create_ah     = mlx5_create_ah,
	.destroy_ah    = mlx5_destroy_ah,
	.attach_mcast  = mlx5_attach_mcast,
	.detach_mcast  = mlx5_detach_mcast,

	.advise_mr = mlx5_advise_mr,
	.alloc_dm = mlx5_alloc_dm,
	.alloc_parent_domain = mlx5_alloc_parent_domain,
	.alloc_td = mlx5_alloc_td,
	.attach_counters_point_flow = mlx5_attach_counters_point_flow,
	.close_xrcd = mlx5_close_xrcd,
	.create_counters = mlx5_create_counters,
	.create_cq_ex = mlx5_create_cq_ex,
	.create_flow = mlx5_create_flow,
	.create_flow_action_esp = mlx5_create_flow_action_esp,
	.create_qp_ex = mlx5_create_qp_ex,
	.create_rwq_ind_table = mlx5_create_rwq_ind_table,
	.create_srq_ex = mlx5_create_srq_ex,
	.create_wq = mlx5_create_wq,
	.dealloc_td = mlx5_dealloc_td,
	.destroy_counters = mlx5_destroy_counters,
	.destroy_flow = mlx5_destroy_flow,
	.destroy_flow_action = mlx5_destroy_flow_action,
	.destroy_rwq_ind_table = mlx5_destroy_rwq_ind_table,
	.destroy_wq = mlx5_destroy_wq,
	.free_dm = mlx5_free_dm,
	.get_srq_num = mlx5_get_srq_num,
	.import_dm = mlx5_import_dm,
	.import_mr = mlx5_import_mr,
	.import_pd = mlx5_import_pd,
	.modify_cq = mlx5_modify_cq,
	.modify_flow_action_esp = mlx5_modify_flow_action_esp,
	.modify_qp_rate_limit = mlx5_modify_qp_rate_limit,
	.modify_wq = mlx5_modify_wq,
	.open_qp = mlx5_open_qp,
	.open_xrcd = mlx5_open_xrcd,
	.post_srq_ops = mlx5_post_srq_ops,
	.query_device_ex = mlx5_query_device_ex,
	.query_ece = mlx5_query_ece,
	.query_rt_values = mlx5_query_rt_values,
	.read_counters = mlx5_read_counters,
	.reg_dm_mr = mlx5_reg_dm_mr,
	.alloc_null_mr = mlx5_alloc_null_mr,
	.free_context = mlx5_free_context,
	.set_ece = mlx5_set_ece,
	.unimport_dm = mlx5_unimport_dm,
	.unimport_mr = mlx5_unimport_mr,
	.unimport_pd = mlx5_unimport_pd,
	.query_qp_data_in_order = mlx5_query_qp_data_in_order,
	.alloc_dmah = mlx5_alloc_dmah,
	.dealloc_dmah = mlx5_dealloc_dmah,
	.reg_mr_ex = mlx5_reg_mr_ex,
	.query_port_speed = mlx5_query_port_speed,
	.dm_export_dmabuf_fd = mlx5_dm_export_dmabuf_fd,
};

        mlx5_dv_context_ops:承接厂家高性能接口,使用厂家内核驱动创建的虚拟文件

static struct mlx5_dv_context_ops mlx5_dv_ctx_ops = {
	.query_device = _mlx5dv_query_device,

	.query_qp_lag_port = _mlx5dv_query_qp_lag_port,
	.modify_qp_lag_port = _mlx5dv_modify_qp_lag_port,

	.modify_qp_udp_sport = _mlx5dv_modify_qp_udp_sport,

	.sched_node_create = _mlx5dv_sched_node_create,
	.sched_leaf_create = _mlx5dv_sched_leaf_create,
	.sched_node_modify = _mlx5dv_sched_node_modify,
	.sched_leaf_modify = _mlx5dv_sched_leaf_modify,
	.sched_node_destroy = _mlx5dv_sched_node_destroy,
	.sched_leaf_destroy = _mlx5dv_sched_leaf_destroy,
	.modify_qp_sched_elem = _mlx5dv_modify_qp_sched_elem,

	.reserved_qpn_alloc = _mlx5dv_reserved_qpn_alloc,
	.reserved_qpn_dealloc = _mlx5dv_reserved_qpn_dealloc,

	.set_context_attr = _mlx5dv_set_context_attr,
	.get_clock_info = _mlx5dv_get_clock_info,
	.init_obj = _mlx5dv_init_obj,
};

        mlx5dv_*接口实现的功能有:

        BlueFlame:mlx5 的低延迟 doorbell 发送机制。用户态可直接把小 WQE 通过 MMIO 写入 NIC 的 BlueFlame buffer,减少 PCIe DMA 读取延迟,实现超低延迟发送;
        DevX:“Device eXtended API”。允许用户态直接创建和操作 mlx5 firmware object(硬件对象),绕过标准 verbs 抽象,直接下发 firmware command;
        flow steering:硬件流表/流量转发引擎。NIC 可以在硬件里按五元组、隧道头、metadata 等匹配报文,并执行转发、修改、丢弃、镜像等动作;
        packet pacing:报文发送速率控制。NIC 按设定速率平滑发包,避免瞬时 burst,常用于 RoCE 拥塞控制、QoS、AI 集群限速;
        hairpin queue网卡内部回环队列。报文无需经过 CPU/host memory,直接在 NIC 内部从一个 queue 转发到另一个 queue,用于 OVS/DPDK/switch offload;
        UMR:Unified Memory Region。mlx5 的高级内存注册机制,可动态组合多个 MR,支持 scatter/gather memory layout,减少频繁 reg_mr 开销;
        dynamic UAR:动态 UAR(User Access Region)分配机制。允许用户态动态申请/释放 doorbell MMIO 区域,提高多线程和高并发 doorbell 的效率;
        ASO:Access Steering Object。mlx5 硬件加速对象,用于 flow meter、connection tracking、状态维护等高级 flow steering 功能;
        direct WQE layout:直接操作 mlx5 WQE(Work Queue Element)硬件格式。应用可绕过通用 ibverbs 封装,直接写硬件 WQE,实现更低延迟、更高性能。

        mlx5dv_*接口使用的虚拟文件系统有:

        /dev/infiniband/uverbsX:核心控制通道(cmd_fd);
        /sys/class/infiniband/mlx5_X/:device discovery;
        /sys/class/infiniband_verbs/uverbsX:uverbs 设备对应关系;
        mmap(uverbs fd):UAR/CQ/BF/doorbell 映射。

总结来讲,RDMA的用户态驱动需要完成的功能是:

  1. 向上:向libibverbs.so注册provider;向libibverbs.so提供verbs_device_ops结构;在rdma标准ABI方面为libibverbs.so提供verbs_context_ops结构;在厂家自定义高性能API方面为用户提供mlx5_dv_context_ops结构;
  2. 向下:在rdma标准ABI方面使用libibverbs.so打开的/dev/infiniband/uverbsX的描述符fd用来进行ioctl/mmap;在厂家自定义高性能API方面既使用/dev/infiniband/uverbsX的描述符fd,又是用mlx5内核驱动(mlx5_ib.ko/mlx5_core.ko)创建的其他虚拟文件
  • mlx5_ib.ko和mlx5_core.ko

        内核态驱动核心任务是使用注册ib_register_device()函数注册ib_device,ib_device里带有ib_device_opsib_device_ops承载所有的rdma操作。

        mlx内核态驱动分为了mlx5_ib和mlx5_core两层,mlx5_ib负责注册ib_device等接口类功能,mlx5_core负责枚举PCIe、读写mlx网卡硬件寄存器等操作,ib向core注册interface,core枚举到PCIe网卡后会创建mlx5_core_dev,并触发interface->add()执行,add()实际就是ib里的ib_register_device()流程了,如下所示:

1、CORE注册PCIe驱动:
————————————————————————————————————————————————————————
module_init()
  ↓
  mlx5_init()
    ↓
    pci_register_driver(&mlx5_core_driver)

2、IB向CORE注册interface:
————————————————————————————————————————————————————————
module_init()
  ↓
  mlx5_ib_init()
    ↓
    mlx5_register_interface(&mlx5_ib_interface)

3、CORE枚举到PCIe网卡:
————————————————————————————————————————————————————————
static struct pci_driver mlx5_core_driver = {
    .name     = "mlx5_core",
    .id_table = mlx5_core_pci_table,
    .probe    = init_one,
    .remove   = remove_one,
};

init_one()                            // 识别到PCIe卡
  ↓
  mlx5_pci_probe()
    ↓
    ├─ pci_enable_device()            // enable PCI
    ├─ dma_set_mask_and_coherent()    // DMA mask
    ├─ ioremap()                      // BAR mapping
    ├─ mlx5_cmd_init()                // FW init
    ├─ mlx5_query_hca_caps()          // FW init
    ├─ mlx5_eq_init()                 // EQ/MSI-X
    ├─ mlx5_init_once()               // UAR
    │
    └  创建 mlx5_core_dev              // 触发执行interface->add()

4、IB执行interface->add():
————————————————————————————————————————————————————————
struct mlx5_interface mlx5_ib_interface = {
    .add    = mlx5_ib_add,
    .remove = mlx5_ib_remove,
};

struct mlx5_ib_dev {
    struct ib_device ib_dev;
    struct mlx5_core_dev *mdev;
};

mlx5_ib_add(struct mlx5_core_dev *mdev)
  ↓
  ├─ ibdev = ib_alloc_device(struct mlx5_ib_dev, ib_dev);    // 把mlx5_core_dev绑定到ib_device
  ├─ ibdev->ib_dev.ops = &mlx5_ib_dev_ops;                   // 注册ops,旧版本api
  ├─ ib_set_device_ops();                                    // 注册ops,新版本api
  ├─ ib_device.uverbs_cmd_mask |= OTRDMA_UVERBS_CMD_MASK;    // 开放哪些ibv_* API
  ├─ ib_device.node_type = RDMA_NODE_IB_CA;                  // 网卡类型:IB/RoCE/iWARP
  ├─ ib_device.phys_port_cnt = OTRDMA_MAX_PORTS;             // RDMA设备有多少个物理端口
  ├─ ib_device.num_comp_vectors = 1;                         // 设备支持多少个completion interrupt vector
  └  ib_register_device(&ibdev->ib_dev, "mlx5_%d", dev);
       ↓
       ├─ 生成 /sys/class/infiniband/mlx5_0
       └  生成 /dev/infiniband/uverbs0

        mlx5_ib_dev_ops内容如下:

static const struct ib_device_ops mlx5_ib_dev_ops = {
    .owner                  = THIS_MODULE,

    /* device/query */
    .query_device           = mlx5_ib_query_device,
    .query_port             = mlx5_ib_query_port,
    .query_gid              = mlx5_ib_query_gid,
    .get_link_layer         = mlx5_ib_port_link_layer,

    /* PD */
    .alloc_pd               = mlx5_ib_alloc_pd,
    .dealloc_pd             = mlx5_ib_dealloc_pd,

    /* MR */
    .reg_user_mr            = mlx5_ib_reg_user_mr,
    .dereg_mr               = mlx5_ib_dereg_mr,
    .alloc_mr               = mlx5_ib_alloc_mr,

    /* CQ */
    .create_cq              = mlx5_ib_create_cq,
    .destroy_cq             = mlx5_ib_destroy_cq,
    .poll_cq                = mlx5_ib_poll_cq,
    .req_notify_cq          = mlx5_ib_arm_cq,

    /* QP */
    .create_qp              = mlx5_ib_create_qp,
    .modify_qp              = mlx5_ib_modify_qp,
    .query_qp               = mlx5_ib_query_qp,
    .destroy_qp             = mlx5_ib_destroy_qp,

    /* SRQ */
    .create_srq             = mlx5_ib_create_srq,
    .modify_srq             = mlx5_ib_modify_srq,
    .query_srq              = mlx5_ib_query_srq,
    .destroy_srq            = mlx5_ib_destroy_srq,

    /* AH */
    .create_ah              = mlx5_ib_create_ah,
    .destroy_ah             = mlx5_ib_destroy_ah,

    /* mmap */
    .mmap                   = mlx5_ib_mmap,

    /* ucontext */
    .alloc_ucontext         = mlx5_ib_alloc_ucontext,
    .dealloc_ucontext       = mlx5_ib_dealloc_ucontext,

    /* MAD */
    .process_mad            = mlx5_ib_process_mad,

    /* flow steering */
    .create_flow            = mlx5_ib_create_flow,
    .destroy_flow           = mlx5_ib_destroy_flow,

    /* counters */
    .alloc_hw_port_stats    = mlx5_ib_alloc_hw_port_stats,
    .get_hw_stats           = mlx5_ib_get_hw_stats,

    /* ODP */
    .invalidate_range       = mlx5_ib_invalidate_range,

    /* XRC */
    .alloc_xrcd             = mlx5_ib_alloc_xrcd,
    .dealloc_xrcd           = mlx5_ib_dealloc_xrcd,

    /* DEVX */
    .devx_obj_create        = mlx5_ib_devx_obj_create,
    .devx_obj_destroy       = mlx5_ib_devx_obj_destroy,
};
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐