跳转到内容
搜索文档

Pull 消费者

最后更新 查看 MarkdownAgent 设置

Pull 消费者允许您从 Cloudflare Workers 之外的任何环境和/或编程语言通过 HTTP 从队列拉取消息。当消息消费速率受上游基础设施或长时间运行的任务限制时,Pull 消费者很有用。

如何选择 push 或 pull 消费者

决定配置 push 消费者还是 pull 消费者取决于您如何使用队列,以及队列消费者上游基础设施的配置。

  • push 消费者 开始是从队列消费的最简单方式。Push 消费者运行在 Workers 上,默认会在消息写入队列时自动扩展并消费消息。
  • 若您需要从 Cloudflare Workers 之外的现有基础设施消费消息,和/或需要仔细控制消息消费速度,请使用 pull 消费者。Pull 消费者必须显式调用以从队列拉取(然后确认)消息,仅在准备好时才这样做。

您可以随时移除并附加新消费者到队列,若需求变化,可从 pull 消费者更改为 push 消费者。

要配置 pull 消费者并从队列接收消息,您需要:

  1. 为队列启用 HTTP pull。
  2. 为 HTTP 客户端创建有效的身份验证令牌。
  3. 从队列拉取消息批次。
  4. 确认和/或重试批次中的消息。

1. 启用 HTTP pull

您可以通过 wrangler CLI 或 Cloudflare 仪表板 启用 HTTP pull 或将队列从 push 更改为 pull。Wrangler 配置文件 中不再支持启用 HTTP pull。

wrangler CLI

您可以使用 wrangler queues consumer http 子命令在任何现有队列上启用 pull 消费者,并提供队列名称。

npx wrangler queues consumer http add $QUEUE-NAME

若您已有 push 消费者,需要先移除。若尝试在已有消费者配置的队列上调用 consumer http addwrangler 将返回错误:

wrangler queues consumer worker remove $QUEUE-NAME $SCRIPT_NAME

2. 消费者身份验证

HTTP Pull 消费者需要具有 com.cloudflare.api.account.queues_readcom.cloudflare.api.account.queues_write 权限的 API 令牌

pull 消费者需要同时具有 read 和 write 权限,因为它需要写入队列状态以确认收到的消息。消费消息会更改队列。

API 令牌在 HTTP 请求的 Authorization 头中以 Bearer 令牌形式呈现,格式为 Authorization: Bearer $YOUR_TOKEN_HERE。以下示例展示如何使用 curl HTTP 客户端传递 API 令牌:

curl "https://api.cloudflare.com/client/v4/accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/pull" \
--header "Authorization: Bearer ${QUEUES_TOKEN}" \
--header "Content-Type: application/json" \
--data '{ "visibility_timeout_ms": 10000, "batch_size": 2 }'

您可以针对单个队列同时身份验证并运行多个并发 pull 消费者。

创建 API 令牌

要创建 API 令牌:

  1. 登录 Cloudflare 仪表板
  2. 前往 My Profile(我的个人资料) > API Tokens
  3. 选择 Create Token(创建令牌)
  4. 滚动到页面底部并选择 Create Custom Token(创建自定义令牌)
  5. 为令牌命名。例如,queue-pull-token
  6. 在 **Permissions(权限)**部分,选择 Account(账户),然后选择 Queues。确保已选择 Edit(read+write)(编辑(读+写))
  7. (可选)选择 **All accounts(默认)**或特定账户以限定令牌范围。
  8. 选择 Continue to summary(继续到摘要),然后选择 Create token(创建令牌)

您需要记下令牌:它只会显示一次。

3. 拉取消息

要拉取消息,请向 Queues REST API 发出 HTTP POST 请求,JSON 编码的正文可选指定 visibility_timeoutbatch_size,或空 JSON 对象({}):

