Skip to content

Distributed Semaphores with Redis

asyncio.Semaphore limits concurrency inside one process. When several replicas share a downstream that allows only N concurrent calls — a licensed API, a fragile legacy service, a GPU pool — the limit has to live somewhere they all see. Measured with Redis 8.10 and redis-py 5.3.1, three processes of ten asyncio workers each, a limit of 5 and 600 jobs holding a slot for 20 ms: a naive INCR/DECR counter never let more than five in, but when two holders were killed with SIGKILL their slots never came back — a new acquirer was still waiting after 10 s with the counter stuck at 2. A leased semaphore built on a sorted set recovered crashed slots after the lease expired (1.82 s with a 2 s lease), but its polling acquirers were badly unfair: median wait 0 ms, p99 1.1–1.8 s, worst 2.4 s, because a worker that had just released usually won the slot back immediately. A ticketed FIFO version fixed fairness (p99 240 ms) but polling left freed slots idle and stretched the run from 2.7 s to 4.6 s. Adding a release notification over pub/sub gave both: 2.84–2.88 s for the run, p99 wait 148–157 ms, never more than five holders. This guide builds that final version.

Prerequisites

1. Give every slot a lease, not a counter

The tempting implementation is a counter: INCR, check it is at most N, DECR when done. It enforces the limit, and it leaks. A process that crashes, is OOM-killed or loses its network between INCR and DECR never decrements, and nothing else knows that the slot is free:

# leaks a slot forever when the holder dies before DECR
while await r.incr("sem:count") > LIMIT:
    await r.decr("sem:count")
    await asyncio.sleep(0.02)

Measured with a limit of 2: two holders were killed with SIGKILL, and a third acquirer timed out after 10 s with sem:count at 2. The fix is to record who holds each slot and until when, so expired holders can be removed by anyone. A sorted set does both: the member is a random token per acquisition, the score is the lease's expiry time. Every attempt to acquire first removes members whose expiry has passed. Take that time from the Redis server with TIME inside a Lua script, not from the client: replicas' clocks differ, and the server's clock is the only one every client agrees on.

Verify: after kill -9 of every holder, a new acquirer gets a slot within one lease period.

Four Redis semaphores under the same load A grid of 4 rows by 5 columns. Four Redis semaphores under the same load design max holders after holder crash run time wait p50 / p99 INCR/DECR counter 5 slot lost forever - - leased sorted set, polling 5 recovered in 1.82 s 2.5-2.7 s 0 ms / 1.1-1.8 s FIFO tickets, polling 5 recovered after lease 4.6 s 212 ms / 240 ms FIFO tickets + release notify 5 recovered in 2.02 s 2.84-2.88 s 119 ms / 148-157 ms Three processes x 10 workers; lower bound for the run is 2.4 s.

2. Queue acquirers in FIFO order

With polling, whichever client asks at the right moment wins. The client that just released is always asking at the right moment — its next loop iteration runs immediately — so it tends to win again, and the others starve. Measured: half of all acquisitions waited 0 ms, while the unluckiest waited 2.4 s.

To make it fair, each acquirer takes a ticket from a counter and adds itself to a queue ordered by ticket; the limit lowest live tickets hold the slots, and everyone else is waiting in line:

NOW = "local t = redis.call('TIME') local now = t[1] * 1000 + math.floor(t[2] / 1000)\n"

ENQUEUE = NOW + """
local ticket = redis.call('INCR', KEYS[3])
redis.call('ZADD', KEYS[1], ticket, ARGV[2])                -- queue: ordered by ticket
redis.call('ZADD', KEYS[2], now + tonumber(ARGV[1]), ARGV[2])  -- lease: ordered by expiry
"""

CHECK = NOW + """
local dead = redis.call('ZRANGEBYSCORE', KEYS[2], '-inf', now)
if #dead > 0 then
  redis.call('ZREM', KEYS[1], unpack(dead))
  redis.call('ZREM', KEYS[2], unpack(dead))
end
if not redis.call('ZSCORE', KEYS[2], ARGV[3]) then return -1 end   -- our entry expired
redis.call('ZADD', KEYS[2], now + tonumber(ARGV[2]), ARGV[3])        -- extend lease
if redis.call('ZRANK', KEYS[1], ARGV[3]) < tonumber(ARGV[1]) then return 1 end
return 0
"""

