帮你快速理解、总结文档立即下载
文档中心>Agent Runtime>Agent 集群>快速入门>批量启动沙箱运行 RL Rollout

批量启动沙箱运行 RL Rollout

最近更新时间:2026-10-08 17:40:01
本文档已由 AI 辅助审校
我的收藏
训练 Coding Agent 时,一轮采样往往要执行多次代码尝试:每次尝试准备独立环境、修改代码、运行测试,再把奖励交给训练系统。如果每次都等待上一个环境启动和执行结束,整轮采样会花更多时间。
本指南演示如何在 Agent 集群中一次提交 10 个独立沙箱,并行完成“修复代码 → 运行测试 → 返回奖励”,同时测量整批环境准备和任务完成的时间。适合正在搭建 RL Rollout、代码评测或批量工具执行平台,希望先验证环境创建和结果收集流程的开发者。

你将完成什么

每个沙箱独立处理同一个代码修复任务:订单金额达到 99 元时应免运费,原函数却使用 total > 99。任务先确认边界测试失败,再将条件修复为 total >= 99,重新测试,成功时返回 reward=1。
一次运行将得到:
10 个隔离的任务环境,每个环境拥有独立的代码和测试目录。
10 条奖励记录,可按 Episode ID 关联到对应沙箱。
整批沙箱全部准备完成的时间、全部奖励返回的时间,以及每个沙箱的启动分布。
为了先测清环境准备耗时,Demo 使用固定修复动作,无需模型密钥或 GPU。接入真实 RL 系统时,再将固定动作替换为模型生成的动作;奖励计算和结果汇总链路可以继续复用。

开始前准备

1. 已有可访问的 Agent 集群,配置好 kubeconfig;本机安装 Python 3.9 或以上版本和 kubectl。
2. 集群中有 cube RuntimeClass,节点可运行该运行时,且能拉取 python:3.12-alpine 镜像。
3. 账号允许查询 RuntimeClass、创建和删除测试 Namespace、ConfigMap、Job,以及读取 Pod、事件和日志。
4. 集群有足够资源并行运行 10 个沙箱。每个任务容器请求 100m CPU、64Mi 内存,限制为 500m CPU、128Mi 内存;还需计入 RuntimeClass 的 Pod overhead。
export KUBECONFIG='/absolute/path/to/your/agc-kubeconfig'
kubectl config current-context
kubectl get nodes
kubectl get runtimeclass cube -o yaml
例如 RuntimeClass 每个 Pod 额外需要 768Mi 内存,10 个任务就需额外预留 7.5GiB,不能只按业务容器内存估算容量。
将文末完整脚本保存为 rl-batch-demo.py。脚本只依赖 Python 标准库和 kubectl,不需要下载内部仓库。

第一步:一次运行 10 个 Episode

在脚本所在目录执行:
python3 rl-batch-demo.py run --count 10 --output rl-parallel-10
脚本会生成独立的 agc-rl-... 命名空间,创建 completions=10、parallelism=10 的 Job。每个 Pod 对应一个 Episode,使用自己的 emptyDir 工作目录;不共享代码修改结果。
运行后会打印命名空间。任务完成后,终端显示统计摘要,完整结果保存到 rl-parallel-10/。已有同名本地目录时脚本会停止,防止覆盖实验记录;再次运行请换一个输出目录。
镜像或 RuntimeClass 名称不同时,通过参数指定:
python3 rl-batch-demo.py run --count 10 \\
--runtime YOUR_RUNTIME_CLASS \\
--image YOUR_PYTHON_IMAGE \\
--output rl-parallel-custom
自有镜像需包含 python 命令,支持以 UID/GID 1000 执行 Python 3.9 以上程序。节点需具备相应拉取权限。

第二步:核对奖励与启动时间

cat rl-parallel-10/results.json
先检查 episodes=10、successful_episodes=10、reward_sum=10。脚本还会检查每个任务的原始边界测试失败、修复后测试通过、容器退出码为 0 且无重启;任一条件不满足都会报错,不把失败任务排除后继续给出成功报告。
再看以下指标,单位均为秒:
指标
回答的问题
batch_to_all_ready_seconds
从 Job 创建到本批最后一个 Episode 的 Python 进程写好代码和测试文件,需要等多久?
pod_to_ready_p50_seconds
从各自 Pod 创建到任务环境准备完成,一半沙箱在多久内可执行?
pod_to_ready_p95_seconds
本次样本中,较慢沙箱的准备时间是多少?
batch_to_all_rewards_seconds
从 Job 创建到最后一条奖励产生,整批计算花了多久?
client_job_complete_seconds
本机提交创建到观察到 Job Complete 的总等待时间是多少?
ready 事件由每个沙箱在代码和测试文件准备完毕后写入日志,比只看 Pod Running 更贴近“可以开始执行任务”。所有环境先后达到 ready 不代表它们始终同时存活;先完成的 Episode 会正常退出。
启动统计包含 Job 控制器创建 Pod、调度、运行时、镜像准备和 Python 初始化,不能解读为纯虚拟机启动耗时。client_job_complete_seconds 还包含客户端请求、任务退出、控制器更新状态和轮询等待。
Kubernetes 创建时间与容器事件来自不同时间源,测试前需保证集群时钟同步。Kubernetes 创建时间为秒级,结果应按约秒级理解;脚本遇到负时延或时间顺序异常会报错。

