Skip to content
Queues

使用 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.

1
import 'jsr:@supabase/functions-js/edge-runtime.d.ts'
2
import { createClient } from 'npm:@supabase/supabase-js@2'
3
4
const supabaseUrl = 'supabaseURL'
5
const supabaseKey = 'supabaseKey'
6
7
const supabase = createClient(supabaseUrl, supabaseKey)
8
const queueName = 'your_queue_name'
9
10
// Type definition for queue messages
11
interface QueueMessage {
12
msg_id: bigint
13
read_ct: number
14
vt: string
15
enqueued_at: string
16
message: any
17
}
18
19
async function processMessage(message: QueueMessage) {
20
//
21
// Do whatever logic you need to with the message content
22
//
23
// Delete the message from the queue
24
const { error: deleteError } = await supabase.schema('pgmq_public').rpc('delete', {
25
queue_name: queueName,
26
msg_id: message.msg_id,
27
})
28
29
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
}
35
36
Deno.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 empty
40
n: 5, // Read 5 messages off the queue
41
})
42
43
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
}
50
51
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
}
58
59
console.log(`Found ${messages.length} messages to process`)
60
61
// Process each message that was read off the queue
62
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
}
69
70
// Return immediately while background processing continues
71
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:

  1. 从队列中读取5条消息
  2. 调用 processMessage 函数
  3. processMessage结束时,消息会从队列中删除
  4. 如果 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.