☰
Cortex BatchAPI 端到端实战:从 FastAPI 服务编写、镜像推送、部署到批量作业提交与日志查看
2026/9/27 7:05:37 网站建设 项目流程
  • 后端
  • 云原生
  • 模型推理服务
  • MLOps
  • 人工智能

【免费下载链接】cortex

Production infrastructure for machine learning at scale

项目地址:https://gitcode.com/gh_mirrors/co/cortex
点击查看免费下载

导读

本文以 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-world
docker run -p 8080:8080 hello-world
curl -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。依次执行:

  1. 登录 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
  1. 创建仓库:
aws ecr create-repository --repository-name hello-world
  1. 打标签并推送:
docker tag hello-world <AWS_ACCOUNT_ID>.dkr.ecr.us-east-1.amazonaws.com/hello-world
docker 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 中:

  1. 数据随请求提交(item_list):items中每个元素(可以是任意类型:对象、列表、字符串等)作为一个样本,按batch_size聚合成批次。每个批次必须小于 256 KiB,且整个请求小于 10 MiB(pkg/operator/endpoints/submit_batch.go 中通过http.MaxBytesReader(w, r.Body, 10<<20)施加 10 MiB 限制)。适合样本数量少、单个样本小、想避免中间存储的场景;
  2. S3 文件路径列表(file_path_lister):通过s3_paths指定文件或前缀,配合includes/excludes过滤,按batch_size聚合路径。适合图片/视频等 S3 目录场景,单个文件代表少量样本;
  3. 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作业正在被拆分成批次并放入队列
runningworker 正在从队列获取批次并执行
succeededworker 无失败地完成了队列中所有条目
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

项目地址:https://gitcode.com/gh_mirrors/co/cortex
点击查看免费下载

相关推荐

上一篇:快速将复杂PDF转Markdown:Marker五分钟上手指南
下一篇:探索高效媒体处理:强大的远程FFmpeg工具——`rffmpeg`

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询