Waiters hold leases too, extended on every check, so a waiter that crashes drops out of the queue instead of blocking everyone behind it. Measured: p99 wait fell from 1.1–1.8 s to 240 ms and the worst case to 240 ms — but the run took 4.6 s instead of 2.7 s, because a freed slot stayed empty until the next waiter's 20 ms poll came round.

Verify: with many waiters, the spread between median and p99 acquisition time is small.

3. Wake waiters when a slot is released

Polling faster wastes Redis round trips; polling slower wastes slots. Instead, publish a message on every release and let waiters re-check immediately, keeping a slow poll only as a fallback for lost messages:

class RedisSemaphore:
    def __init__(self, r, name, limit, ttl=5.0, fallback_poll=0.25):
        self.r, self.limit, self.ttl_ms, self.fallback_poll = r, limit, int(ttl * 1000), fallback_poll
        self.keys = [f"sem:{name}:queue", f"sem:{name}:lease", f"sem:{name}:seq"]
        self.channel = f"sem:{name}:released"
        self._enqueue, self._check = r.register_script(ENQUEUE), r.register_script(CHECK)
        self._changed = asyncio.Event()
        self._listener = None

    async def _listen(self):
        async with self.r.pubsub() as pubsub:
            await pubsub.subscribe(self.channel)
            async for msg in pubsub.listen():
                if msg["type"] == "message":
                    self._changed.set()

    async def acquire(self) -> str:
        if self._listener is None:
            self._listener = asyncio.create_task(self._listen())
        token = uuid.uuid4().hex
        await self._enqueue(keys=self.keys, args=[self.ttl_ms, token])
        try:
            while True:
                self._changed.clear()
                state = await self._check(keys=self.keys[:2], args=[self.limit, self.ttl_ms, token])
                if state == 1:
                    return token
                if state == -1:
                    raise TimeoutError("semaphore lease expired while waiting")
                try:
                    async with asyncio.timeout(self.fallback_poll):
                        await self._changed.wait()
                except TimeoutError:
                    pass                               # pub/sub is fire-and-forget
        except BaseException:
            await self.release(token)                  # leave the queue on cancel or error
            raise

    async def release(self, token: str) -> None:
        async with self.r.pipeline(transaction=True) as p:
            p.zrem(self.keys[0], token).zrem(self.keys[1], token).publish(self.channel, b"")
            await p.execute()

One subscriber per process feeds an asyncio.Event that all local waiters share, so a thousand waiting coroutines cost one Redis connection for notifications, not a thousand. Measured: the run took 2.84–2.88 s against a theoretical minimum of 2.4 s, median wait 119 ms and p99 148–157 ms, and the maximum number of simultaneous holders observed across all three processes was 5.

Verify: under load, the run time approaches jobs × hold time / limit, and the maximum observed holder count equals the limit.

Acquiring a slot with release notification A sequence of 6 messages between 4 participants. Acquiring a slot with release notification waiter Redis holder subscriber ENQUEUE: ticket 42 CHECK: rank 5 -> wait ZREM + PUBLISH released message on released event.set() CHECK: rank 4 -> acquired The fallback poll covers a lost message; it is not the main path.

4. Renew leases for long holds

A lease must be longer than a normal hold, or slots expire under healthy holders and the limit is exceeded. For holds that vary widely, keep the lease short — so crashes are recovered quickly — and renew it from a background task:

async def renew(self, token: str) -> bool:
    return await self._check(keys=self.keys[:2], args=[self.limit, self.ttl_ms, token]) == 1


async def call_with_slot(sem, fn):
    token = await sem.acquire()
    work = asyncio.create_task(fn())
    try:
        while True:
            done, _ = await asyncio.wait({work}, timeout=sem.ttl_ms / 3000)
            if done:
                return work.result()
            if not await sem.renew(token):
                raise RuntimeError("lost the semaphore slot")
    finally:
        work.cancel()                       # no-op if finished; stops orphaned work otherwise
        await sem.release(token)

