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:
- 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.
-
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
rtokens/second, capacityb). 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
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)
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
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)
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
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)
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)
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 totime.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_idwe 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)