index.jsjs
// POST /accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/pull with the timeout & batch size
let resp = await fetch(
	`https://api.cloudflare.com/client/v4/accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/pull`,
	{
		method: "POST",
		headers: {
			"content-type": "application/json",
			authorization: `Bearer ${QUEUES_API_TOKEN}`,
		},
		// Optional - you can provide an empty object '{}' and the defaults will apply.
		body: JSON.stringify({ visibility_timeout_ms: 6000, batch_size: 50 }),
	},
);
index.tsts
// POST /accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/pull with the timeout & batch size
let resp = await fetch(
	`https://api.cloudflare.com/client/v4/accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/pull`,
	{
		method: "POST",
		headers: {
			"content-type": "application/json",
			authorization: `Bearer ${QUEUES_API_TOKEN}`,
		},
		// Optional - you can provide an empty object '{}' and the defaults will apply.
		body: JSON.stringify({ visibility_timeout_ms: 6000, batch_size: 50 }),
	},
);
import json
from workers import fetch

# POST /accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/pull with the timeout & batch size

resp = await fetch(
	f"https://api.cloudflare.com/client/v4/accounts/{CF_ACCOUNT_ID}/queues/{QUEUE_ID}/messages/pull",
	method="POST",
	headers={
		"content-type": "application/json",
		"authorization": f"Bearer {QUEUES_API_TOKEN}",
	}, # Optional - you can provide an empty object '{}' and the defaults will apply.
	body=json.dumps({"visibility_timeout_ms": 6000, "batch_size": 50}),
)

这将返回消息数组(最多为指定的 batch_size),格式如下:

{
	"success": true,
	"errors": [],
	"messages": [],
	"result": {
		"message_backlog_count": 10,
		"messages": [
			{
				"body": "hello",
				"id": "1ad27d24c83de78953da635dc2ea208f",
				"timestamp_ms": 1689615013586,
				"attempts": 2,
				"metadata": {
					"CF-sourceMessageSource": "dash",
					"CF-Content-Type": "json"
				},
				"lease_id": "eyJhbGciOiJkaXIiLCJlbmMiOiJBMjU2Q0JDLUhTNTEyIn0..NXmbr8h6tnKLsxJ_AuexHQ.cDt8oBb_XTSoKUkVKRD_Jshz3PFXGIyu7H1psTO5UwI.smxSvQ8Ue3-ymfkV6cHp5Va7cyUFPIHuxFJA07i17sc"
			},
			{
				"body": "world",
				"id": "95494c37bb89ba8987af80b5966b71a7",
				"timestamp_ms": 1689615013586,
				"attempts": 2,
				"metadata": {
					"CF-sourceMessageSource": "dash",
					"CF-Content-Type": "json"
				},
				"lease_id": "eyJhbGciOiJkaXIiLCJlbmMiOiJBMjU2Q0JDLUhTNTEyIn0..QXPgHfzETsxYQ1Vd-H0hNA.mFALS3lyouNtgJmGSkTzEo_imlur95EkSiH7fIRIn2U.PlwBk14CY_EWtzYB-_5CR1k30bGuPFPUx1Nk5WIipFU"
			}
		]
	}
}

Pull 消费者遵循"短轮询"方式:若有消息可投递,Queues 将立即返回最多为配置 batch_size 的消息响应。若没有消息可投递,Queues 将返回空响应。Queues 不会保持开放连接(通常称为"长轮询")等待消息可投递。

每个消息对象有五个字段:

  1. body - 根据消息发布时使用的内容类型,可能为 base64 编码。
  2. id - 消息的唯一、只读临时标识符。
  3. timestamp_ms - 消息发布到队列的时间(毫秒,自 Unix 纪元 起)。可用于通过从当前时间戳减去来确定消息的年龄。
  4. attempts - 消息被完整投递尝试的次数。当达到 max_retries 值时,消息将不再重新投递,并从队列中永久删除。
  5. lease_id - 消息的编码租约 ID。lease_id 用于显式确认或重试消息。

lease_id 允许 pull 消费者显式确认或重试批次中的部分、全部或不确认任何消息。若消费者未确认或标记消息重试,则消息将在达到 visibility_timeout 后标记为重新投递。lease_id 在此超时后不再有效。

拉取队列时可配置 batch_sizevisibility_timeout

  • batch_size(默认为 5;最大 100)- 每次 pull 返回给消费者的消息数。
  • visibility_timeout(默认为 30 秒;最大 12 小时)- 定义消费者根据 lease_id 确认批次中投递消息的时间。此超时过期后,消息被视为未确认并再次排队重新投递。

并发消费者

