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.
  • NOTIFY fails with “payload string too long” for larger events.
  • Listeners stop receiving anything after running fine for hours, with no error.
  • Through PgBouncer, LISTEN appears 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.

Notifications lost during a listener reconnect Timeline of sixty seconds in which a node's listening connection drops for eight seconds; notifications sent during the gap are lost unless the node reads the table after reconnecting. Notifications lost during a listener reconnect Listener Notifies listening listening 3 lost 0 12 24 36 48 60 seconds drop catch up
The table closes the gap. On reconnect the node reads everything after the last sequence it saw, so notifications are an optimisation, not the source of truth.

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.

Notification as a wake-up, table as the truth Flow from a transaction inserting into sse_events, a trigger sending NOTIFY with the sequence, the node's listener waking up, reading rows after its last sequence, and delivering frames to local connections. Notification as a wake-up, table as the truth INSERT event in the business txn commit NOTIFY sse seq only wake Listener 1 per node catch up SELECT seq > last range read deliver Local hub frames to clients
The listener never trusts the notification's content. It reads the range from the table, so a lost or collapsed notification only delays delivery.

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.

LISTEN/NOTIFY pitfalls and their fixes Matrix of four LISTEN/NOTIFY pitfalls, their symptoms and the fixes used in this guide. LISTEN/NOTIFY pitfalls and their fixes Pitfall Symptom Fix No queue for absent listeners lost events read range after reconnect 8,000-byte payload limit NOTIFY errors send seq only Transaction-mode pooler silent listener direct connection Listener per SSE client connection exhaustion one per node
All four pitfalls vanish when the notification carries only a sequence and the table is the source of truth.

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.