Cloudflare Queues 与 Cloudflare Workers 集成。要发送和接收消息,必须使用 Worker。
可以向 Queue 发送消息的 Worker 是生产者 Worker,可以从 Queue 接收消息的 Worker 是消费者 Worker。同一 Worker 可以同时是生产者和消费者。
未来,我们期望支持其他 API,例如用于发送或接收消息的 HTTP 端点。要报告 bug 或请求功能,请前往 Cloudflare Community Forums ↗。要提供反馈,请前往 #queues ↗ Discord 频道。
这些 API 允许生产者 Worker 向 Queue 发送消息。
向 Queue 写入单条消息的示例:
export default {
async fetch(req, env, ctx) {
await env.MY_QUEUE.send({
url: req.url,
method: req.method,
headers: Object.fromEntries(req.headers),
});
return new Response("Sent!");
},
};interface Env {
readonly MY_QUEUE: Queue;
}
export default {
async fetch(req, env, ctx): Promise<Response> {
await env.MY_QUEUE.send({
url: req.url,
method: req.method,
headers: Object.fromEntries(req.headers),
});
return new Response("Sent!");
},
} satisfies ExportedHandler<Env>;from workers import Response, WorkerEntrypoint
class Default(WorkerEntrypoint):
async def fetch(self, request):
await self.env.MY_QUEUE.send({
"url": request.url,
"method": request.method,
"headers": dict(request.headers),
})
return Response("Sent!")Queues API 还支持一次写入多条消息:
const sendResultsToQueue = async (results, env) => {
const batch = results.map((value) => ({
body: value,
}));
await env.MY_QUEUE.sendBatch(batch);
};const sendResultsToQueue = async (results: Array<unknown>, env: Env) => {
const batch: MessageSendRequest[] = results.map((value) => ({
body: value,
}));
await env.MY_QUEUE.sendBatch(batch);
};async def send_results_to_queue(results, env):
batch = [
{"body": value}
for value in results
]
await env.MY_QUEUE.sendBatch(batch)允许生产者向 Queue 发送消息的绑定(binding)。
interface Queue<Body = unknown> {
send(body: Body, options?: QueueSendOptions): Promise<QueueSendResult>;
sendBatch(messages: Iterable<MessageSendRequest<Body>>, options?: QueueSendBatchOptions): Promise<QueueSendResult>;
metrics(): Promise<QueueMetrics>;
}send(body: unknown, options?: {contentType?: QueuesContentType })Promise<QueueSendResult>- 向 Queue 发送消息。body 可以是 structured clone 算法 ↗ 支持的任何类型,只要大小小于 128 KB。
- Promise 解析时,消息已确认写入磁盘。
- 返回包含队列实时指标的 QueueSendResult。
sendBatch(messages: Iterable<MessageSendRequest<unknown>>, options?: QueueSendBatchOptions)Promise<QueueSendBatchResult>- 向 Queue 发送一批消息。提供的 Iterable ↗ 中每项必须是 structured clone 算法 ↗ 支持的类型。一批最多可包含 100 条消息,但每项限制为 128 KB,数组总大小不能超过 256 KB。
- 可选的
options参数可用于对批次中所有消息应用设置(如delaySeconds)。请参阅 QueueSendBatchOptions。 - Promise 解析时,消息已确认写入磁盘。
metrics()Promise<QueueMetrics>- 返回队列的实时 QueueMetrics。
用于发送消息批次的包装类型。
interface MessageSendRequest<Body = unknown> {
body: Body;
contentType?: QueueContentType;
delaySeconds?: number;
}-
bodyunknown- 消息正文。
- body 可以是 structured clone 算法 ↗ 支持的任何类型,只要大小小于 128 KB。
-
contentTypeQueueContentType- 消息的显式内容类型,以便使用从仪表板列出消息 功能正确预览。可选参数。
- 请参阅 QueuesContentType 了解可能的值。
-
delaySecondsnumber- 在队列内延迟消息的秒数,然后才能投递给消费者。
- 必须是 0 到 86400(24 小时)之间的整数。
向队列发送消息时应用的可选配置。
-
contentTypeQueuesContentType- 消息的显式内容类型,以便使用从仪表板列出消息 功能正确预览。可选参数。
- 目前,此选项供内部使用。未来,
contentType将被其他消费者类型用于显式标记序列化消息,以便以所需类型消费。 - 请参阅 QueuesContentType 了解可能的值。
-
delaySecondsnumber- 在队列内延迟消息的秒数,然后才能投递给消费者。
- 必须是 0 到 86400(24 小时)之间的整数。将此值设置为零将显式阻止消息被延迟,即使队列级别有全局(默认)延迟也是如此。
向队列发送一批消息时应用的可选配置。
-
delaySecondsnumber- 在队列内延迟消息的秒数,然后才能投递给消费者。
- 必须是正整数。
包含有效消息内容类型的联合类型。
// Default: json
type QueuesContentType = "text" | "bytes" | "json" | "v8";- 使用
"json"发送可 JSON 序列化的 JavaScript 对象。此内容类型可从 Cloudflare 仪表板 ↗ 预览。json内容类型为默认值。 - 使用
"text"发送String。此内容类型可使用从仪表板列出消息 功能预览。 - 使用
"bytes"发送ArrayBuffer。此内容类型无法从 Cloudflare 仪表板 ↗ 预览,将显示为 Base64 编码。 - 使用
"v8"发送无法 JSON 序列化但受 structured clone ↗ 支持的 JavaScript 对象(例如Date和Map)。此内容类型无法从 Cloudflare 仪表板 ↗ 预览,将显示为 Base64 编码。
若您指定无效的内容类型,或指定的内容类型与消息内容类型不匹配,发送操作将失败并返回错误。
成功发送操作的结果。
interface QueueSendResult {
metadata: {
metrics: QueueMetrics;
};
}-
metadataobject- 包含发送操作后队列的元数据。
-
metadata.metricsQueueMetrics- 队列的实时指标。请参阅 QueueMetrics。
队列的实时指标。
interface QueueMetrics {
backlogCount: number;
backlogBytes: number;
oldestMessageTimestamp: number;
}-
backlogCountnumber- 队列中当前的消息数量。
-
backlogBytesnumber- 队列中消息的总大小(字节)。
-
oldestMessageTimestampnumber- 队列中最旧消息的时间戳(毫秒,自纪元起)。
这些 API 允许消费者 Worker 从 Queue 消费消息。
要定义消费者 Worker,请向 Worker 的默认导出添加 queue() 函数。这将允许它从 Queue 接收消息。
默认情况下,满足以下所有条件时,批次中的所有消息将被确认:
queue()函数已返回。- 若
queue()函数返回 promise,该 promise 已解析。 - 传递给
waitUntil()的任何 promise 已解析。
若 queue() 函数抛出异常,或它返回的 promise 或传递给 waitUntil() 的任何 promise 被拒绝,则整个批次将被视为失败,并根据消费者的重试设置重试。
export default {
async queue(batch, env, ctx) {
for (const message of batch.messages) {
console.log("Received", message.body);
}
},
};interface Env {
// Add your bindings here
}
export default {
async queue(batch, env, ctx): Promise<void> {
for (const message of batch.messages) {
console.log("Received", message.body);
}
},
} satisfies ExportedHandler<Env>;from workers import WorkerEntrypoint
class Default(WorkerEntrypoint):
async def queue(self, batch):
for message in batch.messages:
print("Received", message)env 和 ctx 字段如 Workers 文档 中所述。
您可以在生产者上使用 Queue<T>,在消费者上使用 ExportedHandler<Env, T> 为队列消息添加类型。
type MyMessage = {
id: string;
};
interface Env {
MY_QUEUE: Queue<MyMessage>;
}
export default {
async queue(batch) {
for (const message of batch.messages) {
console.log(message.body.id);
}
},
} satisfies ExportedHandler<Env, MyMessage>;对于原始消息,使用 Queue<number> 或 satisfies ExportedHandler<Env, number>。若未指定类型,message.body 为 unknown。
或者,可以使用(已弃用的)service worker 语法编写队列消费者:
addEventListener('queue', (event) => {
event.waitUntil(handleMessages(event));
});在 service worker 语法中,event 提供与下面定义的 MessageBatch 相同的字段和方法,以及 waitUntil() ↗。
发送给消费者 Worker 的消息批次。
interface MessageBatch<Body = unknown> {
readonly queue: string;
readonly messages: readonly Message<Body>[];
ackAll(): void;
retryAll(options?: QueueRetryOptions): void;
}-
queuestring- 此批次所属 Queue 的名称。
-
messagesMessage[]- 批次中的消息数组。消息顺序为尽力而为——不保证与发布顺序完全一致。
-
ackAll()void- 将所有消息标记为成功投递,无论
queue()消费者处理程序是否成功返回。
- 将所有消息标记为成功投递,无论
-
retryAll(options?: QueueRetryOptions)void- 将所有消息标记为在下一批次中重试。
- 支持可选的
options对象。
发送给消费者 Worker 的消息。
interface Message<Body = unknown> {
readonly id: string;
readonly timestamp: Date;
readonly body: Body;
readonly attempts: number;
ack(): void;
retry(options?: QueueRetryOptions): void;
}-
idstring- 消息的唯一系统生成 ID。
-
timestampDate- 消息发送时的时间戳。
-
bodyunknown- 消息正文。
- body 可以是 structured clone 算法 ↗ 支持的任何类型,只要大小小于 128 KB。
-
attemptsnumber- 消费者尝试处理此消息的次数。从 1 开始。
-
ack()void- 将消息标记为成功投递,无论
queue()消费者处理程序是否成功返回。
- 将消息标记为成功投递,无论
-
retry(options?: QueueRetryOptions)void- 将消息标记为在下一批次中重试。
- 支持可选的
options对象。
标记消息或一批消息重试时的可选配置。
interface QueueRetryOptions {
delaySeconds?: number;
}-
delaySecondsnumber- 在队列内延迟消息的秒数,然后才能投递给消费者。
- 必须是正整数。
-
Promise 解析时,消息写入磁盘。
- 返回包含队列实时指标的 QueueSendResult。