1. 从单机到集群:为什么你的机器学习项目需要Ray

如果你正在用Python做机器学习,尤其是玩过PyTorch或者TensorFlow,那你肯定经历过这样的时刻:模型越来越大,数据越来越多,跑一次训练从几小时变成了几天。你看着任务管理器里那个满载的CPU或者孤零零的GPU,心里盘算着是不是得升级硬件了。但说实话,自己买一堆显卡或者租昂贵的云服务器,不仅成本高,管理起来也头疼。这时候,分布式计算就成了一个绕不开的话题。

但一提到“分布式”,很多朋友可能就头大了。是不是要学一堆像MPI那样复杂的通信库?是不是要把代码重写得面目全非?是不是要搞一堆配置文件,管理一堆服务器节点?我以前也是这么想的,直到我开始用Ray。简单来说,Ray的目标就是让你用写普通Python脚本的方式,来写分布式程序。你几乎不用改变你原有的数据处理和模型训练的思维,就能让代码跑在多台机器、多个CPU核心、多个GPU上。

我最初接触Ray是因为一个图像分类项目。当时数据集有上千万张图片,用单卡训练一个Epoch就要将近一天,调一次超参数简直是一场噩梦。后来我把数据加载和预处理部分用Ray包装了一下,轻松就分发到了8台机器的CPU上进行并行处理,预处理时间从几个小时缩短到了几分钟。训练部分用上RaySGD后,几行代码就把单卡训练改成了多卡并行。整个过程,我大部分时间还是在写我熟悉的PyTorch代码,Ray就像个隐形的力量倍增器,在背后默默地把计算任务摊开。

所以,Ray到底适合谁?我觉得是这三类人:第一是算法工程师和研究员,你只想专注模型和算法,不想被复杂的集群部署和通信细节困扰;第二是数据科学家,你需要快速处理海量数据,进行特征工程或模型评估;第三是有一定Python基础的开发者,希望让自己写的计算密集型工具(比如模拟、批量推理)跑得更快。如果你符合其中任何一条,那么这篇实战指南就是为你准备的。接下来,我不会讲太多抽象的原理,而是直接带你上手,看看怎么用Ray把分布式机器学习系统“搭”起来。

2. 环境搭建与Ray初体验:5分钟跑起你的第一个分布式任务

理论说再多,不如动手跑一行代码。Ray的安装简单到令人发指,它就是一个纯粹的Python库。打开你的终端,无论用的是conda还是pip,一行命令搞定:

pip install -U ray

如果你的训练涉及到PyTorch,通常也会一起装上:

pip install torch torchvision

安装完成后,我们不用急着配置复杂的集群。Ray的一个巨大优势就是,你可以在自己的笔记本电脑上,以“模拟”分布式的方式开发和调试代码,确认无误后再无缝部署到真正的集群上。这大大降低了学习和试错成本。

让我们写一个最简单的脚本,感受一下Ray的“远程函数”(Remote Function)能力。假设我们有一个计算量很大的函数,比如模拟一个复杂的计算过程:

import ray
import time

# 初始化Ray。这里在本地启动,`num_cpus`指定可用的CPU核心数,你可以根据自己电脑情况调整。
ray.init(num_cpus=4)

# 普通的Python函数
def normal_function(x):
    time.sleep(1)  # 模拟1秒的计算耗时
    return x * x

# 用`@ray.remote`装饰器,把它变成一个可以远程异步执行的任务。
@ray.remote
def remote_function(x):
    time.sleep(1)
    return x * x

# 普通顺序执行
start = time.time()
results = [normal_function(i) for i in range(4)]
print(f"顺序执行耗时: {time.time() - start:.2f}秒, 结果: {results}")

# Ray并行执行
start = time.time()
# 这里会立即返回4个“未来对象”(ObjectRef),任务在后台并行执行
futures = [remote_function.remote(i) for i in range(4)]
# 用`ray.get()`收集所有结果,这里会等待所有任务完成
results = ray.get(futures)
print(f"Ray并行执行耗时: {time.time() - start:.2f}秒, 结果: {results}")

# 关闭Ray
ray.shutdown()

