跳转到内容
搜索文档

JavaScript API

最后更新 查看 MarkdownAgent 设置

Cloudflare Queues 与 Cloudflare Workers 集成。要发送和接收消息,必须使用 Worker。

可以向 Queue 发送消息的 Worker 是生产者 Worker,可以从 Queue 接收消息的 Worker 是消费者 Worker。同一 Worker 可以同时是生产者和消费者。

未来,我们期望支持其他 API,例如用于发送或接收消息的 HTTP 端点。要报告 bug 或请求功能,请前往 Cloudflare Community Forums。要提供反馈,请前往 #queues Discord 频道。

生产者

这些 API 允许生产者 Worker 向 Queue 发送消息。

向 Queue 写入单条消息的示例:

index.jsjs
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!");
	},
};
index.tsts
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 还支持一次写入多条消息:

index.jsjs
const sendResultsToQueue = async (results, env) => {
	const batch = results.map((value) => ({
		body: value,
	}));
	await env.MY_QUEUE.sendBatch(batch);
};
index.tsts
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

允许生产者向 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>

MessageSendRequest

用于发送消息批次的包装类型。

interface MessageSendRequest<Body = unknown> {
  body: Body;
  contentType?: QueueContentType;
  delaySeconds?: number;
}
  • body unknown
  • contentType QueueContentType
  • delaySeconds number
    • 在队列内延迟消息的秒数,然后才能投递给消费者。
    • 必须是 0 到 86400(24 小时)之间的整数。

QueueSendOptions

向队列发送消息时应用的可选配置。

  • contentType QueuesContentType
    • 消息的显式内容类型,以便使用从仪表板列出消息 功能正确预览。可选参数。
    • 目前,此选项供内部使用。未来,contentType 将被其他消费者类型用于显式标记序列化消息,以便以所需类型消费。
    • 请参阅 QueuesContentType 了解可能的值。
  • delaySeconds number
    • 在队列内延迟消息的秒数,然后才能投递给消费者。
    • 必须是 0 到 86400(24 小时)之间的整数。将此值设置为零将显式阻止消息被延迟,即使队列级别有全局(默认)延迟也是如此。

QueueSendBatchOptions

向队列发送一批消息时应用的可选配置。

  • delaySeconds number
    • 在队列内延迟消息的秒数,然后才能投递给消费者。
    • 必须是正整数。

QueuesContentType

包含有效消息内容类型的联合类型。

// Default: json
type QueuesContentType = "text" | "bytes" | "json" | "v8";
  • 使用 "json" 发送可 JSON 序列化的 JavaScript 对象。此内容类型可从 Cloudflare 仪表板 预览。json 内容类型为默认值。
  • 使用 "text" 发送 String。此内容类型可使用从仪表板列出消息 功能预览。
  • 使用 "bytes" 发送 ArrayBuffer。此内容类型无法从 Cloudflare 仪表板 预览,将显示为 Base64 编码。
  • 使用 "v8" 发送无法 JSON 序列化但受 structured clone 支持的 JavaScript 对象(例如 DateMap)。此内容类型无法从 Cloudflare 仪表板 预览,将显示为 Base64 编码。

若您指定无效的内容类型,或指定的内容类型与消息内容类型不匹配,发送操作将失败并返回错误。

QueueSendResult

成功发送操作的结果。

interface QueueSendResult {
	metadata: {
		metrics: QueueMetrics;
	};
}
  • metadata object
    • 包含发送操作后队列的元数据。
  • metadata.metrics QueueMetrics

QueueMetrics

队列的实时指标。

interface QueueMetrics {
	backlogCount: number;
	backlogBytes: number;
	oldestMessageTimestamp: number;
}
  • backlogCount number
    • 队列中当前的消息数量。
  • backlogBytes number
    • 队列中消息的总大小(字节)。
  • oldestMessageTimestamp number
    • 队列中最旧消息的时间戳(毫秒,自纪元起)。

消费者

这些 API 允许消费者 Worker 从 Queue 消费消息。

要定义消费者 Worker,请向 Worker 的默认导出添加 queue() 函数。这将允许它从 Queue 接收消息。

默认情况下,满足以下所有条件时,批次中的所有消息将被确认:

  1. queue() 函数已返回。
  2. queue() 函数返回 promise,该 promise 已解析。
  3. 传递给 waitUntil() 的任何 promise 已解析。

queue() 函数抛出异常,或它返回的 promise 或传递给 waitUntil() 的任何 promise 被拒绝,则整个批次将被视为失败,并根据消费者的重试设置重试。

index.jsjs
export default {
	async queue(batch, env, ctx) {
		for (const message of batch.messages) {
			console.log("Received", message.body);
		}
	},
};
index.tsts
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)

envctx 字段如 Workers 文档 中所述。

TypeScript 消息类型

您可以在生产者上使用 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.bodyunknown

或者,可以使用(已弃用的)service worker 语法编写队列消费者:

addEventListener('queue', (event) => {
	event.waitUntil(handleMessages(event));
});

在 service worker 语法中,event 提供与下面定义的 MessageBatch 相同的字段和方法,以及 waitUntil()

MessageBatch

发送给消费者 Worker 的消息批次。

interface MessageBatch<Body = unknown> {
  readonly queue: string;
  readonly messages: readonly Message<Body>[];
  ackAll(): void;
  retryAll(options?: QueueRetryOptions): void;
}
  • queue string
    • 此批次所属 Queue 的名称。
  • messages Message[]
    • 批次中的消息数组。消息顺序为尽力而为——不保证与发布顺序完全一致。
  • ackAll() void
    • 将所有消息标记为成功投递,无论 queue() 消费者处理程序是否成功返回。
  • retryAll(options?: QueueRetryOptions) void
    • 将所有消息标记为在下一批次中重试。
    • 支持可选的 options 对象。

Message

发送给消费者 Worker 的消息。

interface Message<Body = unknown> {
  readonly id: string;
  readonly timestamp: Date;
  readonly body: Body;
	readonly attempts: number;
  ack(): void;
  retry(options?: QueueRetryOptions): void;
}
  • id string
    • 消息的唯一系统生成 ID。
  • timestamp Date
    • 消息发送时的时间戳。
  • body unknown
  • attempts number
    • 消费者尝试处理此消息的次数。从 1 开始。
  • ack() void
    • 将消息标记为成功投递,无论 queue() 消费者处理程序是否成功返回。
  • retry(options?: QueueRetryOptions) void
    • 将消息标记为在下一批次中重试。
    • 支持可选的 options 对象。

QueueRetryOptions

标记消息或一批消息重试时的可选配置。

interface QueueRetryOptions {
  delaySeconds?: number;
}
  • delaySeconds number
    • 在队列内延迟消息的秒数,然后才能投递给消费者。
    • 必须是正整数。
  • Promise 解析时,消息写入磁盘。

这篇文档对您有帮助吗?