목차

, , , , ,

PGMQ

PGMQ(Postgres Message Queue)는 PostgreSQL 내부의 SQL object만으로 message queue를 제공하는 extension이다. 별도 background worker 없이 queue 생성, JSON message 송수신, visibility timeout, archive, FIFO group, metrics를 SQL function으로 다룬다.

Summary

Installation

Docker

공식 image에는 PGMQ extension이 미리 설치되어 있다. 아래 password는 local test용 예시이며 production에서는 안전한 secret 전달 방식을 사용한다.

docker run -d \
  --name pgmq-postgres \
  -e POSTGRES_PASSWORD=postgres \
  -p 5432:5432 \
  ghcr.io/pgmq/pg18-pgmq:v1.10.0
 
psql 'postgres://postgres:postgres@localhost:5432/postgres'
CREATE EXTENSION pgmq;

PGXN

PostgreSQL server host에서 pg_configPATH에 잡힌 상태로 설치한다.

pgxn install pgmq

Source build

git clone https://github.com/pgmq/pgmq.git
cd pgmq/pgmq-extension
make
sudo make install

설치 후 사용할 database마다 extension을 활성화한다.

SELECT name, default_version, installed_version
FROM pg_available_extensions
WHERE name = 'pgmq';
 
CREATE EXTENSION pgmq;

SQL-only

server filesystem에 extension file을 설치할 수 없는 환경에서는 SQL object만 설치할 수 있다. 아래 unversioned 방식은 fresh installation 전용이므로 upgrade가 필요한 운영 환경에서는 extension 또는 versioned client 방식을 우선 검토한다.

git clone https://github.com/pgmq/pgmq.git
cd pgmq
psql -f pgmq-extension/sql/pgmq.sql 'postgres://USER:PASSWORD@HOST:5432/DBNAME'
공식 upstream은 Debian/Ubuntu APT, RHEL/Fedora DNF/YUM, Homebrew, Windows winget 전용 PGMQ package를 안내하지 않는다. 해당 platform에서도 PostgreSQL server와 build dependency를 준비한 뒤 PGXN/source 방식을 쓰거나 공식 Docker image를 사용한다.

Verification

\dx pgmq
SELECT extversion FROM pg_extension WHERE extname = 'pgmq';
SELECT * FROM pgmq.list_queues();

Usage

SELECT pgmq.create('jobs');
SELECT * FROM pgmq.send('jobs', '{"task":"resize","image_id":42}'::jsonb);
SELECT * FROM pgmq.read('jobs', 30, 1);
SELECT pgmq.archive('jobs', 1);

Commands

Queue lifecycle

Function Purpose Notes
pgmq.create(QUEUE_NAME) 일반 queue 생성 queue 이름은 최대 47자
pgmq.create_non_partitioned(QUEUE_NAME) non-partitioned queue 생성 create()의 명시적 형태
pgmq.create_partitioned(…) partitioned queue 생성 pg_partman 필요
pgmq.create_unlogged(QUEUE_NAME) unlogged queue 생성 성능 우선, crash durability 없음
pgmq.list_queues() queue 목록 조회 partitioned/unlogged 여부 포함
pgmq.purge_queue(QUEUE_NAME) queue의 모든 message 영구 삭제 삭제 건수 반환
pgmq.drop_queue(QUEUE_NAME) queue와 archive table 삭제 destructive
SELECT pgmq.create('jobs');
SELECT pgmq.create_partitioned('events', '100000', '10000000');
SELECT * FROM pgmq.list_queues();
pgmq.purge_queue()는 queue 내용 전체를, pgmq.drop_queue()는 queue와 archive table을 삭제한다. production에서는 대상 이름과 backup/retention 정책을 확인한 뒤 실행한다.

Send

message와 선택적 headersjsonb다. delay에는 seconds 정수 또는 timestamptz를 전달할 수 있다.

-- 단일 message
SELECT * FROM pgmq.send(
  queue_name => 'jobs',
  msg        => '{"task":"resize","image_id":42}'::jsonb
);
 
-- header와 10초 delay
SELECT * FROM pgmq.send(
  queue_name => 'jobs',
  msg        => '{"task":"email","user_id":7}'::jsonb,
  headers    => '{"trace_id":"abc123"}'::jsonb,
  delay      => 10
);
 
