跳转到内容
搜索文档

使用 Queues 将数据存储到 R2

如何使用 Queues 批处理数据并将其存储到 R2 存储桶的示例。

最后更新 查看 MarkdownAgent 设置

以下 Worker 将捕获 JavaScript 错误并将其发送到队列。同一 Worker 将分批接收这些错误并将其存储到 R2 存储桶中的日志文件。

{
	"$schema": "./node_modules/wrangler/config-schema.json",
	"name": "my-worker",
	"queues": {
		"producers": [
			{
				"queue": "my-queue",
				"binding": "ERROR_QUEUE"
			}
		],
		"consumers": [
			{
				"queue": "my-queue",
				"max_batch_size": 100,
				"max_batch_timeout": 30
			}
		]
	},
	"r2_buckets": [
		{
			"bucket_name": "my-bucket",
			"binding": "ERROR_BUCKET"
		}
	]
}
"$schema" = "./node_modules/wrangler/config-schema.json"
name = "my-worker"

[[queues.producers]]
queue = "my-queue"
binding = "ERROR_QUEUE"

[[queues.consumers]]
queue = "my-queue"
max_batch_size = 100
max_batch_timeout = 30

[[r2_buckets]]
bucket_name = "my-bucket"
binding = "ERROR_BUCKET"
interface ErrorMessage {
	message: string;
	stack?: string;
}

interface Env {
	readonly ERROR_QUEUE: Queue<ErrorMessage>;
	readonly ERROR_BUCKET: R2Bucket;
}

export default {
  async fetch(req, env, ctx): Promise<Response> {
    try {
      return doRequest(req);
    } catch (e) {
      const error: ErrorMessage = {
        message: e instanceof Error ? e.message : String(e),
        stack: e instanceof Error ? e.stack : undefined,
      };
      await env.ERROR_QUEUE.send(error);
      return new Response(error.message, { status: 500 });
    }
  },
  async queue(batch, env, ctx): Promise<void> {
    let file = "";
    for (const message of batch.messages) {
      const error = message.body;
      file += error.stack ?? error.message;
      file += "\r\n";
    }
    await env.ERROR_BUCKET.put(`errors/${Date.now()}.log`, file);
  },
} satisfies ExportedHandler<Env, ErrorMessage>;

function doRequest(request: Request): Response {
  if (Math.random() > 0.5) {
    return new Response("Success!");
  }
  throw new Error("Failed!");
}

这篇文档对您有帮助吗?