第三步:用相同任务数观察并行效果

先保存并清理上一轮,再让同样 10 个任务以单并发执行:
python3 rl-batch-demo.py cleanup --output rl-parallel-10
python3 rl-batch-demo.py run --count 10 --parallelism 1 --output rl-sequential-10
cat rl-sequential-10/results.json
比较两组的 batch_to_all_rewards_seconds 和 client_job_complete_seconds,可以看到并行采样对整批任务等待时间的影响。两组任务数保持相同,避免用 5 个串行任务与 10 个并行任务直接比较。
单并发组中,后面的 Episode 需要等待前面的任务结束才会创建。因此该组的 batch_to_all_ready_seconds 含排队时间,不等于单个沙箱启动变慢。并行缩短整批耗时,也不代表单个沙箱相对其他运行时快了同样倍数。

一次实测的参考结果

2026 年 9 月 21 日,使用同一任务脚本和计时口径、固定的 Python 镜像 digest、相同的 Cube 运行时和容器资源规格,重新执行了下表三组测试。每个沙箱为单容器、独立 emptyDir,容器 CPU 限制为 500m、内存限制为 128Mi;可调度范围固定为两个 64 vCPU 节点。
检查项
10 路并行
单并发
逐节点预热后 10 路并行
测试规模
1 轮 × 10 个任务
1 轮 × 10 个任务
3 轮 × 10 个任务
完成任务 / 返回奖励
10 / 10,奖励均为 1
10 / 10,奖励均为 1
30 / 30,奖励均为 1
沙箱创建 RunPodSandbox:P50 / P95
752 / 765 毫秒
676 / 709 毫秒
761 / 835 毫秒
容器创建 CreateContainer:P50 / P95
15 / 23 毫秒
14 / 15 毫秒
14 / 21 毫秒
容器启动 StartContainer:P50 / P95
176 / 211 毫秒
175 / 209 毫秒
180 / 265 毫秒
Python 模块导入:P50 / P95
1.32 / 1.33 秒
1.32 / 1.35 秒
1.33 / 1.37 秒
导入完成至应用 ready:P50 / P95
2.3 / 28.8 毫秒
2.3 / 2.7 毫秒
2.3 / 2.7 毫秒
整批应用环境准备完成(含 Python 初始化)
约 3.7 秒
约 93.1 秒
三轮:3.2 / 3.5 / 3.6 秒
单 Pod 创建至应用 ready:P50 / P95
3.7 / 3.7 秒
3.3 / 3.4 秒
3.4 / 3.5 秒
整批奖励全部产生
约 8.4 秒
约 97.6 秒
三轮:7.9 / 8.2 / 8.2 秒
客户端观察到 Job 完成
约 12.7 秒
约 101.6 秒
三轮:12.9 / 13.0 / 12.2 秒
执行顺序为 10 路并行、单并发,再在两个节点分别执行一次相同任务预热,最后运行三轮 10 路并行。全部 50 个正式任务均命中镜像缓存;前两组也不是冷启动,第三组用于观察显式预热后的重复运行结果,不能据此计算冷启动与热启动的加速比例。
P50 / P95 以各组单个任务为样本计算:前两列各 10 个样本,第三列合并 30 个样本;整批耗时按各轮分别展示。单并发的整批耗时包含等待前一个任务结束的排队时间,各阶段分位数不能直接相加得到整体分位数。
RunPodSandbox、CreateContainer 和 StartContainer 的耗时来自同一宿主机的运行时调用与返回日志。RunPodSandbox 包含沙箱及网络准备,不包含后续容器启动、Python 初始化和任务文件准备。
Python 导入记录代码段墙钟耗时,可能包含运行时、I/O 和调度等待;“导入完成至应用 ready”还包含文件准备、函数定义和少量事件输出。Job/Pod 创建时间为秒级,整批和应用 ready 时间按约秒级理解。
本次三组都使用调度器正常分配节点,没有改变共享运行时配置。具体节点分布、镜像 imageID、逐任务状态和原始日志已随测试结果保留。镜像、节点负载、任务初始化和并发数都会影响结果,小样本 P95 不作为大规模尾延迟或服务承诺。

