使用 Edge Functions 消费 Supabase 队列消息
Learn how to consume Supabase Queue messages server-side with a Supabase Edge Function
本指南帮助你使用 Supabase Edge Function 在服务器端读取和处理队列消息。想了解我们 API 的更多细节,请阅读 Queues API 参考。
🌐 This guide helps you read & process queue messages server-side with a Supabase Edge Function. Read Queues API Reference for more details on our API.
概念 #
🌐 Concepts
Supabase Queues 是一个基于拉取的消息队列,主要由三个组件组成:队列、消息和队列类型。你应该已经熟悉 队列快速入门。
🌐 Supabase Queues is a pull-based Message Queue consisting of three main components: Queues, Messages, and Queue Types. You should already be familiar with the Queues Quickstart.
在 Edge 函数中消费消息 #
🌐 Consuming messages in an Edge Function
这是一个 Supabase Edge 函数,它会从队列中读取 5 条消息,处理每条消息,并在处理完成后删除每条消息。
🌐 This is a Supabase Edge Function that reads 5 messages off the queue, processes each of them, and deletes each message when it is done.
1import 'jsr:@supabase/functions-js/edge-runtime.d.ts'2import { createClient } from 'npm:@supabase/supabase-js@2'34const supabaseUrl = 'supabaseURL'5const supabaseKey = 'supabaseKey'67const supabase = createClient(supabaseUrl, supabaseKey)8const queueName = 'your_queue_name'910// Type definition for queue messages11interface QueueMessage {12 msg_id: bigint13 read_ct: number14 vt: string15 enqueued_at: string16 message: any17}1819async function processMessage(message: QueueMessage) {20 //21 // Do whatever logic you need to with the message content22 //23 // Delete the message from the queue24 const { error: deleteError } = await supabase.schema('pgmq_public').rpc('delete', {25 queue_name: queueName,26 msg_id: message.msg_id,27 })2829 if (deleteError) {30 console.error(`Failed to delete message ${message.msg_id}:`, deleteError)31 } else {32 console.log(`Message ${message.msg_id} deleted from queue`)33 }34}3536Deno.serve(async (req) => {37 const { data: messages, error } = await supabase.schema('pgmq_public').rpc('read', {38 queue_name: queueName,39 sleep_seconds: 0, // Don't wait if queue is empty40 n: 5, // Read 5 messages off the queue41 })4243 if (error) {44 console.error(`Error reading from ${queueName} queue:`, error)45 return new Response(JSON.stringify({ error: error.message }), {46 status: 500,47 headers: { 'Content-Type': 'application/json' },48 })49 }5051 if (!messages || messages.length === 0) {52 console.log('No messages in workflow_messages queue')53 return new Response(JSON.stringify({ message: 'No messages in queue' }), {54 status: 200,55 headers: { 'Content-Type': 'application/json' },56 })57 }5859 console.log(`Found ${messages.length} messages to process`)6061 // Process each message that was read off the queue62 for (const message of messages) {63 try {64 await processMessage(message as QueueMessage)65 } catch (error) {66 console.error(`Error processing message ${message.msg_id}:`, error)67 }68 }6970 // Return immediately while background processing continues71 return new Response(72 JSON.stringify({73 message: `Processing ${messages.length} messages in background`,74 count: messages.length,75 }),76 {77 status: 200,78 headers: { 'Content-Type': 'application/json' },79 }80 )81})每次运行这个 Edge 函数时,它会:
🌐 Every time this Edge Function is run it:
- 从队列中读取5条消息
- 调用
processMessage函数 processMessage结束时,消息会从队列中删除- 如果
processMessage抛出错误,错误会被记录下来。在这种情况下,消息仍然在队列中,所以下一次这个 Edge Function 运行时,它会再次读取这条消息。
你可能会发现这种设置在配合 Supabase Cron 使用时很方便。你可以设置 Cron,让每隔 N 分钟或秒,Edge Function 就会运行并处理队列中的一些消息。
🌐 You might find this kind of setup handy to run with Supabase Cron. You can set up Cron so that every N number of minutes or seconds, the Edge Function will run and process a number of messages off the queue.
同样,你可以在任何时候通过 supabase.functions.invoke 命令调用 Edge 函数。
🌐 Similarly, you can invoke the Edge Function on command at any given time with supabase.functions.invoke.