Using Postgres LISTEN/NOTIFY for SSE Permalink to this section
Part of Kafka & NATS as SSE Event Sources, under Backend Stream Generation & Connection Management.
If the data already lives in PostgreSQL, adding a broker just to push changes to browsers is often unnecessary. LISTEN and NOTIFY let every SSE node learn about changes the moment a transaction commits. Used naively — putting the event itself in the notification — they lose events whenever a listener reconnects and break on payloads over 8,000 bytes. Used as a wake-up signal over an events table, they give a durable, replayable, gap-free SSE source with no extra infrastructure. This guide builds the second version.
Symptom & Developer Intent Permalink to this section
- Events are occasionally missing in the browser, especially around deploys or database failovers.
NOTIFYfails with “payload string too long” for larger events.- Listeners stop receiving anything after running fine for hours, with no error.
- Through PgBouncer,
LISTENappears to work but notifications never arrive. - Each SSE connection opens its own database connection to listen, and the database runs out of connections.
The intent is an SSE source backed only by Postgres that never loses an event, supports Last-Event-ID replay, and uses one database connection per node for listening.
Root Cause Analysis Permalink to this section
NOTIFY delivers to sessions that are listening at the moment the notifying transaction commits. There is no queue for absent listeners. A node whose listening connection dropped — a network blip, a failover, a restart — misses every notification sent while it was reconnecting.
The payload limit is 8,000 bytes by default, and every notification is also held in a shared queue (8 GB by default) until every listener has consumed it — a listener that stops reading can fill it and eventually block NOTIFY for everyone. Transaction-mode connection poolers break LISTEN because it is session state: the pooler hands your session to someone else after the transaction, and the listen registration goes with it.
Step-by-Step Resolution Permalink to this section
Step 1 — Store events in a table with a sequence Permalink to this section
CREATE TABLE sse_events (
seq bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
audience text NOT NULL, -- 'user:42', 'room:7', 'public'
type text NOT NULL,
body jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX ON sse_events (audience, seq);
Insert events in the same transaction as the change they describe, so an event exists if and only if the change committed.
Step 2 — Notify with the sequence only Permalink to this section
CREATE FUNCTION sse_notify() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
PERFORM pg_notify('sse', NEW.seq::text);
RETURN NULL;
END $$;
CREATE TRIGGER sse_events_notify AFTER INSERT ON sse_events
FOR EACH ROW EXECUTE FUNCTION sse_notify();
Payloads stay tiny, so the 8,000-byte limit never matters. Notifications from one transaction are delivered together at commit; identical notifications within a transaction are collapsed, which is harmless because the listener reads by range.
Step 3 — One dedicated listener per node, reading by range Permalink to this section
// listener.js — node-postgres, direct connection (not through a transaction-mode pooler).
import pg from 'pg';
let lastSeq = 0;
let reading = false, again = false;
async function catchUp(pool) {
if (reading) { again = true; return; } // coalesce bursts of notifications
reading = true;
do {
again = false;
const { rows } = await pool.query(
'SELECT seq, audience, type, body FROM sse_events WHERE seq > $1 ORDER BY seq LIMIT 1000',
[lastSeq]);
for (const r of rows) {
lastSeq = Number(r.seq);
hub.deliver(r.audience, `id: ${r.seq}\nevent: ${r.type}\ndata: ${JSON.stringify(r.body)}\n\n`, lastSeq);
}
if (rows.length === 1000) again = true;
} while (again);
reading = false;
}
export async function startListener(pool) {
const client = new pg.Client({ connectionString: process.env.DATABASE_DIRECT_URL });
client.on('notification', () => catchUp(pool));
client.on('error', () => setTimeout(() => startListener(pool), 1000)); // reconnect
await client.connect();
const { rows } = await client.query('SELECT coalesce(max(seq), 0) AS s FROM sse_events');
if (!lastSeq) lastSeq = Number(rows[0].s); // first start: live tail from now
await client.query('LISTEN sse');
await catchUp(pool); // after (re)connect: read the gap
setInterval(() => catchUp(pool), 5000).unref(); // safety net for lost notifications
}
The periodic catchUp is a belt-and-braces measure: even if a notification is lost for a reason not yet imagined, events arrive within five seconds.
Step 4 — Serve Last-Event-ID replay from the same table Permalink to this section
app.get('/api/stream', requireUser, async (req, res) => {
openStream(res);
const after = Number(req.get('Last-Event-ID') ?? 0);
const audiences = [`user:${req.user.id}`, 'public'];
const unregister = hub.register(audiences, res, { buffer: true }); // live first
const { rows } = await pool.query(
'SELECT seq, type, body FROM sse_events WHERE audience = ANY($1) AND seq > $2 ORDER BY seq LIMIT 500',
[audiences, after]);
let high = after;
for (const r of rows) { res.write(`id: ${r.seq}\nevent: ${r.type}\ndata: ${JSON.stringify(r.body)}\n\n`); high = Number(r.seq); }
hub.flushBuffered(res, high); // live frames with seq > high
req.on('close', unregister);
});
Step 5 — Prune the table Permalink to this section
Delete events older than your replay window in small batches, or partition the table by day and drop old partitions:
DELETE FROM sse_events WHERE seq IN (
SELECT seq FROM sse_events WHERE created_at < now() - interval '3 days' LIMIT 10000);
Clients presenting a sequence older than the oldest retained row receive resync.
Validation & Monitoring Permalink to this section
# Watch notifications arrive in psql while inserting from another session.
psql "$DATABASE_DIRECT_URL" -c 'LISTEN sse' -c 'SELECT pg_sleep(30)'
psql "$DATABASE_URL" -c "INSERT INTO sse_events (audience, type, body) VALUES ('public','ping','{}')"
# Queue health: fraction of the notification queue in use (should be ~0).
psql "$DATABASE_URL" -c 'SELECT pg_notification_queue_usage()'
Kill a node’s listening connection with pg_terminate_backend during a burst of inserts and confirm clients receive every event once the listener reconnects. Monitor pg_notification_queue_usage(), the node’s lastSeq lag behind max(seq), and catch-up query duration.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Can LISTEN/NOTIFY handle high event rates?
Thousands of notifications per second are fine; each one is small and the listener reads ranges. For tens of thousands per second, a log-structured broker such as Kafka or JetStream is a better fit.
Why not put the whole event in the notification?
Payloads are limited to 8,000 bytes, and a notification sent while a listener is reconnecting is lost. A sequence number plus a table read has neither problem.
Does LISTEN work through PgBouncer?
Only in session pooling mode. In transaction mode the session is shared, so the listen registration is not tied to your client. Give the listener a direct connection.
Is logical replication a better option?
For capturing every row change across many tables, logical decoding or a change-data-capture tool is more complete. For application-defined events written deliberately, a table plus NOTIFY is simpler and entirely sufficient.