时钟周期

PyTRIO 的异步训练 API 可以让本地数据准备与远端模型计算并行推进。理解时钟周期后,可以更合理地安排 forward_backward_async()optim_step_async() 的提交时机,减少共享计算资源等待本地客户端的空档。

时钟周期是一种调度抽象,没有固定的毫秒数。每个周期的实际耗时由请求队列、批次规模、模型和系统负载共同决定。

理解时钟周期

PyTRIO 服务会把多个 LoRA 训练作业调度到共享 worker pool。worker pool 以同步步伐推进,每一步就是一个时钟周期。一个周期内可以处理来自多个作业的前向传播、反向传播和优化器更新。

不同作业共享计算容量,各自的 LoRA 权重和优化器状态保持隔离。小批次也能利用共享容量;请求需要等到合适的周期边界才能开始执行,因此提交时机也会影响端到端延迟。

PyTRIO 时钟周期与共享 worker pool

让一次更新进入同一个周期

PyTRIO 的训练异步 API 有两层 await

  1. await training_client.forward_backward_async(...) 提交请求并返回 APIFuture
  2. await fwdbwd_future 等待远端计算完成并取得结果。

第一次 await 返回后,本地程序可以继续提交下一个训练请求。利用这段时间先提交 optim_step_async(),前向/反向与优化器更新就可以进入同一次训练更新的调度周期。

先提交同一次更新的两个请求,再等待结果

下面的代码假设已经创建 training_clientbatch 是一组 trio.Datum,并且 adam_params 已配置完成。

逐个提交并等待

下面的写法先等待前向/反向完成,再提交优化器更新:

fwdbwd_future = await training_client.forward_backward_async(
    batch,
    "cross_entropy",
)
fwdbwd_result = await fwdbwd_future

optim_future = await training_client.optim_step_async(adam_params)
await optim_future

在图示的调度中,optim_step 会错过 Cycle N+1,到 Cycle N+2 才开始执行,一次更新跨越三个时钟周期。

先提交,再等待

先按训练顺序提交两个请求,再分别等待结果:

fwdbwd_future = await training_client.forward_backward_async(
    batch,
    "cross_entropy",
)
optim_future = await training_client.optim_step_async(adam_params)

fwdbwd_result = await fwdbwd_future
await optim_future

两个请求在本地等待之前都已进入队列,服务端可以把同一次更新安排进同一个时钟周期。请保持 forward_backward_async() 在前、optim_step_async() 在后的提交顺序。

用流水线填满后续周期

完成一次更新后才提交下一批数据,会让请求队列在两个批次之间出现空档。批次流水线会先提交 Batch N+1,再等待 Batch N 的结果,让后续周期始终有可执行的请求。

用批次流水线保持时钟周期连续

下面的示例使用一批前瞻:队列中最多保留当前批次和下一批次,既能填补客户端空档,也能避免一次提交过多请求。

import pytrio as trio


async def submit_update(training_client, batch, adam_params):
    fwdbwd_future = await training_client.forward_backward_async(
        batch,
        "cross_entropy",
    )
    optim_future = await training_client.optim_step_async(adam_params)
    return fwdbwd_future, optim_future


async def train_pipelined(
    training_client,
    batches: list[list[trio.Datum]],
    learning_rate: float,
):
    if not batches:
        return []

    adam_params = trio.AdamParams(learning_rate=learning_rate)
    current = await submit_update(training_client, batches[0], adam_params)
    results = []

    for next_batch in batches[1:]:
        # 先让下一批进入队列,再等待当前批次。
        following = await submit_update(training_client, next_batch, adam_params)

        fwdbwd_future, optim_future = current
        results.append(await fwdbwd_future)
        await optim_future
        current = following

    fwdbwd_future, optim_future = current
    results.append(await fwdbwd_future)
    await optim_future
    return results

这个模式适合下一批数据已经准备好、并且不依赖当前批次返回值的训练任务。若下一批数据必须根据当前结果生成,应保留批次间的依赖,只优化单次更新内的请求提交顺序。

选择提交方式

场景推荐方式原因
单步调试或短任务同一次更新先提交两个请求,再等待结果逻辑简单,同时避免更新内的空档
独立的 SFT 批次使用一批前瞻流水线下一批可以提前进入队列
下一批依赖当前结果只重叠当前更新内的两个请求保持算法要求的数据依赖
批次数量很多使用有界流水线控制正在处理的 future 数量和本地内存占用

评估吞吐量时,应观察多个训练步骤的端到端耗时。单个时钟周期的 wall-clock 时间会随队列和系统负载变化,图中的周期数量用于解释请求调度关系。

有关 APIFuture 和其他异步方法,请继续阅读异步指南

这篇文档对你有帮助吗?

本页目录