把这段代码保存成 ray_demo.py 并运行。如果你的电脑有4个或更多物理/逻辑核心,你很可能会看到顺序执行大约用了4秒,而Ray并行执行只用了1秒多一点!这就是最直观的加速。@ray.remote 这个装饰器是Ray的魔法之源,它把普通函数变成了一个可以发送到Ray运行时(可能在其他进程、其他机器)执行的“任务”。remote() 调用是异步的,立刻返回一个引用(Future),ray.get() 才真正去取结果。

注意:在本地调试时,ray.init() 不传参数会自动检测所有可用CPU。在生产集群中,你需要连接到集群头节点,例如 ray.init(address='auto')ray.init(address='ray://<head-node-ip>:10001')

这个简单的例子揭示了Ray的核心抽象:任务(Task)。基于这个简单的“远程函数”,Ray构建了整个生态。接下来,我们会用它来处理更实际的问题——分布式数据预处理。

3. 实战核心一:用Ray轻松搞定分布式数据预处理

在真实的机器学习流水线中,数据预处理(Data Preprocessing)往往是时间瓶颈。特别是当你有TB级的图像、文本或视频数据需要做清洗、增强、编码时,单机处理显得力不从心。用Ray来并行化这个过程,逻辑非常清晰。

假设我们有一个包含大量图片文件的目录,我们需要读取每张图片,将其调整为固定尺寸,进行归一化,然后可能还要做一些数据增强(如随机翻转、裁剪)。我们用PyTorch的Dataset和DataLoader来做示范,但你会发现,原生的DataLoader在多机多卡场景下进行数据加载可能成为瓶颈。我们可以用Ray来创建一个分布式的数据加载管道。

思路是:将数据索引(比如文件路径列表)划分成多个分片(Shard),每个分片由一个Ray远程任务(一个“worker”)负责处理。每个worker内部使用标准的PyTorch DataLoader来高效读取和转换自己分片的数据。最后,我们需要一个协调者来从各个worker那里收集处理好的批次(batch)。

下面是一个简化但可运行的示例:

import ray
import torch
from torch.utils.data import Dataset, DataLoader
from torchvision import transforms
from PIL import Image
import os
import time

# 1. 初始化Ray
ray.init(num_cpus=8)  # 假设我们使用8个CPU worker

# 2. 定义一个简单的图片数据集类(模拟)
class ImageDataset(Dataset):
    def __init__(self, file_paths, transform=None):
        self.file_paths = file_paths
        self.transform = transform

    def __len__(self):
        return len(self.file_paths)

    def __getitem__(self, idx):
        img = Image.open(self.file_paths[idx]).convert('RGB')
        if self.transform:
            img = self.transform(img)
        # 这里简单返回一个标签0,实际中应从文件名或标注文件读取
        return img, 0

# 3. 定义一个远程Worker类,每个实例负责一个数据分片
@ray.remote
class DataWorker:
    def __init__(self, file_paths_shard):
        self.file_paths = file_paths_shard
        # 定义数据转换
        self.transform = transforms.Compose([
            transforms.Resize((256, 256)),
            transforms.RandomHorizontalFlip(),
            transforms.ToTensor(),
            transforms.Normalize(mean=[0.485, 0.456, 0.406],
                                 std=[0.229, 0.224, 0.225]),
        ])
        self.dataset = ImageDataset(self.file_paths, self.transform)
        # 每个worker有自己的DataLoader
        self.dataloader = DataLoader(self.dataset, batch_size=32, shuffle=True, num_workers=0) # 注意:Ray worker内不建议再用多子进程
        self.iterator = iter(self.dataloader)

    def next_batch(self):
        """获取下一个批次的数据。当分片数据遍历完,可以重新创建迭代器或返回None。"""
        try:
            batch = next(self.iterator)
            return batch
        except StopIteration:
            self.iterator = iter(self.dataloader)  # 重新开始迭代
            batch = next(self.iterator)
            return batch