保存结果与清理

输出目录包含以下文件:
文件
用途
results.json
汇总和逐 Episode 指标、节点分布、实际镜像 imageID。
rollout-*.log
每个 Episode 的 ready、reward 事件及测试结果。
job.json、pods.json、events.txt
任务状态和故障排查依据。
receipt.json
本次测试的 context、命名空间和所有权信息,用于安全清理。
结果已保存后,清理单并发测试:
python3 rl-batch-demo.py cleanup --output rl-sequential-10
如果没有执行单并发组,就清理 rl-parallel-10。脚本核对测试资源的 UID 和标签后,先删除 Job 及其 Pod,再删除独立命名空间。本机输出目录会保留。Pod 删除后,emptyDir 中的代码文件也随之删除;真实任务需在退出前上传需要保留的补丁、轨迹和奖励。

遇到问题时

脚本异常时,先记录终端显示的命名空间,再检查:
kubectl -n YOUR_TEST_NAMESPACE get job,pods
kubectl -n YOUR_TEST_NAMESPACE get events --sort-by=.metadata.creationTimestamp
kubectl -n YOUR_TEST_NAMESPACE logs YOUR_POD_NAME -c episode
现象
检查方向
沙箱长期 Pending
节点可调度状态、RuntimeClass 调度条件、容器资源与 Pod overhead。
镜像拉取失败
镜像名称、仓库权限和节点网络。
Job Failed 或等待超时
输出目录中的 Pod 状态、事件与各 Episode 日志;脚本不自动重试失败任务。
没有 results.json
任务或结果校验未通过,先处理报错,不将部分完成当作整批成功。
时间顺序异常
集群时钟同步及日志事件是否完整。
清理失败
确认 kubeconfig 仍可访问原 context,检查资源终止状态;不要强制去除共享组件的 finalizer。

接入自己的训练流程

把脚本中的固定代码修复动作替换为真实模型和工具调用,为每个 Episode 准备独立任务输入、工作目录与奖励规则。训练调度器可以按一轮采样的规模设置 completions,按集群可用容量设置 parallelism,收齐奖励和轨迹后交给训练进程。
本 Demo 不更新模型权重。它验证的是 RL 系统中的环境准备、并行执行和奖励回收;先把这条链路跑通,再接入真实模型、重试策略和结果存储。

完整运行脚本

复制下面代码,保存为 rl-batch-demo.py。脚本将体验规模限制在最多 20 个 Episode,便于首次验证;这个限制不代表产品容量上限。
#!/usr/bin/env python3
"""Reproducible, bounded RL rollout environment demo; requires kubectl."""
import argparse
import concurrent.futures
import datetime
import json
import math
import os
from pathlib import Path
import subprocess
import sys
import time
import uuid

WORKER = r'''
import json, os, pathlib, subprocess, time
p = pathlib.Path('/workspace')
os.chdir(p)
p.joinpath('shipping.py').write_text('def shipping(total):\\n return 0 if total > 99 else 10\\n')
p.joinpath('test_shipping.py').write_text("""import unittest
from shipping import shipping
class ShippingTest(unittest.TestCase):
def test_below(self): self.assertEqual(shipping(98), 10)
def test_boundary(self): self.assertEqual(shipping(99), 0)
def test_above(self): self.assertEqual(shipping(100), 0)
if __name__ == '__main__': unittest.main()
""")
def emit(event, **data):
print(json.dumps(dict(event=event, at=time.time(),
episode_id=os.environ['EPISODE_ID'], **data)), flush=True)
def test():
result = subprocess.run(['python', '-B', 'test_shipping.py'],
text=True, capture_output=True, timeout=20)
return result.returncode, result.stderr
emit('ready')
before, before_output = test()
if before == 0:
raise RuntimeError('baseline must fail the boundary test')
# Deterministic stand-in for a coding Agent's action; no external model call.
p.joinpath('shipping.py').write_text(p.joinpath('shipping.py').read_text().replace('total > 99', 'total >= 99'))
after, after_output = test()
emit('reward', reward=int(after == 0), baseline_exit=before,
tests_exit=after, tests_output=after_output)
if after != 0:
raise RuntimeError('patched code failed its tests')
'''


def percentile(values, q):
values = sorted(values)
i = (len(values) - 1) * q
lo, hi = math.floor(i), math.ceil(i)
return values[lo] + (values[hi] - values[lo]) * (i - lo)