您可以有多个 HTTP 客户端并发从同一队列拉取:每个客户端将收到唯一的消息批次,并在 visibility_timeout 过期或这些消息被标记重试之前保留这些消息的"租约"。

标记重试的消息将放回队列,可由任何消费者投递。消息绑定到特定消费者,因为消费者没有身份,且为避免慢速或卡住的消费者阻碍队列中消息的处理。

多个消费者在您有多个上游资源(例如 GPU 基础设施)、希望根据队列积压 自动扩展和/或成本的情况下很有用。

4. 确认消息

Pull 消费者拉取的消息需要被确认或标记重试。

要确认和/或标记消息重试,请根据 Queues REST API 向队列的 /ack 端点发出 HTTP POST 请求,提供要确认和/或重试的 lease_id 对象数组:

index.jsjs
// POST /accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/ack with the lease_ids
let resp = await fetch(
	`https://api.cloudflare.com/client/v4/accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/ack`,
	{
		method: "POST",
		headers: {
			"content-type": "application/json",
			authorization: `Bearer ${QUEUES_API_TOKEN}`,
		},
		// If you have no messages to retry, you can specify an empty array - retries: []
		body: JSON.stringify({
			acks: [
				{ lease_id: "lease_id1" },
				{ lease_id: "lease_id2" },
				{ lease_id: "etc" },
			],
			retries: [{ lease_id: "lease_id4" }],
		}),
	},
);
index.tsts
// POST /accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/ack with the lease_ids
let resp = await fetch(
	`https://api.cloudflare.com/client/v4/accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/ack`,
	{
		method: "POST",
		headers: {
			"content-type": "application/json",
			authorization: `Bearer ${QUEUES_API_TOKEN}`,
		},
		// If you have no messages to retry, you can specify an empty array - retries: []
		body: JSON.stringify({
			acks: [
				{ lease_id: "lease_id1" },
				{ lease_id: "lease_id2" },
				{ lease_id: "etc" },
			],
			retries: [{ lease_id: "lease_id4" }],
		}),
	},
);
import json
from workers import fetch

# POST /accounts/${CF_ACCOUNT_ID}/queues/${QUEUE_ID}/messages/ack with the lease_ids

resp = await fetch(
	f"https://api.cloudflare.com/client/v4/accounts/{CF_ACCOUNT_ID}/queues/{QUEUE_ID}/messages/ack",
	method="POST",
	headers={
		"content-type": "application/json",
		"authorization": f"Bearer {QUEUES_API_TOKEN}",
	}, # If you have no messages to retry, you can specify an empty array - retries: []
	body=json.dumps({
		"acks": [
			{"lease_id": "lease_id1"},
			{"lease_id": "lease_id2"},
			{"lease_id": "etc"},
		],
		"retries": [{"lease_id": "lease_id4"}],
	}),
)

retries 数组中提供 { lease_id: string, delay_seconds: number } 对象,可选择在标记消息重试时指定延迟秒数:

{
	"acks": [
		{ "lease_id": "lease_id1" },
		{ "lease_id": "lease_id2" },
		{ "lease_id": "lease_id3" }
	],
	"retries": [{ "lease_id": "lease_id4", "delay_seconds": 600 }]
}

此外:

  • 若您在消费者中处理这些消息,应在 /ack 端点请求中提供每个 lease_id。若不确认消息,它将被标记为重新投递(放回队列)。
  • 您可以选择标记消息重试:例如,处理消息时出错或上游资源压力。显式标记消息重试将立即将其放回队列,而不是等待(可能很长的)visibility_timeout
  • 您可以在处理一批消息的过程中多次调用 /ack 端点,但我们建议分组确认以减少所需的 API 调用次数。

Queues 旨在对租约 ID 保持宽松:若消费者在 visibility_timeout 达到之后通过租约 ID 确认消息,Queues 仍将接受该确认。若消息在期间被投递给另一个消费者,它也能确认消息而不会出错。

内容类型

向具有外部消费者的队列发布时,您应了解某些内容类型可能被编码以便在 JSON 对象中安全序列化。

对于 jsonbytes 内容类型,这意味着它们将被 base64 编码(RFC 4648)。text 类型将作为纯 UTF-8 编码字符串发送。

您的消费者需要在操作数据之前解码 jsonbytes 类型。

后续步骤

这篇文档对您有帮助吗?