# 4. 主程序逻辑
if __name__ == "__main__":
    # 模拟一堆图片文件路径
    all_image_paths = [f"dummy_image_{i}.jpg" for i in range(10000)]

    # 将数据路径分成4个分片
    num_shards = 4
    shards = [all_image_paths[i::num_shards] for i in range(num_shards)]

    # 创建4个远程DataWorker,每个分配一个分片
    print("正在创建分布式数据workers...")
    workers = [DataWorker.remote(shard) for shard in shards]

    # 模拟训练循环,从所有worker那里轮询获取数据
    num_batches_to_fetch = 100
    start_time = time.time()

    for batch_idx in range(num_batches_to_fetch):
        # 使用round-robin方式从不同worker获取batch,实现数据并行读取
        worker_index = batch_idx % num_shards
        # 异步调用`next_batch`方法
        future = workers[worker_index].next_batch.remote()
        # 同步等待结果(在实际训练中,这个获取可以和模型计算重叠进行)
        batch_data, batch_labels = ray.get(future)

        # 这里应该是你的模型前向传播、损失计算、反向更新等步骤
        # 我们只是模拟一下处理时间
        time.sleep(0.01)
        if batch_idx % 20 == 0:
            print(f"已处理批次 {batch_idx}, 数据形状: {batch_data.shape}")

    end_time = time.time()
    print(f"分布式数据加载与处理 {num_batches_to_fetch} 个批次,总耗时: {end_time - start_time:.2f}秒")

    ray.shutdown()

这个例子虽然简单,但展示了用Ray组织分布式数据流的经典模式:分片 -> 创建远程处理器 -> 异步拉取。在实际更复杂的场景中,你可能会用到Ray Datasets,它是一个更高级的抽象,专门为分布式数据加载和转换设计,支持从Parquet、JSON、CSV等格式读取,并能与RaySGD、Ray Tune等库更好集成。但理解上面这个“手动”模式,能让你更清楚底层发生了什么,遇到问题时也更好调试。

4. 实战核心二:利用RaySGD进行分布式模型训练

数据准备好了,接下来就是重头戏——模型训练。对于PyTorch用户,原生的 DistributedDataParallel (DDP) 功能已经很强大了,但它需要你手动管理进程组、设置环境变量(如MASTER_ADDR, MASTER_PORT),写启动脚本。RaySGD中的 TorchTrainer 把这些繁琐的步骤都封装了起来,让你用几行代码就能启动分布式训练,并且能很好地与Ray的其他部分(如资源管理、容错)结合。

让我们用一个真实的例子,在CIFAR-10数据集上训练一个ResNet-18模型。我将详细解释每一步:

import torch
import torch.nn as nn
from torch.utils.data import DataLoader
from torchvision.datasets import CIFAR10
import torchvision.transforms as transforms
import ray
from ray.util.sgd import TorchTrainer
from ray.util.sgd.torch import TrainingOperator
from torchvision.models import resnet18  # 使用torchvision官方模型

# 1. 定义数据创建函数。这个函数会在每个训练worker上被调用。
def data_creator(config):
    """创建训练和验证数据加载器。"""
    transform_train = transforms.Compose([
        transforms.RandomCrop(32, padding=4),
        transforms.RandomHorizontalFlip(),
        transforms.ToTensor(),
        transforms.Normalize((0.4914, 0.4822, 0.4465), (0.2023, 0.1994, 0.2010)),
    ])
    transform_val = transforms.Compose([
        transforms.ToTensor(),
        transforms.Normalize((0.4914, 0.4822, 0.4465), (0.2023, 0.1994, 0.2010)),
    ])

    train_set = CIFAR10(root="./data", train=True, download=True, transform=transform_train)
    val_set = CIFAR10(root="./data", train=False, download=True, transform=transform_val)

    train_loader = DataLoader(train_set, batch_size=config["batch_size"], shuffle=True, num_workers=2, pin_memory=True)
    val_loader = DataLoader(val_set, batch_size=config["batch_size"], shuffle=False, num_workers=2, pin_memory=True)

    return train_loader, val_loader

# 2. 定义模型创建函数。
def model_creator(config):
    """创建神经网络模型。"""
    model = resnet18(num_classes=10)  # CIFAR-10有10类
    # 调整第一层卷积,因为CIFAR图片是32x32x3,而ResNet原输入是224x224
    model.conv1 = nn.Conv2d(3, 64, kernel_size=3, stride=1, padding=1, bias=False)
    # 移除原来的第一个最大池化层,因为特征图已经很小了
    model.maxpool = nn.Identity()
    return model