def collect(pods, logs, job, count):
origin = datetime.datetime.fromisoformat(job['metadata']['creationTimestamp'].replace('Z', '+00:00')).timestamp()
rows = []
if len(pods) != count:
raise ValueError('Pod count differs from requested episodes; inspect failures/replacements')
for pod, log in zip(pods, logs):
if pod['status']['phase'] != 'Succeeded':
raise ValueError('Episode did not succeed: ' + pod['metadata']['name'])
events = [json.loads(line) for line in log.splitlines() if line.strip()]
ready = [x for x in events if x['event'] == 'ready']
reward = [x for x in events if x['event'] == 'reward']
if len(ready) != 1 or len(reward) != 1:
raise ValueError('Expected one ready and one reward event per episode')
r, w = ready[0], reward[0]
if r['episode_id'] != pod['metadata']['name'] or w['episode_id'] != r['episode_id']:
raise ValueError('Episode identity mismatch')
if w['reward'] != 1 or w['tests_exit'] != 0 or w['baseline_exit'] == 0:
raise ValueError('Reward or baseline verification failed')
created = datetime.datetime.fromisoformat(pod['metadata']['creationTimestamp'].replace('Z', '+00:00')).timestamp()
if not origin <= created <= r['at'] <= w['at']:
raise ValueError('Invalid timestamp order; check clock synchronization')
status = pod['status']['containerStatuses'][0]
if status['restartCount'] != 0 or status['state']['terminated']['exitCode'] != 0:
raise ValueError('Unexpected restart or exit code')
rows.append(dict(pod=r['episode_id'], node=pod['spec']['nodeName'],
image_id=status.get('imageID'), reward=w['reward'],
pod_to_ready_seconds=r['at']-created,
batch_to_ready_seconds=r['at']-origin,
batch_to_reward_seconds=w['at']-origin))
startup = [x['pod_to_ready_seconds'] for x in rows]
return dict(episodes=count, successful_episodes=len(rows), reward_sum=sum(x['reward'] for x in rows),
batch_to_all_ready_seconds=max(x['batch_to_ready_seconds'] for x in rows),
batch_to_all_rewards_seconds=max(x['batch_to_reward_seconds'] for x in rows),
pod_to_ready_p50_seconds=percentile(startup, .5),
pod_to_ready_p95_seconds=percentile(startup, .95),
pod_to_ready_max_seconds=max(startup), pods=rows)