-- batch
SELECT * FROM pgmq.send_batch(
  'jobs',
  ARRAY['{"task":"one"}', '{"task":"two"}']::jsonb[]
);

Read and acknowledge

-- 최대 5개를 읽고 60초간 다른 consumer에게 숨긴다.
SELECT * FROM pgmq.read('jobs', 60, 5);
 
-- queue가 비었으면 최대 5초간 100ms 간격으로 poll한다.
SELECT * FROM pgmq.read_with_poll('jobs', 60, 5, 5, 100);
 
-- 성공 처리 후 영구 삭제
SELECT pgmq.delete('jobs', 42);
 
-- 또는 archive table로 이동
SELECT pgmq.archive('jobs', 42);
 
-- batch acknowledge
SELECT * FROM pgmq.delete('jobs', ARRAY[43, 44]::bigint[]);
SELECT * FROM pgmq.archive('jobs', ARRAY[45, 46]::bigint[]);
vt는 예상 처리 시간보다 길게 잡는다. vt 안에서는 동일 message를 한 consumer에게 전달하지만, 시간이 끝나기 전에 delete/archive하지 않으면 다시 visible 상태가 되므로 consumer 작업은 idempotent하게 설계하는 편이 안전하다.

pgmq.pop()은 읽는 즉시 삭제하므로 consumer가 이후 처리에 실패하면 message를 복구할 수 없다.

SELECT * FROM pgmq.pop('jobs');

Visibility timeout

처리 시간이 늘어난 message는 현재 시점부터 새 vt seconds만큼 visibility timeout을 연장할 수 있다.

SELECT * FROM pgmq.set_vt('jobs', 42, 120);
SELECT * FROM pgmq.set_vt('jobs', ARRAY[43, 44]::bigint[], 120);

FIFO groups

x-pgmq-group header가 같은 message는 group 안에서 순서대로 처리한다. FIFO 조회를 자주 쓰는 queue에는 header GIN index를 만든다.

SELECT * FROM pgmq.send(
  'jobs',
  '{"sequence":1}'::jsonb,
  '{"x-pgmq-group":"customer-42"}'::jsonb
);
 
SELECT pgmq.create_fifo_index('jobs');
SELECT * FROM pgmq.read_grouped('jobs', 60, 10);
SELECT * FROM pgmq.read_grouped_rr('jobs', 60, 10);

Metrics

SELECT * FROM pgmq.metrics('jobs');
SELECT * FROM pgmq.metrics_all();

주요 field는 queue_length, queue_visible_length, newest_msg_age_sec, oldest_msg_age_sec, total_messages, scrape_time이다.

LISTEN / NOTIFY

sporadic traffic에서는 insert notification을 활성화해 불필요한 polling을 줄일 수 있다. channel 이름은 pgmq.q_QUEUE_NAME.INSERT 형태다.

SELECT pgmq.enable_notify_insert('jobs', 250);
LISTEN "pgmq.q_jobs.INSERT";
 
-- 더 이상 notification이 필요하지 않을 때
SELECT pgmq.disable_notify_insert('jobs');

Examples

Safe consumer transaction

아래는 한 message를 읽고 application 처리 결과가 성공했을 때 archive하는 기본 흐름이다. 외부 API 호출처럼 database transaction으로 되돌릴 수 없는 작업에는 별도의 idempotency key와 retry 정책이 필요하다.

BEGIN;
 
SELECT *
FROM pgmq.read('jobs', 60, 1);
 
-- application이 반환된 msg_id=42의 작업을 성공적으로 처리했다고 가정
SELECT pgmq.archive('jobs', 42);
 
COMMIT;

Inspect archive

archive table은 pgmq.a_QUEUE_NAME 형식이다.

SELECT msg_id, read_ct, enqueued_at, archived_at, message, headers
FROM pgmq.a_jobs
ORDER BY archived_at DESC
LIMIT 20;

Help

PGMQ는 독립 CLI가 아니므로 pgmq –help가 없다. 설치된 database에서 실제 function signature와 extension version을 조회한다.

PGMQ function inspection

Troubleshooting

Compatibility

Deprecated / Legacy

See Also

History