# 3. 定义优化器创建函数。
def optimizer_creator(model, config):
    """创建优化器。"""
    return torch.optim.SGD(model.parameters(),
                           lr=config["lr"],
                           momentum=config["momentum"],
                           weight_decay=config["weight_decay"])

# 4. 初始化Ray。在单机多GPU上测试。
ray.init(num_gpus=2, ignore_reinit_error=True)  # 假设机器有2块GPU

# 5. 配置训练参数
trainer_config = {
    "batch_size": 128,      # 每个GPU上的批次大小
    "lr": 0.1,
    "momentum": 0.9,
    "weight_decay": 5e-4,
}

# 6. 创建自定义训练操作符,将上面的创建函数组装起来。
CustomTrainingOperator = TrainingOperator.from_creators(
    model_creator=model_creator,
    optimizer_creator=optimizer_creator,
    data_creator=data_creator,
    loss_creator=nn.CrossEntropyLoss,  # 直接使用损失函数类
    scheduler_creator=None,  # 可以在这里定义学习率调度器创建函数
)

# 7. 实例化TorchTrainer,这是核心。
trainer = TorchTrainer(
    training_operator_cls=CustomTrainingOperator,
    num_workers=2,               # 使用2个worker,通常对应2块GPU
    num_cpus_per_worker=2,       # 每个worker分配2个CPU核心,用于数据加载
    use_gpu=True,                # 使用GPU
    config=trainer_config,
    backend="nccl",              # 使用NCCL后端进行GPU间通信,效率最高
    use_fp16=False,              # 是否使用混合精度训练,可以加速并减少显存占用
)

# 8. 开始训练和验证!
print("开始分布式训练...")
for epoch in range(5):  # 我们只跑5个epoch作为演示
    train_stats = trainer.train()  # 执行一个epoch的训练
    val_stats = trainer.validate() # 在验证集上评估
    print(f"Epoch {epoch+1}: Train Loss = {train_stats['train_loss']:.4f}, "
          f"Val Accuracy = {val_stats['val_accuracy']:.4f}")

# 9. 保存训练好的模型状态
model_state = trainer.state_dict()
torch.save(model_state, "distributed_resnet18_cifar10.pth")
print("模型已保存。")

# 10. 可以加载模型状态继续训练或用于推理
# new_trainer = TorchTrainer(...)
# new_trainer.load_state_dict(model_state)

# 11. 关闭训练器和Ray
trainer.shutdown()
ray.shutdown()

我来拆解一下这里面的关键点:

  • num_workers: 这是并行的训练进程数。如果你有2块GPU,就设为2。每个worker会运行一个完整的训练循环,处理不同的数据子集(数据并行)。
  • backend: 设置为 "nccl" 对于GPU训练是必须的,这是NVIDIA的高性能通信库。
  • use_fp16: 如果设为 True,RaySGD会启用自动混合精度训练,这在Ampere架构及以后的GPU上(如A100, RTX 30系列)能显著加速并节省显存。
  • TrainingOperator: 这是一个非常灵活的抽象。通过 from_creators 方法,你以函数式的方式定义了训练的所有组件。你还可以继承 TrainingOperator 类,重写 train_batchvalidate_batch 等方法,实现完全自定义的训练逻辑。

我自己的经验是,用RaySGD的最大好处是代码干净与生态无缝集成。你的训练脚本几乎就是标准的PyTorch代码,没有杂七杂八的分布式初始化代码。当你需要做超参数搜索时,可以轻松地把这个训练器包装起来,丢给Ray Tune去跑,这是我们接下来要看的。

5. 实战核心三:用Ray Tune实现自动化超参数搜索

手动调参是机器学习中最耗时的“玄学”活动之一。Ray Tune把这个过程自动化、并行化了。它可以同时启动成百上千个训练试验(Trial),每个试验使用不同的超参数组合,并智能地调度资源、提前终止表现不好的试验,帮你快速找到最优配置。

我们把上一节的训练代码,改造成一个可以被Tune调用的“可训练函数”。然后,我们定义要搜索的超参数空间,选择一个调度器(Scheduler)和搜索算法(Search Algorithm),就可以开始了。

import torch
import torch.nn as nn
from torch.utils.data import DataLoader
from torchvision.datasets import CIFAR10
import torchvision.transforms as transforms
from torchvision.models import resnet18
import ray
from ray import tune
from ray.tune.schedulers import ASHAScheduler  # 异步连续减半算法,高效终止差劲的试验
from ray.tune import CLIReporter
import os

