karawaci.kode

← Semua snippet

Python Lanjut Performance

Distributed lock pattern Redis (Redlock)

Distributed lock dengan Redis — pastikan job hanya jalan di satu worker sekaligus. SET NX EX pattern + token validation untuk safety.

Dipublikasikan 21 Juni 2026

Cron job yang sama jalan di 3 server bisa bikin chaos — double-charge user, double-send email, double-update stok. Distributed lock di Redis bikin cuma satu worker yang menang. SET NX EX murah dan cukup buat 95% kasus. Snippet ini async lock dengan auto-refresh untuk task panjang.

Kode

import asyncio
import secrets
import time
from contextlib import asynccontextmanager
from dataclasses import dataclass
from typing import AsyncIterator, Optional

import redis.asyncio as redis

# Lua: hanya hapus kalau token match
UNLOCK_LUA = """
if redis.call("GET", KEYS[1]) == ARGV[1] then
  return redis.call("DEL", KEYS[1])
else
  return 0
end
"""

# Lua: hanya extend kalau token match
REFRESH_LUA = """
if redis.call("GET", KEYS[1]) == ARGV[1] then
  return redis.call("PEXPIRE", KEYS[1], ARGV[2])
else
  return 0
end
"""


@dataclass
class LockToken:
    key: str
    token: str


class LockAcquisitionFailed(Exception):
    pass


class DistributedLock:
    def __init__(self, redis_client: redis.Redis):
        self.r = redis_client
        self._unlock_sha: Optional[str] = None
        self._refresh_sha: Optional[str] = None

    async def _load_scripts(self) -> None:
        if self._unlock_sha is None:
            self._unlock_sha = await self.r.script_load(UNLOCK_LUA)
            self._refresh_sha = await self.r.script_load(REFRESH_LUA)

    async def acquire(
        self,
        key: str,
        ttl_seconds: int = 30,
        wait_timeout: float = 0,
        retry_interval: float = 0.05,
    ) -> LockToken:
        """Acquire lock atau raise LockAcquisitionFailed."""
        await self._load_scripts()
        full_key = f"lock:{key}"
        token = secrets.token_urlsafe(16)
        deadline = time.monotonic() + wait_timeout

        while True:
            ok = await self.r.set(full_key, token, nx=True, ex=ttl_seconds)
            if ok:
                return LockToken(key=full_key, token=token)

            if time.monotonic() >= deadline:
                raise LockAcquisitionFailed(f"Tidak dapat lock '{key}' dalam {wait_timeout}s")

            await asyncio.sleep(retry_interval)

    async def release(self, lock: LockToken) -> bool:
        """Release lock. Return True kalau berhasil hapus (masih punya kita)."""
        await self._load_scripts()
        result = await self.r.evalsha(self._unlock_sha, 1, lock.key, lock.token)
        return bool(result)

    async def refresh(self, lock: LockToken, ttl_seconds: int) -> bool:
        """Extend TTL kalau masih punya kita."""
        await self._load_scripts()
        result = await self.r.evalsha(
            self._refresh_sha, 1, lock.key, lock.token, str(ttl_seconds * 1000)
        )
        return bool(result)

    @asynccontextmanager
    async def lock(
        self,
        key: str,
        ttl_seconds: int = 30,
        wait_timeout: float = 0,
        auto_refresh: bool = True,
    ) -> AsyncIterator[LockToken]:
        """Context manager dengan auto-refresh background."""
        lock_token = await self.acquire(key, ttl_seconds=ttl_seconds, wait_timeout=wait_timeout)
        refresh_task: Optional[asyncio.Task] = None

        if auto_refresh:
            async def _refresher():
                interval = ttl_seconds / 3
                while True:
                    await asyncio.sleep(interval)
                    if not await self.refresh(lock_token, ttl_seconds):
                        return  # lock sudah hilang, stop refresh

            refresh_task = asyncio.create_task(_refresher())

        try:
            yield lock_token
        finally:
            if refresh_task:
                refresh_task.cancel()
                try:
                    await refresh_task
                except asyncio.CancelledError:
                    pass
            await self.release(lock_token)

Pemakaian

# Setup
async def main():
    r = redis.from_url("redis://localhost:6379/0", decode_responses=True)
    locks = DistributedLock(r)

    # Use case 1: cron job kirim notifikasi harian
    async with locks.lock("daily-notif", ttl_seconds=60, wait_timeout=0):
        print("Worker ini yang kirim notif hari ini")
        await kirim_notifikasi_harian()
        # Auto-refresh setiap 20 detik supaya lock tidak expire
        # walau task butuh 5 menit
# Use case 2: prevent double order checkout
async def checkout(user_id: int, cart_id: str, locks: DistributedLock):
    try:
        async with locks.lock(
            f"checkout:user:{user_id}",
            ttl_seconds=10,
            wait_timeout=2,  # tunggu max 2 detik
        ):
            # Hanya satu request checkout per user yang bisa proceed
            await proses_checkout(user_id, cart_id)
    except LockAcquisitionFailed:
        raise ValueError("Checkout sudah diproses, mohon tunggu")
# Use case 3: lock per produk saat update stok
async def update_stok(produk_id: int, qty: int, locks: DistributedLock):
    async with locks.lock(f"produk:{produk_id}:stok", ttl_seconds=5, wait_timeout=3):
        # Read-modify-write yang race-safe
        current = await get_stok(produk_id)
        if current + qty < 0:
            raise ValueError("Stok tidak cukup")
        await set_stok(produk_id, current + qty)

Kapan dipakai

  • Cron job yang harus single-instance di multi-server deployment.
  • Prevent double charge / double order checkout.
  • Update stok inventory yang sering race.
  • Send-once notification yang fan-out ke banyak worker.
  • Migration script yang gak boleh paralel.

Catatan

  • Token wajib — tanpa token, race lock-takeover bikin worker A unlock yang seharusnya punya B.
  • TTL > job duration untuk safety, atau pakai auto-refresh. TTL terlalu pendek + job lambat = lock expired, worker lain masuk, conflict.
  • Lua atomic — check-and-del / check-and-pexpire harus atomic. Tanpa Lua, ada window race.
  • Redlock untuk financial — single-node Redis hilang lock saat failover. Untuk life-critical, pakai 5 node dengan algoritma Redlock (lib redlock-py) atau hindari Redis untuk hal ini.
  • Async cancellation — kalau context manager dibatalkan (TimeoutError), task refresher harus di-cancel juga. Snippet ini handle di finally.
  • Wait timeout 0 = non-blocking. Cocok untuk cron job yang kalau worker lain sudah pegang, skip aja.

Distributed lock = bottleneck. Kalau jadi hot key (puluhan ribu acquire/detik), shard berdasarkan resource ID (lock:produk:{id % 16}) atau pertimbangkan optimistic concurrency control di DB.

# tags

redisdistributed-lockredlockconcurrencyasync

Ditulis oleh Asti Larasati · 21 Juni 2026