本实践将代码执行平台的一次文本统计请求放到 Agent 集群运行:本机程序接收文本,通过 Kubernetes API 创建独立 Pod,等待执行完成,读取 JSON 结果并删除 Pod。平台已有的用户交互和任务编排流程可以调用同一个函数。
准备工作
先完成 Kubernetes 快速入门,在本机配置集群 kubeconfig,并准备
agc-docs-demo 命名空间。示例使用 Python 3 和 kubectl,容器镜像需要包含 Python 3 并支持 UID 1000 运行。设置运行时和镜像:
export AGC_RUNTIME_CLASS='REPLACE_WITH_RUNTIMECLASS'export AGC_PYTHON_IMAGE='REPLACE_WITH_PYTHON_IMAGE'
创建任务执行器
将以下代码保存为
run-agent-task.py:import jsonimport osimport subprocessimport sysimport timeimport uuidNAMESPACE = 'agc-docs-demo'TASK_CODE = 'import json,os; print(json.dumps({"word_count":len(os.environ["INPUT_TEXT"].split())}))'def kubectl(*args, data=None):result = subprocess.run(['kubectl', '-n', NAMESPACE, *args], input=data,text=True, capture_output=True, timeout=30, check=True,)return result.stdoutdef run_task(text):runtime = os.environ['AGC_RUNTIME_CLASS']image = os.environ['AGC_PYTHON_IMAGE']if any(v.startswith('REPLACE_WITH_') or not v for v in (runtime, image)):raise ValueError('请设置 RuntimeClass 和 Python 镜像')name = 'agent-task-' + uuid.uuid4().hex[:12]pod = {'apiVersion': 'v1', 'kind': 'Pod','metadata': {'name': name, 'namespace': NAMESPACE},'spec': {'runtimeClassName': runtime, 'restartPolicy': 'Never','automountServiceAccountToken': False, 'activeDeadlineSeconds': 120,'containers': [{'name': 'runner', 'image': image,'command': ['python3', '-c', TASK_CODE],'env': [{'name': 'INPUT_TEXT', 'value': text}],'resources': {'requests': {'cpu': '100m', 'memory': '64Mi'},'limits': {'cpu': '500m', 'memory': '128Mi'},},'securityContext': {'runAsNonRoot': True, 'runAsUser': 1000,'allowPrivilegeEscalation': False,'capabilities': {'drop': ['ALL']},'seccompProfile': {'type': 'RuntimeDefault'},},}],},}# 创建前记录资源名,响应丢失时也能查询原任务。print('task_pod:', name, file=sys.stderr, flush=True)created = Falsetry:kubectl('create', '-f', '-', data=json.dumps(pod))created = Truedeadline = time.monotonic() + 180while time.monotonic() < deadline:current = json.loads(kubectl('get', 'pod', name, '-o', 'json'))phase = current.get('status', {}).get('phase')if phase in ('Succeeded', 'Failed'):output = kubectl('logs', name, '-c', 'runner')if phase == 'Failed':raise RuntimeError('任务失败:' + output)result = json.loads(output)if type(result.get('word_count')) is not int:raise ValueError('任务返回的 JSON 缺少有效 word_count')return {'task_pod': name, 'result': result}time.sleep(1)raise TimeoutError('等待任务超时:' + name)finally:if created:try:kubectl('delete', 'pod', name, '--wait=true', '--timeout=20s')except (subprocess.SubprocessError, OSError) as error:print('清理失败,请核对 Pod ' + name + ': ' + str(error), file=sys.stderr)if __name__ == '__main__':text = sys.argv[1] if len(sys.argv) > 1 else 'hello agent cluster'print(json.dumps(run_task(text), ensure_ascii=False))
运行:
python3 run-agent-task.py 'hello agent cluster'
预期返回包含任务 Pod 名称和
"word_count": 3 的 JSON。输入文本通过环境变量传递,不拼入 Shell 命令。接入已有平台
在平台的任务执行模块中调用
run_task(text),将返回结果写入平台的任务记录。用户身份和模型调用继续由原平台管理;执行器负责本次 Pod 的创建、等待、结果读取和清理。平台为每次任务分配业务 ID,并关联返回的
task_pod。生产任务需要保存输入引用、执行状态、结果位置和重试次数。大文件使用工作区或外部存储交付,避免放入环境变量。处理失败
情况 | 处理方式 |
创建请求失败或超时 | 用已经输出的 Pod 名称查询是否创建成功,再决定重试或清理。 |
Pod Pending | 查看事件,检查容量、镜像和 RuntimeClass。 |
程序失败或返回格式错误 | 保存异常信息,检查输入和代码;脚本会尝试删除已创建的 Pod。 |
等待任务超时 | 检查执行时长,脚本会尝试结束已创建的 Pod。 |
清理失败 | 按标准错误输出中的 Pod 名称查询和清理残留资源。 |