# 1. 将训练逻辑包装成一个函数,这是Tine要求的格式。
def train_cifar(config, checkpoint_dir=None):
    # 这个函数会在每个Trial(试验)中独立运行
    net = resnet18(num_classes=10)
    net.conv1 = nn.Conv2d(3, 64, kernel_size=3, stride=1, padding=1, bias=False)
    net.maxpool = nn.Identity()

    # 使用config中的超参数
    criterion = nn.CrossEntropyLoss()
    optimizer = torch.optim.SGD(net.parameters(),
                                lr=config["lr"],
                                momentum=config["momentum"],
                                weight_decay=config["weight_decay"])

    # 数据加载(每个试验独立加载)
    transform = transforms.Compose([
        transforms.RandomCrop(32, padding=4),
        transforms.RandomHorizontalFlip(),
        transforms.ToTensor(),
        transforms.Normalize((0.4914, 0.4822, 0.4465), (0.2023, 0.1994, 0.2010)),
    ])
    trainset = CIFAR10(root="./data", train=True, download=True, transform=transform)
    trainloader = DataLoader(trainset, batch_size=config["batch_size"], shuffle=True, num_workers=2)

    # 如果有检查点,加载它(用于从故障中恢复或从暂停的试验继续)
    if checkpoint_dir:
        checkpoint = torch.load(os.path.join(checkpoint_dir, "checkpoint.pth"))
        net.load_state_dict(checkpoint["model_state_dict"])
        optimizer.load_state_dict(checkpoint["optimizer_state_dict"])
        start_epoch = checkpoint["epoch"]
    else:
        start_epoch = 0

    # 训练循环
    for epoch in range(start_epoch, 10):  # 每个试验训练10个epoch
        net.train()
        running_loss = 0.0
        for i, data in enumerate(trainloader, 0):
            inputs, labels = data
            optimizer.zero_grad()
            outputs = net(inputs)
            loss = criterion(outputs, labels)
            loss.backward()
            optimizer.step()
            running_loss += loss.item()

        # 每个epoch结束后,计算验证集精度
        val_transform = transforms.Compose([
            transforms.ToTensor(),
            transforms.Normalize((0.4914, 0.4822, 0.4465), (0.2023, 0.1994, 0.2010)),
        ])
        valset = CIFAR10(root="./data", train=False, download=True, transform=val_transform)
        valloader = DataLoader(valset, batch_size=config["batch_size"], shuffle=False, num_workers=2)
        net.eval()
        correct = 0
        total = 0
        with torch.no_grad():
            for data in valloader:
                images, labels = data
                outputs = net(images)
                _, predicted = torch.max(outputs.data, 1)
                total += labels.size(0)
                correct += (predicted == labels).sum().item()
        val_accuracy = correct / total

        # 创建检查点(Tune要求)
        with tune.checkpoint_dir(epoch) as checkpoint_dir:
            path = os.path.join(checkpoint_dir, "checkpoint.pth")
            torch.save({
                "epoch": epoch + 1,
                "model_state_dict": net.state_dict(),
                "optimizer_state_dict": optimizer.state_dict(),
                "loss": running_loss / len(trainloader),
            }, path)

        # 向Tune报告指标,这是驱动调度器和搜索算法的关键
        tune.report(loss=running_loss / len(trainloader), accuracy=val_accuracy)

# 2. 定义超参数搜索空间
search_space = {
    "lr": tune.loguniform(1e-4, 1e-1),  # 学习率在0.0001到0.1之间对数均匀采样
    "momentum": tune.uniform(0.8, 0.99),
    "weight_decay": tune.loguniform(1e-5, 1e-3),
    "batch_size": tune.choice([64, 128, 256])  # 批次大小也可以作为超参数搜索
}

# 3. 初始化Ray(Tune会自动利用Ray集群)
ray.init(num_cpus=8, num_gpus=2, address='auto')  # 连接到现有集群或本地启动

