Postgres LISTEN/NOTIFY + KMP SharedFlow: No WebSockets
Meta description: Map PostgreSQL LISTEN/NOTIFY channels to KMP SharedFlows for real-time mobile push. Zero broker overhead, sub-100ms latency, with a reconnect strategy for flaky networks.
Tags: kotlin kmp multiplatform android backend
TL;DR
PostgreSQL’s built-in LISTEN/NOTIFY mechanism can drive real-time mobile updates without WebSockets, polling, or a message broker. Map each pg_notify channel to a Kotlin Multiplatform SharedFlow, manage the JDBC connection on a supervised coroutine, and implement exponential backoff reconnection. The result: sub-100ms push latency at a fraction of the infrastructure cost.
Sub-100ms push from a database you already run
Most mobile backends achieve real-time push by introducing a dedicated message broker — Redis Pub/Sub, a managed WebSocket service, or a streaming platform — before checking whether the existing Postgres instance can meet the requirement. For the majority of product workloads, it can.
Postgres has shipped LISTEN/NOTIFY since version 6.4. It’s transactional, has zero external dependencies, and is already running. A backend process calls pg_notify from a trigger or application query; any client holding a LISTEN connection receives the payload asynchronously, outside the normal request/response cycle. This is the primitive worth reaching for first.
How LISTEN/NOTIFY actually works
Consider a trigger that fires on every insert into an events table:
CREATE OR REPLACE FUNCTION notify_event_insert()
RETURNS trigger AS $$
BEGIN
PERFORM pg_notify(
'event_channel',
row_to_json(NEW)::text
);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER on_event_insert
AFTER INSERT ON events
FOR EACH ROW EXECUTE FUNCTION notify_event_insert();
Any client subscribed to event_channel receives the serialized row the moment the transaction commits. No polling interval, no separate write path.
Mapping pg_notify channels to KMP SharedFlows
On the Kotlin Multiplatform side, a SharedFlow fits naturally. It’s hot, it multicasts, and it survives subscriber churn without dropping the upstream connection.
// commonMain
class PgNotifyChannel(private val channelName: String) {
private val _events = MutableSharedFlow<String>(
replay = 0,
extraBufferCapacity = 64,
onBufferOverflow = BufferOverflow.DROP_OLDEST
)
val events: SharedFlow<String> = _events.asSharedFlow()
suspend fun emit(payload: String) = _events.emit(payload)
}
DROP_OLDEST is a deliberate choice: a slow subscriber should never back-pressure the notification pipeline. Size extraBufferCapacity to your expected burst volume, not the average.
On the backend (JVM), a dedicated coroutine holds the JDBC connection in LISTEN mode and forwards payloads into the flow:
// jvmMain / backend service
fun listenOnChannel(scope: CoroutineScope, channel: PgNotifyChannel) {
scope.launch(Dispatchers.IO + SupervisorJob()) {
val conn = dataSource.connection
conn.createStatement().execute("LISTEN ${channel.channelName}")
val pgConn = conn.unwrap(PGConnection::class.java)
while (isActive) {
val notifications = pgConn.getNotifications(1000) ?: continue
notifications.forEach { channel.emit(it.parameter) }
}
}
}
SupervisorJob() ensures a failure in this coroutine does not cascade to the parent scope — critical for a long-lived connection managed alongside other application work.
Connection lifecycle: treat it as a first-class resource
A persistent JDBC connection holding a LISTEN is stateful. Network interruptions, Postgres restarts, and load balancer idle timeouts will kill it silently. Blind retry loops cause thundering-herd problems on outage recovery. Use supervised reconnection instead.
| Strategy | Behaviour | Risk |
|---|---|---|
| Immediate retry | Reconnects instantly | Thundering herd on outage |
| Fixed interval (e.g. 5s) | Simple, predictable | Slow recovery under load |
| Exponential backoff + jitter | Scales gracefully under load | Slightly more complex |
| Circuit breaker + backoff | Production-grade, observable | Most implementation effort |
Exponential backoff with full jitter is the minimum bar for any mobile-facing system:
suspend fun withReconnect(block: suspend () -> Unit) {
var attempt = 0
while (true) {
try {
block()
attempt = 0
} catch (e: Exception) {
val delay = minOf(30_000L, (500L shl attempt)) + Random.nextLong(500)
delay(delay)
attempt++
}
}
}
On Android and iOS, gate reconnection attempts on device connectivity state. There’s no point cycling through backoff attempts while the device is offline — surface NetworkCallback (Android) or NWPathMonitor (iOS) via expect/actual declarations and skip reconnect until the path is satisfied.
Latency and cost comparison
The numbers below reflect single-node, LAN-condition benchmarks typical of a staging environment; production latency varies with connection count, payload size, and network topology. Read them as order-of-magnitude comparisons, not hard numbers.
| Approach | Infra added | Typical latency | Ops overhead |
|---|---|---|---|
| Polling (5s interval) | None | 0–5s | Minimal |
| WebSocket server | Load balancer, sticky sessions | ~50–150ms | Medium |
| Redis Pub/Sub | Redis cluster | ~5–20ms | High |
| Postgres LISTEN/NOTIFY | None | ~50–150ms | Minimal |
LISTEN/NOTIFY is latency-competitive with a WebSocket layer while adding nothing to your infrastructure. Redis Pub/Sub wins on raw latency at scale, but that tradeoff only matters once listener counts start stressing Postgres connection limits — typically in the tens of thousands.
Conclusion
Most teams reach for a message broker before they’ve tested whether Postgres already does the job. That’s usually a mistake. LISTEN/NOTIFY has been in Postgres since 6.4, costs nothing to run, and delivers latency that mobile users won’t notice. Paired with SharedFlow and a proper reconnection strategy, it handles real-time push for most product workloads without adding a new infrastructure tier to operate and monitor.
The three things worth getting right:
- Benchmark
pg_notifybefore deploying a broker. It typically meets latency requirements below tens of thousands of concurrent listeners. - Supervise the
LISTENconnection withSupervisorJob(), instrument reconnect attempts with metrics, and gate reconnection on device connectivity state on mobile. - Use
SharedFlowwithDROP_OLDEST. Size the buffer to your expected burst volume and validate that assumption under load.