Blog
Real-Time Notifications with PostgreSQL LISTEN/NOTIFY and SvelteKit: Complete Guide
Learn how to implement real-time notifications using PostgreSQL LISTEN/NOTIFY with SvelteKit. Step-by-step tutorial for building scalable real-time features without external services.
Why PostgreSQL LISTEN/NOTIFY for Real-Time Features?
In today’s web applications, users expect real-time updates — live chat messages, instant notifications, collaborative editing, and live dashboards. While many developers reach for external services like Socket.io, Pusher, or Firebase, there’s a powerful built-in solution that’s often overlooked: PostgreSQL’s LISTEN/NOTIFY.
PostgreSQL LISTEN/NOTIFY provides a native, scalable, and cost-effective way to implement real-time features without adding external dependencies. When combined with SvelteKit’s reactive architecture and Server-Sent Events (SSE), you get a powerful combination for building modern, real-time applications.
The Power of Native Database Notifications
PostgreSQL LISTEN/NOTIFY offers several compelling advantages:
- Zero External Dependencies: No need for Redis, Socket.io, or third-party services
- Database-Level Reliability: Notifications are transactional and ACID-compliant
- Scalable Architecture: Works seamlessly with your existing PostgreSQL setup
- Cost-Effective: No additional infrastructure costs
- Type Safety: Full TypeScript support with proper typing
Real-World Use Cases
This approach is perfect for:
- Live Chat Applications: Instant message delivery
- Collaborative Tools: Real-time document editing
- Dashboard Updates: Live metrics and analytics
- Notification Systems: User alerts and updates
- Game State Management: Multiplayer game synchronization
Understanding PostgreSQL LISTEN/NOTIFY
Before diving into implementation, let’s understand how PostgreSQL LISTEN/NOTIFY works.
How LISTEN/NOTIFY Works
PostgreSQL LISTEN/NOTIFY is a publish-subscribe messaging system built into the database:
- LISTEN: Clients subscribe to specific channels
- NOTIFY: Database sends messages to subscribed clients
- Channels: Named message queues (e.g., ‘user_notifications’, ‘chat_messages’)
- Payload: Optional JSON data with each notification
Architecture Overview
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ SvelteKit │ │ PostgreSQL │ │ Database │
│ Frontend │◄──►│ LISTEN/NOTIFY │◄──►│ Triggers │
│ (SSE) │ │ Channels │ │ & Functions │
└─────────────────┘ └─────────────────┘ └─────────────────┘ Key Benefits Over External Services
| Feature | External Services | PostgreSQL LISTEN/NOTIFY |
|---|---|---|
| Cost | Monthly fees | Free with existing DB |
| Latency | Network overhead | Direct database connection |
| Reliability | Service dependency | Database ACID guarantees |
| Scalability | Service limits | Database capacity |
| Setup Complexity | API integration | Native PostgreSQL |
Prerequisites and Setup
Let’s set up our development environment for building real-time notifications.
Required Dependencies
First, ensure you have the necessary packages:
# Install Postgres.js
bun install postgres Database Schema Setup
Create the foundation for our notification system:
-- Create notifications table
CREATE TABLE notifications (
id SERIAL PRIMARY KEY,
user_id INTEGER NOT NULL,
type VARCHAR(50) NOT NULL,
title VARCHAR(255) NOT NULL,
message TEXT,
data JSONB,
read_at TIMESTAMP,
created_at TIMESTAMP DEFAULT NOW()
);
-- Create notification channels table
CREATE TABLE notification_channels (
id SERIAL PRIMARY KEY,
name VARCHAR(100) UNIQUE NOT NULL,
description TEXT
);
-- Insert default channels
INSERT INTO notification_channels (name, description) VALUES
('user_notifications', 'User-specific notifications'),
('chat_messages', 'Real-time chat messages'),
('system_alerts', 'System-wide alerts'); SvelteKit Project Structure
Organize your project for real-time features:
src/
├── lib/
│ ├── server/
│ │ ├── db/
│ │ │ ├── index.ts # Database connection
│ │ │ ├── schema.ts # Database schema
│ │ │ └── notifications.ts # Notification queries
│ │ └── realtime/
│ │ ├── postgres.ts # PostgreSQL LISTEN/NOTIFY
│ │ └── sse.ts # Server-Sent Events handling
│ └── components/
│ ├── NotificationList.svelte
│ └── RealTimeChat.svelte
├── routes/
│ ├── api/
│ │ ├── notifications/
│ │ │ └── +server.ts
│ │ └── realtime/
│ │ └── +server.ts
│ └── notifications/
│ └── +page.svelte Step-by-Step Implementation
Let’s build a complete real-time notification system step by step.
Step 1: Database Connection Setup
Create a robust database connection with LISTEN/NOTIFY support:
// src/lib/server/db/index.ts
import postgres from 'postgres';
import { env } from '$env/dynamic/private';
class DatabaseManager {
private sql: postgres.Sql;
private listeners: Map<string, (payload: any) => void> = new Map();
constructor() {
this.sql = postgres(env.DATABASE_URL, {
max: 20,
idle_timeout: 30,
connect_timeout: 20,
onnotice: () => {} // Suppress notice messages
});
}
async listen(channel: string, callback: (payload: any) => void): Promise<void> {
try {
await this.sql.listen(channel, (payload) => {
const parsedPayload = payload ? JSON.parse(payload) : null;
callback(parsedPayload);
});
this.listeners.set(channel, callback);
} catch (error) {
console.error(`Error listening to channel ${channel}:`, error);
throw error;
}
}
async notify(channel: string, payload: any): Promise<void> {
try {
const payloadStr = JSON.stringify(payload);
await this.sql.notify(channel, payloadStr);
} catch (error) {
console.error(`Error notifying channel ${channel}:`, error);
throw error;
}
}
async close(): Promise<void> {
await this.sql.end();
}
// Get the sql instance for direct queries
get sql() {
return this.sql;
}
}
export const db = new DatabaseManager(); Step 2: Notification Service Layer
Create a service layer for managing notifications:
// src/lib/server/db/notifications.ts
import { db } from './index';
export interface Notification {
id: number;
user_id: number;
type: string;
title: string;
message?: string;
data?: any;
read_at?: Date;
created_at: Date;
}
export class NotificationService {
async createNotification(
notification: Omit<Notification, 'id' | 'created_at'>
): Promise<Notification> {
try {
const [newNotification] = await db.sql`
INSERT INTO notifications (user_id, type, title, message, data)
VALUES (${notification.user_id}, ${notification.type}, ${notification.title}, ${notification.message}, ${notification.data})
RETURNING *
`;
// Send real-time notification
await db.notify('user_notifications', {
user_id: notification.user_id,
notification: newNotification
});
return newNotification;
} catch (error) {
console.error('Error creating notification:', error);
throw error;
}
}
async getUserNotifications(userId: number, limit = 50): Promise<Notification[]> {
try {
return await db.sql`
SELECT * FROM notifications
WHERE user_id = ${userId}
ORDER BY created_at DESC
LIMIT ${limit}
`;
} catch (error) {
console.error('Error fetching notifications:', error);
throw error;
}
}
async markAsRead(notificationId: number, userId: number): Promise<void> {
try {
await db.sql`
UPDATE notifications
SET read_at = NOW()
WHERE id = ${notificationId} AND user_id = ${userId}
`;
} catch (error) {
console.error('Error marking notification as read:', error);
throw error;
}
}
}
export const notificationService = new NotificationService(); Step 3: Server-Sent Events (SSE) for Real-Time Updates
Create an SSE endpoint to bridge PostgreSQL notifications to the frontend:
// src/routes/api/realtime/+server.ts
import { json } from '@sveltejs/kit';
import type { RequestHandler } from './$types';
import { db } from '$lib/server/db';
export const GET: RequestHandler = async ({ request }) => {
const url = new URL(request.url);
const userId = url.searchParams.get('userId');
const channel = url.searchParams.get('channel') || 'user_notifications';
if (!userId) {
return json({ error: 'User ID required' }, { status: 400 });
}
// Set up Server-Sent Events
const stream = new ReadableStream({
start(controller) {
const sendEvent = (data: any) => {
const event = `data: ${JSON.stringify(data)}\n\n`;
controller.enqueue(new TextEncoder().encode(event));
};
// Listen to PostgreSQL notifications
db.listen(channel, (payload) => {
if (payload && payload.user_id === parseInt(userId)) {
sendEvent({
type: 'notification',
data: payload
});
}
});
// Send initial connection message
sendEvent({
type: 'connected',
channel,
userId
});
// Keep connection alive
const heartbeat = setInterval(() => {
sendEvent({ type: 'heartbeat' });
}, 30000);
// Cleanup on close
request.signal.addEventListener('abort', () => {
clearInterval(heartbeat);
controller.close();
});
}
});
return new Response(stream, {
headers: {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
Connection: 'keep-alive',
'Access-Control-Allow-Origin': '*'
}
});
}; Step 4: Frontend Real-Time Components
Create reactive Svelte components for real-time updates:
// src/lib/components/NotificationList.svelte
import { onMount, onDestroy } from 'svelte';
import type { Notification } from '$lib/server/db/notifications';
export let userId: number;
export let notifications: Notification[] = [];
let eventSource: EventSource;
let unreadCount = 0;
onMount(() => {
// Connect to real-time notifications
eventSource = new EventSource(`/api/realtime?userId=${userId}&channel=user_notifications`);
eventSource.onmessage = (event) => {
const data = JSON.parse(event.data);
if (data.type === 'notification') {
notifications = [data.data.notification, ...notifications];
unreadCount++;
// Show browser notification
if (Notification.permission === 'granted') {
new Notification(data.data.notification.title, {
body: data.data.notification.message,
icon: '/favicon.ico'
});
}
}
};
eventSource.onerror = (error) => {
console.error('EventSource error:', error);
// Implement reconnection logic
setTimeout(() => {
if (eventSource.readyState === EventSource.CLOSED) {
onMount();
}
}, 5000);
};
});
onDestroy(() => {
if (eventSource) {
eventSource.close();
}
});
async function markAsRead(notificationId: number) {
try {
await fetch(`/api/notifications/${notificationId}/read`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ userId })
});
notifications = notifications.map(n =>
n.id === notificationId ? { ...n, read_at: new Date() } : n
);
unreadCount = Math.max(0, unreadCount - 1);
} catch (error) {
console.error('Error marking notification as read:', error);
}
}
function requestNotificationPermission() {
if (Notification.permission === 'default') {
Notification.requestPermission();
}
} <!-- src/lib/components/NotificationList.svelte -->
<div class="notifications-container">
<div class="notifications-header">
<h2>
Notifications {#if unreadCount > 0}({unreadCount} unread){/if}
</h2>
<button on:click={requestNotificationPermission} class="permission-btn">
Enable Notifications
</button>
</div>
<div class="notifications-list">
{#each notifications as notification (notification.id)}
<div
class="notification-item {notification.read_at ? 'read' : 'unread'}"
on:click={() => markAsRead(notification.id)}
>
<div class="notification-header">
<span class="notification-type">{notification.type}</span>
<span class="notification-time">
{new Date(notification.created_at).toLocaleString()}
</span>
</div>
<h3 class="notification-title">{notification.title}</h3>
{#if notification.message}
<p class="notification-message">{notification.message}</p>
{/if}
{#if notification.data}
<div class="notification-data">
<pre>{JSON.stringify(notification.data, null, 2)}</pre>
</div>
{/if}
</div>
{:else}
<p class="no-notifications">No notifications yet</p>
{/each}
</div>
</div>
<style>
.notifications-container {
max-width: 600px;
margin: 0 auto;
padding: 20px;
}
.notifications-header {
display: flex;
justify-content: space-between;
align-items: center;
margin-bottom: 20px;
}
.permission-btn {
padding: 8px 16px;
background: #007bff;
color: white;
border: none;
border-radius: 4px;
cursor: pointer;
}
.notification-item {
padding: 15px;
margin-bottom: 10px;
border-radius: 8px;
cursor: pointer;
transition: background-color 0.2s;
}
.notification-item.unread {
background-color: #f8f9fa;
border-left: 4px solid #007bff;
}
.notification-item.read {
background-color: #ffffff;
border: 1px solid #e9ecef;
}
.notification-item:hover {
background-color: #e9ecef;
}
.notification-header {
display: flex;
justify-content: space-between;
margin-bottom: 8px;
}
.notification-type {
font-size: 0.8em;
color: #6c757d;
text-transform: uppercase;
}
.notification-time {
font-size: 0.8em;
color: #6c757d;
}
.notification-title {
margin: 0 0 8px 0;
font-size: 1.1em;
font-weight: 600;
}
.notification-message {
margin: 0;
color: #495057;
}
.notification-data {
margin-top: 10px;
padding: 10px;
background-color: #f8f9fa;
border-radius: 4px;
}
.notification-data pre {
margin: 0;
font-size: 0.8em;
color: #6c757d;
}
.no-notifications {
text-align: center;
color: #6c757d;
font-style: italic;
}
</style> Step 5: API Endpoints for Notifications
Create RESTful endpoints for notification management:
// src/routes/api/notifications/+server.ts
import { json } from '@sveltejs/kit';
import type { RequestHandler } from './$types';
import { notificationService } from '$lib/server/db/notifications';
export const GET: RequestHandler = async ({ url }) => {
const userId = url.searchParams.get('userId');
const limit = parseInt(url.searchParams.get('limit') || '50');
if (!userId) {
return json({ error: 'User ID required' }, { status: 400 });
}
try {
const notifications = await notificationService.getUserNotifications(parseInt(userId), limit);
return json(notifications);
} catch (error) {
console.error('Error fetching notifications:', error);
return json({ error: 'Failed to fetch notifications' }, { status: 500 });
}
};
export const POST: RequestHandler = async ({ request }) => {
try {
const body = await request.json();
const notification = await notificationService.createNotification(body);
return json(notification);
} catch (error) {
console.error('Error creating notification:', error);
return json({ error: 'Failed to create notification' }, { status: 500 });
}
}; // src/routes/api/notifications/[id]/read/+server.ts
import { json } from '@sveltejs/kit';
import type { RequestHandler } from './$types';
import { notificationService } from '$lib/server/db/notifications';
export const POST: RequestHandler = async ({ params, request }) => {
try {
const { userId } = await request.json();
const notificationId = parseInt(params.id);
await notificationService.markAsRead(notificationId, userId);
return json({ success: true });
} catch (error) {
console.error('Error marking notification as read:', error);
return json({ error: 'Failed to mark notification as read' }, { status: 500 });
}
}; Advanced Features and Optimizations
Database Triggers for Automatic Notifications
Create database triggers to automatically send notifications on data changes:
-- Function to handle notification triggers
CREATE OR REPLACE FUNCTION notify_user_change()
RETURNS TRIGGER AS $$
BEGIN
-- Send notification when user data changes
PERFORM pg_notify(
'user_notifications',
json_build_object(
'user_id', NEW.id,
'type', 'user_update',
'title', 'Profile Updated',
'message', 'Your profile has been updated',
'data', json_build_object(
'field', TG_ARGV[0],
'old_value', OLD,
'new_value', NEW
)
)::text
);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
-- Trigger for user table updates
CREATE TRIGGER user_notification_trigger
AFTER UPDATE ON users
FOR EACH ROW
EXECUTE FUNCTION notify_user_change('profile'); Channel Management and Filtering
Implement sophisticated channel management:
// src/lib/server/realtime/channel-manager.ts
export class ChannelManager {
private channels: Map<string, Set<string>> = new Map();
subscribe(userId: string, channel: string): void {
if (!this.channels.has(channel)) {
this.channels.set(channel, new Set());
}
this.channels.get(channel)!.add(userId);
}
unsubscribe(userId: string, channel: string): void {
const subscribers = this.channels.get(channel);
if (subscribers) {
subscribers.delete(userId);
if (subscribers.size === 0) {
this.channels.delete(channel);
}
}
}
getSubscribers(channel: string): string[] {
return Array.from(this.channels.get(channel) || []);
}
broadcast(channel: string, message: any): void {
const subscribers = this.getSubscribers(channel);
// Send to all subscribers
subscribers.forEach((userId) => {
// Implementation depends on your SSE setup
});
}
} Performance Optimizations
Implement connection pooling and caching:
// src/lib/server/db/connection-pool.ts
import postgres from 'postgres';
class ConnectionPool {
private sql: postgres.Sql;
private notificationClients: Map<string, any> = new Map();
constructor() {
this.sql = postgres(process.env.DATABASE_URL!, {
max: 20,
idle_timeout: 30,
connect_timeout: 20
});
}
async getNotificationClient(channel: string) {
if (!this.notificationClients.has(channel)) {
// Postgres.js handles connection pooling automatically
await this.sql.listen(channel, () => {});
this.notificationClients.set(channel, true);
}
return this.sql;
}
async cleanup() {
await this.sql.end();
}
} Real-World Implementation Examples
Live Chat Application
// Chat message notification
async function sendChatMessage(roomId: string, userId: number, message: string) {
const notification = await notificationService.createNotification({
user_id: userId,
type: 'chat_message',
title: 'New Message',
message: message,
data: {
room_id: roomId,
sender_id: userId,
timestamp: new Date()
}
});
// Notify all users in the room
await db.notify(`chat_room_${roomId}`, {
type: 'new_message',
message: notification
});
} Collaborative Document Editing
// Document change notification
async function notifyDocumentChange(documentId: string, userId: number, change: any) {
await db.notify(`document_${documentId}`, {
type: 'document_change',
user_id: userId,
change: change,
timestamp: new Date()
});
} Live Dashboard Updates
// Dashboard metrics update
async function updateDashboardMetrics(metrics: any) {
await db.notify('dashboard_metrics', {
type: 'metrics_update',
data: metrics,
timestamp: new Date()
});
} Testing and Debugging
Testing Real-Time Features
// src/lib/server/db/__tests__/notifications.test.ts
import { describe, it, expect, beforeEach, afterEach } from 'bun:test';
import { notificationService } from '../notifications';
import { db } from '../index';
describe('NotificationService', () => {
beforeEach(async () => {
// Setup test database
});
afterEach(async () => {
// Cleanup test data
});
it('should create and send real-time notification', async () => {
const notification = await notificationService.createNotification({
user_id: 1,
type: 'test',
title: 'Test Notification',
message: 'This is a test'
});
expect(notification).toBeDefined();
expect(notification.user_id).toBe(1);
});
}); Debugging Tools
// src/lib/server/realtime/debug.ts
export class NotificationDebugger {
static logNotification(channel: string, payload: any) {
console.log(`[${new Date().toISOString()}] ${channel}:`, payload);
}
static monitorChannels() {
// Monitor all active channels
const channels = ['user_notifications', 'chat_messages', 'system_alerts'];
channels.forEach((channel) => {
db.listen(channel, (payload) => {
this.logNotification(channel, payload);
});
});
}
} Production Deployment Considerations
Environment Configuration
// src/lib/server/config.ts
export const config = {
database: {
url: process.env.DATABASE_URL,
maxConnections: parseInt(process.env.DB_MAX_CONNECTIONS || '20'),
idleTimeout: parseInt(process.env.DB_IDLE_TIMEOUT || '30000')
},
notifications: {
maxRetries: parseInt(process.env.NOTIFICATION_MAX_RETRIES || '3'),
retryDelay: parseInt(process.env.NOTIFICATION_RETRY_DELAY || '1000'),
heartbeatInterval: parseInt(process.env.HEARTBEAT_INTERVAL || '30000')
},
security: {
corsOrigins: process.env.CORS_ORIGINS?.split(',') || ['http://localhost:5173'],
maxPayloadSize: parseInt(process.env.MAX_PAYLOAD_SIZE || '1024')
}
}; Monitoring and Health Checks
// src/routes/api/health/+server.ts
import { json } from '@sveltejs/kit';
import type { RequestHandler } from './$types';
import { db } from '$lib/server/db';
export const GET: RequestHandler = async () => {
try {
// Test database connection using Postgres.js
const result = await db.sql`SELECT 1`;
return json({
status: 'healthy',
database: 'connected',
timestamp: new Date().toISOString()
});
} catch (error) {
return json(
{
status: 'unhealthy',
error: error.message,
timestamp: new Date().toISOString()
},
{ status: 500 }
);
}
}; Best Practices and Security
Security Considerations
- Input Validation: Always validate notification payloads
- Rate Limiting: Implement rate limiting for notification creation
- Authentication: Verify user permissions before sending notifications
- Payload Size: Limit notification payload size to prevent abuse
// src/lib/server/middleware/notification-validation.ts
export function validateNotificationPayload(payload: any): boolean {
const maxSize = 1024; // 1KB limit
const payloadStr = JSON.stringify(payload);
if (payloadStr.length > maxSize) {
return false;
}
// Validate required fields
if (!payload.user_id || !payload.type || !payload.title) {
return false;
}
return true;
} Performance Best Practices
- Connection Pooling: Postgres.js handles connection pooling automatically
- Batch Notifications: Group multiple notifications when possible
- Channel Optimization: Use specific channels for different notification types
- Cleanup: Properly close connections and clean up resources
Error Handling and Resilience
// src/lib/server/realtime/error-handler.ts
export class NotificationErrorHandler {
static async withRetry<T>(
operation: () => Promise<T>,
maxRetries: number = 3,
delay: number = 1000
): Promise<T> {
let lastError: Error;
for (let i = 0; i < maxRetries; i++) {
try {
return await operation();
} catch (error) {
lastError = error;
if (i < maxRetries - 1) {
await new Promise((resolve) => setTimeout(resolve, delay * (i + 1)));
}
}
}
throw lastError!;
}
} Troubleshooting Common Issues
Connection Issues
Problem: SSE connections dropping frequently Solution: Implement automatic reconnection with exponential backoff
class ReconnectingEventSource {
private url: string;
private eventSource: EventSource | null = null;
private reconnectAttempts = 0;
private maxReconnectAttempts = 5;
constructor(url: string) {
this.url = url;
this.connect();
}
private connect() {
this.eventSource = new EventSource(this.url);
this.eventSource.onopen = () => {
this.reconnectAttempts = 0;
};
this.eventSource.onerror = () => {
this.eventSource?.close();
this.scheduleReconnect();
};
}
private scheduleReconnect() {
if (this.reconnectAttempts < this.maxReconnectAttempts) {
const delay = Math.pow(2, this.reconnectAttempts) * 1000;
setTimeout(() => {
this.reconnectAttempts++;
this.connect();
}, delay);
}
}
} Database Performance Issues
Problem: High database connection usage Solution: Postgres.js handles connection pooling automatically, but you can configure it
Memory Leaks
Problem: Memory usage growing over time Solution: Proper cleanup of event listeners and database connections
Conclusion
PostgreSQL LISTEN/NOTIFY with SvelteKit and Postgres.js provides a powerful, scalable, and cost-effective solution for real-time features. By leveraging the database’s built-in notification system with the modern Postgres.js client, you can build sophisticated real-time applications without external dependencies.
The key benefits of this approach include:
- Zero External Dependencies: Everything runs on your existing PostgreSQL database
- ACID Compliance: Notifications are transactional and reliable
- Scalable Architecture: Grows with your database capacity
- Cost Effective: No additional infrastructure costs
- Type Safety: Full TypeScript support throughout the stack
- Modern Client: Postgres.js provides a clean, modern API
Whether you’re building a chat application, collaborative tools, or live dashboards, PostgreSQL LISTEN/NOTIFY with SvelteKit and Postgres.js provides the foundation for modern, real-time web applications.
Start with the basic implementation and gradually add advanced features like channel management, database triggers, and sophisticated error handling. The modular approach allows you to scale your real-time features as your application grows.
Happy coding with real-time notifications! 🚀