---
title: "Postgres source and sink"
description: "Stream PostgreSQL 17+ tables into queues and queues into tables from inside the broker: setup, the change events, how a JSON message maps onto columns, what exactly-once means here, and where it stops."
---

> Queen MQ documentation, for AI agents
> Complete self-contained summary of Queen MQ: https://queenmq.com/llms-brief.txt
> Fetch that first when the question is about the product rather than about this page.
> Index of all pages: https://queenmq.com/llms.txt

# Postgres source and sink

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:

```bash
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:latest
```

A table with a nested JSON column, and two rows in it:

```bash
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);
SQL
```

Create the source:

```bash
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/null
```

Every node re-reads the connector documents every 5 seconds, so give it a few seconds:

```bash
curl -s localhost:6632/api/v1/connectors/orders-src | jq -c '.status | {phase, pointerLsn, messages}'
```

```text
{"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`:

```bash
curl -s 'localhost:6632/api/v1/pop/queue/orders?batch=10&partitions=10&consumerGroup=reader&subscriptionMode=all&autoAck=true' \
  | jq -c '.messages[] | {partition, data}'
```

```text
{"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.

```bash
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}'
```

```text
{"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:

```bash
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/null
```

A few seconds later the copy holds what the table holds:

```bash
docker exec pg psql -U postgres -c 'SELECT id, customer, status, total FROM orders_copy ORDER BY id;'
```

```text
 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:

- `bool` becomes `true`/`false`. Integers and `numeric` become 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). `NaN` and `Infinity` become strings.
- `json` and `jsonb` are embedded as JSON, nested objects and arrays included.
- `timestamptz` becomes `"2026-10-03T06:14:39.365472+00:00"` and `timestamp`
  `"2026-10-03T06:14:39.365472"`; both connections run in UTC.
- Arrays of the built-in types become JSON arrays. `bytea` becomes 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 `DEFAULT` on an
  insert. It is never set to NULL unless the message says `null`.
- PostgreSQL converts the values itself (`jsonb_populate_record`), from the exact characters the
  producer wrote: `"12.50"` into a `numeric(12,2)` is 12.50, a nested object into a `jsonb` column
  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.city` does not reach a `city` column by itself. Keep
  the subtree in a `jsonb` column 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:

```bash
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/null
```

A few seconds later:

```bash
docker exec pg psql -U postgres -c 'SELECT id, customer, lines, total, customer_tier, city, queen_offset FROM invoices;'
```

```text
 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:

```json
{
  "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.

1. Set `wal_level = logical` and restart the server. Raise `max_replication_slots` and
   `max_wal_senders` if other consumers already use slots (one slot per source).
2. Give the source a role with `LOGIN REPLICATION` and `SELECT` on 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:

```sql
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
```

```json
"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.
3. 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.
4. 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 with `slot_invalidated`. Alert on
   `queen_pg_source_slot_lag_bytes` well below the cap.
5. For a primary with a standby, the source creates its slot as a failover slot. Turn on
   PostgreSQL 17's slot synchronization (`sync_replication_slots` on the standby, with the
   prerequisites the PostgreSQL documentation lists, and `synchronized_standby_slots` on the
   primary), and the slot survives a promotion. Without it a failover loses the slot, and the source
   stops with `slot_lost` rather 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:

```sql
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](/guides/exactly-once/) 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 FULL` on 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](/concepts/guarantees/). The connectors are tested against
  PostgreSQL 17 and 18 by their own suites, and end to end (`test/pgconn`) with `kill -9`,
  three-node takeovers, slot loss and resync, poison messages and a 200,000-row transaction killed
  halfway.

## Next

- [Charge a card once](/guides/exactly-once/): exactly-once effects inside Queen, the same idea the
  connectors use, applied to your own workers.
- [Consuming](/concepts/consuming/): consumer groups, leases and the dead-letter queue the sink
  relies on.
- [Guarantees](/concepts/guarantees/): what the broker underneath promises, and how it was tested.

Source: https://queenmq.com/guides/postgres/index.mdx
