DEV Community

Timevolt
Timevolt

Posted on

Building a Real-Time Notification System at Scale: The Gandalf Approach

The Quest Begins (The "Why")

I still remember the first time our notification service started to feel like a dragon hoarding gold. We were pushing out real‑time alerts for a fledgling social app—likes, comments, follow requests—everything had to arrive within a second or users would think the app was broken. At first we just fired off a database write for every event, assuming our modest PostgreSQL could keep up.

Spoiler: it couldn’t.

During a peak hour, the write queue grew longer than a line at Comic‑Con, latency spiked to 12 seconds, and our monitoring lit up like a Christmas tree. Users started complaining that they missed important updates, and the support team was drowning in tickets. We needed a way to throttle the flood without losing the real‑time feel that made our product special.

The quest was clear: design a system that could smooth out bursts, protect our backend, and still deliver notifications as fast as possible.

The Revelation (The Insight)

After a few sleepless nights (and way too much coffee), the breakthrough hit me while I was re‑watching The Lord of the Rings: Gandalf doesn’t try to stop the flood; he guides the water.

In other words, instead of trying to block every notification outright, we let a controlled amount through, store the excess for later, and release it smoothly when the system can handle it. The concrete shape of that idea is a token bucket rate limiter combined with a sliding‑window counter backed by Redis, with a tiny in‑memory fallback for sub‑millisecond bursts.

Here’s why this beats the usual alternatives:

Approach Pros Cons
Global DB write per event Simple, strong consistency Overwhelms DB, high latency
Fixed‑window counter (e.g., INCR + EXPIRE) Easy to implement Allows burst at window edge → thundering herd
Pure token bucket (local) Low latency, smooth No shared state → limits per instance, not per service
Token bucket + Redis + local fallback ✅ Shared, smooth rate
✅ Burst absorption
✅ Fast path via local cache
✅ Graceful degradation if Redis hiccups
Slightly more moving parts (but still simple)

The magic lies in the dual‑layer:

  1. Fast path – each service instance keeps a tiny in‑memory bucket (say, 10 tokens). If a token is available, we send the notification immediately and decrement locally.
  2. Slow path – when the local bucket is empty, we ask Redis for a token via a Lua script that implements a true token bucket (refilling rate r tokens/second, capacity b). If Redis grants a token, we forward the event; otherwise we drop it (or optionally persist to a dead‑letter queue for retry).

The ASCII diagram below shows the flow:

+----------------+       +-------------------+       +----------------+
|   Event Source | --->  |  Local Token Bucket| --->  |   Redis Lua    |
| (web/worker)   |       |   (in‑memory)     |       |   Token Bucket |
+----------------+       +-------------------+       +----------------+
          |                       |                         |
          |   token? (yes)        |   token? (yes)          |   token? (no)
          v                       v                         v
   Send Notification   Send Notification   Drop / Queue for Retry
Enter fullscreen mode Exit fullscreen mode

The critical insight is that the local bucket absorbs micro‑bursts (think of a sudden spike when a post goes viral) without hitting Redis at all, while Redis guarantees a globally fair rate over longer periods. If Redis becomes unreachable, the local bucket continues to work for a short grace period, preventing a total outage—something a pure Redis limiter can’t do.

Wielding the Power (Code & Examples)

The Struggle: Naïve per‑event DB write

# pseudo‑code – what we started with
def send_notification(user_id, payload):
    # Every notification writes a row to the notifications table
    db.insert("notifications", {
        "user_id": user_id,
        "payload": payload,
        "created_at": now()
    })
    # Then we push to a websocket/fan‑out service
    push_to_ws(user_id, payload)
Enter fullscreen mode Exit fullscreen mode

Under load, db.insert became the bottleneck. Latency climbed, and we wasted CPU cycles on retries.

The Victory: Token bucket with local + Redis fallback

First, we set up a tiny in‑memory bucket per process (using threading.local or a simple class).

import time
import threading

class LocalBucket:
    def __init__(self, capacity, fill_rate):
        self.capacity = capacity          # max tokens
        self.fill_rate = fill_rate        # tokens per second
        self._tokens = capacity
        self._timestamp = time.monotonic()
        self._lock = threading.Lock()

    def consume(self, tokens=1):
        """Try to consume tokens; return True if successful."""
        with self._lock:
            now = time.monotonic()
            # refill based on elapsed time
            elapsed = now - self._timestamp
            self._tokens = min(self.capacity,
                               self._tokens + elapsed * self.fill_rate)
            self._timestamp = now
            if self._tokens >= tokens:
                self._tokens -= tokens
                return True
            return False
