训练 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-contextkubectl get nodeskubectl 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-10python3 rl-batch-demo.py run --count 10 --parallelism 1 --output rl-sequential-10cat 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,podskubectl -n YOUR_TEST_NAMESPACE get events --sort-by=.metadata.creationTimestampkubectl -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 argparseimport concurrent.futuresimport datetimeimport jsonimport mathimport osfrom pathlib import Pathimport subprocessimport sysimport timeimport uuidWORKER = r'''import json, os, pathlib, subprocess, timep = 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 unittestfrom shipping import shippingclass 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.stderremit('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) * qlo, 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.parallelismif 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'); returnlive = 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() + 120while 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() + 120while 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); returnoutput.mkdir(parents=True, exist_ok=False)run_id = uuid.uuid4().hex[:12]ns = 'agc-rl-' + run_idreceipt = 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 + 270while 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': breakif time.monotonic()>deadline: breaktime.sleep(1)observed = time.monotonic()-startpods = 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 logwith 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()