vercel-workers 版本演进全解析:Python 侧 Vercel Queues 与 Worker Services 的 SDK 能力图谱
【免费下载链接】vercelDevelop. Preview. Ship.项目地址: https://gitcode.com/gh_mirrors/ve/vercel
导读
vercel-workers是 Vercel 官方开源仓库中面向 Python 的队列与 Worker 服务 SDK,提供send()消息发布、@subscribe消费原语,以及 Celery、Dramatiq、Django tasks 三类任务框架适配器。本文以仓库内 python/vercel-workers/CHANGELOG.md 为骨架,逐版本梳理 0.0.10 → 0.0.25 的功能演进脉络,并结合源码与示例验证每项变更背后的真实实现,帮助读者快速判断该 SDK 当前具备哪些能力、底层如何工作,以及如何在自己的 Python 服务中正确使用。
版本脉络:一个从"初始导入"到"企业级队列 SDK"的演进史
vercel-workers的 CHANGELOG 记录了两个阶段:0.0.10 的 Initial Release与随后持续迭代的 0.0.11 ~ 0.0.25(每个版本均以patch形式推进)。当前包版本为 0.0.25,声明于 pyproject.toml 中,要求requires-python = ">=3.12"。
0.0.10:Initial Release
- 将
vercel-workersPython 包初始导入 monorepo; - 提供Celery、Django tasks、Dramatiq三套适配器,用于对接 Vercel Queues。
这一定位延续至今,在 README.md 中有明确描述:Python SDK for Vercel Queues and Vercel Worker Services,核心原语即send()与@subscribe。
核心 API 层演进:从client到_queue/的内部重构
0.0.20 / 0.0.21:框架逻辑下沉与队列 SDK 拆分
- 0.0.20:将 framework-specific 逻辑重构进
vercel-workers; - 0.0.21:把 Python 队列 SDK 重构进
_queue/目录。
从当前源码结构看,该重构的结果是 src/vercel/workers/_queue/ 成为队列能力的底层实现区,包含:
| 文件 | 职责 |
|---|---|
client.py | 发送消息的同步/异步客户端(send、send_async),HTTP 调用 Vercel Queue Service V3 API |
subscribe.py | @subscribe装饰器、Subscription注册表、Ack/RetryAfter指令、payload 类型校验 |
callback.py | 回调解析:CloudEvent v1beta 与 v2beta 两种格式的解析与识别 |
receive.py | 消息接收、可见性超时处理 |
send.py | 发送请求构建:URL、鉴权头、幂等键、保留/延迟头 |
types.py | WorkerJSONEncoder、MessageMetadata、SendMessageResult等类型定义 |
exceptions.py | 全套错误类型 |
而面向用户的公共入口 src/vercel/workers/init.py 重新导出QueueClient、AsyncQueueClient、send、subscribe、Ack、RetryAfter、WorkerJSONEncoder、get_wsgi_app、get_asgi_app等。公共 API 保持稳定,内部实现被模块化——这是该 SDK 演进中最重要的架构决策。
0.0.22:新增 QueueClient 与 AsyncQueueClient
QueueClient与AsyncQueueClient在 src/vercel/workers/_queue/client.py 中实现,为需要以"客户端对象"方式管理队列连接的用户提供面向对象接口,与函数式send()/send_async()并行。
0.0.18:从公共 API 移除consumer
consumer不再作为公开 API 暴露。当前公开导出列表(见init.py 的__all__)中确实已无consumer,用户配置消费组改为在vercel.json的 worker service 声明中完成。
send() 发送链路:参数、请求头与底层实现
0.0.21:retention/delay 支持 timedelta
0.0.21 用retention与delay取代了retention_seconds与delay_seconds,并支持datetime.timedelta传参,例如retention=timedelta(hours=6)。在 src/vercel/workers/_queue/send.py 中可见完整处理逻辑:
_resolve_duration_alias()同时兼容新旧参数名,但同一语义不能同时传两个参数,否则抛TypeError;_duration_to_seconds()校验必须为非负有限数值,timedelta会被转换为秒;- 最终通过
Vqs-Retention-Seconds、Vqs-Delay-Seconds请求头发送。
0.0.12:send API 增加 headers,subscribe() 支持 topic filter
- 0.0.12 为 send API 增加
headers参数(自定义请求头会与Authorization、Content-Type合并); subscribe()支持 topic 过滤。
当前subscribe支持三种形式(见 src/vercel/workers/_queue/subscribe.py):
@subscribe # 不限定 topic def worker(message, metadata): ... @subscribe(topic="events") # 精确匹配 topic def billing_worker(message, metadata): ... @subscribe(topic=("user-*", lambda t: t.startswith("user-"))) # 自定义谓词过滤 def user_worker(message, metadata): ...0.0.19:部署固定(deployment pinning)对齐 TypeScript SDK
0.0.19 将队列部署固定行为与 TypeScript SDK 对齐,区分三种状态:
- 自动固定:默认通过
VERCEL_DEPLOYMENT_ID环境变量将消息固定到当前部署; - 显式部署 ID:调用时传入
deployment_id="dpl_xxx"; - 显式取消固定:传入
deployment_id=None。
resolve_deployment_id()(send.py)的实现要点:
- 开发模式下(
VERCEL_WORKERS_IN_PROCESS=1或VERCEL_QUEUE_TOKEN="vc-dev-token")永不发送部署 ID,与vercel dev行为一致; - 未显式传参且环境变量缺失时抛出
RuntimeError,提示可传显式deployment_id或deployment_id=None显式退出固定; - 固定后通过
Vqs-Deployment-Id请求头随消息发送。
0.0.13:支持非标准库类型的编码
0.0.13 为 send 支持了常见非标准库类型作为参数编码。其实现是 src/vercel/workers/_queue/types.py 中的WorkerJSONEncoder:
class WorkerJSONEncoder(json.JSONEncoder): def default(self, o): match o: case UUID(): return str(o) case datetime() | date(): return o.isoformat() case Decimal(): return float(o) case _: return super().default(o)即UUID→ 字符串、datetime/date→ ISO 格式、Decimal→ 浮点数,用户也可通过send(json_encoder=...)传入自定义编码器。
消息消费:回调协议从 v1beta 到 v2beta
0.0.14 → 0.0.17:从 Queues V3 API 迁移到 v2beta triggers
- 0.0.14:python workers 迁移至Queues V3 API(对应
get_queue_base_path()默认的/api/v3/topic端点); - 0.0.17:python workers 迁移至v2beta triggers + 私有路由。
0.0.23:v2beta 元数据回调的兜底实现
0.0.23 为 v2betametadata-only 回调增加receive_message_by_id兜底:当回调仅携带消息元数据(队列名、消费组、消息 ID 等头信息)而不含完整 payload 时,SDK 会按消息 ID 主动拉取完整消息后再分发给订阅者。该逻辑位于 src/vercel/workers/_queue/callback.py 与client.py的_handle_queue_callback()中:
- 通过
Ce-Type: com.vercel.queue.v2beta头识别 v2beta 回调(is_v2beta_callback); - 从
Ce-Vqsqueuename、Ce-Vqsconsumergroup、Ce-Vqsmessageid、Ce-Vqsreceipthandle等头解析队列上下文(parse_v2beta_callback); - 若回调中无 receipt/payload,则调用
resolve_v2beta_message以receive_message_by_id补全; - 兼容 v1beta CloudEvent 结构(
parse_cloudevent,type == "com.vercel.queue.v1beta"),两种格式共用同一套分发管线。
可见性超时与自动续期
回调处理遵循 Node 侧ConsumerGroupOptions的默认值(源码注释明确说明"Mirror the Node defaults"):
VQS_VISIBILITY_TIMEOUT:可见性超时,默认30 秒;VQS_VISIBILITY_REFRESH_INTERVAL:自动续期间隔,默认10 秒。
处理期间通过VisibilityExtender后台任务持续刷新消息可见性,防止任务执行中消息被重复投递。
显式重试与确认指令:Ack / RetryAfter
0.0.19:worker 可返回或抛出 RetryAfter / Ack
0.0.19 为 Python worker 增加显式重试与确认指令。在 subscribe.py 中实现:
class Ack(Exception): ... class RetryAfter(Exception): def __init__(self, delay, reason=None): # delay 支持 int 秒或 timedelta;负值会被截断为 0invoke_subscriptions()的语义:
- 返回或抛出
RetryAfter(delay)→ 消息按 delay 秒后重试(通过change_visibility设置可见性); - 返回或抛出
Ack→ 立即确认(删除消息); - 返回其他任何值(含
None)→ 视为处理成功并确认。
examples/basic/worker.py 展示了实际用法:返回None即确认消息,需要重试则return RetryAfter(60)。
0.0.25:支持 per-actor Dramatiq 重试选项
最新版本 0.0.25 让 Dramatiq worker 能遵循每个 actor 级别配置的重试选项,而非仅使用全局默认。这补齐了 Dramatiq 适配器与上游 Dramatiq 生态(actor 级max_retries、retry_when等)的对接能力。
0.0.15:Dramatiq 中间件与序列化修复
- 处理 Dramatiq 中间件(middlewares);
- 修复 UUID/Decimal/datetime 的序列化问题(即前文
WorkerJSONEncoder的覆盖场景)。
认证与运行环境:token 解析与进程内开发模式
send()的 token 解析顺序(send.py 的get_queue_token):
- 显式
token=参数; VERCEL_QUEUE_TOKEN环境变量;- Vercel OIDC token(
vercel.oidc.get_vercel_oidc_token),用于部署环境自动认证。
队列端点解析(get_queue_base_url):
- 优先
VERCEL_QUEUE_BASE_URL; - 其次若设置
VERCEL_REGION,则路由到https://{region}.vercel-queue.com(如iad1); - 兜底
https://vercel-queue.com;路径默认/api/v3/topic。
进程内开发模式
设置VERCEL_WORKERS_IN_PROCESS=1时,send()会在当前进程内直接调用匹配的@subscribe处理器(_send_in_process),不访问队列服务,模拟 TypeScript 侧的本地 dev 体验——无持久化、无可见性超时、无重试。若未注册任何订阅或没有匹配 topic 的订阅,会抛出带可用 topic 列表的明确错误,便于快速排查配置不匹配。
安装与快速上手
安装
pip install vercel-workers按需安装适配器 extras:
pip install "vercel-workers[celery]" pip install "vercel-workers[dramatiq]" pip install "vercel-workers[django]"依赖见 pyproject.toml:httpx>=0.27.0、anyio>=4.0.0、pydantic>=2.7.0、python-dotenv、vercel>=0.3.7。
Worker Service 部署形态(vercel.json)
{ "projectSettings": { "framework": "services" }, "experimentalServices": { "web": { "framework": "fastapi", "entrypoint": "main.py", "routePrefix": "/" }, "worker": { "type": "worker", "entrypoint": "worker.py", "topic": "default", "consumer": "default" } } }worker.py需要暴露 worker 定义(@subscribe函数、Celeryapp或 Dramatiqbroker),并导入任务模块以完成处理器注册。参考 examples/basic/vercel.json 的完整示例,其中topic使用topics数组形式声明,且与main.py中send()的队列名保持一一对应。
完整示例(FastAPI 生产者 + @subscribe Worker)
生产者侧 examples/basic/main.py:
import worker # noqa: F401 # 导入以注册 @subscribe 处理器 from fastapi import FastAPI from pydantic import BaseModel from vercel.workers import send QUEUE_NAME = "default" # 必须与 vercel.json 中 worker 的 topic 一致 app = FastAPI() @app.post("/enqueue") def enqueue_job(body: EnqueueRequest): result = send(QUEUE_NAME, {"message": body.message}) return {"queued": True, "messageId": result["messageId"], "queue": QUEUE_NAME}消费侧 examples/basic/worker.py:
from vercel.workers import MessageMetadata, RetryAfter, subscribe @subscribe(topic="default") def process_message(message, metadata: MessageMetadata) -> RetryAfter | None: print("Received message from queue:", message) return None # None = 确认;RetryAfter(60) = 60 秒后重试错误处理体系
exceptions.py 与公共导出提供了一整套结构化错误类型,便于精确处理队列调用失败:
- 基础类
VQSError(含status_code、可选retry_after); - 4xx 类:
BadRequestError(400)、UnauthorizedError(401)、ForbiddenError(403)、DuplicateIdempotencyKeyError(409,幂等键冲突)、InvalidLimitError; - 队列/消息状态类:
QueueEmptyError、MessageNotFoundError、MessageNotAvailableError、MessageCorruptedError、MessageLockedError; - 其他:
InternalServerError(5xx)、ThrottledError(限流)、TokenResolutionError(token 无法解析)。
回调处理管线对VQSError会按其status_code原样回传 HTTP 状态码,并附上type与可选retryAfter字段,保证消费者侧错误语义一致。
环境变量速查表
| 环境变量 | 作用 | 默认值 |
|---|---|---|
VERCEL_QUEUE_TOKEN | 队列鉴权 token(部署外运行必需) | 无 |
VERCEL_QUEUE_BASE_URL | 队列服务端点覆盖 | https://vercel-queue.com或按 region |
VERCEL_REGION | 按 region 路由队列端点 | 无 |
VERCEL_QUEUE_BASE_PATH | API 路径前缀 | /api/v3/topic |
VERCEL_DEPLOYMENT_ID | 自动部署固定 | 无 |
VERCEL_WORKERS_IN_PROCESS | 启用进程内开发模式(1/true/yes) | 关闭 |
VQS_VISIBILITY_TIMEOUT | 可见性超时(秒) | 30 |
VQS_VISIBILITY_REFRESH_INTERVAL | 可见性续期间隔(秒) | 10 |
在 Vercel 之外运行本地开发时,至少需要设置VERCEL_QUEUE_TOKEN(可选VERCEL_QUEUE_BASE_URL),这是 README.md 明确说明的运行前提。
测试与示例资产
仓库提供了覆盖各能力面的测试与示例,便于读者对照验证:
- 测试:tests/test_client_and_callback.py、tests/test_celery_adapter.py、tests/test_dramatiq_adapter.py、tests/test_django_adapter.py、tests/test_runtime_bridge.py;
- 示例:examples/basic(FastAPI +
@subscribe)、examples/celery、examples/dramatiq、examples/django。
小结
从 0.0.10 到 0.0.25,vercel-workers完成了三件关键事:一是架构上把队列 SDK 下沉到_queue/并保持公共 API 稳定;二是消息协议上从 v1beta CloudEvent 演进到 v2beta triggers,并补齐 metadata-only 回调的按 ID 拉取兜底;三是消费语义上引入Ack/RetryAfter显式指令与 deployment pinning,使其行为与 TypeScript SDK 对齐。对于需要在 Vercel 上用 Python 构建异步任务系统的开发者,该 SDK 当前已具备生产可用的发布、消费、重试、鉴权与多框架适配能力,可直接参照上述配置与示例落地。
【免费下载链接】vercelDevelop. Preview. Ship.项目地址: https://gitcode.com/gh_mirrors/ve/vercel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考