The broker can stream a PostgreSQL table into a queue, and a queue into a table, with no other process running. The source reads committed changes through logical replication and turns each one into a message, in a partition named after the row’s key, in commit order. The sink is a consumer group the broker runs for you, and it writes every message into a table. Both are set up with one HTTP call, both run on every node of a cluster, and both give exactly-once effects: a crash, a failover or a redelivery never writes a change twice and never loses one.
Try it
PostgreSQL 17 with logical decoding, and a broker on the same Docker network. The broker needs
QUEEN_ENCRYPTION_KEY, because it stores the database password sealed with that key:
docker network create queen-pg
docker run -d --name pg --network queen-pg -e POSTGRES_PASSWORD=secret \
postgres:17 -c wal_level=logical
docker run -d --name queen --network queen-pg --platform linux/amd64 -p 6632:6632 \
-v queen-data:/var/lib/queen/raft \
-e QUEEN_ENCRYPTION_KEY=$(openssl rand -hex 32) \
ghcr.io/queen-mq/queen:latestA table with a nested JSON column, and two rows in it:
docker exec -i pg psql -U postgres <<'SQL'
CREATE TABLE orders (
id bigint PRIMARY KEY,
customer jsonb NOT NULL,
status text NOT NULL,
total numeric(12,2) NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
);
INSERT INTO orders (id, customer, status, total) VALUES
(1, '{"name": "Ada", "tier": "gold"}', 'new', 42.00),
(2, '{"name": "Linus", "tier": "silver"}', 'new', 17.50);
SQLCreate the source:
curl -s -X PUT localhost:6632/api/v1/connectors/orders-src -H 'content-type: application/json' -d '{
"kind": "source",
"connection": {"host": "pg", "database": "postgres", "user": "postgres", "password": "secret", "sslMode": "disable"},
"source": {"tables": [{"table": "public.orders", "queue": "orders"}]}
}' > /dev/nullEvery node re-reads the connector documents every 5 seconds, so give it a few seconds:
curl -s localhost:6632/api/v1/connectors/orders-src | jq -c '.status | {phase, pointerLsn, messages}'{"phase":"streaming","pointerLsn":"0/1526990","messages":2}Within a few seconds the source created a publication and a replication slot, copied the two rows
that were already there (the snapshot) and started streaming. The queue orders holds one message
per row, each in the partition named after the row’s primary key. Read them as any Queen consumer
would. A new consumer group starts at the end of a queue unless it asks for subscriptionMode=all:
curl -s 'localhost:6632/api/v1/pop/queue/orders?batch=10&partitions=10&consumerGroup=reader&subscriptionMode=all&autoAck=true' \
| jq -c '.messages[] | {partition, data}'{"partition":"1","data":{"op":"r","table":"public.orders","key":{"id":1},"after":{"id":1,"customer":{"name":"Ada","tier":"gold"},"status":"new","total":42.00,"created_at":"2026-10-03T06:14:39.365472+00:00"},"before":null,"lsn":"0/1526990","ts":"2026-10-03T06:14:43.797508Z","seq":0}}
{"partition":"2","data":{"op":"r","table":"public.orders","key":{"id":2},"after":{"id":2,"customer":{"name":"Linus","tier":"silver"},"status":"new","total":17.50,"created_at":"2026-10-03T06:14:39.365472+00:00"},"before":null,"lsn":"0/1526990","ts":"2026-10-03T06:14:43.797508Z","seq":1}}r is a row read by the snapshot. Now change the table: one update, then an insert and a delete in
one transaction.
docker exec -i pg psql -U postgres -q <<'SQL'
UPDATE orders SET status = 'paid' WHERE id = 1;
BEGIN;
INSERT INTO orders (id, customer, status, total) VALUES (3, '{"name": "Grace", "tier": "gold"}', 'new', 99.90);
DELETE FROM orders WHERE id = 2;
COMMIT;
SQL
curl -s 'localhost:6632/api/v1/pop/queue/orders?batch=10&partitions=10&consumerGroup=reader&autoAck=true' \
| jq -c '.messages[] | {partition, transactionId, data}'{"partition":"1","transactionId":"pg:00c0b669:0000000001570D68:0","data":{"op":"u","table":"public.orders","key":{"id":1},"after":{"id":1,"customer":{"name":"Ada","tier":"gold"},"status":"paid","total":42.00,"created_at":"2026-10-03T06:14:39.365472+00:00"},"before":null,"lsn":"0/1570D68","xid":753,"ts":"2026-10-03T06:15:03.583036Z","seq":0}}
{"partition":"2","transactionId":"pg:00c0b669:0000000001570EA0:1","data":{"op":"d","table":"public.orders","key":{"id":2},"after":null,"before":{"id":2},"lsn":"0/1570EA0","xid":754,"ts":"2026-10-03T06:15:03.591716Z","seq":1}}
{"partition":"3","transactionId":"pg:00c0b669:0000000001570EA0:0","data":{"op":"c","table":"public.orders","key":{"id":3},"after":{"id":3,"customer":{"name":"Grace","tier":"gold"},"status":"new","total":99.90,"created_at":"2026-10-03T06:15:03.591416+00:00"},"before":null,"lsn":"0/1570EA0","xid":754,"ts":"2026-10-03T06:15:03.591716Z","seq":0}}The insert and the delete came from one PostgreSQL transaction, so they carry the same lsn and
xid, and they reached Queen in one Queen transaction: no consumer could see one without the
other. Each transactionId is built from the change’s place in the WAL, so it is the same however
many times the source has to send it.
Now the other direction. A sink in cdc mode applies these events to a second table:
docker exec pg psql -U postgres -c 'CREATE TABLE orders_copy (LIKE orders INCLUDING ALL);'
curl -s -X PUT localhost:6632/api/v1/connectors/orders-copy -H 'content-type: application/json' -d '{
"kind": "sink",
"connection": {"host": "pg", "database": "postgres", "user": "postgres", "password": "secret", "sslMode": "disable"},
"sink": {"queue": "orders", "table": "public.orders_copy", "mode": "cdc"}
}' > /dev/nullA few seconds later the copy holds what the table holds:
docker exec pg psql -U postgres -c 'SELECT id, customer, status, total FROM orders_copy ORDER BY id;' id | customer | status | total
----+-----------------------------------+--------+-------
1 | {"name": "Ada", "tier": "gold"} | paid | 42.00
3 | {"name": "Grace", "tier": "gold"} | new | 99.90
(2 rows)How it works
Both connectors run inside the broker process, on threads of their own (QUEEN_PG_THREADS, 2 by
default). They reach the broker in memory, through the same transactions, consumer groups and KV
store every client uses, so a source’s messages are ordinary Queen messages and a sink’s progress
shows up as an ordinary consumer group.
The source. Every node of a cluster runs every source, and one of them owns it at a time,
chosen through a lease in KV. When the owner dies, another node takes over once the lease runs out
(10 seconds by default, QUEEN_PG_LEASE_TTL_MS); when it stops on SIGTERM it gives the lease back
at once. The owner reads the replication slot with PostgreSQL’s built-in pgoutput plugin and
pushes each committed transaction (or several small ones together) as one Queen transaction. That
transaction carries the messages and, in the same commit, the WAL position they reach (the
pointer, a KV entry written only if nobody else moved it first). Queen commits both or neither.
Only after Queen has committed does the source tell PostgreSQL it may recycle that WAL. After a
crash, the next owner reads the pointer and restarts the slot there, and drops anything at or below
it: those changes are already in Queen. That is the whole exactly-once argument for the source: the
position commits with the data, on the side where the data lands.
The snapshot reads each table in key order, 2,000 rows at a time (snapshotChunkRows), while the
stream keeps running. Before and after each chunk the source writes a marker into the WAL. Rows that
changed between the two markers come from the stream instead of the chunk, so a consumer never sees
an older version of a row after a newer one. No long transaction holds back VACUUM, and a crash
repeats at most one chunk.
A transaction too large for one Queen transaction (more than maxBundleMessages changes or
maxBundleBytes of payload) is pushed in pieces, and the pointer records how far into it each piece
got. Every change still arrives exactly once; the pieces are no longer atomic.
The sink. A sink is a consumer group (pg- followed by the connector name, unless you name
one). Every node runs workers workers for it, each with its own PostgreSQL connection, and Queen’s
partition leases spread the partitions over all of them. Each batch is one PostgreSQL transaction.
It locks the batch’s rows in a progress table (queen.sink_progress: sink, partition id, last
applied offset), drops every message at or below its partition’s last offset, writes the rest,
moves the offsets forward, and commits. Only then does it ack the batch in Queen.
A lost ack, an expired lease or a crash brings messages back, and the progress rows turn them into no-ops. Two nodes that end up with the same partition queue on the same progress row, and the second finds the work already done. That is the sink’s half of the argument: the offsets commit with the rows.
A message PostgreSQL refuses (a constraint, a value of the wrong type, a payload the mode cannot
read) is tried maxAttempts times, then dead-lettered with the database’s error, and its partition
carries on after it. A connection that drops, a deadlock or a serialization failure is retried with
backoff and never dead-letters anything.
The change events
Every message the source writes is one JSON object:
| Field | |
|---|---|
op |
c insert, u update, d delete, r a row read by the snapshot, t truncate (only with onTruncate: "emit") |
table |
schema.name |
key |
the key columns and their values: from the new row, or the old one for a delete |
after |
the new row (c, u, r); null for a delete |
before |
the old key for a delete or for an update that changed the key; the whole old row under REPLICA IDENTITY FULL; otherwise null |
unchanged |
names of columns PostgreSQL did not send (see below); absent when empty |
lsn |
the commit position of the source transaction (for r, the snapshot chunk’s position) |
xid |
the source transaction id (absent for r) |
ts |
the commit time (for r, when the chunk was read), UTC, microseconds |
seq |
the change’s index within its transaction (or its chunk) |
Values are converted from PostgreSQL’s own text output, the same way for the snapshot and the
stream, so an r event and a c event for the same row carry identical values:
boolbecomestrue/false. Integers andnumericbecome JSON numbers with the digits PostgreSQL printed, so a 30-digit numeric or a bigint above 2^53 stays exact (as long as your consumer does not parse numbers as doubles).NaNandInfinitybecome strings.jsonandjsonbare embedded as JSON, nested objects and arrays included.timestamptzbecomes"2026-10-03T06:14:39.365472+00:00"andtimestamp"2026-10-03T06:14:39.365472"; both connections run in UTC.- Arrays of the built-in types become JSON arrays.
byteabecomes a"\\x..."string. - Everything else (dates, intervals,
uuid,inet, enums, composites, ranges) becomes the string PostgreSQL prints.
Partitions. By default a row’s messages go to the partition named after its key: the value of
a one-column key as printed ("42"), or the values of a composite key joined with a vertical bar
(a vertical bar or backslash inside a value is escaped with a backslash, NULL is written \N). A
name longer than 128 bytes is replaced by ~ and the first 32 hex digits of its SHA-256. One
partition per row gives every row its own ordered stream. Set partitionBy to a list of columns to
group rows (all orders of a customer, say), or to "single" for one partition named all that
keeps the commit order of the whole table. An update that changes a row’s partition is written as a
d to the old partition and a c to the new one.
Large values. PostgreSQL does not log a TOASTed value (a large one, usually over 2 kB) that an
update left unchanged, so the event lists that column in unchanged instead of carrying it, and a
consumer keeps the value it already has. The source fills the two places where a consumer could not
have it. A row first seen through such an update during the snapshot gets a fill event at the end of
its chunk: an ordinary u whose after holds just the key and the missing columns, with every other
column in unchanged. A row whose key changed gets the missing values read back from the table, so
its c is complete; that value can be slightly newer than the change, and the next event corrects
it. Under REPLICA IDENTITY FULL PostgreSQL logs the old row whole, and unchanged is always empty.
Message ids. pg:<epoch>:<commit LSN in 16 hex digits>:<seq> for a change (a key change’s
second event adds .1), pg:<epoch>:s:<position>:<seq> for a snapshot row and
pg:<epoch>:s:<position>:f<seq> for a fill. The epoch is new each time the source starts over from
a new snapshot, so ids never repeat across a resync.
How a message becomes a row
The sink never creates or alters the table it writes: you create it, and the sink checks it at
start. The only table it creates is its progress table. In append and upsert mode it maps the
message onto the table like this, and in cdc mode it maps the event’s after the same way:
- Each top-level key of the message goes to the column of exactly that name. Keys with no column are ignored. Matching is case-sensitive.
- A column the message does not name keeps its value on an update and gets its
DEFAULTon an insert. It is never set to NULL unless the message saysnull. - PostgreSQL converts the values itself (
jsonb_populate_record), from the exact characters the producer wrote:"12.50"into anumeric(12,2)is 12.50, a nested object into ajsonbcolumn is stored as it is, a nested object into a composite-type column fills it field by field, a JSON array into an array column fills it element by element. A value the column’s type cannot take is a data error, and the message ends up in the dead-letter queue. - Only top-level keys match:
customer.address.citydoes not reach acitycolumn by itself. Keep the subtree in ajsonbcolumn and let PostgreSQL derive the fields with generated columns, which the sink never writes.
For example, a table that keeps two subtrees whole and derives two fields from one of them:
docker exec -i pg psql -U postgres <<'SQL'
CREATE TABLE invoices (
id bigint PRIMARY KEY,
customer jsonb,
lines jsonb,
total numeric(12,2),
customer_tier text GENERATED ALWAYS AS (customer->>'tier') STORED,
city text GENERATED ALWAYS AS (customer->'address'->>'city') STORED,
queen_offset bigint
);
SQL
curl -s -X PUT localhost:6632/api/v1/connectors/invoices-sink -H 'content-type: application/json' -d '{
"kind": "sink",
"connection": {"host": "pg", "database": "postgres", "user": "postgres", "password": "secret", "sslMode": "disable"},
"sink": {"queue": "invoices", "table": "public.invoices", "mode": "upsert", "metadata": {"offset": "queen_offset"}}
}' > /dev/null
curl -s -X POST localhost:6632/api/v1/push -H 'content-type: application/json' -d '{"items":[
{"queue":"invoices","partition":"7","payload":{"id":7,"customer":{"name":"Ada","tier":"gold","address":{"city":"Turin"}},"lines":[{"sku":"A1","qty":2}],"total":"12.50","note":"no column, ignored"}},
{"queue":"invoices","partition":"7","payload":{"id":7,"total":15}}
]}' > /dev/nullA few seconds later:
docker exec pg psql -U postgres -c 'SELECT id, customer, lines, total, customer_tier, city, queen_offset FROM invoices;' id | customer | lines | total | customer_tier | city | queen_offset
----+---------------------------------------------------------------+---------------------------+-------+---------------+-------+--------------
7 | {"name": "Ada", "tier": "gold", "address": {"city": "Turin"}} | [{"qty": 2, "sku": "A1"}] | 15.00 | gold | Turin | 1
(1 row)The second message named only id and total, so it changed only total. note has no column
and was dropped. metadata writes the message’s own facts into columns you name: partition,
partitionId, offset, transactionId, createdAt, and payload, which stores the whole
message in one jsonb column (the way to keep messages that are not objects, or to keep them
whole).
The four modes.
| Mode | Each message | Key |
|---|---|---|
append |
inserted as a new row; a run of messages is one INSERT |
none needed |
upsert |
inserted, or the row with its key updated, with the columns it names | key, else the table’s primary key; it needs a unique index |
cdc |
a source event applied: c/r/u upsert after (an update keeps its unchanged columns and moves the row when its key changed), d deletes, t truncates |
as upsert |
sql |
your statement, run once per message with params |
none |
In sql mode each parameter is a path into the message ($ for all of it, $.a.b, $.lines[0])
or one of @partition, @partitionId, @offset, @transactionId, @createdAt. Every value is
passed as text (objects and arrays as their JSON, a missing path as NULL), and the statement casts
it. Because the statement runs in the same transaction as the progress row, a statement that is not
idempotent runs exactly once per message:
{
"kind": "sink",
"connection": {"host": "pg", "database": "postgres", "user": "postgres", "password": "secret"},
"sink": {
"queue": "payments",
"table": "public.accounts",
"mode": "sql",
"statement": "UPDATE accounts SET balance = balance + $1::numeric WHERE id = $2::bigint",
"params": ["$.amount", "$.accountId"]
}
}With an upsert or cdc sink, keep the table’s key aligned with the partition. Queen orders
messages within a partition, not across partitions, so two partitions writing the same row would
race.
Setting up PostgreSQL
The source needs PostgreSQL 17 or later; the sink works with any version.
-
Set
wal_level = logicaland restart the server. Raisemax_replication_slotsandmax_wal_sendersif other consumers already use slots (one slot per source). -
Give the source a role with
LOGIN REPLICATIONandSELECTon its tables. The source manages its publication (managePublication: true) only if that role owns the tables. If it does not, have the owner create the publication and name it in the document:CREATE ROLE queen_cdc LOGIN REPLICATION PASSWORD '...'; GRANT SELECT ON public.orders TO queen_cdc; CREATE PUBLICATION queen_orders FOR TABLE public.orders; -- as the table owner"source": {"publication": "queen_orders", "managePublication": false, "tables": [...]}A publication you manage must publish whole tables: one with a row filter or a column list is refused, because the snapshot would not apply it.
-
Every source table needs a key the source can read from the WAL: a primary key, or a unique index set as replica identity (
ALTER TABLE t REPLICA IDENTITY USING INDEX t_key). A table without one is refused before the publication is touched: once published, PostgreSQL would reject the application’s own updates and deletes on it. -
Cap what a slot can hold:
max_slot_wal_keep_size. A slot keeps every WAL file the source has not confirmed, so a source that stays down for days fills the database’s disk. With a cap, PostgreSQL drops the slot instead, and the source stops withslot_invalidated. Alert onqueen_pg_source_slot_lag_byteswell below the cap. -
For a primary with a standby, the source creates its slot as a failover slot. Turn on PostgreSQL 17’s slot synchronization (
sync_replication_slotson the standby, with the prerequisites the PostgreSQL documentation lists, andsynchronized_standby_slotson the primary), and the slot survives a promotion. Without it a failover loses the slot, and the source stops withslot_lostrather than skip the gap.
The sink’s role needs INSERT, UPDATE and DELETE on its table (TRUNCATE for cdc with
truncates), and either CREATE on the database for the progress table’s schema or a progress table
created for it:
CREATE SCHEMA IF NOT EXISTS queen;
CREATE TABLE IF NOT EXISTS queen.sink_progress (
sink text NOT NULL,
partition_id bigint NOT NULL,
partition text NOT NULL,
last_offset bigint NOT NULL,
updated_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (sink, partition_id)
);Managed PostgreSQL. On Google Cloud SQL, logical decoding is the cloudsql.logical_decoding
flag (it restarts the instance), the postgres user can grant REPLICATION, and the smallest
max_slot_wal_keep_size Cloud SQL accepts is 100 GB. Cloud SQL has no superuser, so either the
source’s role owns the tables or the owner creates the publication.
Reference
The connector document. PUT /api/v1/connectors/<name> takes the whole document; a connector
name is lowercase letters, digits, - and _, 48 characters at most, starting with a letter or a
digit. Values outside a range are refused, never adjusted.
connection, for both kinds:
| Field | Default | |
|---|---|---|
host, port |
port 5432 |
Or url instead of the fields: postgres://user:password@host:port/database?sslmode=require |
database, user |
||
password |
Write-only: sealed with QUEEN_ENCRYPTION_KEY and never returned; a read says passwordSet: true. Omit it to keep the stored one, send "" to remove it |
|
sslMode |
prefer |
disable, prefer, require (encrypted, certificate not checked) or verify-full (certificate and host name checked) |
sslRootCert |
PEM, for verify-full with a private CA |
|
connectTimeoutMs |
10000 |
1,000 to 120,000; bounds the whole connection setup |
source:
| Field | Default | |
|---|---|---|
tables |
1 to 256 of {"table": "schema.name", "queue": "...", "partitionBy": ...} |
|
tables[].partitionBy |
"key" |
"key", "single", or a list of columns from the table’s key (any columns under REPLICA IDENTITY FULL) |
slot, publication |
queen_ + name |
- becomes _; lowercase, 63 characters at most |
managePublication |
true |
Create and update the publication to match tables |
snapshot |
"initial" |
"never" streams from now on, without copying existing rows |
snapshotChunkRows |
2000 |
100 to 100,000 |
maxBundleMessages |
1000 |
1 to 50,000 messages per Queen transaction |
maxBundleBytes |
4194304 |
64 KiB to 64 MiB per Queen transaction |
lingerMs |
20 |
0 to 5,000: how long to wait for more changes before a commit |
heartbeatSeconds |
10 |
1 to 3,600: a heartbeat written into the WAL when the tables are idle, so the slot keeps moving |
onTruncate |
"skip" |
"emit" sends t events, allowed only when every table is partitionBy: "single" |
sink:
| Field | Default | |
|---|---|---|
queue |
The queue to consume | |
table |
schema.name |
|
mode |
append, upsert, cdc or sql |
|
consumerGroup |
pg- + name |
|
subscriptionMode |
"all" |
"new" skips what the queue already holds when the group is first created |
key |
primary key | upsert and cdc only |
statement, params |
sql only; the statement is 64 KiB at most |
|
metadata |
append and upsert only, see above |
|
batch |
500 |
1 to 10,000 messages per pop and per PostgreSQL transaction |
workers |
1 |
1 to 32 per node |
leaseSeconds |
60 |
5 to 3,600; renewed while a batch is still running |
maxAttempts |
3 |
1 to 100 tries before a message is dead-lettered |
progressTable |
queen.sink_progress |
Must not be the target table |
createProgressTable |
true |
The API. Reads need read access, every write needs an admin token.
| Route | |
|---|---|
GET /api/v1/connectors |
Your connectors: each document, this node’s status and, for a source, its state (the pointer and the lease) |
GET /api/v1/connectors/<name> |
One of them, or 404 connector_not_found |
PUT /api/v1/connectors/<name> |
Create or replace; 200 with the stored document. 400 invalid_connector names the field, 400 encryption_required when the broker has no key, 409 kind_change when a source would become a sink or the other way round, 409 conflict when another write won |
DELETE /api/v1/connectors/<name> |
A sink is gone at once (its consumer group and progress rows stay). A source answers 202: its owner drops the slot, and the publication if it manages it, then the connector is gone. ?dropSlot=false removes it at once and leaves the slot to you |
POST /api/v1/connectors/<name>/resync |
A source only: drop the pointer and the slot and start over from a new snapshot, under a new epoch. 202 |
Settings, one set per node:
| Variable | Default | |
|---|---|---|
QUEEN_PG_CONNECTORS |
true |
Run connectors on this node |
QUEEN_PG_THREADS |
2 |
1 to 16 |
QUEEN_PG_RELOAD_MS |
5000 |
500 to 60,000: how often a node re-reads the documents |
QUEEN_PG_LEASE_TTL_MS |
10000 |
3,000 to 120,000: how long a dead source owner keeps its source |
QUEEN_PG_SHUTDOWN_GRACE_MS |
10000 |
0 to 120,000: how long a SIGTERM waits for the work in flight |
QUEEN_PG_ALLOW_PRIVATE_NETWORKS |
true, false behind the built-in proxy |
Whether a connector may reach loopback, private and link-local addresses |
Operating it
Status. GET /api/v1/connectors and the pg block of /status show each connector’s
phase: streaming, snapshot, connecting, standby (another node owns this source),
waiting_for_slot (another process holds the slot) or error for a source; running or error
for a sink; disabled for either. The status is this node’s view. For a source, state.lease.node
names the node that runs it and state.pointer shows how far it got, from any node.
Errors. An error carries a code and a sentence naming the fix. The ones that need you:
| Code | What happened | What to do |
|---|---|---|
version, wal_level |
PostgreSQL older than 17, or wal_level is not logical |
Upgrade, or set it and restart |
replica_identity, table_missing, partition_by |
A table has no usable key, does not exist, or partitionBy names a column the WAL does not carry |
Fix the table or the document |
publication |
The publication is missing a table, filters one, or the role cannot manage it | Fix the publication, or let the source manage it |
slot_lost, slot_invalidated |
The slot is gone (dropped, a failover without failover slots) or PostgreSQL dropped what it held (max_slot_wal_keep_size) |
POST .../resync |
slot_ahead, system_changed |
Someone else consumed the slot, or the database is a different cluster (a restore) | POST .../resync |
egress |
The host is on a private network and QUEEN_PG_ALLOW_PRIVATE_NETWORKS is false |
Use a public address, or allow it |
no_key, progress_table, statement, metadata_column |
The sink’s key has no unique index, its progress table cannot be created, its statement does not prepare, or a metadata column is not writable | Fix the table or the document |
unseal |
This node cannot decrypt the password | Give every node the same QUEEN_ENCRYPTION_KEY |
A resync re-reads every table and writes each row that still exists as an r event, under a new
epoch. Changes made while the slot was gone arrive as those rows’ current values, but a row deleted
in that window has no event at all, so a sink’s copy keeps it. After a resync, remove such rows from
the copies, or rebuild a copy from scratch.
Metrics on /metrics/prometheus, labelled connector (and tenant):
queen_pg_source_transactions_total, queen_pg_source_messages_total,
queen_pg_source_bundles_total, queen_pg_source_snapshot_rows_total,
queen_pg_source_slot_lag_bytes, queen_pg_source_pointer_lsn, queen_pg_sink_applied_total,
queen_pg_sink_skipped_total (redeliveries recognized and dropped), queen_pg_sink_dlq_total,
queen_pg_sink_batches_total, queen_pg_connector_errors_total and queen_pg_connector_owner
(1 on the node running the work).
Limits
- Exactly-once covers what the connectors write: each change lands in Queen once, and each message lands in the table once. What your own consumers do with a message is theirs; give them a Queen transaction for that.
- A source transaction larger than one Queen transaction is pushed in pieces, which are not atomic.
And a source transaction that spans several partitions reaches a sink’s table as several
PostgreSQL transactions, one per batch, because partitions are consumed independently. For the
table’s whole commit order, use
partitionBy: "single", at the throughput of one partition. - Schema changes are not captured. A new column shows up as a new field in the next event; a table added to a running source is streamed from then on but not copied (resync to copy it).
- A consumer that joins a queue halfway, with no snapshot behind it, can meet an update whose large
unchanged columns it has never seen.
REPLICA IDENTITY FULLon such tables removes the case. - Truncates are skipped by default: they cannot be ordered against rows spread over many partitions.
- Exactly-once holds while the replication slot lives. Losing it (a failover without failover slots, a slot dropped by hand, a restore from backup) needs a resync, and a resync cannot report the rows deleted while the slot was gone.
- A sink cannot apply what retention already deleted. Keep queue retention longer than the longest time a sink may stay down; retention is off by default.
- Connections use a password: no client certificates, no IAM tokens, no
verify-ca. - Not part of the Jepsen campaign. The connectors are tested against
PostgreSQL 17 and 18 by their own suites, and end to end (
test/pgconn) withkill -9, three-node takeovers, slot loss and resync, poison messages and a 200,000-row transaction killed halfway.
Next
- Charge a card once: exactly-once effects inside Queen, the same idea the connectors use, applied to your own workers.
- Consuming: consumer groups, leases and the dead-letter queue the sink relies on.
- Guarantees: what the broker underneath promises, and how it was tested.