Building a Real-Time Notification System with FastAPI and Redis Streams

Explore the architecture and implementation of a robust notification system using FastAPI and Redis Streams.

Blog cover image
2101050's avatar
2101050
54 views

Building a Real-Time Notification System with FastAPI and Redis Streams

In today's digital landscape, real-time notifications are crucial for enhancing user engagement and experience. This blog post explores how to build a robust notification system using FastAPI, Redis Streams, and Server-Sent Events (SSE). We'll address complex scenarios such as multi-Pod consumption, message deduplication, weak network conditions, offline users, and more.

System Architecture Overview

The system leverages Redis Streams for message queuing and FastAPI for handling HTTP requests and SSE connections. Redis Streams provide a powerful mechanism for managing real-time data flows, while FastAPI offers a modern, asynchronous web framework for building APIs.

Key Components

  1. Message Producer (FastAPI Endpoint)
  2. SSE Push Service (FastAPI Endpoint)
  3. Redis Streams for Message Queuing

Message Producer

The message producer is responsible for sending notifications to Redis Streams. Each message is tagged with a unique identifier and timestamp to ensure deduplication and proper ordering.

Python
1@app.post("/notify/{user_id}") 2async def send_notification(user_id: str, message: dict): 3 message.update({ 4 "user_id": user_id, 5 "timestamp": int(time.time() * 1000), 6 "message_id": str(uuid.uuid4()) 7 }) 8 stream_id = redis_client.xadd( 9 f"user_notifications:{user_id}", 10 message, 11 maxlen=10000, 12 id="*" 13 ) 14 return {"status": "success", "stream_id": stream_id} 15

Example

A user action, such as a new message or alert, triggers this endpoint, which then writes the notification to the Redis Stream associated with the user.

SSE Push Service

The SSE push service streams notifications to clients in real-time. It uses Redis consumer groups to manage message distribution and ensure each message is processed by only one consumer.

Python
1async def event_generator(user_id: str, last_id: str = None): 2 group_name = "notifications_group" 3 consumer_name = f"consumer_{os.getpid()}" 4 5 try: 6 redis_client.xgroup_create( 7 f"user_notifications:{user_id}", 8 group_name, 9 id="0", 10 mkstream=True 11 ) 12 except redis.exceptions.ResponseError as e: 13 if "BUSYGROUP" not in str(e): 14 raise 15 16 while True: 17 try: 18 messages = redis_client.xreadgroup( 19 group_name, consumer_name, 20 {f"user_notifications:{user_id}": ">"}, 21 count=1, block=5000, noack=False 22 ) 23 24 if messages: 25 stream, message_list = messages[0] 26 for message_id, message_data in message_list: 27 yield f"data: {json.dumps(message_data)}\n\n" 28 redis_client.xack( 29 f"user_notifications:{user_id}", 30 group_name, 31 message_id 32 ) 33 else: 34 yield ": heartbeat\n\n" 35 36 except (redis.exceptions.ConnectionError, asyncio.CancelledError): 37 await asyncio.sleep(1) 38 continue 39 40@app.get("/notifications/{user_id}") 41async def stream_notifications(user_id: str, request: Request): 42 return StreamingResponse( 43 event_generator(user_id), 44 media_type="text/event-stream" 45 ) 46

Example

When a client connects to the /notifications/{user_id} endpoint, they receive a continuous stream of notifications. The server sends periodic heartbeats to maintain the connection during idle periods.

Handling Complex Scenarios

Multi-Pod Consumption

Redis consumer groups ensure that each message is processed by only one consumer, even in a multi-Pod environment. This prevents duplicate processing and allows for efficient load balancing.

Pod Failures and Recovery

If a Pod becomes unresponsive, messages it holds remain in a "pending" state. To address this, implement a recovery mechanism that reassigns pending messages to active consumers.

Python
1async def recover_pending_messages(): 2 while True: 3 try: 4 user_streams = redis_client.keys("user_notifications:*") 5 for stream in user_streams: 6 pending = redis_client.xpending(...) 7 # Logic to reassign pending messages 8 except Exception as e: 9 # Handle exceptions 10 pass 11

Pod Scaling

Redis Streams and consumer groups handle dynamic scaling well. New consumers can join the group and start processing messages without disrupting the system.

Conclusion

By combining FastAPI, Redis Streams, and SSE, you can build a scalable and efficient real-time notification system. This architecture handles complex scenarios like multi-Pod environments, message deduplication, and offline user management, ensuring a seamless user experience.

Recommended Articles

Discover more articles you might find interesting

Implementing LangGraph REST API with FastAPI
Technical Insights

Implementing LangGraph REST API with FastAPI

This guide provides a comprehensive implementation plan for building a LangGraph REST API using FastAPI, covering environment setup, agent definitions, endpoint creation, testing, and deployment.

2101050
Jun 18
153
Read More
DeepSite v2 Practical Guide
Technical Insights

DeepSite v2 Practical Guide

A comprehensive guide to DeepSite v2, covering its features, installation, and advanced workflows.

2101050
Jun 21
112
Read More
Fastify OpenTelemetry: Logging, Metrics, and Tracing in Practice
Technical Insights

Fastify OpenTelemetry: Logging, Metrics, and Tracing in Practice

Learn how to implement logging, metrics, and tracing in Fastify using OpenTelemetry.

2101050
Jul 11
106
Read More
Creating Diverse Logo Designs with Flux Model and ComfyUI
Technical Insights

Creating Diverse Logo Designs with Flux Model and ComfyUI

Learn to leverage the Flux model and ComfyUI for unique logo designs through effective prompts and examples.

2101050
Jan 10
93
Read More
Formatting Dates in TypeScript to UTC
Technical Insights

Formatting Dates in TypeScript to UTC

A guide on how to format dates in TypeScript to the specific format YYYY-MM-DDTHH:mm:ss+00:00.

2101050
Dec 19
83
Read More
Implementing a Custom Chat Model with LangChain
Technical Insights

Implementing a Custom Chat Model with LangChain

This guide provides a comprehensive blueprint for creating a custom chat model by subclassing LangChain's BaseChatModel, including configuration, method overrides, and error handling.

2101050
Jun 17
78
Read More