Rate Limiter -- компонент, ограничивающий количество запросов клиента за единицу времени. Защита от DDoS, злоупотреблений API и обеспечение справедливого использования ресурсов.
Шаг 1: Требования
Функциональные требования
Ограничение запросов по разным ключам (IP, user_id, API key)
Поддержка разных правил (100 req/min, 1000 req/hour)
Возврат HTTP 429 (Too Many Requests) при превышении лимита
Наиболее распространённый алгоритм. Токены добавляются с фиксированной скоростью. Каждый запрос забирает один токен.
<?php
declare(strict_types=1);
final class TokenBucketLimiter
{
public function __construct(
private readonly \Redis $redis,
) {}
/**
* @param string $key Unique client identifier
* @param int $capacity Max tokens (burst size)
* @param float $rate Tokens added per second
* @return array{allowed: bool, remaining: int, retryAfter: float}
*/
public function allow(string $key, int $capacity, float $rate): array
{
$now = microtime(true);
$redisKey = "token_bucket:{$key}";
// Lua script for atomicity
$script = <<<'LUA'
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call('HMGET', key, 'tokens', 'last_refill')
local tokens = tonumber(data[1])
local lastRefill = tonumber(data[2])
-- Initialize if first request
if tokens == nil then
tokens = capacity
lastRefill = now
end
-- Add tokens based on elapsed time
local elapsed = now - lastRefill
local newTokens = elapsed * rate
tokens = math.min(capacity, tokens + newTokens)
lastRefill = now
-- Try to consume one token
local allowed = 0
if tokens >= 1 then
tokens = tokens - 1
allowed = 1
end
-- Save state
redis.call('HMSET', key, 'tokens', tokens, 'last_refill', lastRefill)
redis.call('EXPIRE', key, math.ceil(capacity / rate) + 1)
-- Calculate retry-after if not allowed
local retryAfter = 0
if allowed == 0 then
retryAfter = (1 - tokens) / rate
end
return {allowed, math.floor(tokens), tostring(retryAfter)}
LUA;
$result = $this->redis->eval($script, [$redisKey, $capacity, $rate, $now], 1);
return [
'allowed' => (bool) $result[0],
'remaining' => (int) $result[1],
'retryAfter' => (float) $result[2],
];
}
}
package ratelimit
import (
"context"
"fmt"
"strconv"
"time"
"github.com/redis/go-redis/v9"
)
// TokenBucketResult holds the result of a rate limit check.
type TokenBucketResult struct {
Allowed bool
Remaining int
RetryAfter float64
}
// TokenBucketLimiter implements the token bucket algorithm using Redis.
type TokenBucketLimiter struct {
rdb *redis.Client
}
func NewTokenBucketLimiter(rdb *redis.Client) *TokenBucketLimiter {
return &TokenBucketLimiter{rdb: rdb}
}
// Allow checks if a request is allowed under the token bucket algorithm.
func (l *TokenBucketLimiter) Allow(ctx context.Context, key string, capacity int, rate float64) (TokenBucketResult, error) {
now := float64(time.Now().UnixMicro()) / 1e6
redisKey := "token_bucket:" + key
script := redis.NewScript(`
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call('HMGET', key, 'tokens', 'last_refill')
local tokens = tonumber(data[1])
local lastRefill = tonumber(data[2])
if tokens == nil then
tokens = capacity
lastRefill = now
end
local elapsed = now - lastRefill
tokens = math.min(capacity, tokens + elapsed * rate)
lastRefill = now
local allowed = 0
if tokens >= 1 then
tokens = tokens - 1
allowed = 1
end
redis.call('HMSET', key, 'tokens', tokens, 'last_refill', lastRefill)
redis.call('EXPIRE', key, math.ceil(capacity / rate) + 1)
local retryAfter = 0
if allowed == 0 then
retryAfter = (1 - tokens) / rate
end
return {allowed, math.floor(tokens), tostring(retryAfter)}
`)
res, err := script.Run(ctx, l.rdb, []string{redisKey},
capacity, rate, fmt.Sprintf("%.6f", now),
).Slice()
if err != nil {
return TokenBucketResult{}, fmt.Errorf("token bucket eval: %w", err)
}
retryAfter, _ := strconv.ParseFloat(res[2].(string), 64)
return TokenBucketResult{
Allowed: res[0].(int64) == 1,
Remaining: int(res[1].(int64)),
RetryAfter: retryAfter,
}, nil
}
using System.Globalization;
using StackExchange.Redis;
namespace RateLimit;
// TokenBucketResult holds the result of a rate limit check.
public readonly record struct TokenBucketResult(bool Allowed, int Remaining, double RetryAfter);
// TokenBucketLimiter implements the token bucket algorithm using Redis.
public sealed class TokenBucketLimiter(IDatabase db)
{
// Lua script for atomicity
private const string Script =
"""
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call('HMGET', key, 'tokens', 'last_refill')
local tokens = tonumber(data[1])
local lastRefill = tonumber(data[2])
if tokens == nil then
tokens = capacity
lastRefill = now
end
local elapsed = now - lastRefill
tokens = math.min(capacity, tokens + elapsed * rate)
lastRefill = now
local allowed = 0
if tokens >= 1 then
tokens = tokens - 1
allowed = 1
end
redis.call('HMSET', key, 'tokens', tokens, 'last_refill', lastRefill)
redis.call('EXPIRE', key, math.ceil(capacity / rate) + 1)
local retryAfter = 0
if allowed == 0 then
retryAfter = (1 - tokens) / rate
end
return {allowed, math.floor(tokens), tostring(retryAfter)}
""";
// AllowAsync checks if a request is allowed under the token bucket algorithm.
public async Task<TokenBucketResult> AllowAsync(string key, int capacity, double rate)
{
// Wall clock, not Stopwatch: the timestamp is shared between nodes via Redis
var now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds() / 1000.0;
var redisKey = $"token_bucket:{key}";
var result = (RedisResult[])(await db.ScriptEvaluateAsync(
Script,
[redisKey],
[capacity, rate, now.ToString("F6", CultureInfo.InvariantCulture)]))!;
return new TokenBucketResult(
Allowed: (int)result[0] == 1,
Remaining: (int)result[1],
RetryAfter: double.Parse((string)result[2]!, CultureInfo.InvariantCulture));
}
}
import time
from dataclasses import dataclass
import redis.asyncio as redis
# Lua script for atomicity
TOKEN_BUCKET_SCRIPT = """
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call('HMGET', key, 'tokens', 'last_refill')
local tokens = tonumber(data[1])
local lastRefill = tonumber(data[2])
if tokens == nil then
tokens = capacity
lastRefill = now
end
local elapsed = now - lastRefill
tokens = math.min(capacity, tokens + elapsed * rate)
lastRefill = now
local allowed = 0
if tokens >= 1 then
tokens = tokens - 1
allowed = 1
end
redis.call('HMSET', key, 'tokens', tokens, 'last_refill', lastRefill)
redis.call('EXPIRE', key, math.ceil(capacity / rate) + 1)
local retryAfter = 0
if allowed == 0 then
retryAfter = (1 - tokens) / rate
end
return {allowed, math.floor(tokens), tostring(retryAfter)}
"""
@dataclass(frozen=True, slots=True)
class TokenBucketResult:
allowed: bool
remaining: int
retry_after: float
class TokenBucketLimiter:
def __init__(self, client: redis.Redis) -> None:
# register_script caches the SHA and falls back to EVAL on NOSCRIPT
self._script = client.register_script(TOKEN_BUCKET_SCRIPT)
async def allow(self, key: str, capacity: int, rate: float) -> TokenBucketResult:
# time.time(), not monotonic: the timestamp is shared between nodes via Redis
now = time.time()
redis_key = f"token_bucket:{key}"
allowed, remaining, retry_after = await self._script(
keys=[redis_key],
args=[capacity, rate, f"{now:.6f}"],
)
return TokenBucketResult(
allowed=bool(allowed),
remaining=int(remaining),
retry_after=float(retry_after),
)
**Плюсы:** допускает burst, гибкий, интуитивный.
**Минусы:** нужно хранить состояние для каждого клиента.
4.2 Leaky Bucket
Запросы обрабатываются с постоянной скоростью, лишние отбрасываются.
<?php
declare(strict_types=1);
final class LeakyBucketLimiter
{
public function __construct(
private readonly \Redis $redis,
) {}
/**
* @param string $key Unique client identifier
* @param int $capacity Queue size (max waiting requests)
* @param float $leakRate Requests processed per second
*/
public function allow(string $key, int $capacity, float $leakRate): array
{
$now = microtime(true);
$redisKey = "leaky_bucket:{$key}";
$script = <<<'LUA'
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local leakRate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call('HMGET', key, 'water', 'last_leak')
local water = tonumber(data[1]) or 0
local lastLeak = tonumber(data[2]) or now
-- Leak water based on time passed
local elapsed = now - lastLeak
local leaked = elapsed * leakRate
water = math.max(0, water - leaked)
lastLeak = now
local allowed = 0
if water < capacity then
water = water + 1
allowed = 1
end
redis.call('HMSET', key, 'water', water, 'last_leak', lastLeak)
redis.call('EXPIRE', key, math.ceil(capacity / leakRate) + 1)
return {allowed, capacity - math.floor(water)}
LUA;
$result = $this->redis->eval($script, [$redisKey, $capacity, $leakRate, $now], 1);
return [
'allowed' => (bool) $result[0],
'remaining' => (int) $result[1],
];
}
}
package ratelimit
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// LeakyBucketResult holds the result of a leaky bucket check.
type LeakyBucketResult struct {
Allowed bool
Remaining int
}
// LeakyBucketLimiter implements the leaky bucket algorithm using Redis.
type LeakyBucketLimiter struct {
rdb *redis.Client
}
func NewLeakyBucketLimiter(rdb *redis.Client) *LeakyBucketLimiter {
return &LeakyBucketLimiter{rdb: rdb}
}
// Allow checks if a request is allowed under the leaky bucket algorithm.
func (l *LeakyBucketLimiter) Allow(ctx context.Context, key string, capacity int, leakRate float64) (LeakyBucketResult, error) {
now := float64(time.Now().UnixMicro()) / 1e6
redisKey := "leaky_bucket:" + key
script := redis.NewScript(`
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local leakRate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call('HMGET', key, 'water', 'last_leak')
local water = tonumber(data[1]) or 0
local lastLeak = tonumber(data[2]) or now
local elapsed = now - lastLeak
water = math.max(0, water - elapsed * leakRate)
lastLeak = now
local allowed = 0
if water < capacity then
water = water + 1
allowed = 1
end
redis.call('HMSET', key, 'water', water, 'last_leak', lastLeak)
redis.call('EXPIRE', key, math.ceil(capacity / leakRate) + 1)
return {allowed, capacity - math.floor(water)}
`)
res, err := script.Run(ctx, l.rdb, []string{redisKey},
capacity, leakRate, fmt.Sprintf("%.6f", now),
).Slice()
if err != nil {
return LeakyBucketResult{}, fmt.Errorf("leaky bucket eval: %w", err)
}
return LeakyBucketResult{
Allowed: res[0].(int64) == 1,
Remaining: int(res[1].(int64)),
}, nil
}
using System.Globalization;
using StackExchange.Redis;
namespace RateLimit;
// LeakyBucketResult holds the result of a leaky bucket check.
public readonly record struct LeakyBucketResult(bool Allowed, int Remaining);
// LeakyBucketLimiter implements the leaky bucket algorithm using Redis.
public sealed class LeakyBucketLimiter(IDatabase db)
{
private const string Script =
"""
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local leakRate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call('HMGET', key, 'water', 'last_leak')
local water = tonumber(data[1]) or 0
local lastLeak = tonumber(data[2]) or now
-- Leak water based on time passed
local elapsed = now - lastLeak
water = math.max(0, water - elapsed * leakRate)
lastLeak = now
local allowed = 0
if water < capacity then
water = water + 1
allowed = 1
end
redis.call('HMSET', key, 'water', water, 'last_leak', lastLeak)
redis.call('EXPIRE', key, math.ceil(capacity / leakRate) + 1)
return {allowed, capacity - math.floor(water)}
""";
/// <param name="capacity">Queue size (max waiting requests)</param>
/// <param name="leakRate">Requests processed per second</param>
public async Task<LeakyBucketResult> AllowAsync(string key, int capacity, double leakRate)
{
var now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds() / 1000.0;
var redisKey = $"leaky_bucket:{key}";
var result = (RedisResult[])(await db.ScriptEvaluateAsync(
Script,
[redisKey],
[capacity, leakRate, now.ToString("F6", CultureInfo.InvariantCulture)]))!;
return new LeakyBucketResult(
Allowed: (int)result[0] == 1,
Remaining: (int)result[1]);
}
}
import time
from dataclasses import dataclass
import redis.asyncio as redis
LEAKY_BUCKET_SCRIPT = """
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local leakRate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call('HMGET', key, 'water', 'last_leak')
local water = tonumber(data[1]) or 0
local lastLeak = tonumber(data[2]) or now
-- Leak water based on time passed
local elapsed = now - lastLeak
water = math.max(0, water - elapsed * leakRate)
lastLeak = now
local allowed = 0
if water < capacity then
water = water + 1
allowed = 1
end
redis.call('HMSET', key, 'water', water, 'last_leak', lastLeak)
redis.call('EXPIRE', key, math.ceil(capacity / leakRate) + 1)
return {allowed, capacity - math.floor(water)}
"""
@dataclass(frozen=True, slots=True)
class LeakyBucketResult:
allowed: bool
remaining: int
class LeakyBucketLimiter:
def __init__(self, client: redis.Redis) -> None:
self._script = client.register_script(LEAKY_BUCKET_SCRIPT)
async def allow(
self,
key: str,
capacity: int, # queue size (max waiting requests)
leak_rate: float, # requests processed per second
) -> LeakyBucketResult:
now = time.time()
redis_key = f"leaky_bucket:{key}"
allowed, remaining = await self._script(
keys=[redis_key],
args=[capacity, leak_rate, f"{now:.6f}"],
)
return LeakyBucketResult(allowed=bool(allowed), remaining=int(remaining))
**Плюсы:** стабильный output rate.
**Минусы:** burst трафик теряется, не подходит для API.
4.3 Fixed Window Counter
Простейший алгоритм: счётчик на фиксированный интервал.
<?php
declare(strict_types=1);
final class FixedWindowLimiter
{
public function __construct(
private readonly \Redis $redis,
) {}
/**
* @param string $key Unique client identifier
* @param int $limit Max requests per window
* @param int $windowSec Window size in seconds
*/
public function allow(string $key, int $limit, int $windowSec): array
{
$window = (int) (time() / $windowSec);
$redisKey = "fixed_window:{$key}:{$window}";
$script = <<<'LUA'
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local windowSec = tonumber(ARGV[2])
local current = tonumber(redis.call('GET', key) or '0')
if current < limit then
redis.call('INCR', key)
redis.call('EXPIRE', key, windowSec)
return {1, limit - current - 1}
else
return {0, 0}
end
LUA;
$result = $this->redis->eval($script, [$redisKey, $limit, $windowSec], 1);
return [
'allowed' => (bool) $result[0],
'remaining' => (int) $result[1],
'resetAt' => ($window + 1) * $windowSec,
];
}
}
package ratelimit
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// FixedWindowResult holds the result of a fixed window check.
type FixedWindowResult struct {
Allowed bool
Remaining int
ResetAt int64
}
// FixedWindowLimiter implements the fixed window counter algorithm.
type FixedWindowLimiter struct {
rdb *redis.Client
}
func NewFixedWindowLimiter(rdb *redis.Client) *FixedWindowLimiter {
return &FixedWindowLimiter{rdb: rdb}
}
// Allow checks if a request is allowed under the fixed window algorithm.
func (l *FixedWindowLimiter) Allow(ctx context.Context, key string, limit int, windowSec int) (FixedWindowResult, error) {
now := time.Now().Unix()
window := now / int64(windowSec)
redisKey := fmt.Sprintf("fixed_window:%s:%d", key, window)
script := redis.NewScript(`
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local windowSec = tonumber(ARGV[2])
local current = tonumber(redis.call('GET', key) or '0')
if current < limit then
redis.call('INCR', key)
redis.call('EXPIRE', key, windowSec)
return {1, limit - current - 1}
else
return {0, 0}
end
`)
res, err := script.Run(ctx, l.rdb, []string{redisKey}, limit, windowSec).Slice()
if err != nil {
return FixedWindowResult{}, fmt.Errorf("fixed window eval: %w", err)
}
return FixedWindowResult{
Allowed: res[0].(int64) == 1,
Remaining: int(res[1].(int64)),
ResetAt: (window + 1) * int64(windowSec),
}, nil
}
using StackExchange.Redis;
namespace RateLimit;
// FixedWindowResult holds the result of a fixed window check.
public readonly record struct FixedWindowResult(bool Allowed, int Remaining, long ResetAt);
// FixedWindowLimiter implements the fixed window counter algorithm.
public sealed class FixedWindowLimiter(IDatabase db)
{
private const string Script =
"""
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local windowSec = tonumber(ARGV[2])
local current = tonumber(redis.call('GET', key) or '0')
if current < limit then
redis.call('INCR', key)
redis.call('EXPIRE', key, windowSec)
return {1, limit - current - 1}
else
return {0, 0}
end
""";
/// <param name="limit">Max requests per window</param>
/// <param name="windowSec">Window size in seconds</param>
public async Task<FixedWindowResult> AllowAsync(string key, int limit, int windowSec)
{
var now = DateTimeOffset.UtcNow.ToUnixTimeSeconds();
var window = now / windowSec;
var redisKey = $"fixed_window:{key}:{window}";
var result = (RedisResult[])(await db.ScriptEvaluateAsync(
Script,
[redisKey],
[limit, windowSec]))!;
return new FixedWindowResult(
Allowed: (int)result[0] == 1,
Remaining: (int)result[1],
ResetAt: (window + 1) * windowSec);
}
}
import time
from dataclasses import dataclass
import redis.asyncio as redis
FIXED_WINDOW_SCRIPT = """
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local windowSec = tonumber(ARGV[2])
local current = tonumber(redis.call('GET', key) or '0')
if current < limit then
redis.call('INCR', key)
redis.call('EXPIRE', key, windowSec)
return {1, limit - current - 1}
else
return {0, 0}
end
"""
@dataclass(frozen=True, slots=True)
class FixedWindowResult:
allowed: bool
remaining: int
reset_at: int
class FixedWindowLimiter:
def __init__(self, client: redis.Redis) -> None:
self._script = client.register_script(FIXED_WINDOW_SCRIPT)
async def allow(
self,
key: str,
limit: int, # max requests per window
window_sec: int, # window size in seconds
) -> FixedWindowResult:
window = int(time.time()) // window_sec
redis_key = f"fixed_window:{key}:{window}"
allowed, remaining = await self._script(
keys=[redis_key],
args=[limit, window_sec],
)
return FixedWindowResult(
allowed=bool(allowed),
remaining=int(remaining),
reset_at=(window + 1) * window_sec,
)
**Плюсы:** простой, мало памяти.
**Минусы:** edge case на границе окна (можно получить 2x трафика).
4.4 Sliding Window Log
Точный, но ресурсоёмкий: хранит timestamp каждого запроса.
<?php
declare(strict_types=1);
final class SlidingWindowLogLimiter
{
public function __construct(
private readonly \Redis $redis,
) {}
public function allow(string $key, int $limit, int $windowSec): array
{
$now = microtime(true);
$redisKey = "sliding_log:{$key}";
$windowStart = $now - $windowSec;
$script = <<<'LUA'
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local now = tonumber(ARGV[2])
local windowStart = tonumber(ARGV[3])
local windowSec = tonumber(ARGV[4])
-- Remove expired entries
redis.call('ZREMRANGEBYSCORE', key, '-inf', windowStart)
-- Count current requests in window
local count = redis.call('ZCARD', key)
if count < limit then
-- Add current request
redis.call('ZADD', key, now, now .. ':' .. math.random(1000000))
redis.call('EXPIRE', key, windowSec)
return {1, limit - count - 1}
else
redis.call('EXPIRE', key, windowSec)
return {0, 0}
end
LUA;
$result = $this->redis->eval(
$script,
[$redisKey, $limit, $now, $windowStart, $windowSec],
1,
);
return [
'allowed' => (bool) $result[0],
'remaining' => (int) $result[1],
];
}
}
package ratelimit
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// SlidingWindowLogResult holds the result of a sliding window log check.
type SlidingWindowLogResult struct {
Allowed bool
Remaining int
}
// SlidingWindowLogLimiter implements the sliding window log algorithm.
type SlidingWindowLogLimiter struct {
rdb *redis.Client
}
func NewSlidingWindowLogLimiter(rdb *redis.Client) *SlidingWindowLogLimiter {
return &SlidingWindowLogLimiter{rdb: rdb}
}
// Allow checks if a request is allowed under the sliding window log algorithm.
func (l *SlidingWindowLogLimiter) Allow(ctx context.Context, key string, limit int, windowSec int) (SlidingWindowLogResult, error) {
now := float64(time.Now().UnixMicro()) / 1e6
redisKey := "sliding_log:" + key
windowStart := now - float64(windowSec)
script := redis.NewScript(`
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local now = tonumber(ARGV[2])
local windowStart = tonumber(ARGV[3])
local windowSec = tonumber(ARGV[4])
redis.call('ZREMRANGEBYSCORE', key, '-inf', windowStart)
local count = redis.call('ZCARD', key)
if count < limit then
redis.call('ZADD', key, now, now .. ':' .. math.random(1000000))
redis.call('EXPIRE', key, windowSec)
return {1, limit - count - 1}
else
redis.call('EXPIRE', key, windowSec)
return {0, 0}
end
`)
res, err := script.Run(ctx, l.rdb, []string{redisKey},
limit, fmt.Sprintf("%.6f", now), fmt.Sprintf("%.6f", windowStart), windowSec,
).Slice()
if err != nil {
return SlidingWindowLogResult{}, fmt.Errorf("sliding window log eval: %w", err)
}
return SlidingWindowLogResult{
Allowed: res[0].(int64) == 1,
Remaining: int(res[1].(int64)),
}, nil
}
using System.Globalization;
using StackExchange.Redis;
namespace RateLimit;
// SlidingWindowLogResult holds the result of a sliding window log check.
public readonly record struct SlidingWindowLogResult(bool Allowed, int Remaining);
// SlidingWindowLogLimiter implements the sliding window log algorithm.
public sealed class SlidingWindowLogLimiter(IDatabase db)
{
private const string Script =
"""
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local now = tonumber(ARGV[2])
local windowStart = tonumber(ARGV[3])
local windowSec = tonumber(ARGV[4])
-- Remove expired entries
redis.call('ZREMRANGEBYSCORE', key, '-inf', windowStart)
local count = redis.call('ZCARD', key)
if count < limit then
redis.call('ZADD', key, now, now .. ':' .. math.random(1000000))
redis.call('EXPIRE', key, windowSec)
return {1, limit - count - 1}
else
redis.call('EXPIRE', key, windowSec)
return {0, 0}
end
""";
public async Task<SlidingWindowLogResult> AllowAsync(string key, int limit, int windowSec)
{
var now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds() / 1000.0;
var redisKey = $"sliding_log:{key}";
var windowStart = now - windowSec;
var result = (RedisResult[])(await db.ScriptEvaluateAsync(
Script,
[redisKey],
[
limit,
now.ToString("F6", CultureInfo.InvariantCulture),
windowStart.ToString("F6", CultureInfo.InvariantCulture),
windowSec,
]))!;
return new SlidingWindowLogResult(
Allowed: (int)result[0] == 1,
Remaining: (int)result[1]);
}
}
import time
from dataclasses import dataclass
import redis.asyncio as redis
SLIDING_LOG_SCRIPT = """
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local now = tonumber(ARGV[2])
local windowStart = tonumber(ARGV[3])
local windowSec = tonumber(ARGV[4])
-- Remove expired entries
redis.call('ZREMRANGEBYSCORE', key, '-inf', windowStart)
local count = redis.call('ZCARD', key)
if count < limit then
redis.call('ZADD', key, now, now .. ':' .. math.random(1000000))
redis.call('EXPIRE', key, windowSec)
return {1, limit - count - 1}
else
redis.call('EXPIRE', key, windowSec)
return {0, 0}
end
"""
@dataclass(frozen=True, slots=True)
class SlidingWindowLogResult:
allowed: bool
remaining: int
class SlidingWindowLogLimiter:
def __init__(self, client: redis.Redis) -> None:
self._script = client.register_script(SLIDING_LOG_SCRIPT)
async def allow(self, key: str, limit: int, window_sec: int) -> SlidingWindowLogResult:
now = time.time()
redis_key = f"sliding_log:{key}"
window_start = now - window_sec
allowed, remaining = await self._script(
keys=[redis_key],
args=[limit, f"{now:.6f}", f"{window_start:.6f}", window_sec],
)
return SlidingWindowLogResult(allowed=bool(allowed), remaining=int(remaining))
### 4.5 Sliding Window Counter (гибридный)
Комбинация Fixed Window и Sliding Window. Оптимальный баланс точности и ресурсов.
<?php
declare(strict_types=1);
final class SlidingWindowCounterLimiter
{
public function __construct(
private readonly \Redis $redis,
) {}
public function allow(string $key, int $limit, int $windowSec): array
{
$now = time();
$currentWindow = (int) ($now / $windowSec);
$previousWindow = $currentWindow - 1;
$currentKey = "swc:{$key}:{$currentWindow}";
$previousKey = "swc:{$key}:{$previousWindow}";
$script = <<<'LUA'
local currentKey = KEYS[1]
local previousKey = KEYS[2]
local limit = tonumber(ARGV[1])
local windowSec = tonumber(ARGV[2])
local elapsedRatio = tonumber(ARGV[3])
local currentCount = tonumber(redis.call('GET', currentKey) or '0')
local previousCount = tonumber(redis.call('GET', previousKey) or '0')
-- Weighted count: previous window * remaining ratio + current count
local remainingRatio = 1 - elapsedRatio
local estimatedCount = math.floor(previousCount * remainingRatio) + currentCount
if estimatedCount < limit then
redis.call('INCR', currentKey)
redis.call('EXPIRE', currentKey, windowSec * 2)
return {1, limit - estimatedCount - 1}
else
return {0, 0}
end
LUA;
$elapsedInCurrentWindow = $now % $windowSec;
$elapsedRatio = $elapsedInCurrentWindow / $windowSec;
$result = $this->redis->eval(
$script,
[$currentKey, $previousKey, $limit, $windowSec, $elapsedRatio],
2,
);
return [
'allowed' => (bool) $result[0],
'remaining' => (int) $result[1],
'resetAt' => ($currentWindow + 1) * $windowSec,
];
}
}
package ratelimit
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// SlidingWindowCounterResult holds the result of a sliding window counter check.
type SlidingWindowCounterResult struct {
Allowed bool
Remaining int
ResetAt int64
}
// SlidingWindowCounterLimiter implements the hybrid sliding window counter.
type SlidingWindowCounterLimiter struct {
rdb *redis.Client
}
func NewSlidingWindowCounterLimiter(rdb *redis.Client) *SlidingWindowCounterLimiter {
return &SlidingWindowCounterLimiter{rdb: rdb}
}
// Allow checks if a request is allowed under the sliding window counter.
func (l *SlidingWindowCounterLimiter) Allow(ctx context.Context, key string, limit int, windowSec int) (SlidingWindowCounterResult, error) {
now := time.Now().Unix()
currentWindow := now / int64(windowSec)
previousWindow := currentWindow - 1
currentKey := fmt.Sprintf("swc:%s:%d", key, currentWindow)
previousKey := fmt.Sprintf("swc:%s:%d", key, previousWindow)
elapsedInWindow := now % int64(windowSec)
elapsedRatio := float64(elapsedInWindow) / float64(windowSec)
script := redis.NewScript(`
local currentKey = KEYS[1]
local previousKey = KEYS[2]
local limit = tonumber(ARGV[1])
local windowSec = tonumber(ARGV[2])
local elapsedRatio = tonumber(ARGV[3])
local currentCount = tonumber(redis.call('GET', currentKey) or '0')
local previousCount = tonumber(redis.call('GET', previousKey) or '0')
local remainingRatio = 1 - elapsedRatio
local estimatedCount = math.floor(previousCount * remainingRatio) + currentCount
if estimatedCount < limit then
redis.call('INCR', currentKey)
redis.call('EXPIRE', currentKey, windowSec * 2)
return {1, limit - estimatedCount - 1}
else
return {0, 0}
end
`)
res, err := script.Run(ctx, l.rdb, []string{currentKey, previousKey},
limit, windowSec, fmt.Sprintf("%.6f", elapsedRatio),
).Slice()
if err != nil {
return SlidingWindowCounterResult{}, fmt.Errorf("sliding window counter eval: %w", err)
}
return SlidingWindowCounterResult{
Allowed: res[0].(int64) == 1,
Remaining: int(res[1].(int64)),
ResetAt: (currentWindow + 1) * int64(windowSec),
}, nil
}
using System.Globalization;
using StackExchange.Redis;
namespace RateLimit;
// SlidingWindowCounterResult holds the result of a sliding window counter check.
public readonly record struct SlidingWindowCounterResult(bool Allowed, int Remaining, long ResetAt);
// SlidingWindowCounterLimiter implements the hybrid sliding window counter.
public sealed class SlidingWindowCounterLimiter(IDatabase db)
{
private const string Script =
"""
local currentKey = KEYS[1]
local previousKey = KEYS[2]
local limit = tonumber(ARGV[1])
local windowSec = tonumber(ARGV[2])
local elapsedRatio = tonumber(ARGV[3])
local currentCount = tonumber(redis.call('GET', currentKey) or '0')
local previousCount = tonumber(redis.call('GET', previousKey) or '0')
-- Weighted count: previous window * remaining ratio + current count
local remainingRatio = 1 - elapsedRatio
local estimatedCount = math.floor(previousCount * remainingRatio) + currentCount
if estimatedCount < limit then
redis.call('INCR', currentKey)
redis.call('EXPIRE', currentKey, windowSec * 2)
return {1, limit - estimatedCount - 1}
else
return {0, 0}
end
""";
public async Task<SlidingWindowCounterResult> AllowAsync(string key, int limit, int windowSec)
{
var now = DateTimeOffset.UtcNow.ToUnixTimeSeconds();
var currentWindow = now / windowSec;
var previousWindow = currentWindow - 1;
var currentKey = $"swc:{key}:{currentWindow}";
var previousKey = $"swc:{key}:{previousWindow}";
var elapsedRatio = (double)(now % windowSec) / windowSec;
var result = (RedisResult[])(await db.ScriptEvaluateAsync(
Script,
[currentKey, previousKey],
[limit, windowSec, elapsedRatio.ToString("F6", CultureInfo.InvariantCulture)]))!;
return new SlidingWindowCounterResult(
Allowed: (int)result[0] == 1,
Remaining: (int)result[1],
ResetAt: (currentWindow + 1) * windowSec);
}
}
import time
from dataclasses import dataclass
import redis.asyncio as redis
SLIDING_COUNTER_SCRIPT = """
local currentKey = KEYS[1]
local previousKey = KEYS[2]
local limit = tonumber(ARGV[1])
local windowSec = tonumber(ARGV[2])
local elapsedRatio = tonumber(ARGV[3])
local currentCount = tonumber(redis.call('GET', currentKey) or '0')
local previousCount = tonumber(redis.call('GET', previousKey) or '0')
-- Weighted count: previous window * remaining ratio + current count
local remainingRatio = 1 - elapsedRatio
local estimatedCount = math.floor(previousCount * remainingRatio) + currentCount
if estimatedCount < limit then
redis.call('INCR', currentKey)
redis.call('EXPIRE', currentKey, windowSec * 2)
return {1, limit - estimatedCount - 1}
else
return {0, 0}
end
"""
@dataclass(frozen=True, slots=True)
class SlidingWindowCounterResult:
allowed: bool
remaining: int
reset_at: int
class SlidingWindowCounterLimiter:
def __init__(self, client: redis.Redis) -> None:
self._script = client.register_script(SLIDING_COUNTER_SCRIPT)
async def allow(self, key: str, limit: int, window_sec: int) -> SlidingWindowCounterResult:
now = int(time.time())
current_window = now // window_sec
previous_window = current_window - 1
current_key = f"swc:{key}:{current_window}"
previous_key = f"swc:{key}:{previous_window}"
elapsed_ratio = (now % window_sec) / window_sec
allowed, remaining = await self._script(
keys=[current_key, previous_key],
args=[limit, window_sec, f"{elapsed_ratio:.6f}"],
)
return SlidingWindowCounterResult(
allowed=bool(allowed),
remaining=int(remaining),
reset_at=(current_window + 1) * window_sec,
)
## Сравнение алгоритмов
Алгоритм
Память
Точность
Burst
Сложность
Token Bucket
O(1)
Высокая
Да
Средняя
Leaky Bucket
O(1)
Высокая
Нет
Средняя
Fixed Window
O(1)
Низкая (edge)
Да (2x)
Низкая
Sliding Log
O(N)
Идеальная
Нет
Высокая
Sliding Counter
O(1)
Высокая
Частично
Средняя
Шаг 5: Middleware приложения
<?php
declare(strict_types=1);
final class RateLimitMiddleware
{
public function __construct(
private readonly TokenBucketLimiter $limiter,
private readonly RateLimitConfig $config,
) {}
public function handle(Request $request, callable $next): Response
{
$key = $this->resolveKey($request);
$rule = $this->config->getRuleForPath($request->getPathInfo());
$result = $this->limiter->allow($key, $rule->capacity, $rule->rate);
if (!$result['allowed']) {
return new Response(
json_encode(['error' => 'Rate limit exceeded']),
429,
[
'Content-Type' => 'application/json',
'X-RateLimit-Limit' => (string) $rule->capacity,
'X-RateLimit-Remaining' => '0',
'X-RateLimit-Reset' => (string) ($result['resetAt'] ?? time() + 60),
'Retry-After' => (string) ceil($result['retryAfter']),
],
);
}
$response = $next($request);
// Add rate limit headers to successful responses
return $response->withHeaders([
'X-RateLimit-Limit' => (string) $rule->capacity,
'X-RateLimit-Remaining' => (string) $result['remaining'],
'X-RateLimit-Reset' => (string) ($result['resetAt'] ?? time() + 60),
]);
}
private function resolveKey(Request $request): string
{
// Priority: API Key > User ID > IP
if ($apiKey = $request->headers->get('X-API-Key')) {
return "api_key:{$apiKey}";
}
if ($userId = $request->attributes->get('user_id')) {
return "user:{$userId}";
}
return "ip:{$request->getClientIp()}";
}
}
final readonly class RateLimitRule
{
public function __construct(
public int $capacity, // Max requests (burst)
public float $rate, // Refill rate per second
public string $keyType, // ip, user, api_key
) {}
}
final class RateLimitConfig
{
/** @var array<string, RateLimitRule> */
private array $rules = [];
public function addRule(string $pathPattern, RateLimitRule $rule): void
{
$this->rules[$pathPattern] = $rule;
}
public function getRuleForPath(string $path): RateLimitRule
{
foreach ($this->rules as $pattern => $rule) {
if (preg_match($pattern, $path)) {
return $rule;
}
}
// Default rule
return new RateLimitRule(
capacity: 100,
rate: 10.0,
keyType: 'ip',
);
}
}
using System.Globalization;
using System.Security.Claims;
using System.Text.RegularExpressions;
using StackExchange.Redis;
namespace RateLimit;
// RateLimitRule defines a rate limit rule for a path pattern.
public sealed record RateLimitRule(
int Capacity, // Max requests (burst)
double Rate, // Refill rate per second
string KeyType); // ip, user, api_key
// RateLimitConfig stores rate limit rules per path pattern.
public sealed class RateLimitConfig
{
private readonly List<(Regex Pattern, RateLimitRule Rule)> _rules = [];
public void AddRule(string pathPattern, RateLimitRule rule) =>
_rules.Add((new Regex(pathPattern, RegexOptions.Compiled), rule));
public RateLimitRule GetRuleForPath(string path)
{
foreach (var (pattern, rule) in _rules)
{
if (pattern.IsMatch(path))
{
return rule;
}
}
// Default rule
return new RateLimitRule(Capacity: 100, Rate: 10.0, KeyType: "ip");
}
}
// RateLimitMiddleware enforces rate limits on the ASP.NET Core pipeline.
public sealed class RateLimitMiddleware(
RequestDelegate next,
TokenBucketLimiter limiter,
RateLimitConfig config)
{
public async Task InvokeAsync(HttpContext context)
{
var key = ResolveKey(context);
var rule = config.GetRuleForPath(context.Request.Path);
TokenBucketResult result;
try
{
result = await limiter.AllowAsync(key, rule.Capacity, rule.Rate);
}
catch (RedisException)
{
// Fail-open: allow request if rate limiter is down
await next(context);
return;
}
var resetAt = DateTimeOffset.UtcNow.AddSeconds(60).ToUnixTimeSeconds();
var headers = context.Response.Headers;
headers["X-RateLimit-Limit"] = rule.Capacity.ToString(CultureInfo.InvariantCulture);
headers["X-RateLimit-Remaining"] = result.Remaining.ToString(CultureInfo.InvariantCulture);
headers["X-RateLimit-Reset"] = resetAt.ToString(CultureInfo.InvariantCulture);
if (!result.Allowed)
{
headers["Retry-After"] = ((int)Math.Ceiling(result.RetryAfter))
.ToString(CultureInfo.InvariantCulture);
context.Response.StatusCode = StatusCodes.Status429TooManyRequests;
await context.Response.WriteAsJsonAsync(new { error = "Rate limit exceeded" });
return;
}
await next(context);
}
private static string ResolveKey(HttpContext context)
{
// Priority: API Key > User ID > IP
if (context.Request.Headers.TryGetValue("X-API-Key", out var apiKey) && apiKey.Count > 0)
{
return $"api_key:{apiKey}";
}
var userId = context.User.FindFirstValue(ClaimTypes.NameIdentifier);
if (!string.IsNullOrEmpty(userId))
{
return $"user:{userId}";
}
return $"ip:{context.Connection.RemoteIpAddress}";
}
}
import math
import re
import time
from dataclasses import dataclass
from redis.exceptions import RedisError
from starlette.middleware.base import BaseHTTPMiddleware, RequestResponseEndpoint
from starlette.requests import Request
from starlette.responses import JSONResponse, Response
from starlette.types import ASGIApp
@dataclass(frozen=True, slots=True)
class RateLimitRule:
capacity: int # max requests (burst)
rate: float # refill rate per second
key_type: str # ip, user, api_key
class RateLimitConfig:
def __init__(self) -> None:
self._rules: list[tuple[re.Pattern[str], RateLimitRule]] = []
def add_rule(self, path_pattern: str, rule: RateLimitRule) -> None:
self._rules.append((re.compile(path_pattern), rule))
def rule_for_path(self, path: str) -> RateLimitRule:
for pattern, rule in self._rules:
if pattern.search(path):
return rule
# Default rule
return RateLimitRule(capacity=100, rate=10.0, key_type="ip")
class RateLimitMiddleware(BaseHTTPMiddleware):
def __init__(self, app: ASGIApp, limiter: TokenBucketLimiter, config: RateLimitConfig) -> None:
super().__init__(app)
self._limiter = limiter
self._config = config
async def dispatch(self, request: Request, call_next: RequestResponseEndpoint) -> Response:
key = self._resolve_key(request)
rule = self._config.rule_for_path(request.url.path)
try:
result = await self._limiter.allow(key, rule.capacity, rule.rate)
except RedisError:
# Fail-open: allow request if rate limiter is down
return await call_next(request)
headers = {
"X-RateLimit-Limit": str(rule.capacity),
"X-RateLimit-Remaining": str(result.remaining),
"X-RateLimit-Reset": str(int(time.time()) + 60),
}
if not result.allowed:
return JSONResponse(
{"error": "Rate limit exceeded"},
status_code=429,
headers=headers | {"Retry-After": str(math.ceil(result.retry_after))},
)
response = await call_next(request)
response.headers.update(headers)
return response
@staticmethod
def _resolve_key(request: Request) -> str:
# Priority: API Key > User ID > IP
api_key = request.headers.get("X-API-Key")
if api_key:
return f"api_key:{api_key}"
user_id = request.headers.get("X-User-Id")
if user_id:
return f"user:{user_id}"
client = request.client
return f"ip:{client.host if client else 'unknown'}"
## Шаг 6: Distributed Rate Limiting
Проблема: несколько серверов
Client -> LB -> Server 1 (limit: 50 of 100) -> Redis
└──-> Server 2 (limit: 50 of 100) -> Redis (same key!)
Redis решает проблему: все серверы используют один Redis для состояния.
Race conditions
Lua-скрипты в Redis атомарны -- нет race conditions. Это главная причина использования Lua, а не отдельных GET/SET.
Синхронизация при Redis Cluster
<?php
declare(strict_types=1);
final class DistributedRateLimiter
{
/** @var \Redis[] */
private array $redisNodes;
public function __construct(array $redisNodes)
{
$this->redisNodes = $redisNodes;
}
/**
* Use consistent hashing to route to the correct Redis node
*/
public function allow(string $key, int $limit, int $windowSec): array
{
$node = $this->getNodeForKey($key);
return $this->checkLimit($node, $key, $limit, $windowSec);
}
private function getNodeForKey(string $key): \Redis
{
$hash = crc32($key);
$index = $hash % count($this->redisNodes);
return $this->redisNodes[$index];
}
private function checkLimit(\Redis $redis, string $key, int $limit, int $windowSec): array
{
// Same Lua script as before, executed on the specific node
// ...
return ['allowed' => true, 'remaining' => $limit];
}
}
using System.IO.Hashing;
using System.Text;
using StackExchange.Redis;
namespace RateLimit;
// DistributedRateLimiter routes rate limit checks to the correct Redis node.
public sealed class DistributedRateLimiter
{
private readonly IReadOnlyList<FixedWindowLimiter> _limiters;
public DistributedRateLimiter(IReadOnlyList<IDatabase> nodes) =>
_limiters = [.. nodes.Select(node => new FixedWindowLimiter(node))];
// Routes the key to the appropriate Redis node via consistent hashing.
public Task<FixedWindowResult> AllowAsync(string key, int limit, int windowSec) =>
GetLimiterForKey(key).AllowAsync(key, limit, windowSec);
private FixedWindowLimiter GetLimiterForKey(string key)
{
// Crc32 lives in the System.IO.Hashing package -- the BCL has no built-in CRC32
var hash = Crc32.HashToUInt32(Encoding.UTF8.GetBytes(key));
return _limiters[(int)(hash % (uint)_limiters.Count)];
}
}
import zlib
import redis.asyncio as redis
class DistributedRateLimiter:
"""Routes rate limit checks to the correct Redis node."""
def __init__(self, nodes: list[redis.Redis]) -> None:
if not nodes:
raise ValueError("at least one Redis node is required")
self._limiters = [FixedWindowLimiter(node) for node in nodes]
async def allow(self, key: str, limit: int, window_sec: int) -> FixedWindowResult:
# Consistent hashing picks the node that owns this key
return await self._limiter_for_key(key).allow(key, limit, window_sec)
def _limiter_for_key(self, key: str) -> FixedWindowLimiter:
# zlib.crc32 already returns an unsigned 32-bit value
index = zlib.crc32(key.encode()) % len(self._limiters)
return self._limiters[index]
## Возможные вопросы интервьюера
Как обрабатывать ситуацию, когда Redis недоступен?
Fail-open: пропускаем все запросы (лучше доступность)
Fail-closed: блокируем все (лучше безопасность)
Local fallback: in-memory rate limiter на каждом сервере
Как реализовать разные лимиты для разных тарифов?
Конфигурация правил в базе данных
Ключ кэша включает tier: rate_limit:{user_id}:{tier}
Как избежать hot key в Redis?
Шардирование по ключам
Local cache для часто проверяемых клиентов
Pipelining для batch проверок
Rate limiting на уровне API Gateway vs Application?
Gateway: грубая защита (IP-based), DDoS
Application: точная (user-based), бизнес-правила
Как тестировать rate limiter?
Unit-тесты с mock Redis
Integration-тесты с реальным Redis
Load testing для проверки точности при высоком QPS