消费者并发允许从队列处理消息的消费者 Worker 自动水平扩展,以跟上向队列写入消息的速率。
在许多系统中,向队列写入消息的速率很容易超过单个消费者读取和处理这些消息的速率。这通常是因为消费者可能在解析消息内容、写入存储或数据库,或调用第三方(上游)API。
请注意,队列生产者始终可扩展,最高可达每条队列支持的最大每秒消息数限制。
默认情况下,所有队列均启用并发。队列消费者将自动扩展至最大并发调用数,以管理队列积压和/或错误率。
处理一批消息后,Queues 将检查是否应调整并发消费者数量。为队列调用的并发消费者数量将根据多个因素自动扩展,包括:
- 队列中的消息数量(积压)及其增长速度。
- 失败(与成功)调用的比率。失败调用是指
queue()处理程序返回未捕获异常而非void(无返回值)的情况。 - 为该消费者设置的
max_concurrency值。
在可能的情况下,Queues 会优化以防止积压呈指数增长,以尽量减少队列中消息在消息保留限制之前被处理之前达到该限制的情况。
若您以每秒 100 条消息的速度写入队列,而单个并发消费者处理 100 条消息的批次需要 5 秒,则正在处理的消息数量将以超过消费者处理能力的速度持续增长。
在此场景中,Queues 会注意到积压增长,并将并发消费者 Worker 调用数扩展至(大约)五(5)个的稳定状态,直到传入消息速率下降、消费者处理消息更快,或消费者开始产生错误。
若消费者没有自动扩展,可能的原因包括:
max_concurrency已设置为 1。- 消费者 Worker 返回错误而非处理消息。检查消费者以确保其健康。
- 正在处理一批消息。Queues 仅在处理完整一批消息后才检查是否应自动扩展,因此在处理批次期间不会自动扩展。考虑减小批大小或重构消费者以更快处理消息。
若您的工作流受上游 API 和/或系统限制,您可能希望积压增长,以换取增加的整体延迟,从而避免压垮上游系统。
您可以通过两种方式配置消费者 Worker 的并发:
- 在 Cloudflare 仪表板中设置并发
- 通过 Wrangler 配置文件 设置并发
要从仪表板配置消费者 Worker 的并发设置:
-
在 Cloudflare 仪表板中,前往 Queues 页面。
Go to Queues ↗ -
选择队列 > Settings(设置)。
-
在 Consumer details(消费者详情)下选择 Edit Consumer(编辑消费者)。
-
将 **Maximum consumer invocations(最大消费者调用数)**设置为
1到250之间的值。此值表示队列可用的最大并发消费者调用数。
要移除固定最大值,选择 auto (recommended)(自动(推荐))。
请注意,若向队列写入消息的速度超过处理速度,消息最终可能达到该队列设置的最大保留期。达到该限制的单个消息将从队列中过期并删除。
在 Wrangler 配置文件 中设置并发
要为给定队列设置固定的最大并发消费者调用数,请在 Wrangler 文件中配置 max_concurrency:
{
"queues": {
"consumers": [
{
"queue": "my-queue",
"max_concurrency": 1
}
]
}
}[[queues.consumers]]
queue = "my-queue"
max_concurrency = 1要移除限制,从给定队列的 [[queues.consumers]] 配置中删除 max_concurrency 设置,并调用 npx wrangler deploy 推送配置更新。
当多个消费者 Worker 被调用时,每次 Worker 调用都会产生 CPU 时间成本。
- 若您打算处理写入队列的所有消息,总体有效成本相同,即使启用了并发也是如此。
- 启用并发只是提前产生这些成本,并有助于防止消息达到消息保留限制。
消费者的计费遵循 Workers 标准使用模型,即开发者按请求和请求中使用的 CPU 时间计费。
处理一批消息需要 2 秒的消费者 Worker,无论是并发(更快)还是单独(更慢)处理 5000 万(50,000,000)条消息,总体成本都相同。