PGMQ(Postgres Message Queue)는 PostgreSQL 내부의 SQL object만으로 message queue를 제공하는 extension이다. 별도 background worker 없이 queue 생성, JSON message 송수신, visibility timeout, archive, FIFO group, metrics를 SQL function으로 다룬다.
pgmq.q_QUEUE_NAME table에 저장되며 명시적으로 delete 또는 archive할 때까지 유지된다.pgmq.delete() 또는 pgmq.archive()를 호출해야 한다.공식 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;
PostgreSQL server host에서 pg_config가 PATH에 잡힌 상태로 설치한다.
pgxn install pgmq
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;
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'
winget 전용 PGMQ package를 안내하지 않는다. 해당 platform에서도 PostgreSQL server와 build dependency를 준비한 뒤 PGXN/source 방식을 쓰거나 공식 Docker image를 사용한다.
\dx pgmq SELECT extversion FROM pg_extension WHERE extname = 'pgmq'; SELECT * FROM pgmq.list_queues();
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);
pgmq.create(QUEUE_NAME) queue 생성pgmq.send(QUEUE_NAME, MESSAGE) JSON message 전송pgmq.read(QUEUE_NAME, VT_SECONDS, QTY) message 읽기pgmq.archive(QUEUE_NAME, MSG_ID) 처리한 message 보관| 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 정책을 확인한 뒤 실행한다.
message와 선택적 headers는 jsonb다. 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[] );
-- 최대 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');
처리 시간이 늘어난 message는 현재 시점부터 새 vt seconds만큼 visibility timeout을 연장할 수 있다.
SELECT * FROM pgmq.set_vt('jobs', 42, 120); SELECT * FROM pgmq.set_vt('jobs', ARRAY[43, 44]::bigint[], 120);
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);
read_grouped()는 가장 오래된 available group을 우선 채운다.read_grouped_rr()는 여러 group을 round-robin 방식으로 읽는다.read_grouped_with_poll(), read_grouped_rr_with_poll()이 있다.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이다.
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');
아래는 한 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;
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;
PGMQ는 독립 CLI가 아니므로 pgmq –help가 없다. 설치된 database에서 실제 function signature와 extension version을 조회한다.
function pgmq… does not exist: 현재 database에 CREATE EXTENSION pgmq를 실행했는지, \dx pgmq와 \df pgmq.*로 확인한다.vt 안에 delete/archive하지 못한 것이다. 처리 시간과 retry 정책에 맞춰 vt를 늘리거나 set_vt()로 연장한다.vt가 미래인지 확인하고 metrics()의 visible length를 함께 본다.pg_partman extension 설치 및 활성화 여부를 확인한다.pgmq.create_fifo_index()를 적용했는지 확인한다.pg_partman에 의존한다.\df+ pgmq.* 결과를 우선한다.pgmq.detach_archive(QUEUE_NAME)pgmq.drop_queue(QUEUE_NAME, PARTITIONED)