push subscription由Pub/Sub控制并发上限,容易在长时间模型调用时超配额;pull subscription配合--concurrency和--max-instances可精确控制模型调用速率,保护共享配额。
推送订阅把你的 worker 变成一个 HTTP 服务器,Pub/Sub 向它发起调用。这在任务耗时不超过 ack 超时时限时非常方便,但一旦处理时间超过确认期限,Pub/Sub 就会开始把同一条消息投递给第二个实例——而第一个实例还在和模型通信。
Push 模式及其对并发的影响
使用 pull 订阅时,worker 自己决定一次拿多少任务。使用 push 订阅则由 Pub/Sub 来决定,而你唯一的调节杠杆是服务自身的最大实例数以及单实例并发数。对于一个调用模型的 worker,这件事的意义比普通 Web 服务大得多,因为你需要保护的并不是 CPU——而是所有运行实例共享的 provider 配额。
Cloud Run 默认的单实例并发数很高,这个设定适合那些花大量时间等待 I/O 的服务,但如果每个飞行中的请求都会占用一个 provider 槽位,那就不适用了。把 --concurrency 设置为你希望单个实例同时调用模型的次数,把 --max-instances 设为上限。两个数的乘积就是你向 provider 发起的实际请求速率,这个乘法值得找个地方记下来。
处理器收到的信封
Pub/Sub 会把消息包装一下。你的处理器收到的 JSON body 里有一个 message 对象,其中包含 base64 编码的数据、messageId、publishTime、可选的 attributes、可选的 orderingKey,以及——如果你配置了死信 topic——一个 deliveryAttempt 计数器。Google 同时记录了 id 和时间戳字段的驼峰命名和蛇形命名两种拼写方式,所以要防御性读取,不要假设只有某一种。
deliveryAttempt 字段是值得使用的那个。它是你拥有的最廉价的信号,能告诉你这条消息之前已经被尝试过,在这个工作负载上意味着上一次尝试可能已经产生了一笔计费调用。
import base64, json, os
from flask import Flask, request
app = Flask(__name__)
@app.post("/")
def handle():
envelope = request.get_json(silent=True) or {}
message = envelope.get("message", {})
payload = json.loads(base64.b64decode(message["data"]).decode())
attempt = envelope.get("deliveryAttempt", 1)
if already_done(payload["job_id"]):
return "", 204 # ack: nothing to redo
try:
write_result(payload["job_id"], call_model(payload["prompt"]))
except RetryableError:
return "", 500 # nack: Pub/Sub redelivers
return "", 204
状态码不是一种约定,而是协议本身。Google 文档说明返回 102、200、201、202 或 204 会确认消息,而任何其他状态码都是否定确认,会导致重新投递。路由错误的 404 是 nack。缺少调用者权限的 403 是 nack,而且是永久性的,以完整重试速率重试——这就是一个权限错误如何变成一张账单的原因。
Ack 超时与 Cloud Run 请求超时
同一个请求上有两个独立的时钟在运行,而且两个都可以终止它。Google 文档说明订阅确认期限默认为 10 秒,最小 10 秒,最大 600 秒,并指出你无法修改通过 push 订阅接收到的单条消息的期限——pull 订阅者可用的逐消息扩展技巧在你这里用不了。另外,Cloud Run 的请求超时默认 5 分钟(300 秒),可以延长到 60 分钟(3600 秒)。
所以 push 投递的模型调用的上限是 600 秒,无论你在 Cloud Run 上设置的是什么。把 push 触发服务的 Cloud Run 超时提高到十分钟以上什么都买不到,只会得到一个更长的窗口——在此期间 Pub/Sub 已经放弃并重新投递了。把 ack 超时设置为略高于你最坏情况的调用时间,把 Cloud Run 超时设置为略高于 ack 超时,然后把 600 秒作为一个硬性的架构上限:超过这个上限,答案应该是 pull 订阅或 Cloud Tasks,而不是一个更大的数字。
期限范围、保留期和死信数字来自 Google 的 Pub/Sub 订阅属性文档,超时数字来自 Cloud Run 的请求超时页面,两者均为 2026 年 8 月读取。Google:订阅属性
还有两个默认值值得主动设置而不是继承。消息保留期默认为 7 天,文档记录的范围是 10 分钟到 31 天——一次故障后重放一周累积的模型任务本身就是一起事件。死信 topic 默认 5 次投递尝试,可配置为 5 到 100 之间的任意数字,这个设置的作用是阻止一条永久有毒的消息被重试到保留窗口过期。
对私有服务进行身份验证
部署服务时加上 --no-allow-unauthenticated。然后给 push 订阅一个 oidcToken,指定一个服务账号,并授予该服务账号在服务上的 Cloud Run Invoker 角色,在处理器前面验证产生的 bearer token——或者让 Cloud Run 自己的 IAM 检查来做这件事,这也是把服务保持为私有的原因。
有一项授权人们容易漏掉:Pub/Sub 自己的服务代理需要获得以你指定的服务账号身份来铸造 token 的权限,也就是说需要在那个账号上拥有 Service Account Token Creator 角色。代理的地址是项目特定的,从你项目的 IAM 页面读取 Google 托管的账号,而不是从博客文章里复制一个地址。
创建 topic,然后部署 Cloud Run 服务,加上 --no-allow-unauthenticated、显式的 --concurrency、显式的 --max-instances,以及高于你最坏情况模型调用的 --timeout。
为 push 创建一个服务账号,授予它在服务上的 roles/run.invoker,并授予 Pub/Sub 服务代理在该服务账号上的 Service Account Token Creator 角色。
创建 push 订阅,加上 --push-auth-service-account、按调用时间调整大小的 --ack-deadline,以及带有 --max-delivery-attempts 的死信 topic。
发布一条消息并确认出现了一行结果。然后发布一条会超过 ack 超时的消息,确认你可以在自己的日志中看到重复投递——你需要亲眼看过一次这种事发生。
加上基于 messageId 或你自己的 job id 的幂等性检查,并确认重复现在不会产生任何费用。
gcloud run deploy inference-worker \
--image europe-docker.pkg.dev/PROJECT/repo/worker:1 \
--no-allow-unauthenticated --concurrency 4 --max-instances 20 --timeout 420
gcloud pubsub subscriptions create inference-push \
--topic inference-jobs \
--push-endpoint https://inference-worker-xxxx.europe-west1.run.app/ \
--push-auth-service-account [email protected] \
--ack-deadline 400 \
--dead-letter-topic inference-dead \
--max-delivery-attempts 5
Google Cloud Tasks for Rate-Limited Model API Calls
Long Polling and Push Delivery for a Queue Feeding a Model Pipeline
Idempotency Keys for a Queued Model Request That Might Retry