- 后端
- 云原生
- 模型推理服务
- MLOps
- 人工智能
【免费下载链接】cortex
Production infrastructure for machine learning at scale
导读
本文以 Cortex 开源仓库中的 Batch API 实战示例文档为核心,完整演示"从零搭建一个可分布式执行批量处理任务的 Batch API"的端到端流程:先编写基于 FastAPI 的处理逻辑,再制作镜像并推送到 AWS ECR,随后编写 cortex.yaml 配置并执行cortex deploy部署,最后通过 POST 请求提交批量作业、用cortex get/cortex logs查看状态与日志。读完本文,你将掌握 Batch API 的定义、部署、作业提交与观测的完整链路,并能结合仓库源码理解 enqueuer/dequeuer 背后的分批与队列机制。
一、Batch API 是什么:一次讲清适用场景与核心能力
在进入示例之前,先明确 Batch API 的定位。仓库文档 docs/workloads/batch/batch.md 定义:Batch APIs 按需运行分布式、可容错的批处理作业,非常适合将工作负载拆散并分发到一组专用 worker 上执行,例如对一批图片批量跑推理。
Batch API 的关键特性包括:
- 将一个批量作业(batch job)分发到多个 worker 并行处理;
- 没有批量作业时自动缩容到 0(scale to 0),不占用资源;
- 所有批次处理完成后,自动触发
/on-job-complete钩子; - 每个批次至少被尝试一次(at-least-once 语义);
- 失败的批次会被转投到死信队列(dead letter queue);
- 能从失败与 Spot 实例终止中自动恢复。
从工作流程看(详见 batch.md):部署 Batch API 后,Cortex 会创建一个接收作业提交的端点;提交作业后返回一个 Job ID,并异步触发一个 Batch Job。Batch Job 启动时先部署一个enqueuer进程,把作业数据拆成批次推入SQS FIFO 队列;入队完成后,Cortex 初始化指定数量的 worker pod,并为每个 pod 挂载dequeuer sidecar,由它从队列取批次、向你的 pod 发起 HTTP 请求;当队列被清空后,作业标记为完成,worker pod 被终止、SQS 队列被删除。作业期间可以通过 GET 请求查看作业状态以及已完成/失败批次数量等指标。
二、第一步:用 FastAPI 定义一个 Batch API 处理器
示例文档 docs/workloads/batch/example.md 中的处理器是一个标准的 FastAPI 应用,核心是两个端点:
# main.py from fastapi import FastAPI from typing import List app = FastAPI() @app.post("/") def handle_batch(batch: List[int]): print(batch) @app.post("/on-job-complete") def on_job_complete(): print("done")两个端点的职责如下:
POST /(handle_batch):worker 的 dequeuer sidecar 每从队列取出一个批次,就会向本端点发起一次 HTTP 请求,请求体即该批次的 items。所以这个函数每批次被调用一次,在这里实现你的核心业务逻辑(比如对一批图片跑推理)。POST /on-job-complete:所有批次处理完毕、队列清空后触发一次。注意它只在整个作业中触发一次(跨所有 worker),适合做结果汇总、写回 S3 等收尾工作。
仓库里的完整示例可以参考 test/apis/batch/sum/app/main.py,它展示了更贴近生产的使用方式:在startup事件中读取挂载到容器内的作业规格/cortex/spec/job.json(包含job_id、config等),通过/healthz就绪探针暴露健康状态,并在on_job_complete中把汇总结果通过 boto3 写入 S3。该示例还说明了config字段的典型用法:提交作业时传入dest_s3_dir,容器内据此计算输出位置。
三、第二步:编写 Dockerfile 并本地验证
定义好处理器后,为其编写容器镜像:
FROM python:3.8-slim RUN pip install --no-cache-dir fastapi uvicorn COPY main.py / CMD uvicorn --host 0.0.0.0 --port 8080 main:app在本地依次构建镜像、运行容器并验证接口是否正常工作:
docker build . -t hello-worlddocker run -p 8080:8080 hello-worldcurl -X POST -H "Content-Type: application/json" -d '[1,2,3,4]' localhost:8080本地请求体是一个 JSON 数组[1,2,3,4],对应handle_batch接收的参数类型List[int],可在容器日志中看到打印结果。这一步确认应用本身可用,再进行镜像推送。
四、第三步:把镜像推送到 AWS ECR
Cortex 集群中的 worker pod 需要从镜像仓库拉取镜像,因此需要先把镜像推送到你的 AWS ECR。依次执行:
- 登录 ECR:
aws ecr get-login-password --region us-east-1 | docker login --username AWS --password-stdin <AWS_ACCOUNT_ID>.dkr.ecr.us-east-1.amazonaws.com- 创建仓库:
aws ecr create-repository --repository-name hello-world- 打标签并推送:
docker tag hello-world <AWS_ACCOUNT_ID>.dkr.ecr.us-east-1.amazonaws.com/hello-worlddocker push <AWS_ACCOUNT_ID>.dkr.ecr.us-east-1.amazonaws.com/hello-world仓库中另有 ECR 相关辅助脚本(见 dev/delete_ecr_repos.py),可参考其使用的 AWS API 了解 ECR 仓库管理方式。
五、第四步:编写 cortex.yaml 部署配置
在项目根目录创建cortex.yaml,声明一个名为hello-world的 BatchAPI:
# cortex.yaml - name: hello-world kind: BatchAPI pod: containers: - name: api image: <AWS_ACCOUNT_ID>.dkr.ecr.us-east-1.amazonaws.com/hello-world command: ["uvicorn", "--host", "0.0.0.0", "--port", "8080", "main:app"]注意这里显式指定了command覆盖镜像中的CMD。关于该配置文件的完整字段,请以 docs/workloads/batch/configuration.md 为准,下面摘录最核心的字段并补充默认值:
name(必填):API 名称;kind(必填):Batch API 必须为"BatchAPI";pod.port:请求发送到的端口,默认 8080,会以环境变量$CORTEX_PORT导出。仓库示例 test/apis/batch/sum/cortex_cpu.yaml 中即使用"$(CORTEX_PORT)"引用它;pod.containers(至少一个容器):name(必填)、image(必填);command:entrypoint(不经 shell 执行),可用$(CORTEX_PORT)形式引用环境变量;args:entrypoint 参数,默认无;env:环境变量字典;compute:资源请求,cpu默认 200m(一个 CPU 单位对应一个虚拟 CPU,支持小数与m后缀)、gpu默认 0、inf(Inferentia 芯片)默认 0、mem默认 Null(支持 K/M/G/T 及二进制 Ki/Mi/Gi/Ti 后缀)、shm默认 Null(如64Mi、1Gi);readiness_probe/liveness_probe:HTTP GET、TCP socket 或 exec 探针,以及initial_delay_seconds(默认 0)、timeout_seconds(默认 1)、period_seconds(默认 10)、success_threshold(默认 1)、failure_threshold(默认 3)。仓库的 sum 示例就配置了指向/healthz的readiness_probe;
node_groups:可运行的节点组列表,默认所有节点组均可;networking.endpoint:API 端点,默认与 API 同名。
六、第五步:部署并获取端点
执行部署命令,Cortex 会解析 cortex.yaml 并在集群中创建对应的 Batch API:
cortex deploy部署完成后获取 API 信息,其中endpoint字段就是后续提交作业要用的地址:
cortex get hello-world七、第六步:提交批量作业(三种数据来源)
Batch API 通过HTTP POST提交作业,提交后异步执行并立即返回 Job ID。示例文档中的请求格式为:
curl -X POST -H "Content-Type: application/json" -d '{"workers": 2, "item_list": {"items": [1,2,3,4], "batch_size": 2}}' http://***.amazonaws.com/hello-world该请求声明使用 2 个 worker,items共 4 个样本按batch_size: 2拆成 2 个批次,每个批次会调用一次handle_batch。作业提交的完整 schema 与三种数据来源方式定义在 docs/workloads/batch/jobs.md 与 pkg/operator/schema/job_submission.go 中:
- 数据随请求提交(
item_list):items中每个元素(可以是任意类型:对象、列表、字符串等)作为一个样本,按batch_size聚合成批次。每个批次必须小于 256 KiB,且整个请求小于 10 MiB(pkg/operator/endpoints/submit_batch.go 中通过http.MaxBytesReader(w, r.Body, 10<<20)施加 10 MiB 限制)。适合样本数量少、单个样本小、想避免中间存储的场景; - S3 文件路径列表(
file_path_lister):通过s3_paths指定文件或前缀,配合includes/excludes过滤,按batch_size聚合路径。适合图片/视频等 S3 目录场景,单个文件代表少量样本; - S3 中的换行分隔 JSON 文件(
delimited_files):逐行解析 S3 JSON 文件,每行一个 JSON 对象作为一个样本,按batch_size拆批。适合单个文件包含大量样本需要拆分的场景。
三种方式均可选timeout(提交后多少秒强制终止作业)、sqs_dead_letter_queue(指定死信队列 ARN 与max_receive_count,批次被 worker 处理超过该次数后转投死信队列)以及任意的config字典(作业专属参数)。提交成功的响应包含job_id、workers、sqs_url、timeout、created_time等字段。此外,整个作业规格会被写入容器内的/cortex/spec/job.json,方便容器启动时读取(sum 示例的 startup 逻辑即依赖这一点)。
需要提一下提交入口的实现:submit_batch.go 中的SubmitBatchJob会先校验 API 的 kind 必须是BatchAPI,反序列化作业提交请求,然后调用batchapi.SubmitJob完成入队;它还支持dryRun=true查询参数,在真正提交前输出将被处理的目标文件列表并提示 "validations passed",可用于核对过滤结果。
从底层实现看,pkg/enqueuer/enqueuer.go 中的Enqueue按三种来源分别调用enqueueItems、enqueueS3Paths、enqueueS3FileContents完成分批与写入 SQS FIFO 队列,最后还会向队列投递一条带job_complete消息属性的占位消息用于标记作业完成,并将批次总数写入 S3(UploadBatchCount)。
八、第七步:查看日志
作业运行中或结束后,可查看指定 Job 的日志:
cortex logs hello-world <JOB_ID>日志会包含 worker 容器内应用自身打印的输出(例如handle_batch的 print 结果)。sum 示例的handle_batch中会打印从/cortex/spec/job.json读到的作业规格,on_job_complete中会打印汇总结果,这些都能在cortex logs里直接看到。
九、作业状态、指标与终止(补充实战能力)
除了示例文档,批处理作业的日常运维还包括状态查询与停止(详见 jobs.md):
- 查看作业状态:
cortex get <batch_api_name> <job_id>,或对端点发起GET <batch_api_endpoint>?jobID=<jobID>。响应中的job_status包含status、batches_in_queue(队列中剩余批次)、worker_counts(pending/initializing/running/succeeded/failed/stalled,其中 stalled 表示卡在 pending 超过 10 分钟)以及start_time/end_time;metrics字段给出succeeded(成功批次数)、failed(失败尝试数)与avg_time_per_batch(每个批次平均处理时间,仅统计成功尝试); - 停止作业:
cortex delete <batch_api_name> <job_id>,或对端点发起DELETE <batch_api_endpoint>?jobID=<jobID>,响应为{"message":"stopped job <job_id>"}。
作业生命周期中可能出现的状态定义在 docs/workloads/batch/statuses.md 中:
| 状态 | 含义 |
|---|---|
| enqueuing | 作业正在被拆分成批次并放入队列 |
| running | worker 正在从队列获取批次并执行 |
| succeeded | worker 无失败地完成了队列中所有条目 |
| failed while enqueuing | 入队阶段发生失败,需查看作业日志 |
| completed with failures | 队列处理完毕,但部分批次未能成功处理并抛出了异常 |
| worker error | 一个或多个 worker 发生不可恢复错误,导致作业失败 |
| out of memory | 一个或多个 worker 内存耗尽导致作业失败 |
| timed out | 作业在达到指定 timeout 后被终止 |
| stopped | 作业被用户停止,或 Batch API 被删除 |
补充一点底层机制:worker pod 侧的 dequeuer 实现在 pkg/dequeuer/dequeuer.go 中,它使用 SQS 长轮询(WaitTimeSeconds10 秒)与 30 秒的 visibility timeout 保证批次消费的可靠性,并按提交时指定的workers数量启动等量的消费协程——这正是"分布式并行 + 至少一次尝试"语义的直接来源。
十、完整示例速览:仓库中的 sum BatchAPI
仓库 test/apis/batch/sum 目录提供了一个完整的可复现样例,可作为理解整套流程的最佳参考:
- app/main.py:FastAPI 处理器,startup 时读取 job spec、
/healthz就绪探针、handle_batch累加求和、on_job_complete把结果写入 S3; - cortex_cpu.yaml:BatchAPI 配置,使用
$(CORTEX_PORT)启动 uvicorn 并配置就绪探针与计算资源(cpu: 200m、mem: 256Mi); - sample.json:两个样本列表,作为
item_list.items的数据来源; - submit.py:通过
cortex.client(env_name)获取端点后,以{"workers": 1, "item_list": {"items": ..., "batch_size": 1}, "config": {"dest_s3_dir": ...}}提交作业的参考实现。
总结
从示例文档出发,本文走完了 Batch API 的完整生命周期:定义 FastAPI 处理器与on-job-complete钩子 → 构建镜像并推送 ECR → 编写 cortex.yaml →cortex deploy部署 → POST 提交作业 →cortex get/cortex logs观测状态与日志。深入仓库源码可以看到,Cortex 用 enqueuer 将三种数据来源(请求内联数据、S3 路径列表、S3 换行分隔 JSON)拆成批次写入 SQS FIFO 队列,再由挂载在每个 worker pod 上的 dequeuer sidecar 并行消费并调用你的handle_batch,最终以on-job-complete钩子收尾。掌握这套模式后,你可以将任意可拆分的计算任务(尤其是批量推理场景)接入 Cortex,获得开箱即用的分布式调度、故障恢复与资源弹性。
- 后端
- 云原生
- 模型推理服务
- MLOps
- 人工智能
【免费下载链接】cortex
Production infrastructure for machine learning at scale
相关推荐
Bark 推送服务端部署完全指南:Docker、Compose、手动部署与批量推送优化
Bark 推送服务端部署完全指南:Docker、Compose、手动部署与批量推送优化 本指南围绕 Bark(iOS 自定义推送 App)配套服务端 bark
开发工具移动开发Qwen3部署实战:从本地推理到云端服务
Qwen3部署实战:从本地推理到云端服务 本文全面介绍了Qwen3大语言模型的多种部署方案,涵盖了从本地CPU推理优化到云端高性能服务的完整技术栈。详细讲解了T
人工智能大模型Qwen模型评测示例工程本地部署教程Apache APISIX syslog 插件实战指南:批量推送请求/响应日志到 Syslog 服务器
Apache APISIX syslog 插件实战指南:批量推送请求/响应日志到 Syslog 服务器 导读 syslog 是 Apache APISIX 内置
API网关后端云原生微服务
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考