# 4. 配置并运行超参数搜索
analysis = tune.run(
    train_cifar,
    metric="accuracy",  # 我们关注验证集准确率
    mode="max",         # 目标是最大化准确率
    config=search_space,
    num_samples=20,     # 总共尝试20组不同的超参数组合
    resources_per_trial={"cpu": 2, "gpu": 0.5},  # 每个试验分配2个CPU和0.5个GPU(意味着2个试验共享1块GPU)
    scheduler=ASHAScheduler(
        metric="accuracy",
        mode="max",
        max_t=10,           # 每个试验最多运行10个epoch
        grace_period=2,     # 至少运行2个epoch后才可能被提前终止
        reduction_factor=2  # 每次减半,保留一半表现最好的试验
    ),
    progress_reporter=CLIReporter(
        metric_columns=["loss", "accuracy", "training_iteration"],
        max_report_frequency=30,  # 每30秒报告一次进度
    ),
    local_dir="./ray_results",  # 所有试验的日志和检查点保存目录
    name="cifar10_tune_exp",    # 实验名称
)

# 5. 分析结果
print("最佳试验配置:", analysis.best_config)
print("最佳试验最终准确率:", analysis.best_result["accuracy"])

# 你可以获取DataFrame进行更详细的分析
df = analysis.results_df
print(df[["config/lr", "config/momentum", "accuracy", "loss"]].head())

# 6. 恢复并测试最佳模型
best_trial = analysis.get_best_trial("accuracy", "max", "last")
best_checkpoint = analysis.get_best_checkpoint(best_trial, "accuracy", "max")
print(f"最佳模型检查点路径: {best_checkpoint.path}")

# 你可以加载这个检查点,用于后续的模型服务或进一步训练
ray.shutdown()

运行这段代码,你会在终端看到一个动态更新的表格,显示所有正在运行的试验的状态、损失和准确率。ASHAScheduler会非常“无情”地提前终止那些看起来没希望的试验,把计算资源集中给更有潜力的组合。这就是为什么Ray Tune能如此高效。

提示:resources_per_trial 的设置是性能调优的关键。如果你的模型很小,GPU显存充足,可以设置 "gpu": 1 让每个试验独占一块GPU。如果模型大或者你想同时跑更多试验,可以设置分数GPU(如0.5),Ray会负责在同一块GPU上调度多个试验(注意这可能增加显存溢出风险)。CPU数量要保证足够每个试验的数据加载。

6. 模型部署与服务化:用Ray Serve让模型上线

模型训练和调优好了,最终要落地提供服务。你可能想过用Flask或FastAPI写一个简单的HTTP服务,但当你面临高并发、需要动态扩缩容、或者要同时部署多个模型版本(A/B测试)时,自己从头搭建这套系统就复杂了。Ray Serve 就是专门解决这个问题的。

Ray Serve允许你将模型(不管是PyTorch、TensorFlow、Scikit-learn还是任意Python函数)包装成一个“部署”(Deployment),然后以高性能、可扩展的方式提供服务。它自动处理请求的批处理(Batching)、多副本(Replicas)以增加吞吐量、以及负载均衡。

我们来部署上一节得到的最佳图像分类模型。假设我们已经从Ray Tune的实验中拿到了最好的模型权重文件。

import ray
from ray import serve
import torch
import torch.nn as nn
from torchvision import transforms
from PIL import Image
import io
import requests
from torchvision.models import resnet18

# 1. 定义我们的模型类,它将被包装成Serve的部署
class ImageClassifier:
    def __init__(self):
        # 加载模型架构
        self.model = resnet18(num_classes=10)
        self.model.conv1 = nn.Conv2d(3, 64, kernel_size=3, stride=1, padding=1, bias=False)
        self.model.maxpool = nn.Identity()
        # 加载训练好的权重(这里假设权重文件在当前目录)
        state_dict = torch.load("distributed_resnet18_cifar10.pth", map_location="cpu")
        # 注意:RaySGD保存的state_dict可能包含额外前缀,需要处理
        # 如果是用上一节TorchTrainer保存的,直接加载即可
        self.model.load_state_dict(state_dict)
        self.model.eval()

        # 定义图像预处理变换(必须与训练时一致)
        self.transform = transforms.Compose([
            transforms.Resize((32, 32)),
            transforms.ToTensor(),
            transforms.Normalize((0.4914, 0.4822, 0.4465), (0.2023, 0.1994, 0.2010)),
        ])
        # CIFAR-10的类别名称
        self.classes = ('plane', 'car', 'bird', 'cat', 'deer', 'dog', 'frog', 'horse', 'ship', 'truck')

    # 这个方法是处理单个请求的
    async def __call__(self, request):
        # 从HTTP请求中获取图片数据
        image_data = await request.body()
        image = Image.open(io.BytesIO(image_data)).convert('RGB')
        # 预处理
        input_tensor = self.transform(image).unsqueeze(0)  # 增加batch维度

        # 推理
        with torch.no_grad():
            output = self.model(input_tensor)
            probabilities = torch.nn.functional.softmax(output[0], dim=0)
            predicted_class_idx = torch.argmax(probabilities).item()

        # 返回JSON格式的结果
        return {
            "predicted_class": self.classes[predicted_class_idx],
            "class_index": predicted_class_idx,
            "confidence": probabilities[predicted_class_idx].item()
        }