def main():
parser = argparse.ArgumentParser()
parser.add_argument('action', choices=['run', 'cleanup'])
parser.add_argument('--output', required=True)
parser.add_argument('--count', type=int, default=10)
parser.add_argument('--parallelism', type=int)
parser.add_argument('--image', default='python:3.12-alpine')
parser.add_argument('--runtime', default='cube')
args = parser.parse_args()
output = Path(args.output)
parallel = args.count if args.parallelism is None else args.parallelism
if not 1 <= parallel <= args.count <= 20:
parser.error('require 1 <= parallelism <= count <= 20')
if args.action == 'cleanup':
receipt = json.loads((output/'receipt.json').read_text())
context = receipt['context']
else:
context = subprocess.check_output(['kubectl', 'config', 'current-context'], text=True).strip()
def kubectl(*parts, data=None):
return subprocess.check_output(['kubectl', '--context', context, '--request-timeout=30s', *parts],
input=json.dumps(data) if data is not None else None, text=True, timeout=60)
def create(obj):
return json.loads(kubectl('create', '-f', '-', '-o', 'json', data=obj))
if args.action == 'cleanup':
ns = receipt['namespace']
raw = kubectl('get', 'namespace', ns, '--ignore-not-found', '-o', 'json')
if not raw.strip():
print('Namespace already absent'); return
live = json.loads(raw)
if live['metadata']['uid'] != receipt['uid'] or live['metadata'].get('labels', {}).get('agc-demo-run') != receipt['run_id']:
raise RuntimeError('Ownership check failed; refusing cleanup')
raw_job = kubectl('-n', ns, 'get', 'job', 'rollout', '--ignore-not-found', '-o', 'json')
if raw_job.strip():
expected_uid = receipt.get('job_uid')
if expected_uid is None and (output/'job.json').exists():
expected_uid = json.loads((output/'job.json').read_text())['metadata']['uid']
if json.loads(raw_job)['metadata']['uid'] != expected_uid:
raise RuntimeError('Job ownership check failed; refusing cleanup')
kubectl('-n', ns, 'delete', 'job', 'rollout', '--cascade=foreground', '--wait=false')
deadline = time.monotonic() + 120
while kubectl('-n', ns, 'get', 'job', 'rollout', '--ignore-not-found', '-o', 'name').strip():
if time.monotonic() > deadline: raise TimeoutError('Job or Pods still terminating')
time.sleep(2)
kubectl('delete', 'namespace', ns, '--wait=false')
deadline = time.monotonic() + 120
while kubectl('get', 'namespace', ns, '--ignore-not-found', '-o', 'name').strip():
if time.monotonic() > deadline: raise TimeoutError('Namespace still terminating; inspect finalizers')
time.sleep(2)
print('Deleted test namespace:', ns); return
output.mkdir(parents=True, exist_ok=False)
run_id = uuid.uuid4().hex[:12]
ns = 'agc-rl-' + run_id
receipt = dict(context=context, namespace=ns, run_id=run_id)
# Save intended namespace before any cluster mutation so failures remain traceable.
(output/'receipt.json').write_text(json.dumps(receipt, indent=2))
kubectl('get', 'runtimeclass', args.runtime)
created = create(dict(apiVersion='v1', kind='Namespace', metadata=dict(name=ns, labels={'agc-demo-run':run_id})))
receipt['uid'] = created['metadata']['uid']
(output/'receipt.json').write_text(json.dumps(receipt, indent=2))
print('Test namespace:', ns, flush=True)
create(dict(apiVersion='v1', kind='ConfigMap', metadata=dict(name='episode', namespace=ns), data={'episode.py':WORKER}))
job = dict(apiVersion='batch/v1', kind='Job', metadata=dict(name='rollout', namespace=ns), spec=dict(
completions=args.count, parallelism=parallel, backoffLimit=0, activeDeadlineSeconds=240,
template=dict(metadata=dict(labels={'app':'rl-batch-demo'}), spec=dict(
runtimeClassName=args.runtime, restartPolicy='Never', automountServiceAccountToken=False,
securityContext=dict(runAsUser=1000, runAsGroup=1000, fsGroup=1000, runAsNonRoot=True),
containers=[dict(name='episode', image=args.image, imagePullPolicy='IfNotPresent',
command=['python','-B','/demo/episode.py'],
env=[dict(name='EPISODE_ID',valueFrom=dict(fieldRef=dict(fieldPath='metadata.name')))],
securityContext=dict(allowPrivilegeEscalation=False, capabilities=dict(drop=['ALL']), seccompProfile=dict(type='RuntimeDefault')),
resources=dict(requests=dict(cpu='100m',memory='64Mi'),limits=dict(cpu='500m',memory='128Mi')),
volumeMounts=[dict(name='script',mountPath='/demo',readOnly=True),dict(name='workspace',mountPath='/workspace')])],
volumes=[dict(name='script',configMap=dict(name='episode')),dict(name='workspace',emptyDir={})]))))
(output/'job-submitted.json').write_text(json.dumps(job, indent=2))
start = time.monotonic()
created_job = create(job)
receipt['job_uid'] = created_job['metadata']['uid']
(output/'receipt.json').write_text(json.dumps(receipt, indent=2))
deadline = start + 270
while True:
live = json.loads(kubectl('-n',ns,'get','job','rollout','-o','json'))
conditions = {c['type']:c['status'] for c in live.get('status',{}).get('conditions',[])}
if conditions.get('Complete')=='True' or conditions.get('Failed')=='True': break
if time.monotonic()>deadline: break
time.sleep(1)
observed = time.monotonic()-start
pods = json.loads(kubectl('-n',ns,'get','pods','-l','job-name=rollout','-o','json'))['items']
(output/'job.json').write_text(json.dumps(live, indent=2))
(output/'pods.json').write_text(json.dumps(pods, indent=2))
(output/'events.txt').write_text(kubectl('-n',ns,'get','events','--sort-by=.metadata.creationTimestamp'))
def get_log(pod):
name = pod['metadata']['name']
try: log = kubectl('-n',ns,'logs',name,'-c','episode')
except subprocess.SubprocessError: log = ''
(output/(name+'.log')).write_text(log)
return log
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as pool:
logs = list(pool.map(get_log,pods))
if conditions.get('Complete')!='True':
raise RuntimeError('Job failed or timed out; inspect saved status/events/logs, then cleanup')
report = collect(pods, logs, live, args.count)
report.update(context=context, namespace=ns, image=args.image, runtime=args.runtime,
parallelism=parallel, cache_state='uncontrolled', client_job_complete_seconds=observed)
(output/'results.json').write_text(json.dumps(report, indent=2))
print(json.dumps({k:v for k,v in report.items() if k!='pods'},indent=2))

if __name__=='__main__':
main()