From 5deecaa44ce7307e01532b56a837be9b7fa0a866 Mon Sep 17 00:00:00 2001 From: Oliver Rice Date: Wed, 11 Sep 2024 15:12:56 -0500 Subject: [PATCH] Add docs for pgmq (#29232) * add docs for pgmq * make prettier happy * Update apps/docs/content/guides/database/extensions/pgmq.mdx Co-authored-by: Charis <26616127+charislam@users.noreply.github.com> * create extension -> create schema * pgmq is not relocatable --------- Co-authored-by: Charis <26616127+charislam@users.noreply.github.com> --- .../NavigationMenu.constants.ts | 4 + .../guides/database/extensions/pgmq.mdx | 707 ++++++++++++++++++ 2 files changed, 711 insertions(+) create mode 100644 apps/docs/content/guides/database/extensions/pgmq.mdx diff --git a/apps/docs/components/Navigation/NavigationMenu/NavigationMenu.constants.ts b/apps/docs/components/Navigation/NavigationMenu/NavigationMenu.constants.ts index c4444fe8ad8..d074fe8e7a0 100644 --- a/apps/docs/components/Navigation/NavigationMenu/NavigationMenu.constants.ts +++ b/apps/docs/components/Navigation/NavigationMenu/NavigationMenu.constants.ts @@ -902,6 +902,10 @@ export const database: NavMenuConstant = { name: 'PostGIS: Geo queries', url: '/guides/database/extensions/postgis', }, + { + name: 'pgmq: Queues', + url: '/guides/database/extensions/pgmq', + }, { name: 'pgsodium (pending deprecation): Encryption Features', url: '/guides/database/extensions/pgsodium', diff --git a/apps/docs/content/guides/database/extensions/pgmq.mdx b/apps/docs/content/guides/database/extensions/pgmq.mdx new file mode 100644 index 00000000000..bd33c51dad7 --- /dev/null +++ b/apps/docs/content/guides/database/extensions/pgmq.mdx @@ -0,0 +1,707 @@ +--- +id: 'pgmq' +title: 'pgmq: Queues' +description: 'pgmq: Managed queues in Postgres' +--- + +pgmq is a lightweight message queue built on Postgres. + +## Features + +- Lightweight - No background worker or external dependencies, just Postgres functions packaged in an extension +- "exactly once" delivery of messages to a consumer within a visibility timeout +- API parity with AWS SQS and RSMQ +- Messages stay in the queue until explicitly removed +- Messages can be archived, instead of deleted, for long-term retention and replayability + +## Enable the extension + +```sql +create extension pgmq; +``` + +## Usage [#get-usage] + +### Sending Messages + +#### send + +Send a single message to a queue. + +{/* prettier-ignore */} +```sql +pgmq.send( + queue_name text, + msg jsonb, + delay integer default 0 +) +returns setof bigint +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :------ | :----------------------------------------------------------------- | +| queue_name | text | The name of the queue | +| msg | jsonb | The message to send to the queue | +| delay | integer | Time in seconds before the message becomes visible. Defaults to 0. | + +Example: + +{/* prettier-ignore */} +```sql +select * from pgmq.send('my_queue', '{"hello": "world"}'); + send +------ + 4 +``` + +--- + +#### send_batch + +Send 1 or more messages to a queue. + +{/* prettier-ignore */} +```sql +pgmq.send_batch( + queue_name text, + msgs jsonb[], + delay integer default 0 +) +returns setof bigint +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :------ | :------------------------------------------------------------------ | +| queue_name | text | The name of the queue | +| msgs | jsonb[] | Array of messages to send to the queue | +| delay | integer | Time in seconds before the messages becomes visible. Defaults to 0. | + +{/* prettier-ignore */} +```sql +select * from pgmq.send_batch( + 'my_queue', + array[ + '{"hello": "world_0"}'::jsonb, + '{"hello": "world_1"}'::jsonb + ] +); + send_batch +------------ + 1 + 2 +``` + +--- + +### Reading Messages + +#### read + +Read 1 or more messages from a queue. The VT specifies the delay in seconds between reading and the message becoming invisible to other consumers. + +{/* prettier-ignore */} +```sql +pgmq.read( + queue_name text, + vt integer, + qty integer +) + +returns setof pgmq.message_record +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :------ | :-------------------------------------------------------------- | +| queue_name | text | The name of the queue | +| vt | integer | Time in seconds that the message become invisible after reading | +| qty | integer | The number of messages to read from the queue. Defaults to 1 | + +Example: + +{/* prettier-ignore */} +```sql +select * from pgmq.read('my_queue', 10, 2); + msg_id | read_ct | enqueued_at | vt | message +--------+---------+-------------------------------+-------------------------------+---------------------- + 1 | 1 | 2023-10-28 19:14:47.356595-05 | 2023-10-28 19:17:08.608922-05 | {"hello": "world_0"} + 2 | 1 | 2023-10-28 19:14:47.356595-05 | 2023-10-28 19:17:08.608974-05 | {"hello": "world_1"} +(2 rows) +``` + +--- + +#### read_with_poll + +Same as read(). Also provides convenient long-poll functionality. +When there are no messages in the queue, the function call will wait for `max_poll_seconds` in duration before returning. +If messages reach the queue during that duration, they will be read and returned immediately. + +{/* prettier-ignore */} +```sql + pgmq.read_with_poll( + queue_name text, + vt integer, + qty integer, + max_poll_seconds integer default 5, + poll_interval_ms integer default 100 +) +returns setof pgmq.message_record +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------------- | :------ | :-------------------------------------------------------------------------- | +| queue_name | text | The name of the queue | +| vt | integer | Time in seconds that the message become invisible after reading. | +| qty | integer | The number of messages to read from the queue. Defaults to 1. | +| max_poll_seconds | integer | Time in seconds to wait for new messages to reach the queue. Defaults to 5. | +| poll_interval_ms | integer | Milliseconds between the internal poll operations. Defaults to 100. | + +Example: + +{/* prettier-ignore */} +```sql +select * from pgmq.read_with_poll('my_queue', 1, 1, 5, 100); + msg_id | read_ct | enqueued_at | vt | message +--------+---------+-------------------------------+-------------------------------+-------------------- + 1 | 1 | 2023-10-28 19:09:09.177756-05 | 2023-10-28 19:27:00.337929-05 | {"hello": "world"} +``` + +--- + +#### pop + +Reads a single message from a queue and deletes it upon read. + +Note: utilization of pop() results in at-most-once delivery semantics if the consuming application does not guarantee processing of the message. + +{/* prettier-ignore */} +```sql +pgmq.pop(queue_name text) +returns setof pgmq.message_record +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :--- | :-------------------- | +| queue_name | text | The name of the queue | + +Example: + +{/* prettier-ignore */} +```sql +pgmq=# select * from pgmq.pop('my_queue'); + msg_id | read_ct | enqueued_at | vt | message +--------+---------+-------------------------------+-------------------------------+-------------------- + 1 | 2 | 2023-10-28 19:09:09.177756-05 | 2023-10-28 19:27:00.337929-05 | {"hello": "world"} +``` + +--- + +### Deleting/Archiving Messages + +#### delete (single) + +Deletes a single message from a queue. + +{/* prettier-ignore */} +```sql +pgmq.delete (queue_name text, msg_id: bigint) +returns boolean +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :----- | :---------------------------------- | +| queue_name | text | The name of the queue | +| msg_id | bigint | Message ID of the message to delete | + +Example: + +{/* prettier-ignore */} +```sql +select pgmq.delete('my_queue', 5); + delete +-------- + t +``` + +--- + +#### delete (batch) + +Delete one or many messages from a queue. + +{/* prettier-ignore */} +```sql +pgmq.delete (queue_name text, msg_ids: bigint[]) +returns setof bigint +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :------- | :----------------------------- | +| queue_name | text | The name of the queue | +| msg_ids | bigint[] | Array of message IDs to delete | + +Examples: + +Delete two messages that exist. + +{/* prettier-ignore */} +```sql +select * from pgmq.delete('my_queue', array[2, 3]); + delete +-------- + 2 + 3 +``` + +Delete two messages, one that exists and one that does not. Message `999` does not exist. + +```sql +select * from pgmq.delete('my_queue', array[6, 999]); + delete +-------- + 6 +``` + +--- + +#### purge_queue + +Permanently deletes all messages in a queue. Returns the number of messages that were deleted. + +```text +purge_queue(queue_name text) +returns bigint +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :--- | :-------------------- | +| queue_name | text | The name of the queue | + +Example: + +Purge the queue when it contains 8 messages; + +{/* prettier-ignore */} +```sql +select * from pgmq.purge_queue('my_queue'); + purge_queue +------------- + 8 +``` + +--- + +#### archive (single) + +Removes a single requested message from the specified queue and inserts it into the queue's archive. + +{/* prettier-ignore */} +```sql +pgmq.archive(queue_name text, msg_id bigint) +returns boolean +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :----- | :----------------------------------- | +| queue_name | text | The name of the queue | +| msg_id | bigint | Message ID of the message to archive | + +Returns +Boolean value indicating success or failure of the operation. + +Example; remove message with ID 1 from queue `my_queue` and archive it: + +{/* prettier-ignore */} +```sql +select * from pgmq.archive('my_queue', 1); + archive +--------- + t +``` + +--- + +#### archive (batch) + +Deletes a batch of requested messages from the specified queue and inserts them into the queue's archive. +Returns an array of message ids that were successfully archived. + +```text +pgmq.archive(queue_name text, msg_ids bigint[]) +RETURNS SETOF bigint +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :------- | :------------------------------ | +| queue_name | text | The name of the queue | +| msg_ids | bigint[] | Array of message IDs to archive | + +Examples: + +Delete messages with ID 1 and 2 from queue `my_queue` and move to the archive. + +{/* prettier-ignore */} +```sql +select * from pgmq.archive('my_queue', array[1, 2]); + archive +--------- + 1 + 2 +``` + +Delete messages 4, which exists and 999, which does not exist. + +{/* prettier-ignore */} +```sql +select * from pgmq.archive('my_queue', array[4, 999]); + archive +--------- + 4 +``` + +--- + +### Queue Management + +#### create + +Create a new queue. + +{/* prettier-ignore */} +```sql +pgmq.create(queue_name text) +returns void +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :--- | :-------------------- | +| queue_name | text | The name of the queue | + +Example: + +{/* prettier-ignore */} +```sql +select from pgmq.create('my_queue'); + create +-------- +``` + +--- + +#### create_partitioned + +Create a partitioned queue. + +{/* prettier-ignore */} +```sql +pgmq.create_partitioned ( + queue-ue_name text, + partition_interval text default '10000'::text, + retention_interval text default '100000'::text +) +returns void +``` + +**Parameters:** + +| Parameter | Type | Description | +| :----------------- | :--- | :-------------------- | +| queue_name | text | The name of the queue | +| partition_interval | text | The name of the queue | +| retention_interval | text | The name of the queue | + +Example: + +Create a queue with 100,000 messages per partition, and will retain 10,000,000 messages on old partitions. Partitions greater than this will be deleted. + +{/* prettier-ignore */} +```sql +select * from pgmq.create_partitioned( + 'my_partitioned_queue', + '100000', + '10000000' +); + create_partitioned +-------------------- +``` + +--- + +#### create_unlogged + +Creates an unlogged table. This is useful when write throughput is more important that durability. +See Postgres documentation for [unlogged tables](https://www.postgresql.org/docs/current/sql-createtable.html#SQL-CREATETABLE-UNLOGGED) for more information. + +{/* prettier-ignore */} +```sql +pgmq.create_unlogged(queue_name text) +returns void +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :--- | :-------------------- | +| queue_name | text | The name of the queue | + +Example: + +{/* prettier-ignore */} +```sql +select pgmq.create_unlogged('my_unlogged'); + create_unlogged +----------------- +``` + +--- + +#### detach_archive + +Drop the queue's archive table as a member of the PGMQ extension. Useful for preventing the queue's archive table from being drop when `drop extension pgmq` is executed. +This does not prevent the further archives() from appending to the archive table. + +{/* prettier-ignore */} +```sql +pgmq.detach_archive(queue_name text) +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :--- | :-------------------- | +| queue_name | text | The name of the queue | + +Example: + +{/* prettier-ignore */} +```sql +select * from pgmq.detach_archive('my_queue'); + detach_archive +---------------- +``` + +--- + +#### drop_queue + +Deletes a queue and its archive table. + +{/* prettier-ignore */} +```sql +pgmq.drop_queue(queue_name text) +returns boolean +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :--- | :-------------------- | +| queue_name | text | The name of the queue | + +Example: + +{/* prettier-ignore */} +```sql +select * from pgmq.drop_queue('my_unlogged'); + drop_queue +------------ + t +``` + +### Utilities + +#### set_vt + +Sets the visibility timeout of a message to a specified time duration in the future. Returns the record of the message that was updated. + +{/* prettier-ignore */} +```sql +pgmq.set_vt( + queue_name text, + msg_id bigint, + vt_offset integer +) +returns pgmq.message_record +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :------ | :-------------------------------------------------------------------- | +| queue_name | text | The name of the queue | +| msg_id | bigint | ID of the message to set visibility time | +| vt_offset | integer | Duration from now, in seconds, that the message's VT should be set to | + +Example: + +Set the visibility timeout of message 1 to 30 seconds from now. + +```sql +select * from pgmq.set_vt('my_queue', 11, 30); + msg_id | read_ct | enqueued_at | vt | message +--------+---------+-------------------------------+-------------------------------+---------------------- + 1 | 0 | 2023-10-28 19:42:21.778741-05 | 2023-10-28 19:59:34.286462-05 | {"hello": "world_0"} +``` + +--- + +#### list_queues + +List all the queues that currently exist. + +{/* prettier-ignore */} +```sql +list_queues() +RETURNS TABLE( + queue_name text, + created_at timestamp with time zone, + is_partitioned boolean, + is_unlogged boolean +) +``` + +Example: + +{/* prettier-ignore */} +```sql +select * from pgmq.list_queues(); + queue_name | created_at | is_partitioned | is_unlogged +----------------------+-------------------------------+----------------+------------- + my_queue | 2023-10-28 14:13:17.092576-05 | f | f + my_partitioned_queue | 2023-10-28 19:47:37.098692-05 | t | f + my_unlogged | 2023-10-28 20:02:30.976109-05 | f | t +``` + +--- + +#### metrics + +Get metrics for a specific queue. + +{/* prettier-ignore */} +```sql +pgmq.metrics(queue_name: text) +returns table( + queue_name text, + queue_length bigint, + newest_msg_age_sec integer, + oldest_msg_age_sec integer, + total_messages bigint, + scrape_time timestamp with time zone +) +``` + +**Parameters:** + +| Parameter | Type | Description | +| :--------- | :--- | :-------------------- | +| queue_name | text | The name of the queue | + +**Returns:** + +| Attribute | Type | Description | +| :----------------- | :----------------------- | :------------------------------------------------------------------------ | +| queue_name | text | The name of the queue | +| queue_length | bigint | Number of messages currently in the queue | +| newest_msg_age_sec | integer \| null | Age of the newest message in the queue, in seconds | +| oldest_msg_age_sec | integer \| null | Age of the oldest message in the queue, in seconds | +| total_messages | bigint | Total number of messages that have passed through the queue over all time | +| scrape_time | timestamp with time zone | The current timestamp | + +Example: + +{/* prettier-ignore */} +```sql +select * from pgmq.metrics('my_queue'); + queue_name | queue_length | newest_msg_age_sec | oldest_msg_age_sec | total_messages | scrape_time +------------+--------------+--------------------+--------------------+----------------+------------------------------- + my_queue | 16 | 2445 | 2447 | 35 | 2023-10-28 20:23:08.406259-05 +``` + +--- + +#### metrics_all + +Get metrics for all existing queues. + +```text +pgmq.metrics_all() +RETURNS TABLE( + queue_name text, + queue_length bigint, + newest_msg_age_sec integer, + oldest_msg_age_sec integer, + total_messages bigint, + scrape_time timestamp with time zone +) +``` + +**Returns:** + +| Attribute | Type | Description | +| :----------------- | :----------------------- | :------------------------------------------------------------------------ | +| queue_name | text | The name of the queue | +| queue_length | bigint | Number of messages currently in the queue | +| newest_msg_age_sec | integer \| null | Age of the newest message in the queue, in seconds | +| oldest_msg_age_sec | integer \| null | Age of the oldest message in the queue, in seconds | +| total_messages | bigint | Total number of messages that have passed through the queue over all time | +| scrape_time | timestamp with time zone | The current timestamp | + +{/* prettier-ignore */} +```sql +select * from pgmq.metrics_all(); + queue_name | queue_length | newest_msg_age_sec | oldest_msg_age_sec | total_messages | scrape_time +----------------------+--------------+--------------------+--------------------+----------------+------------------------------- + my_queue | 16 | 2563 | 2565 | 35 | 2023-10-28 20:25:07.016413-05 + my_partitioned_queue | 1 | 11 | 11 | 1 | 2023-10-28 20:25:07.016413-05 + my_unlogged | 1 | 3 | 3 | 1 | 2023-10-28 20:25:07.016413-05 +``` + +### Types + +#### message_record + +The complete representation of a message in a queue. + +| Attribute Name | Type | Description | +| :------------- | :----------------------- | :--------------------------------------------------------------------- | +| msg_id | bigint | Unique ID of the message | +| read_ct | bigint | Number of times the message has been read. Increments on read(). | +| enqueued_at | timestamp with time zone | time that the message was inserted into the queue | +| vt | timestamp with time zone | Timestamp when the message will become available for consumers to read | +| message | jsonb | The message payload | + +Example: + +{/* prettier-ignore */} +```sql + msg_id | read_ct | enqueued_at | vt | message +--------+---------+-------------------------------+-------------------------------+-------------------- + 1 | 1 | 2023-10-28 19:06:19.941509-05 | 2023-10-28 19:06:27.419392-05 | {"hello": "world"} +``` + +## Resources + +- Official Docs: [pgmq/api](https://tembo.io/pgmq/#creating-a-queue)