# 2. 启动Ray和Serve
ray.init(address="auto")  # 连接到现有Ray集群
serve.start(detached=True)  # 以分离模式启动Serve,这样脚本退出服务仍在运行

# 3. 创建并部署我们的模型服务
# `num_replicas`指定启动多少个副本实例来处理请求,可以实现水平扩展。
# `ray_actor_options`指定每个副本需要的资源。
serve.create_backend("classifier:v1", ImageClassifier, num_replicas=2,
                     ray_actor_options={"num_cpus": 1, "num_gpus": 0.1})  # 每个副本使用0.1个GPU

# 4. 将后端(Backend)挂载到一个端点(Endpoint),客户端通过这个URL访问
serve.create_endpoint("classify", backend="classifier:v1", route="/classify")

print("模型服务已部署!端点: /classify")

# 现在服务已经在运行了。我们可以写一个简单的客户端测试脚本(可以放在另一个文件):
def test_client(image_path):
    with open(image_path, "rb") as f:
        image_bytes = f.read()
    resp = requests.post("http://localhost:8000/classify", data=image_bytes)
    print(resp.json())

# 假设有一张猫的图片叫`cat.jpg`
# test_client("cat.jpg")

# 5. 服务管理:我们可以动态更新、缩放或删除部署
# 例如,将副本数扩展到4个以应对更高流量:
# serve.update_backend_config("classifier:v1", {"num_replicas": 4})

# 部署新版本的模型(蓝绿部署):
# serve.create_backend("classifier:v2", ImageClassifierNew, ...)
# serve.set_traffic("classify", {"classifier:v1": 0.5, "classifier:v2": 0.5})

# 脚本到这里可以结束,Serve服务会在后台持续运行。
# 要停止服务,可以调用 `serve.shutdown()`。

这段代码展示了Ray Serve的核心流程:定义模型类 -> 启动后端 -> 创建端点。一旦部署完成,你的模型就通过HTTP REST API对外提供服务了。Ray Serve会自动在多台机器(如果Ray集群是多节点的)上分布这些副本,并提供负载均衡。

我在实际项目中用Ray Serve部署过推荐系统模型,最大的感受是省心。当流量高峰时,我可以通过Ray的监控指标,动态增加副本数;当发布新模型时,可以用流量拆分做平滑的A/B测试。所有这些操作,都有简单的Python API或Dashboard来完成,不需要去折腾Kubernetes的YAML文件或者复杂的网关配置。

7. 从单机到云集群:Ray Cluster实战配置指南

到目前为止,我们都在单机(可能是多GPU)上操作。但Ray真正的威力在于能轻松扩展到成百上千台机器组成的集群。你可以在本地开发调试,然后几乎不改代码就能部署到云上(AWS、GCP、Azure、私有云)或Kubernetes中。Ray自带了一个集群启动器(Cluster Launcher)和一套配置范式,让这个过程变得简单。

这里我以在AWS EC2上启动一个Ray集群为例,说明核心步骤。同样逻辑也适用于其他环境。

第一步:准备集群配置文件 创建一个YAML文件,比如 ray-cluster.yaml。这个文件描述了集群的形态。

# 一个最小化的AWS集群配置示例
cluster_name: ray-ml-cluster

provider:
    type: aws
    region: us-west-2
    availability_zone: us-west-2a