Measured with a 2 s lease, renewed every 0.67 s: a 2.5 s job — longer than the lease — completed with its slot in 2.50 s. When the lease was deleted behind the holder's back, the next renewal failed and the call raised "lost the semaphore slot" after 0.67 s, with the work cancelled; when the caller itself was cancelled, the work was cancelled and the queue was empty afterwards. When holders were killed instead, the next acquirer got a slot after 2.02 s — one lease period. That is the trade-off to tune: the lease bounds how long a crashed holder blocks a slot, and renewal keeps the lease from bounding legitimate work. A semaphore, like a lock, cannot stop a holder that was paused past its lease from still believing it holds a slot; if exceeding the limit briefly would cause damage rather than load, apply the fencing approach from implementing a Redis lock with fencing tokens; the renewal loop itself is covered in renewing lock leases with a heartbeat task.

Verify: a hold longer than the lease succeeds with renewals, and a killed holder's slot frees within one lease.

5. Size the limit per downstream and observe it

The semaphore's limit is a statement about the downstream's capacity, so name semaphores after what they protect and keep their configuration with that dependency. Observe it from Redis directly — the queue and lease sets are inspectable state:

async def semaphore_stats(r, name: str, limit: int) -> dict:
    queued = await r.zcard(f"sem:{name}:queue")
    return {"holders": min(queued, limit), "waiting": max(0, queued - limit)}

A persistent waiting count means the limit or the downstream is too small; waiters that time out mean callers' deadlines are shorter than the queue. Combine it with an in-process asyncio.Semaphore in front of the distributed one, so a single replica cannot flood Redis with thousands of queued tickets, as in limiting concurrent requests with asyncio.Semaphore. For rate limits rather than concurrency limits — requests per second rather than requests at once — use a token bucket instead; the two answer different questions.

Verify: dashboards show holders and waiters per semaphore, and a local semaphore caps each replica's share.

Which concurrency limit does this need? A decision on What must be limited with 4 outcomes. Which concurrency limit does this need? What must be limited? calls at once, one process asyncio.Semaphore no network calls at once, all replicas fair Redis semaphore p99 wait 157 ms exceeding it causes damage add fencing tokens leases cannot stop paused holders calls per second rate limiter, not a semaphore different question A distributed semaphore limits concurrency, not throughput.

Verification

A distributed semaphore is working when:

  • Crashed holders' slots are recovered within one lease period.
  • Acquirers are served in FIFO order, with p99 wait close to the median.
  • Releases wake waiters immediately, with polling only as a fallback.
  • Long holds renew their leases, and losing a renewal stops the work.

Diagnostic Hook: sample ZCARD sem:<name>:queue every few seconds and alert when it stays above the limit for minutes. A queue that never drains while the downstream is idle means slots are held by leases that should have expired — a holder renewing without working — or by a limit set lower than anyone intended.

Pitfalls & edge cases

  • Counter semaphores. Measured: two crashed holders' slots never returned.
  • Polling acquirers without a queue. Measured: p99 wait 1.8 s while the median was 0 ms.
  • Client clocks for expiry. Use Redis TIME inside the script.
  • Leases shorter than holds. Without renewal, healthy holders lose slots and the limit is exceeded.

Frequently Asked Questions

How do I implement a distributed semaphore in Redis?

Use a sorted set of leased tokens: each acquisition adds a token with an expiry score, every attempt first removes expired tokens, and a slot is free while fewer than N tokens remain. A FIFO ticket queue and a release notification over pub/sub made it fair and fast in testing.

Why not just use INCR and DECR?

A holder that crashes before DECR leaks its slot forever. In testing, two killed holders left the counter at 2 and a new acquirer still waited after 10 s.

How long should the semaphore lease be?

Long enough to cover a normal hold or renewal interval, short enough to recover crashed slots quickly. With a 2 s lease, a crashed slot was recovered in 2.02 s; longer holds renewed the lease every 0.67 s.

Is a Redis semaphore safe against paused processes?

No more than a Redis lock: a holder paused past its lease may still act. If exceeding the limit causes damage, use fencing tokens at the protected resource.