Enter fullscreen mode Exit fullscreen mode

Each service instance creates a bucket tuned to its expected traffic:

# 10 tokens burst, refill at 2 tokens/sec → allows short spikes
LOCAL_BUCKET = LocalBucket(capacity=10, fill_rate=2.0)
Enter fullscreen mode Exit fullscreen mode

When the local bucket is empty, we ask Redis for a token via a Lua script that implements the same token bucket logic atomically.

-- redis_token_bucket.lua
-- KEYS[1] = bucket key (e.g., "notif:bucket:user_id")
-- ARGV[1] = fill_rate (tokens per second)
-- ARGV[2] = capacity (max tokens)
-- ARGV[3] = now (current unix time float)
-- returns 1 if token granted, 0 otherwise

local last = redis.call("HGET", KEYS[1], "last")
local tokens = redis.call("HGET", KEYS[1], "tokens")

if not last then
    last = tonumber(ARGV[3])
    tokens = tonumber(ARGV[2])
else
    last = tonumber(last)
    tokens = tonumber(tokens)
    -- refill
    local delta = tonumber(ARGV[3]) - last
    tokens = math.min(tonumber(ARGV[2]), tokens + delta * tonumber(ARGV[1]))
end

local granted = 0
if tokens >= 1 then
    tokens = tokens - 1
    granted = 1
end

redis.call("HMSET", KEYS[1], "tokens", tokens, "last", tonumber(ARGV[3]))
return granted
Enter fullscreen mode Exit fullscreen mode

In Python we call it like this:

import redis

r = redis.Redis(host='redis.example.com', port=6379, db=0)

def try_redis_token(user_id):
    key = f"notif:bucket:{user_id}"
    now = time.time()
    granted = r.eval(
        open("redis_token_bucket.lua").read(),
        1, key,
        2.0,          # fill_rate
        10.0,         # capacity
        now
    )
    return bool(granted)
Enter fullscreen mode Exit fullscreen mode

Finally, the notification flow:

def send_notification(user_id, payload):
    # 1️⃣ Fast path – local bucket
    if LOCAL_BUCKET.consume():
        push_to_ws(user_id, payload)
        return

    # 2️⃣ Slow path – ask Redis
    if try_redis_token(user_id):
        push_to_ws(user_id, payload)
    else:
        # Optionally persist to a dead‑letter queue for later retry
        dead_letter_queue.enqueue(user_id, payload)
Enter fullscreen mode Exit fullscreen mode

Traps to avoid (the “boss fights” on our quest):

  • Clock drift – using time.time() for refill can cause negative deltas if the system clock jumps. We switched to time.monotonic() for the local bucket and passed a Unix float from the server to Redis (still monotonic enough for short intervals).
  • Starvation – a global limiter with a low rate could idle active users while letting a bursty user hog tokens. By tying the bucket to user_id we guarantee fairness per user.
  • Redis failure – if Redis is down, the local bucket will eventually empty and we start dropping notifications. We mitigate this by caching a short‑term fallback queue (e.g., in‑memory list) that we flush when Redis returns.

Why This New Power Matters

Adopting this dual‑layer token bucket turned our notification pipeline from a brittle, bottleneck‑prone beast into a smooth‑flowing river.

  • Latency dropped from seconds to sub‑100 ms for the vast majority of messages, because most notifications never touch Redis.
  • Throughput rose by ~3× on the same hardware – the DB only sees the occasional persisted fallback, not every event.
  • Operational simplicity – we still have a single source of truth (the Redis script) for the global rate, but we get the resiliency of a local cache.

Most importantly, our users stopped missing alerts. The real‑time feel that made our app addictive is back, and the team can now ship new notification types without fear of taking down the service.


Your Turn

Grab a toy project (maybe a chat app or a recommendation feed) and try slapping a local + Redis token bucket in front of your event producer. Start with a modest burst capacity (5‑10 tokens) and a refill rate that matches your average traffic. Watch how the latency curve flattens, and when you feel comfortable, experiment with persisting dropped events to a dead‑letter queue for a “retry later” boss level.

What’s the biggest surprise you hit when you first saw the system absorb a spike without breaking a sweat? Drop a comment—I’d love to hear your war stories!

Happy building, and may your notifications always flow like a well‑guarded river. 🚀

Top comments (0)