auth:
    ssh_user: ubuntu  # 根据你选择的AMI调整,例如Amazon Linux 2可能是ec2-user

# 定义头节点(主节点)
head_node:
    InstanceType: m5.2xlarge  # 8 vCPU, 32 GB内存
    ImageId: ami-0c55b159cbfafe1f0  # 一个预装了Ray的深度学习AMI,你需要选择或自己制作
    KeyName: your-aws-key-pair-name  # 你的EC2密钥对名称
    IamInstanceProfile:
        Arn: arn:aws:iam::YOUR_ACCOUNT_ID:instance-profile/your-ray-instance-profile  # 需要有适当权限的IAM角色
    SecurityGroupIds:
        - sg-xxxxxxxx  # 你的安全组ID,需要开放6379(Ray GCS)、8265(Dashboard)、10001(Serve)等端口
    BlockDeviceMappings:
        - DeviceName: /dev/sda1
          Ebs:
              VolumeSize: 100  # 根卷大小,GB

# 定义工作节点
worker_nodes:
    min_workers: 2
    max_workers: 10  # 支持自动伸缩
    InstanceType: g4dn.2xlarge  # 配备NVIDIA T4 GPU的实例
    ImageId: ami-0c55b159cbfafe1f0
    KeyName: your-aws-key-pair-name
    IamInstanceProfile:
        Arn: arn:aws:iam::YOUR_ACCOUNT_ID:instance-profile/your-ray-instance-profile
    SecurityGroupIds:
        - sg-xxxxxxxx

# Ray的安装和启动命令
setup_commands:
    - pip install -U ray[default] torch torchvision  # 确保所有节点安装必要库

# 文件同步:将本地代码同步到集群头节点
file_mounts: {
    "/home/ubuntu/my_ray_project": "/path/to/your/local/project",  # 本地路径: 远程路径
}

# 初始化命令,在头节点运行
head_setup_commands: []
head_start_ray_commands:
    - ray stop  # 停止可能存在的Ray进程
    - ulimit -n 65536 && ray start --head --port=6379 --object-manager-port=8076 --autoscaling-config=~/ray_bootstrap_config.yaml --dashboard-host=0.0.0.0

# 工作节点启动命令
worker_start_ray_commands:
    - ray stop
    - ulimit -n 65536 && ray start --address=$RAY_HEAD_IP:6379 --object-manager-port=8076

第二步:启动集群 在本地电脑上,确保已安装Ray和AWS CLI,并且配置了AWS凭证。然后运行:

ray up ray-cluster.yaml

这个命令会自动创建头节点,并按照配置启动指定数量的工作节点。它会处理所有网络配置、安全组规则、软件安装等繁琐工作。

第三步:提交任务到集群 集群运行后,你可以通过端口转发访问Web Dashboard(默认在8265端口),监控集群资源和使用情况。提交任务有两种主要方式:

  1. SSH到头节点手动运行

    ray attach ray-cluster.yaml
    # 这会SSH到头部节点,并进入一个终端
    cd my_ray_project
    python my_distributed_training_script.py
    
  2. 使用ray submit提交脚本(推荐):

    ray submit ray-cluster.yaml my_distributed_training_script.py
    

    这个命令会将本地的脚本文件自动同步到集群头节点,并在那里执行。你的脚本中应该使用 ray.init(address='auto') 来连接到集群。

第四步:管理集群

  • 扩展/收缩ray exec ray-cluster.yaml 'ray add-worker ...' 或直接修改YAML中的 min_workers/max_workers,Ray的自动伸缩器会根据任务负载调整节点数量。
  • 关闭集群ray down ray-cluster.yaml。这会终止所有EC2实例,注意保存好数据和日志

我踩过的一个坑是网络和权限。确保安全组允许集群节点之间所有端口的内部通信(或者至少是Ray使用的端口范围),并且IAM角色有创建EC2、访问S3(如果你数据在S3上)等权限。第一次搭建可能会花些时间在配置上,但一旦配好,以后启动一个包含几十个GPU的集群也就是一条命令的事,非常方便。

从单机脚本到云上分布式集群,Ray提供了一条平滑的路径。它没有消灭分布式系统的复杂性,而是把这些复杂性封装成了简单的Python接口和声明式的配置文件。这让开发者能更专注于机器学习任务本身,而不是基础设施。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