api-database-redis
Redis in-memory data store patterns with ioredis and node-redis -- caching, sessions, rate limiting, pub/sub, streams, queues, transactions, cluster
Install
npx skills add https://github.com/agents-inc/skills/tree/main/dist/plugins/api-database-redis/skills/api-database-redis
claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install agents-inc-skills@llmmart
git clone https://github.com/agents-inc/skills.git
The skills CLI installs just this skill, for any of its supported agents. Claude Code installs the whole agents-inc/skills collection as a plugin from our marketplace. Git is the plain clone.
Skill manifest
Redis Patterns
Quick Guide: Use Redis as an in-memory data store for caching, session management, rate limiting, pub/sub messaging, and job queues. Use ioredis (v5.x) as the primary client for its superior TypeScript support, Cluster/Sentinel integration, auto-pipelining, and Lua scripting. Use node-redis (v5.x) only when you need Redis Stack modules (JSON, Search, TimeSeries). Always set
maxRetriesPerRequest: nullfor BullMQ workers, use separate connections for Pub/Sub subscribers, and define Lua scripts viadefineCommandfor atomic multi-step operations.
<critical_requirements>
CRITICAL: Before Using This Skill
All code must follow project conventions in CLAUDE.md (kebab-case, named exports, import ordering,
import type, named constants)
(You MUST use a SEPARATE Redis connection for Pub/Sub subscribers -- a subscribed connection enters a special mode and cannot execute other commands)
(You MUST set maxRetriesPerRequest: null on any ioredis connection passed to BullMQ -- BullMQ requires infinite retries and will throw if this is not set)
(You MUST use Lua scripts (defineCommand or eval) for any operation requiring atomicity across multiple Redis commands -- separate commands are NOT atomic even in a pipeline)
(You MUST handle the error event on every Redis client instance -- unhandled errors crash the Node.js process)
</critical_requirements>
Examples
- Core Patterns -- ioredis/node-redis connection, error handling, reconnection, cluster, sentinel, pipelining, transactions
- Caching Patterns -- Cache-aside, write-through, invalidation, stampede prevention, multi-key pipeline
- Data Structures -- Strings, hashes, lists, sets, sorted sets with typed helpers
- Sessions -- Express connect-redis (node-redis required for v9+), Hono manual middleware
- Pub/Sub -- Publish/subscribe, event broadcasting, pattern subscriptions
- Rate Limiting -- Sliding window (Lua), token bucket (Lua), middleware integration
- Queues & Locks -- BullMQ job queues, Redis Streams with consumer groups, distributed locks
Additional resources:
- reference.md -- Command cheat sheet, connection options, anti-patterns, production checklist
Auto-detection: Redis, ioredis, node-redis, createClient, RedisStore, BullMQ, Queue, Worker, pub/sub, MULTI, EXEC, pipeline, Lua script, defineCommand, xadd, xread, cache-aside, rate limit, session store, connect-redis, Redis.Cluster, Sentinel
When to use:
- Caching database queries or API responses (cache-aside, write-through)
- Session storage for Express/Hono/Fastify applications
- Distributed rate limiting (sliding window, token bucket)
- Real-time messaging with Pub/Sub
- Background job processing with BullMQ queues
- Leaderboards, counters, and real-time analytics with sorted sets
- Distributed locks and atomic operations with Lua scripts
Key patterns covered:
- ioredis connection setup, configuration, and error handling
- Data structures (strings, hashes, lists, sets, sorted sets, streams)
- Cache-aside and write-through caching with TTL management
- Session storage with connect-redis
- Rate limiting with Lua scripts (sliding window, token bucket)
- Pub/Sub messaging with separate connections
- Redis Streams for persistent message queues
- BullMQ for job queues with retries and scheduling
- Pipelining and transactions (MULTI/EXEC)
- Lua scripting for atomic operations
- Cluster mode and Sentinel for high availability
When NOT to use:
- Primary database for relational data (use your relational database)
- Document storage with complex queries (use a document database)
- Large binary file storage (use S3/object storage)
- Data that must survive total memory loss without persistence configured
<decision_framework>
Decision Framework
Which Redis Client?
Which Redis client should I use?
├─ Need Redis Stack modules (JSON, Search, TimeSeries)? -> node-redis (v5.x)
├─ Using BullMQ for job queues? -> ioredis (BullMQ requires it)
├─ Need Cluster or Sentinel support? -> ioredis (built-in, battle-tested)
├─ Need auto-pipelining? -> ioredis (enableAutoPipelining option)
└─ General caching/sessions/pub-sub? -> ioredis (recommended default)
Which Caching Strategy?
How should I cache this data?
├─ Read-heavy, tolerates brief staleness? -> Cache-aside with TTL
├─ Needs strong consistency after writes? -> Write-through (update DB + invalidate cache)
├─ Write-heavy, can tolerate brief data loss? -> Write-behind (async cache update)
└─ Data changes rarely? -> Cache-aside with long TTL + manual invalidation
Which Data Structure?
What Redis data structure should I use?
├─ Simple key-value (cache, sessions)? -> Strings (GET/SET)
├─ Object with multiple fields? -> Hashes (HSET/HGET)
├─ Ordered ranking/leaderboard? -> Sorted Sets (ZADD/ZRANGE)
├─ Queue (FIFO/LIFO)? -> Lists (LPUSH/RPOP)
├─ Unique collection (tags, categories)? -> Sets (SADD/SMEMBERS)
├─ Persistent message log with consumers? -> Streams (XADD/XREAD)
└─ Rate limiting (sliding window)? -> Sorted Sets + Lua script
Which Messaging Pattern?
How should I implement real-time messaging?
├─ Fire-and-forget broadcast? -> Pub/Sub (no persistence)
├─ Need message persistence and replay? -> Streams with consumer groups
├─ Need reliable job processing with retries? -> BullMQ (built on Redis)
└─ Need request-reply pattern? -> Pub/Sub with correlation IDs
Atomicity Decision
Do I need atomicity across multiple commands?
├─ YES -> Are the commands on the same key?
│ ├─ YES -> Use a single atomic command (INCR, SETNX, etc.)
│ └─ NO -> Use Lua script (defineCommand)
├─ NO, but I want batching -> Use pipeline (non-atomic, single round-trip)
└─ Need optimistic locking? -> Use WATCH + MULTI/EXEC
</decision_framework>
<red_flags>
RED FLAGS
High Priority Issues:
- Using the same connection for Pub/Sub subscribe and regular commands -- subscribed connections cannot execute non-pub/sub commands
- Missing
maxRetriesPerRequest: nullon BullMQ connections -- BullMQ throws immediately without this setting - Using
KEYScommand in production -- blocks the entire Redis server while scanning all keys - No
errorevent handler on Redis client -- unhandled errors crash the Node.js process - Storing large objects (> 1 MB) in Redis -- degrades performance and wastes memory; store a reference and fetch from object storage
Medium Priority Issues:
- Missing TTL on cached keys -- causes unbounded memory growth until Redis runs out of memory
- Using
delwith many keys instead ofunlink--delblocks Redis;unlinkfrees memory asynchronously - Not using pipelining for batch operations -- each command is a separate network round-trip
- Serializing/deserializing complex objects without error handling -- malformed JSON in cache crashes on parse
- Sharing a single Redis connection across BullMQ Queue and Worker -- each needs its own connection
Common Mistakes:
- Assuming pipeline commands are atomic -- pipelines batch for network efficiency but do not provide atomicity (use MULTI/EXEC or Lua)
- Forgetting that
hgetallreturns an empty object{}for non-existent keys (notnull) -- checkObject.keys(result).length === 0 - Using
MULTI/EXECwithoutWATCHfor conditional updates -- transactions execute unconditionally unless you WATCH keys first - Not handling
nullreturns fromGET-- cache misses returnnull, notundefined - Connecting to Redis without TLS in production -- credentials sent in plaintext over the network
Gotchas & Edge Cases:
- Redis
HGETALLreturns all values as strings -- numbers stored withHSETcome back as strings, requiring explicit parsing EXPIREresets when a key is overwritten withSET-- if youSETa key that already has a TTL, the TTL is removed unless you includeEX/PXin theSETcommand- Pub/Sub messages are fire-and-forget -- if no subscriber is listening when a message is published, it is lost forever (use Streams for persistence)
- Redis Cluster does not support multi-key operations across different hash slots -- use
{hash-tag}prefix to force related keys to the same slot WATCHis connection-scoped -- concurrent requests sharing a connection will interfere with each other's WATCH state- ioredis auto-pipelining does not work with
WATCH/MULTIor blocking commands (BRPOP,BLPOP,XREAD BLOCK)
</red_flags>
<critical_reminders>
CRITICAL REMINDERS
All code must follow project conventions in CLAUDE.md (kebab-case, named exports, import ordering,
import type, named constants)
(You MUST use a SEPARATE Redis connection for Pub/Sub subscribers -- a subscribed connection enters a special mode and cannot execute other commands)
(You MUST set maxRetriesPerRequest: null on any ioredis connection passed to BullMQ -- BullMQ requires infinite retries and will throw if this is not set)
(You MUST use Lua scripts (defineCommand or eval) for any operation requiring atomicity across multiple Redis commands -- separate commands are NOT atomic even in a pipeline)
(You MUST handle the error event on every Redis client instance -- unhandled errors crash the Node.js process)
Failure to follow these rules will cause pub/sub failures, BullMQ connection errors, race conditions, and application crashes.
</critical_reminders>
Files (skills)
-
examples
-
caching.md 6.5 KB
# Redis -- Caching Pattern Examples > Cache-aside, write-through, invalidation, stampede prevention, and multi-key cache patterns. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [core.md](core.md) -- Connection setup, pipelining, transactions - [data-structures.md](data-structures.md) -- Hashes, lists, sets, sorted sets - [sessions.md](sessions.md) -- Session storage patterns --- ## Cache-Aside with Generic Helper ```typescript import type Redis from "ioredis"; const DEFAULT_CACHE_TTL_SECONDS = 300; // 5 minutes interface CacheResult<T> { data: T; fromCache: boolean; } async function cacheAside<T>( redis: Redis, key: string, fetcher: () => Promise<T>, ttlSeconds: number = DEFAULT_CACHE_TTL_SECONDS, ): Promise<CacheResult<T>> { try { const cached = await redis.get(key); if (cached !== null) { return { data: JSON.parse(cached) as T, fromCache: true }; } } catch (err) { // Cache read failure -- proceed to fetch from source console.error(`Cache read failed for ${key}:`, (err as Error).message); } const data = await fetcher(); // Non-blocking cache write redis.set(key, JSON.stringify(data), "EX", ttlSeconds).catch((err) => { console.error(`Cache write failed for ${key}:`, err.message); }); return { data, fromCache: false }; } export { cacheAside }; export type { CacheResult }; ``` **Why good:** Graceful degradation on cache read failure, non-blocking write prevents cache failure from blocking response, typed return distinguishes cache hit from miss --- ## Write-Through Cache ```typescript import type Redis from "ioredis"; const PRODUCT_CACHE_PREFIX = "cache:product:"; const PRODUCT_CACHE_TTL_SECONDS = 600; // 10 minutes async function writeThrough<T extends { id: string }>( redis: Redis, entity: T, dbWriter: (entity: T) => Promise<T>, ): Promise<T> { // 1. Write to database first (source of truth) const saved = await dbWriter(entity); // 2. Update cache with fresh data const key = `${PRODUCT_CACHE_PREFIX}${saved.id}`; await redis.set(key, JSON.stringify(saved), "EX", PRODUCT_CACHE_TTL_SECONDS); return saved; } export { writeThrough }; ``` **When to use:** Write-heavy applications where consistency matters more than cache hit rate. **When not to use:** Read-heavy applications with rare writes -- cache-aside with TTL is simpler and sufficient. --- ## Cache Invalidation (Write-Through with Delete) ```typescript import type Redis from "ioredis"; const PRODUCT_CACHE_PREFIX = "cache:product:"; const PRODUCT_LIST_CACHE_KEY = "cache:products:list"; const PRODUCT_CACHE_TTL = 600; // 10 minutes async function updateProduct( redis: Redis, productId: string, updates: Partial<Product>, ): Promise<Product> { // 1. Update database (source of truth) const updated = await db .update(products) .set(updates) .where(eq(products.id, productId)) .returning(); // 2. Invalidate specific cache entry await redis.del(`${PRODUCT_CACHE_PREFIX}${productId}`); // 3. Invalidate list cache (stale after update) await redis.del(PRODUCT_LIST_CACHE_KEY); return updated[0]; } async function deleteProduct(redis: Redis, productId: string): Promise<void> { await db.delete(products).where(eq(products.id, productId)); // Invalidate specific cache key and list cache await redis.del(`${PRODUCT_CACHE_PREFIX}${productId}`); await redis.del(PRODUCT_LIST_CACHE_KEY); } export { updateProduct, deleteProduct }; ``` **Why good:** Database is always updated first (source of truth), both specific and list caches invalidated, uses explicit key deletion instead of `KEYS` pattern scan --- ## Multi-Key Cache with Pipeline ```typescript import type Redis from "ioredis"; const CACHE_PREFIX = "cache:user:"; const CACHE_TTL_SECONDS = 300; async function getMultipleFromCache( redis: Redis, userIds: string[], ): Promise<Map<string, string | null>> { const pipeline = redis.pipeline(); for (const id of userIds) { pipeline.get(`${CACHE_PREFIX}${id}`); } const results = await pipeline.exec(); if (!results) { throw new Error("Pipeline execution returned null"); } const cache = new Map<string, string | null>(); for (let i = 0; i < userIds.length; i++) { const [err, value] = results[i]; cache.set(userIds[i], err ? null : (value as string | null)); } return cache; } async function setMultipleInCache( redis: Redis, entries: Array<{ id: string; data: unknown }>, ): Promise<void> { const pipeline = redis.pipeline(); for (const entry of entries) { pipeline.set( `${CACHE_PREFIX}${entry.id}`, JSON.stringify(entry.data), "EX", CACHE_TTL_SECONDS, ); } await pipeline.exec(); } export { getMultipleFromCache, setMultipleInCache }; ``` **Why good:** Pipeline reduces multiple GET/SET calls to single round-trip, error handling per result, typed Map return --- ## Cache Stampede Prevention (Singleflight) Prevents multiple processes from regenerating the same cache entry simultaneously. ```typescript import type Redis from "ioredis"; const LOCK_TTL_SECONDS = 10; const LOCK_RETRY_DELAY_MS = 50; const LOCK_MAX_RETRIES = 20; async function cacheAsideWithLock<T>( redis: Redis, key: string, fetcher: () => Promise<T>, ttlSeconds: number, ): Promise<T> { // Check cache first const cached = await redis.get(key); if (cached !== null) { return JSON.parse(cached) as T; } // Try to acquire lock const lockKey = `lock:${key}`; const lockAcquired = await redis.set( lockKey, "1", "EX", LOCK_TTL_SECONDS, "NX", ); if (lockAcquired === "OK") { try { // We have the lock -- fetch and populate cache const data = await fetcher(); await redis.set(key, JSON.stringify(data), "EX", ttlSeconds); return data; } finally { await redis.del(lockKey); } } // Another process is fetching -- wait and retry cache read for (let i = 0; i < LOCK_MAX_RETRIES; i++) { await new Promise((resolve) => setTimeout(resolve, LOCK_RETRY_DELAY_MS)); const result = await redis.get(key); if (result !== null) { return JSON.parse(result) as T; } } // Lock expired, cache still empty -- fall back to direct fetch return fetcher(); } export { cacheAsideWithLock }; ``` **Why good:** NX flag prevents multiple processes from acquiring the lock, lock has TTL to prevent deadlocks if process crashes, fallback to direct fetch if lock holder fails --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
core.md 9.3 KB
# Redis -- Core Patterns > Connection setup, error handling, pipelining, transactions, cluster, sentinel, and fundamental configuration. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [caching.md](caching.md) -- Cache-aside, write-through, invalidation - [data-structures.md](data-structures.md) -- Strings, hashes, lists, sets, sorted sets - [queues.md](queues.md) -- BullMQ, Redis Streams, distributed locks --- ## ioredis Connection with Error Handling ```typescript import Redis from "ioredis"; const RETRY_DELAY_BASE_MS = 50; const RETRY_DELAY_MAX_MS = 2000; function createRedisClient(): Redis { const url = process.env.REDIS_URL; if (!url) { throw new Error("REDIS_URL environment variable is required"); } const client = new Redis(url, { maxRetriesPerRequest: 3, retryStrategy(times) { const delay = Math.min(times * RETRY_DELAY_BASE_MS, RETRY_DELAY_MAX_MS); return delay; }, lazyConnect: true, }); client.on("error", (err) => { console.error("Redis connection error:", err.message); }); client.on("connect", () => { console.log("Redis connected"); }); return client; } export { createRedisClient }; ``` **Why good:** Environment variable validation, named constants for retry delays, `lazyConnect` prevents connection before ready, error event handler prevents process crash, reconnection strategy with exponential backoff capped at max delay ```typescript // ❌ Bad Example - No error handling, hardcoded config import Redis from "ioredis"; const redis = new Redis("redis://localhost:6379"); // No error handler -- unhandled errors crash the process // No retry strategy -- uses default which may not suit your needs // Hardcoded connection string -- leaks in version control ``` **Why bad:** Missing error event handler crashes Node.js process on connection failure, hardcoded URL prevents environment-specific configuration, no retry strategy customization --- ## node-redis Connection (Alternative) Use node-redis when you need Redis Stack modules (JSON, Search, TimeSeries) or connect-redis v9+ for session storage. ```typescript import { createClient } from "redis"; async function createNodeRedisClient() { const url = process.env.REDIS_URL; if (!url) { throw new Error("REDIS_URL environment variable is required"); } const client = createClient({ url }); client.on("error", (err) => { console.error("Redis client error:", err.message); }); await client.connect(); return client; } export { createNodeRedisClient }; ``` **When to use:** Redis Stack modules (JSON, Search, TimeSeries), connect-redis v9+ session storage (v9 dropped ioredis support). --- ## Pipelining (Non-Atomic Batching) Batch multiple commands to reduce network round-trips. ```typescript import type Redis from "ioredis"; const USER_KEY_PREFIX = "user:"; const USER_TTL_SECONDS = 3600; async function cacheMultipleUsers( redis: Redis, users: Array<{ id: string; name: string; email: string }>, ): Promise<void> { const pipeline = redis.pipeline(); for (const user of users) { const key = `${USER_KEY_PREFIX}${user.id}`; pipeline.hset(key, { name: user.name, email: user.email }); pipeline.expire(key, USER_TTL_SECONDS); } const results = await pipeline.exec(); if (!results) { throw new Error("Pipeline execution returned null"); } // Check for errors in pipeline results for (const [err] of results) { if (err) { throw new Error(`Pipeline command failed: ${err.message}`); } } } export { cacheMultipleUsers }; ``` **Why good:** Single network round-trip for all commands, error checking on each result, named constants for prefix and TTL --- ## Transactions (MULTI/EXEC -- Atomic) ```typescript import type Redis from "ioredis"; const BALANCE_KEY_PREFIX = "balance:"; async function transferBalance( redis: Redis, fromUserId: string, toUserId: string, amount: number, ): Promise<boolean> { const fromKey = `${BALANCE_KEY_PREFIX}${fromUserId}`; const toKey = `${BALANCE_KEY_PREFIX}${toUserId}`; // WATCH for optimistic locking await redis.watch(fromKey); const currentBalance = await redis.get(fromKey); if (!currentBalance || parseFloat(currentBalance) < amount) { await redis.unwatch(); return false; // Insufficient balance } // MULTI/EXEC -- atomic execution const results = await redis .multi() .decrby(fromKey, amount) .incrby(toKey, amount) .exec(); // results is null if WATCH detected a change (optimistic lock failure) if (!results) { return false; // Retry needed -- another client modified the key } return true; } export { transferBalance }; ``` **Why good:** WATCH provides optimistic locking, MULTI/EXEC ensures atomicity, null check handles concurrent modification, clear return value for retry logic ```typescript // ❌ Bad Example - Non-atomic balance transfer await redis.decrby("balance:user1", 100); await redis.incrby("balance:user2", 100); // If the process crashes between these two commands, // money disappears from user1 but never reaches user2 ``` **Why bad:** Two separate commands are not atomic, crash between them causes data inconsistency, no optimistic locking for concurrent access --- ## Cluster Mode ioredis supports Redis Cluster for horizontal scaling and high availability. ```typescript import Redis from "ioredis"; const CLUSTER_RETRY_BASE_MS = 100; const CLUSTER_RETRY_MAX_MS = 2000; const MAX_REDIRECTIONS = 16; const cluster = new Redis.Cluster( [ { host: "redis-node-1", port: 6379 }, { host: "redis-node-2", port: 6379 }, { host: "redis-node-3", port: 6379 }, ], { clusterRetryStrategy(times) { return Math.min(times * CLUSTER_RETRY_BASE_MS, CLUSTER_RETRY_MAX_MS); }, maxRedirections: MAX_REDIRECTIONS, scaleReads: "slave", // Read from replicas, write to master redisOptions: { password: process.env.REDIS_PASSWORD, }, }, ); cluster.on("error", (err) => { console.error("Cluster error:", err.message); }); // Use cluster exactly like a regular Redis client await cluster.set("key", "value"); const value = await cluster.get("key"); export { cluster }; ``` **Why good:** Multiple seed nodes for discovery, `scaleReads: "slave"` offloads reads to replicas, retry strategy with backoff, password from environment variable --- ## Sentinel Setup ```typescript import Redis from "ioredis"; const sentinel = new Redis({ sentinels: [ { host: "sentinel-1", port: 26379 }, { host: "sentinel-2", port: 26379 }, { host: "sentinel-3", port: 26379 }, ], name: "mymaster", // Sentinel group name sentinelRetryStrategy(times) { return Math.min(times * 10, 1000); }, failoverDetector: true, // Detect failover and reconnect automatically }); sentinel.on("error", (err) => { console.error("Sentinel error:", err.message); }); export { sentinel }; ``` **Why good:** Multiple sentinel nodes for redundancy, `failoverDetector: true` for automatic master failover handling, retry strategy for sentinel connectivity --- ## Auto-Pipelining ioredis can automatically batch commands issued during the same event loop tick: ```typescript import Redis from "ioredis"; const redis = new Redis(process.env.REDIS_URL!, { enableAutoPipelining: true, }); // These three commands are automatically batched into one pipeline const [name, email, role] = await Promise.all([ redis.get("user:name"), redis.get("user:email"), redis.get("user:role"), ]); export { redis }; ``` **When to use:** High-throughput applications issuing many independent commands per request. Auto-pipelining reduces network round-trips without changing application code. **When NOT to use:** When commands depend on each other's results (sequential logic), or when using WATCH/MULTI for transactions. --- ## Scanning Instead of KEYS Never use `KEYS` in production -- it blocks the Redis server while scanning all keys. ```typescript import type Redis from "ioredis"; const SCAN_BATCH_SIZE = 100; async function findKeysByPattern( redis: Redis, pattern: string, ): Promise<string[]> { const allKeys: string[] = []; const stream = redis.scanStream({ match: pattern, count: SCAN_BATCH_SIZE, }); return new Promise((resolve, reject) => { stream.on("data", (keys: string[]) => { allKeys.push(...keys); }); stream.on("end", () => resolve(allKeys)); stream.on("error", (err) => reject(err)); }); } export { findKeysByPattern }; ``` **Why good:** `scanStream` iterates incrementally without blocking Redis, `count` is a hint for batch size (not a guarantee), stream-based API handles large keyspaces --- ## Key Expiration Strategies ```typescript const TTL_SHORT_SECONDS = 60; // 1 minute -- volatile data const TTL_MEDIUM_SECONDS = 300; // 5 minutes -- API response cache const TTL_LONG_SECONDS = 3600; // 1 hour -- user profiles const TTL_SESSION_SECONDS = 86400; // 24 hours -- sessions // SET with TTL (preferred -- atomic) await redis.set("key", "value", "EX", TTL_MEDIUM_SECONDS); // SET with millisecond TTL await redis.set("key", "value", "PX", 500); // SET only if key doesn't exist (distributed lock pattern) const acquired = await redis.set("lock:resource", "owner-id", "EX", 30, "NX"); // Returns "OK" if lock acquired, null if already locked // Update TTL on existing key await redis.expire("key", TTL_LONG_SECONDS); ``` --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
data-structures.md 5 KB
# Redis -- Data Structure Examples > Typed helpers for strings, hashes, lists, sets, and sorted sets. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [caching.md](caching.md) -- Cache-aside and write-through patterns - [core.md](core.md) -- Connection setup, pipelining, transactions - [rate-limiting.md](rate-limiting.md) -- Sorted sets for sliding window rate limiting --- ## Strings (Key-Value) ```typescript import type Redis from "ioredis"; const DEFAULT_TTL_SECONDS = 3600; // 1 hour async function setWithTTL( redis: Redis, key: string, value: string, ttlSeconds: number = DEFAULT_TTL_SECONDS, ): Promise<void> { await redis.set(key, value, "EX", ttlSeconds); } async function getOrNull(redis: Redis, key: string): Promise<string | null> { return redis.get(key); } export { setWithTTL, getOrNull }; ``` --- ## Hashes (Object-like) ```typescript import type Redis from "ioredis"; interface UserProfile { name: string; email: string; role: string; } const USER_KEY_PREFIX = "user:"; const USER_TTL_SECONDS = 1800; // 30 minutes async function setUserProfile( redis: Redis, userId: string, profile: UserProfile, ): Promise<void> { const key = `${USER_KEY_PREFIX}${userId}`; await redis.hset(key, profile); await redis.expire(key, USER_TTL_SECONDS); } async function getUserProfile( redis: Redis, userId: string, ): Promise<UserProfile | null> { const key = `${USER_KEY_PREFIX}${userId}`; const data = await redis.hgetall(key); if (!data || Object.keys(data).length === 0) { return null; } return data as UserProfile; } export { setUserProfile, getUserProfile }; ``` **Why good:** Key prefix separates concerns, TTL prevents stale data, null check on empty hash response, typed return **Gotcha:** `hgetall` returns an empty object `{}` for non-existent keys, not `null`. Always check `Object.keys(data).length === 0`. All values come back as strings -- numbers need explicit `parseInt`/`parseFloat`. --- ## Sorted Sets (Leaderboards, Rankings) ```typescript import type Redis from "ioredis"; const LEADERBOARD_KEY = "leaderboard:global"; const TOP_PLAYERS_COUNT = 10; async function updateScore( redis: Redis, playerId: string, score: number, ): Promise<void> { await redis.zadd(LEADERBOARD_KEY, score, playerId); } async function getTopPlayers( redis: Redis, ): Promise<Array<{ playerId: string; score: number }>> { // ZREVRANGE returns highest scores first // Note: ZREVRANGE is deprecated since Redis 6.2 in favor of ZRANGE ... REV, // but ioredis provides typed support for zrevrange const results = await redis.zrevrange( LEADERBOARD_KEY, 0, TOP_PLAYERS_COUNT - 1, "WITHSCORES", ); const players: Array<{ playerId: string; score: number }> = []; for (let i = 0; i < results.length; i += 2) { players.push({ playerId: results[i], score: parseFloat(results[i + 1]), }); } return players; } async function getPlayerRank( redis: Redis, playerId: string, ): Promise<number | null> { // ZREVRANK returns 0-based rank (highest score = rank 0) const rank = await redis.zrevrank(LEADERBOARD_KEY, playerId); return rank !== null ? rank + 1 : null; // Convert to 1-based } export { updateScore, getTopPlayers, getPlayerRank }; ``` **Why good:** Named constants for key and count, ZREVRANGE for descending order, WITHSCORES returns scores alongside members, 1-based rank conversion for user display --- ## Lists (Queues, Recent Items) ```typescript import type Redis from "ioredis"; const RECENT_ITEMS_KEY = "recent:items"; const MAX_RECENT_ITEMS = 50; async function addRecentItem(redis: Redis, item: string): Promise<void> { await redis .pipeline() .lpush(RECENT_ITEMS_KEY, item) .ltrim(RECENT_ITEMS_KEY, 0, MAX_RECENT_ITEMS - 1) .exec(); } async function getRecentItems(redis: Redis): Promise<string[]> { return redis.lrange(RECENT_ITEMS_KEY, 0, MAX_RECENT_ITEMS - 1); } export { addRecentItem, getRecentItems }; ``` **Why good:** Pipeline groups push and trim into single round-trip, LTRIM caps list size preventing unbounded growth, named constants for key and limit --- ## Sets (Unique Collections) ```typescript import type Redis from "ioredis"; const TAG_KEY_PREFIX = "tags:"; async function addTags( redis: Redis, entityId: string, tags: string[], ): Promise<void> { if (tags.length === 0) return; const key = `${TAG_KEY_PREFIX}${entityId}`; await redis.sadd(key, ...tags); } async function getTags(redis: Redis, entityId: string): Promise<string[]> { return redis.smembers(`${TAG_KEY_PREFIX}${entityId}`); } async function getCommonTags( redis: Redis, entityId1: string, entityId2: string, ): Promise<string[]> { return redis.sinter( `${TAG_KEY_PREFIX}${entityId1}`, `${TAG_KEY_PREFIX}${entityId2}`, ); } export { addTags, getTags, getCommonTags }; ``` **Why good:** Sets automatically deduplicate, SINTER finds common elements without application logic, spread operator for variadic SADD --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
pub-sub.md 5.6 KB
# Redis -- Pub/Sub Examples > Publish/subscribe messaging, event broadcasting, and pattern subscriptions. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [core.md](core.md) -- Connection setup (separate connections required for pub/sub) - [queues.md](queues.md) -- BullMQ and Redis Streams for persistent messaging --- ## Basic Pub/Sub with Separate Connections ```typescript // ✅ Good Example - Separate connections for pub/sub import Redis from "ioredis"; const NOTIFICATION_CHANNEL = "notifications"; function createPubSubClients() { const url = process.env.REDIS_URL!; // Publisher can be your regular Redis client const publisher = new Redis(url); publisher.on("error", (err) => { console.error("Publisher error:", err.message); }); // Subscriber MUST be a separate connection const subscriber = new Redis(url); subscriber.on("error", (err) => { console.error("Subscriber error:", err.message); }); return { publisher, subscriber }; } // Subscribe to channels async function setupSubscriber(subscriber: Redis): Promise<void> { await subscriber.subscribe(NOTIFICATION_CHANNEL); subscriber.on("message", (channel, message) => { const data = JSON.parse(message); console.log(`Received on ${channel}:`, data); handleNotification(data); }); } // Publish messages async function publishNotification( publisher: Redis, notification: { userId: string; type: string; message: string }, ): Promise<number> { // Returns number of subscribers that received the message return publisher.publish(NOTIFICATION_CHANNEL, JSON.stringify(notification)); } export { createPubSubClients, setupSubscriber, publishNotification }; ``` **Why good:** Separate connections for pub and sub (required by Redis protocol), error handlers on both, typed notification payload, publish returns subscriber count for observability ```typescript // ❌ Bad Example - Using same connection for pub and sub const redis = new Redis(); await redis.subscribe("channel"); await redis.set("key", "value"); // ERROR: connection is in subscriber mode ``` **Why bad:** A subscribed connection enters a special mode and cannot execute non-pub/sub commands -- `set` will throw an error --- ## Pattern Subscriptions ```typescript // Subscribe to all channels matching a pattern await subscriber.psubscribe("notifications:*"); subscriber.on("pmessage", (pattern, channel, message) => { // pattern: "notifications:*" // channel: "notifications:user:123" (actual channel) // message: the published data console.log(`Pattern ${pattern} matched channel ${channel}`); }); ``` --- ## Event Broadcasting System ```typescript import Redis from "ioredis"; // Event types interface UserEvent { type: "user:created" | "user:updated" | "user:deleted"; userId: string; timestamp: number; data?: Record<string, unknown>; } const EVENT_CHANNEL_PREFIX = "events:"; function createEventBus() { const publisher = new Redis(process.env.REDIS_URL!); const subscriber = new Redis(process.env.REDIS_URL!); publisher.on("error", (err) => console.error("Event publisher error:", err.message), ); subscriber.on("error", (err) => console.error("Event subscriber error:", err.message), ); type EventHandler = (event: UserEvent) => void | Promise<void>; const handlers = new Map<string, EventHandler[]>(); // Subscribe to event patterns async function subscribe( pattern: string, handler: EventHandler, ): Promise<void> { const channel = `${EVENT_CHANNEL_PREFIX}${pattern}`; if (!handlers.has(channel)) { handlers.set(channel, []); await subscriber.psubscribe(channel); } handlers.get(channel)!.push(handler); } // Handle incoming messages subscriber.on("pmessage", async (_pattern, channel, message) => { const event = JSON.parse(message) as UserEvent; // Find all matching handlers for (const [handlerPattern, handlerList] of handlers) { if (channelMatchesPattern(channel, handlerPattern)) { for (const handler of handlerList) { try { await handler(event); } catch (err) { console.error(`Event handler error on ${channel}:`, err); } } } } }); // Publish events async function publish(event: UserEvent): Promise<number> { const channel = `${EVENT_CHANNEL_PREFIX}${event.type}`; return publisher.publish(channel, JSON.stringify(event)); } // Cleanup async function close(): Promise<void> { await subscriber.punsubscribe(); await publisher.quit(); await subscriber.quit(); } return { subscribe, publish, close }; } function channelMatchesPattern(channel: string, pattern: string): boolean { const regex = new RegExp( "^" + pattern.replace(/\*/g, ".*").replace(/\?/g, ".") + "$", ); return regex.test(channel); } export { createEventBus }; export type { UserEvent }; ``` #### Usage ```typescript const eventBus = createEventBus(); // Subscribe to all user events await eventBus.subscribe("user:*", async (event) => { console.log(`User event: ${event.type} for ${event.userId}`); }); // Subscribe to specific event await eventBus.subscribe("user:created", async (event) => { await sendWelcomeEmail(event.userId); }); // Publish an event await eventBus.publish({ type: "user:created", userId: "user-123", timestamp: Date.now(), data: { name: "Alice" }, }); ``` **Why good:** Typed event payloads, pattern-based subscriptions, error isolation per handler, proper cleanup with `close()`, separate pub/sub connections --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
queues.md 10.7 KB
# Redis -- Queues & Locks Examples > BullMQ job queues, Redis Streams with consumer groups, and distributed locks. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [core.md](core.md) -- Connection setup (`maxRetriesPerRequest: null` for BullMQ) - [pub-sub.md](pub-sub.md) -- Fire-and-forget messaging (vs persistent queues) - [rate-limiting.md](rate-limiting.md) -- Lua scripts for atomic operations --- ## BullMQ: Complete Email Queue ```typescript import { Queue, Worker, QueueEvents, type Job } from "bullmq"; import Redis from "ioredis"; // Job data types interface EmailJobData { to: string; subject: string; templateId: string; variables: Record<string, string>; } interface EmailJobResult { messageId: string; sentAt: string; } // Constants const QUEUE_NAME = "emails"; const MAX_ATTEMPTS = 3; const BACKOFF_DELAY_MS = 1000; const WORKER_CONCURRENCY = 5; const STALLED_INTERVAL_MS = 30000; const COMPLETED_RETENTION = 1000; const FAILED_RETENTION = 5000; // Connection factory -- each BullMQ component needs its own connection function createConnection(): Redis { return new Redis(process.env.REDIS_URL!, { maxRetriesPerRequest: null, // REQUIRED for BullMQ }); } // Queue const emailQueue = new Queue<EmailJobData, EmailJobResult>(QUEUE_NAME, { connection: createConnection(), defaultJobOptions: { attempts: MAX_ATTEMPTS, backoff: { type: "exponential", delay: BACKOFF_DELAY_MS }, removeOnComplete: { count: COMPLETED_RETENTION }, removeOnFail: { count: FAILED_RETENTION }, }, }); // Worker const emailWorker = new Worker<EmailJobData, EmailJobResult>( QUEUE_NAME, async (job: Job<EmailJobData, EmailJobResult>) => { const { to, subject, templateId, variables } = job.data; // Report progress await job.updateProgress(10); // Render template const html = await renderTemplate(templateId, variables); await job.updateProgress(50); // Send email const result = await sendEmail({ to, subject, html }); await job.updateProgress(100); return { messageId: result.id, sentAt: new Date().toISOString(), }; }, { connection: createConnection(), concurrency: WORKER_CONCURRENCY, stalledInterval: STALLED_INTERVAL_MS, }, ); // Event handlers emailWorker.on("completed", (job, result) => { console.log(`Email sent to ${job.data.to} (messageId: ${result.messageId})`); }); emailWorker.on("failed", (job, err) => { console.error( `Email to ${job?.data.to} failed after ${job?.attemptsMade} attempts:`, err.message, ); }); // QueueEvents for monitoring const emailEvents = new QueueEvents(QUEUE_NAME, { connection: createConnection(), }); emailEvents.on("waiting", ({ jobId }) => { console.log(`Email job ${jobId} is waiting`); }); export { emailQueue, emailWorker, emailEvents }; ``` **Why good:** `maxRetriesPerRequest: null` is required for BullMQ, typed job data with generics, exponential backoff for retries, cleanup policies prevent unbounded Redis memory growth, separate connection per Queue/Worker (BullMQ requirement), progress reporting ```typescript // ❌ Bad Example - BullMQ without required config import { Queue, Worker } from "bullmq"; import Redis from "ioredis"; const connection = new Redis(); // Missing maxRetriesPerRequest: null const queue = new Queue("emails", { connection }); // BullMQ will throw: "maxRetriesPerRequest must be null" ``` **Why bad:** BullMQ requires `maxRetriesPerRequest: null` -- without it, ioredis gives up retrying after a set number of attempts, but BullMQ expects to retry forever --- ## BullMQ: Adding Jobs with Scheduling ```typescript const REPORT_DELAY_MS = 60000; // 1 minute const HIGH_PRIORITY = 1; const NORMAL_PRIORITY = 5; // Immediate high-priority email await emailQueue.add( "transactional", { to: "user@example.com", subject: "Password Reset", templateId: "password-reset", variables: { resetLink: "https://..." }, }, { priority: HIGH_PRIORITY }, ); // Delayed email await emailQueue.add( "reminder", { to: "user@example.com", subject: "Complete your profile", templateId: "profile-reminder", variables: { userName: "Alice" }, }, { delay: REPORT_DELAY_MS, priority: NORMAL_PRIORITY }, ); // Repeating email (cron schedule) await emailQueue.add( "daily-digest", { to: "user@example.com", subject: "Your Daily Digest", templateId: "daily-digest", variables: {}, }, { repeat: { pattern: "0 9 * * *", tz: "America/New_York" } }, ); // Bulk add await emailQueue.addBulk([ { name: "welcome", data: { to: "a@example.com", subject: "Welcome", templateId: "welcome", variables: {}, }, }, { name: "welcome", data: { to: "b@example.com", subject: "Welcome", templateId: "welcome", variables: {}, }, }, ]); ``` --- ## BullMQ: Graceful Shutdown ```typescript async function shutdown(): Promise<void> { console.log("Shutting down workers..."); // Close worker first (stop accepting new jobs) await emailWorker.close(); // Close event listener await emailEvents.close(); // Close queue last await emailQueue.close(); console.log("All BullMQ connections closed"); } process.on("SIGTERM", shutdown); process.on("SIGINT", shutdown); export { shutdown }; ``` --- ## Redis Streams: Order Processing Pipeline ```typescript import Redis from "ioredis"; const STREAM_KEY = "stream:orders"; const GROUP_NAME = "order-processors"; const BLOCK_MS = 5000; const BATCH_SIZE = 10; const CLAIM_MIN_IDLE_MS = 30000; // Producer async function addOrderEvent( redis: Redis, event: { orderId: string; action: string; data: Record<string, unknown> }, ): Promise<string> { return redis.xadd( STREAM_KEY, "*", "orderId", event.orderId, "action", event.action, "data", JSON.stringify(event.data), ); } // Consumer with pending message recovery async function startConsumer( redis: Redis, consumerName: string, ): Promise<void> { // Ensure group exists try { await redis.xgroup("CREATE", STREAM_KEY, GROUP_NAME, "0", "MKSTREAM"); } catch (err) { if (!(err instanceof Error) || !err.message.includes("BUSYGROUP")) { throw err; } } // Process pending messages first (messages claimed but not ACKed) await processPending(redis, consumerName); // Then process new messages while (true) { const results = await redis.xreadgroup( "GROUP", GROUP_NAME, consumerName, "COUNT", String(BATCH_SIZE), "BLOCK", String(BLOCK_MS), "STREAMS", STREAM_KEY, ">", ); if (!results) continue; for (const [, messages] of results) { for (const [id, fields] of messages) { const event = parseStreamFields(fields); try { await processOrder(event); await redis.xack(STREAM_KEY, GROUP_NAME, id); } catch (err) { console.error(`Failed to process ${id}:`, err); // Will be retried via pending recovery } } } } } // Claim and reprocess messages stuck in pending state async function processPending( redis: Redis, consumerName: string, ): Promise<void> { const pending = await redis.xpending( STREAM_KEY, GROUP_NAME, "-", "+", String(BATCH_SIZE), ); for (const [id, , idleTime] of pending) { if (Number(idleTime) > CLAIM_MIN_IDLE_MS) { const claimed = await redis.xclaim( STREAM_KEY, GROUP_NAME, consumerName, CLAIM_MIN_IDLE_MS, id, ); for (const [claimedId, fields] of claimed) { try { const event = parseStreamFields(fields); await processOrder(event); await redis.xack(STREAM_KEY, GROUP_NAME, claimedId); } catch (err) { console.error(`Failed to reprocess ${claimedId}:`, err); } } } } } function parseStreamFields(fields: string[]): { orderId: string; action: string; data: Record<string, unknown>; } { const obj: Record<string, string> = {}; for (let i = 0; i < fields.length; i += 2) { obj[fields[i]] = fields[i + 1]; } return { orderId: obj.orderId, action: obj.action, data: JSON.parse(obj.data), }; } export { addOrderEvent, startConsumer }; ``` **Why good:** MKSTREAM creates stream if it doesn't exist, processes pending messages on startup for crash recovery, XCLAIM reclaims stuck messages from dead consumers, XACK confirms processing, field parsing handles Redis stream key-value pairs --- ## Distributed Lock with SET NX ```typescript import type Redis from "ioredis"; import crypto from "node:crypto"; const LOCK_DEFAULT_TTL_MS = 10000; // 10 seconds const LOCK_RETRY_DELAY_MS = 100; interface LockOptions { ttlMs?: number; retryCount?: number; retryDelayMs?: number; } async function acquireLock( redis: Redis, resource: string, options: LockOptions = {}, ): Promise<string | null> { const { ttlMs = LOCK_DEFAULT_TTL_MS, retryCount = 3, retryDelayMs = LOCK_RETRY_DELAY_MS, } = options; const lockKey = `lock:${resource}`; const lockValue = crypto.randomUUID(); // Unique owner ID for (let i = 0; i <= retryCount; i++) { const acquired = await redis.set(lockKey, lockValue, "PX", ttlMs, "NX"); if (acquired === "OK") { return lockValue; // Return owner ID for safe release } if (i < retryCount) { await new Promise((resolve) => setTimeout(resolve, retryDelayMs)); } } return null; // Failed to acquire } // Release lock only if we still own it (Lua script for atomicity) async function releaseLock( redis: Redis, resource: string, lockValue: string, ): Promise<boolean> { const lockKey = `lock:${resource}`; const result = await redis.eval( ` if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('DEL', KEYS[1]) else return 0 end `, 1, lockKey, lockValue, ); return result === 1; } export { acquireLock, releaseLock }; ``` #### Usage with Resource Protection ```typescript async function processExclusiveTask(redis: Redis, taskId: string) { const lockValue = await acquireLock(redis, `task:${taskId}`, { ttlMs: 30000, }); if (!lockValue) { console.log(`Task ${taskId} is already being processed`); return; } try { // Do exclusive work here await performTask(taskId); } finally { // Always release in finally block await releaseLock(redis, `task:${taskId}`, lockValue); } } ``` **Why good:** UUID lock value ensures only the owner can release, Lua script makes GET+DEL atomic (prevents releasing someone else's lock), PX for millisecond TTL precision, retry logic for contention, finally block ensures release even on error --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
rate-limiting.md 6.4 KB
# Redis -- Rate Limiting Examples > Sliding window and token bucket rate limiters with Lua scripts, plus middleware integration. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [core.md](core.md) -- Connection setup, Lua scripting basics - [sessions.md](sessions.md) -- Session storage (often paired with rate limiting) - [queues.md](queues.md) -- BullMQ for throttled job processing --- ## Sliding Window Rate Limiter (Lua Script) ```typescript import Redis from "ioredis"; const redis = new Redis(process.env.REDIS_URL!); redis.on("error", (err) => console.error("Redis error:", err.message)); // Atomic sliding window rate limiter redis.defineCommand("slidingWindowRateLimit", { numberOfKeys: 1, lua: ` local key = KEYS[1] local max_requests = tonumber(ARGV[1]) local window_ms = tonumber(ARGV[2]) local now = tonumber(ARGV[3]) local window_start = now - window_ms -- Remove expired entries outside the window redis.call('ZREMRANGEBYSCORE', key, '-inf', window_start) -- Count requests in current window local current = redis.call('ZCARD', key) if current < max_requests then -- Add this request (score = timestamp, member = unique ID) redis.call('ZADD', key, now, now .. ':' .. math.random(1000000)) redis.call('PEXPIRE', key, window_ms) return {1, max_requests - current - 1} -- allowed, remaining else -- Get TTL until oldest entry expires local oldest = redis.call('ZRANGE', key, 0, 0, 'WITHSCORES') local retry_after = 0 if #oldest > 0 then retry_after = tonumber(oldest[2]) + window_ms - now end return {0, 0, retry_after} -- denied, remaining=0, retry_after_ms end `, }); // TypeScript type declaration declare module "ioredis" { interface RedisCommander<Context> { slidingWindowRateLimit( key: string, maxRequests: string, windowMs: string, nowMs: string, ): Promise<[number, number, number?]>; } } interface RateLimitResult { allowed: boolean; remaining: number; retryAfterMs?: number; } const DEFAULT_MAX_REQUESTS = 100; const DEFAULT_WINDOW_MS = 60000; // 1 minute async function checkRateLimit( identifier: string, maxRequests: number = DEFAULT_MAX_REQUESTS, windowMs: number = DEFAULT_WINDOW_MS, ): Promise<RateLimitResult> { const key = `ratelimit:${identifier}`; const now = Date.now(); const [allowed, remaining, retryAfterMs] = await redis.slidingWindowRateLimit( key, String(maxRequests), String(windowMs), String(now), ); return { allowed: allowed === 1, remaining, retryAfterMs: retryAfterMs ?? undefined, }; } export { checkRateLimit }; export type { RateLimitResult }; ``` **Why good:** Entire check-and-update is atomic via Lua, sorted set naturally orders by timestamp, PEXPIRE auto-cleans stale keys, returns remaining count and retry-after for HTTP headers --- ## Rate Limiting Middleware (Express/Hono) ```typescript import type { Request, Response, NextFunction } from "express"; import { checkRateLimit } from "./rate-limiter"; const API_RATE_LIMIT = 100; const API_RATE_WINDOW_MS = 60000; // 1 minute function rateLimitMiddleware( maxRequests: number = API_RATE_LIMIT, windowMs: number = API_RATE_WINDOW_MS, ) { return async (req: Request, res: Response, next: NextFunction) => { const identifier = req.ip ?? "unknown"; const result = await checkRateLimit(identifier, maxRequests, windowMs); // Set standard rate limit headers res.setHeader("X-RateLimit-Limit", maxRequests); res.setHeader("X-RateLimit-Remaining", result.remaining); if (!result.allowed) { res.setHeader( "Retry-After", Math.ceil((result.retryAfterMs ?? windowMs) / 1000), ); res.status(429).json({ error: "Too many requests" }); return; } next(); }; } export { rateLimitMiddleware }; ``` **Why good:** Standard `X-RateLimit-*` headers for client awareness, `Retry-After` header for 429 responses, configurable limits per route --- ## Token Bucket Rate Limiter (Lua Script) Allows burst traffic up to bucket capacity, then refills at a steady rate. ```typescript import Redis from "ioredis"; const redis = new Redis(process.env.REDIS_URL!); redis.on("error", (err) => console.error("Redis error:", err.message)); redis.defineCommand("tokenBucket", { numberOfKeys: 1, lua: ` local key = KEYS[1] local capacity = tonumber(ARGV[1]) local refill_rate = tonumber(ARGV[2]) -- tokens per second local now = tonumber(ARGV[3]) local requested = tonumber(ARGV[4]) local bucket = redis.call('HMGET', key, 'tokens', 'last_refill') local tokens = tonumber(bucket[1]) or capacity local last_refill = tonumber(bucket[2]) or now -- Refill tokens based on elapsed time local elapsed = (now - last_refill) / 1000 tokens = math.min(capacity, tokens + (elapsed * refill_rate)) if tokens >= requested then tokens = tokens - requested redis.call('HMSET', key, 'tokens', tokens, 'last_refill', now) redis.call('EXPIRE', key, math.ceil(capacity / refill_rate) + 1) return {1, math.floor(tokens)} -- allowed, remaining tokens else redis.call('HMSET', key, 'tokens', tokens, 'last_refill', now) redis.call('EXPIRE', key, math.ceil(capacity / refill_rate) + 1) return {0, math.floor(tokens)} -- denied, remaining tokens end `, }); declare module "ioredis" { interface RedisCommander<Context> { tokenBucket( key: string, capacity: string, refillRate: string, nowMs: string, requested: string, ): Promise<[number, number]>; } } const BUCKET_CAPACITY = 50; const REFILL_RATE = 10; // tokens per second async function checkTokenBucket( identifier: string, tokensRequested: number = 1, ): Promise<{ allowed: boolean; remaining: number }> { const key = `tokenbucket:${identifier}`; const now = Date.now(); const [allowed, remaining] = await redis.tokenBucket( key, String(BUCKET_CAPACITY), String(REFILL_RATE), String(now), String(tokensRequested), ); return { allowed: allowed === 1, remaining }; } export { checkTokenBucket }; ``` **When to use sliding window:** Strict, evenly distributed rate limiting (API endpoints, login attempts). **When to use token bucket:** Allow burst traffic up to capacity (file uploads, batch operations). --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
redis.md 25 KB
# Redis Practical Examples > Caching patterns, session storage, rate limiting, pub/sub messaging, and job queue examples for Redis with ioredis. --- ## Caching Patterns ### Cache-Aside with Generic Helper ```typescript import type Redis from "ioredis"; const DEFAULT_CACHE_TTL_SECONDS = 300; // 5 minutes interface CacheResult<T> { data: T; fromCache: boolean; } async function cacheAside<T>( redis: Redis, key: string, fetcher: () => Promise<T>, ttlSeconds: number = DEFAULT_CACHE_TTL_SECONDS, ): Promise<CacheResult<T>> { try { const cached = await redis.get(key); if (cached !== null) { return { data: JSON.parse(cached) as T, fromCache: true }; } } catch (err) { // Cache read failure -- proceed to fetch from source console.error(`Cache read failed for ${key}:`, (err as Error).message); } const data = await fetcher(); // Non-blocking cache write redis.set(key, JSON.stringify(data), "EX", ttlSeconds).catch((err) => { console.error(`Cache write failed for ${key}:`, err.message); }); return { data, fromCache: false }; } export { cacheAside }; export type { CacheResult }; ``` ### Write-Through Cache ```typescript import type Redis from "ioredis"; const PRODUCT_CACHE_PREFIX = "cache:product:"; const PRODUCT_CACHE_TTL_SECONDS = 600; // 10 minutes async function writeThrough<T extends { id: string }>( redis: Redis, entity: T, dbWriter: (entity: T) => Promise<T>, ): Promise<T> { // 1. Write to database first (source of truth) const saved = await dbWriter(entity); // 2. Update cache with fresh data const key = `${PRODUCT_CACHE_PREFIX}${saved.id}`; await redis.set(key, JSON.stringify(saved), "EX", PRODUCT_CACHE_TTL_SECONDS); return saved; } export { writeThrough }; ``` ### Multi-Key Cache with Pipeline ```typescript import type Redis from "ioredis"; const CACHE_PREFIX = "cache:user:"; const CACHE_TTL_SECONDS = 300; async function getMultipleFromCache( redis: Redis, userIds: string[], ): Promise<Map<string, string | null>> { const pipeline = redis.pipeline(); for (const id of userIds) { pipeline.get(`${CACHE_PREFIX}${id}`); } const results = await pipeline.exec(); if (!results) { throw new Error("Pipeline execution returned null"); } const cache = new Map<string, string | null>(); for (let i = 0; i < userIds.length; i++) { const [err, value] = results[i]; cache.set(userIds[i], err ? null : (value as string | null)); } return cache; } async function setMultipleInCache( redis: Redis, entries: Array<{ id: string; data: unknown }>, ): Promise<void> { const pipeline = redis.pipeline(); for (const entry of entries) { pipeline.set( `${CACHE_PREFIX}${entry.id}`, JSON.stringify(entry.data), "EX", CACHE_TTL_SECONDS, ); } await pipeline.exec(); } export { getMultipleFromCache, setMultipleInCache }; ``` ### Cache Stampede Prevention (Singleflight) ```typescript import type Redis from "ioredis"; const LOCK_TTL_SECONDS = 10; const LOCK_RETRY_DELAY_MS = 50; const LOCK_MAX_RETRIES = 20; // Prevent multiple processes from regenerating the same cache entry simultaneously async function cacheAsideWithLock<T>( redis: Redis, key: string, fetcher: () => Promise<T>, ttlSeconds: number, ): Promise<T> { // Check cache first const cached = await redis.get(key); if (cached !== null) { return JSON.parse(cached) as T; } // Try to acquire lock const lockKey = `lock:${key}`; const lockAcquired = await redis.set( lockKey, "1", "EX", LOCK_TTL_SECONDS, "NX", ); if (lockAcquired === "OK") { try { // We have the lock -- fetch and populate cache const data = await fetcher(); await redis.set(key, JSON.stringify(data), "EX", ttlSeconds); return data; } finally { await redis.del(lockKey); } } // Another process is fetching -- wait and retry cache read for (let i = 0; i < LOCK_MAX_RETRIES; i++) { await new Promise((resolve) => setTimeout(resolve, LOCK_RETRY_DELAY_MS)); const result = await redis.get(key); if (result !== null) { return JSON.parse(result) as T; } } // Lock expired, cache still empty -- fall back to direct fetch return fetcher(); } export { cacheAsideWithLock }; ``` **Why good:** NX flag prevents multiple processes from acquiring the lock, lock has TTL to prevent deadlocks if process crashes, fallback to direct fetch if lock holder fails --- ## Session Storage ### Express Session with connect-redis ```typescript import express from "express"; import session from "express-session"; import RedisStore from "connect-redis"; import Redis from "ioredis"; const SESSION_SECRET = process.env.SESSION_SECRET; const SESSION_TTL_SECONDS = 86400; // 24 hours const SESSION_PREFIX = "sess:"; const COOKIE_MAX_AGE_MS = 86400000; // 24 hours if (!SESSION_SECRET) { throw new Error("SESSION_SECRET environment variable is required"); } const redisClient = new Redis(process.env.REDIS_URL!); redisClient.on("error", (err) => { console.error("Redis session store error:", err.message); }); const app = express(); app.use( session({ store: new RedisStore({ client: redisClient, prefix: SESSION_PREFIX, ttl: SESSION_TTL_SECONDS, }), secret: SESSION_SECRET, resave: false, // Don't save session if unmodified saveUninitialized: false, // Don't create session until something is stored cookie: { secure: process.env.NODE_ENV === "production", httpOnly: true, // Prevent XSS access to cookie maxAge: COOKIE_MAX_AGE_MS, sameSite: "lax", // CSRF protection }, }), ); export { app }; ``` **Why good:** `resave: false` prevents race conditions with parallel requests, `saveUninitialized: false` avoids empty sessions, secure cookie in production, httpOnly prevents XSS, sameSite prevents CSRF ### Hono Session Middleware (Manual) ```typescript import type Redis from "ioredis"; import { createMiddleware } from "hono/factory"; import crypto from "node:crypto"; const SESSION_PREFIX = "session:"; const SESSION_TTL_SECONDS = 86400; const SESSION_COOKIE_NAME = "sid"; function sessionMiddleware(redis: Redis) { return createMiddleware(async (c, next) => { const sessionId = c.req.cookie(SESSION_COOKIE_NAME) ?? crypto.randomUUID(); const key = `${SESSION_PREFIX}${sessionId}`; const raw = await redis.get(key); const session = raw ? JSON.parse(raw) : {}; c.set("session", session); c.set("sessionId", sessionId); await next(); // Save session after response const updatedSession = c.get("session"); await redis.set( key, JSON.stringify(updatedSession), "EX", SESSION_TTL_SECONDS, ); // Set cookie if new session if (!c.req.cookie(SESSION_COOKIE_NAME)) { c.header( "Set-Cookie", `${SESSION_COOKIE_NAME}=${sessionId}; Path=/; HttpOnly; SameSite=Lax; Max-Age=${SESSION_TTL_SECONDS}`, ); } }); } export { sessionMiddleware }; ``` --- ## Rate Limiting ### Sliding Window Rate Limiter (Lua Script) ```typescript import Redis from "ioredis"; const redis = new Redis(process.env.REDIS_URL!); redis.on("error", (err) => console.error("Redis error:", err.message)); // Atomic sliding window rate limiter redis.defineCommand("slidingWindowRateLimit", { numberOfKeys: 1, lua: ` local key = KEYS[1] local max_requests = tonumber(ARGV[1]) local window_ms = tonumber(ARGV[2]) local now = tonumber(ARGV[3]) local window_start = now - window_ms -- Remove expired entries outside the window redis.call('ZREMRANGEBYSCORE', key, '-inf', window_start) -- Count requests in current window local current = redis.call('ZCARD', key) if current < max_requests then -- Add this request (score = timestamp, member = unique ID) redis.call('ZADD', key, now, now .. ':' .. math.random(1000000)) redis.call('PEXPIRE', key, window_ms) return {1, max_requests - current - 1} -- allowed, remaining else -- Get TTL until oldest entry expires local oldest = redis.call('ZRANGE', key, 0, 0, 'WITHSCORES') local retry_after = 0 if #oldest > 0 then retry_after = tonumber(oldest[2]) + window_ms - now end return {0, 0, retry_after} -- denied, remaining=0, retry_after_ms end `, }); // TypeScript type declaration declare module "ioredis" { interface RedisCommander<Context> { slidingWindowRateLimit( key: string, maxRequests: string, windowMs: string, nowMs: string, ): Promise<[number, number, number?]>; } } interface RateLimitResult { allowed: boolean; remaining: number; retryAfterMs?: number; } const DEFAULT_MAX_REQUESTS = 100; const DEFAULT_WINDOW_MS = 60000; // 1 minute async function checkRateLimit( identifier: string, maxRequests: number = DEFAULT_MAX_REQUESTS, windowMs: number = DEFAULT_WINDOW_MS, ): Promise<RateLimitResult> { const key = `ratelimit:${identifier}`; const now = Date.now(); const [allowed, remaining, retryAfterMs] = await redis.slidingWindowRateLimit( key, String(maxRequests), String(windowMs), String(now), ); return { allowed: allowed === 1, remaining, retryAfterMs: retryAfterMs ?? undefined, }; } export { checkRateLimit }; export type { RateLimitResult }; ``` ### Rate Limiting Middleware (Express/Hono) ```typescript import type { Request, Response, NextFunction } from "express"; import { checkRateLimit } from "./rate-limiter"; const API_RATE_LIMIT = 100; const API_RATE_WINDOW_MS = 60000; // 1 minute function rateLimitMiddleware( maxRequests: number = API_RATE_LIMIT, windowMs: number = API_RATE_WINDOW_MS, ) { return async (req: Request, res: Response, next: NextFunction) => { const identifier = req.ip ?? "unknown"; const result = await checkRateLimit(identifier, maxRequests, windowMs); // Set standard rate limit headers res.setHeader("X-RateLimit-Limit", maxRequests); res.setHeader("X-RateLimit-Remaining", result.remaining); if (!result.allowed) { res.setHeader( "Retry-After", Math.ceil((result.retryAfterMs ?? windowMs) / 1000), ); res.status(429).json({ error: "Too many requests" }); return; } next(); }; } export { rateLimitMiddleware }; ``` ### Token Bucket Rate Limiter (Lua Script) ```typescript import Redis from "ioredis"; const redis = new Redis(process.env.REDIS_URL!); redis.on("error", (err) => console.error("Redis error:", err.message)); // Token bucket allows burst traffic up to bucket capacity redis.defineCommand("tokenBucket", { numberOfKeys: 1, lua: ` local key = KEYS[1] local capacity = tonumber(ARGV[1]) local refill_rate = tonumber(ARGV[2]) -- tokens per second local now = tonumber(ARGV[3]) local requested = tonumber(ARGV[4]) local bucket = redis.call('HMGET', key, 'tokens', 'last_refill') local tokens = tonumber(bucket[1]) or capacity local last_refill = tonumber(bucket[2]) or now -- Refill tokens based on elapsed time local elapsed = (now - last_refill) / 1000 tokens = math.min(capacity, tokens + (elapsed * refill_rate)) if tokens >= requested then tokens = tokens - requested redis.call('HMSET', key, 'tokens', tokens, 'last_refill', now) redis.call('EXPIRE', key, math.ceil(capacity / refill_rate) + 1) return {1, math.floor(tokens)} -- allowed, remaining tokens else redis.call('HMSET', key, 'tokens', tokens, 'last_refill', now) redis.call('EXPIRE', key, math.ceil(capacity / refill_rate) + 1) return {0, math.floor(tokens)} -- denied, remaining tokens end `, }); declare module "ioredis" { interface RedisCommander<Context> { tokenBucket( key: string, capacity: string, refillRate: string, nowMs: string, requested: string, ): Promise<[number, number]>; } } const BUCKET_CAPACITY = 50; const REFILL_RATE = 10; // tokens per second async function checkTokenBucket( identifier: string, tokensRequested: number = 1, ): Promise<{ allowed: boolean; remaining: number }> { const key = `tokenbucket:${identifier}`; const now = Date.now(); const [allowed, remaining] = await redis.tokenBucket( key, String(BUCKET_CAPACITY), String(REFILL_RATE), String(now), String(tokensRequested), ); return { allowed: allowed === 1, remaining }; } export { checkTokenBucket }; ``` **When to use sliding window:** Strict, evenly distributed rate limiting (API endpoints, login attempts). **When to use token bucket:** Allow burst traffic up to capacity (file uploads, batch operations). --- ## Pub/Sub Messaging ### Event Broadcasting System ```typescript import Redis from "ioredis"; // Event types interface UserEvent { type: "user:created" | "user:updated" | "user:deleted"; userId: string; timestamp: number; data?: Record<string, unknown>; } const EVENT_CHANNEL_PREFIX = "events:"; function createEventBus() { const publisher = new Redis(process.env.REDIS_URL!); const subscriber = new Redis(process.env.REDIS_URL!); publisher.on("error", (err) => console.error("Event publisher error:", err.message), ); subscriber.on("error", (err) => console.error("Event subscriber error:", err.message), ); type EventHandler = (event: UserEvent) => void | Promise<void>; const handlers = new Map<string, EventHandler[]>(); // Subscribe to event patterns async function subscribe( pattern: string, handler: EventHandler, ): Promise<void> { const channel = `${EVENT_CHANNEL_PREFIX}${pattern}`; if (!handlers.has(channel)) { handlers.set(channel, []); await subscriber.psubscribe(channel); } handlers.get(channel)!.push(handler); } // Handle incoming messages subscriber.on("pmessage", async (_pattern, channel, message) => { const event = JSON.parse(message) as UserEvent; // Find all matching handlers for (const [handlerPattern, handlerList] of handlers) { if (channelMatchesPattern(channel, handlerPattern)) { for (const handler of handlerList) { try { await handler(event); } catch (err) { console.error(`Event handler error on ${channel}:`, err); } } } } }); // Publish events async function publish(event: UserEvent): Promise<number> { const channel = `${EVENT_CHANNEL_PREFIX}${event.type}`; return publisher.publish(channel, JSON.stringify(event)); } // Cleanup async function close(): Promise<void> { await subscriber.punsubscribe(); await publisher.quit(); await subscriber.quit(); } return { subscribe, publish, close }; } function channelMatchesPattern(channel: string, pattern: string): boolean { const regex = new RegExp( "^" + pattern.replace(/\*/g, ".*").replace(/\?/g, ".") + "$", ); return regex.test(channel); } export { createEventBus }; export type { UserEvent }; ``` #### Usage ```typescript const eventBus = createEventBus(); // Subscribe to all user events await eventBus.subscribe("user:*", async (event) => { console.log(`User event: ${event.type} for ${event.userId}`); }); // Subscribe to specific event await eventBus.subscribe("user:created", async (event) => { await sendWelcomeEmail(event.userId); }); // Publish an event await eventBus.publish({ type: "user:created", userId: "user-123", timestamp: Date.now(), data: { name: "Alice" }, }); ``` --- ## Job Queues with BullMQ ### Complete Email Queue Example ```typescript import { Queue, Worker, QueueEvents, type Job } from "bullmq"; import Redis from "ioredis"; // Job data types interface EmailJobData { to: string; subject: string; templateId: string; variables: Record<string, string>; } interface EmailJobResult { messageId: string; sentAt: string; } // Constants const QUEUE_NAME = "emails"; const MAX_ATTEMPTS = 3; const BACKOFF_DELAY_MS = 1000; const WORKER_CONCURRENCY = 5; const STALLED_INTERVAL_MS = 30000; const COMPLETED_RETENTION = 1000; const FAILED_RETENTION = 5000; // Connection factory -- each BullMQ component needs its own connection function createConnection(): Redis { return new Redis(process.env.REDIS_URL!, { maxRetriesPerRequest: null, // REQUIRED for BullMQ }); } // Queue const emailQueue = new Queue<EmailJobData, EmailJobResult>(QUEUE_NAME, { connection: createConnection(), defaultJobOptions: { attempts: MAX_ATTEMPTS, backoff: { type: "exponential", delay: BACKOFF_DELAY_MS }, removeOnComplete: { count: COMPLETED_RETENTION }, removeOnFail: { count: FAILED_RETENTION }, }, }); // Worker const emailWorker = new Worker<EmailJobData, EmailJobResult>( QUEUE_NAME, async (job: Job<EmailJobData, EmailJobResult>) => { const { to, subject, templateId, variables } = job.data; // Report progress await job.updateProgress(10); // Render template const html = await renderTemplate(templateId, variables); await job.updateProgress(50); // Send email const result = await sendEmail({ to, subject, html }); await job.updateProgress(100); return { messageId: result.id, sentAt: new Date().toISOString(), }; }, { connection: createConnection(), concurrency: WORKER_CONCURRENCY, stalledInterval: STALLED_INTERVAL_MS, }, ); // Event handlers emailWorker.on("completed", (job, result) => { console.log(`Email sent to ${job.data.to} (messageId: ${result.messageId})`); }); emailWorker.on("failed", (job, err) => { console.error( `Email to ${job?.data.to} failed after ${job?.attemptsMade} attempts:`, err.message, ); }); // QueueEvents for monitoring const emailEvents = new QueueEvents(QUEUE_NAME, { connection: createConnection(), }); emailEvents.on("waiting", ({ jobId }) => { console.log(`Email job ${jobId} is waiting`); }); export { emailQueue, emailWorker, emailEvents }; ``` ### Adding Different Job Types ```typescript const WELCOME_DELAY_MS = 0; const REMINDER_DELAY_MS = 86400000; // 24 hours const HIGH_PRIORITY = 1; const NORMAL_PRIORITY = 5; // Immediate high-priority email await emailQueue.add( "transactional", { to: "user@example.com", subject: "Password Reset", templateId: "password-reset", variables: { resetLink: "https://..." }, }, { priority: HIGH_PRIORITY }, ); // Delayed email await emailQueue.add( "reminder", { to: "user@example.com", subject: "Complete your profile", templateId: "profile-reminder", variables: { userName: "Alice" }, }, { delay: REMINDER_DELAY_MS, priority: NORMAL_PRIORITY }, ); // Repeating email (daily digest) await emailQueue.add( "daily-digest", { to: "user@example.com", subject: "Your Daily Digest", templateId: "daily-digest", variables: {}, }, { repeat: { pattern: "0 9 * * *", tz: "America/New_York" }, }, ); // Bulk add await emailQueue.addBulk([ { name: "welcome", data: { to: "a@example.com", subject: "Welcome", templateId: "welcome", variables: {}, }, }, { name: "welcome", data: { to: "b@example.com", subject: "Welcome", templateId: "welcome", variables: {}, }, }, { name: "welcome", data: { to: "c@example.com", subject: "Welcome", templateId: "welcome", variables: {}, }, }, ]); ``` ### Graceful Shutdown ```typescript async function shutdown(): Promise<void> { console.log("Shutting down workers..."); // Close worker first (stop accepting new jobs) await emailWorker.close(); // Close event listener await emailEvents.close(); // Close queue last await emailQueue.close(); console.log("All BullMQ connections closed"); } process.on("SIGTERM", shutdown); process.on("SIGINT", shutdown); export { shutdown }; ``` --- ## Distributed Lock ### Simple Lock with SET NX ```typescript import type Redis from "ioredis"; import crypto from "node:crypto"; const LOCK_DEFAULT_TTL_MS = 10000; // 10 seconds const LOCK_RETRY_DELAY_MS = 100; interface LockOptions { ttlMs?: number; retryCount?: number; retryDelayMs?: number; } async function acquireLock( redis: Redis, resource: string, options: LockOptions = {}, ): Promise<string | null> { const { ttlMs = LOCK_DEFAULT_TTL_MS, retryCount = 3, retryDelayMs = LOCK_RETRY_DELAY_MS, } = options; const lockKey = `lock:${resource}`; const lockValue = crypto.randomUUID(); // Unique owner ID for (let i = 0; i <= retryCount; i++) { const acquired = await redis.set(lockKey, lockValue, "PX", ttlMs, "NX"); if (acquired === "OK") { return lockValue; // Return owner ID for safe release } if (i < retryCount) { await new Promise((resolve) => setTimeout(resolve, retryDelayMs)); } } return null; // Failed to acquire } // Release lock only if we still own it (Lua script for atomicity) async function releaseLock( redis: Redis, resource: string, lockValue: string, ): Promise<boolean> { const lockKey = `lock:${resource}`; const result = await redis.eval( ` if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('DEL', KEYS[1]) else return 0 end `, 1, lockKey, lockValue, ); return result === 1; } export { acquireLock, releaseLock }; ``` #### Usage with Resource Protection ```typescript async function processExclusiveTask(redis: Redis, taskId: string) { const lockValue = await acquireLock(redis, `task:${taskId}`, { ttlMs: 30000, }); if (!lockValue) { console.log(`Task ${taskId} is already being processed`); return; } try { // Do exclusive work here await performTask(taskId); } finally { // Always release in finally block await releaseLock(redis, `task:${taskId}`, lockValue); } } ``` **Why good:** UUID lock value ensures only the owner can release, Lua script makes GET+DEL atomic (prevents releasing someone else's lock), PX for millisecond TTL precision, retry logic for contention, finally block ensures release even on error --- ## Redis Streams with Consumer Groups ### Order Processing Pipeline ```typescript import Redis from "ioredis"; const STREAM_KEY = "stream:orders"; const GROUP_NAME = "order-processors"; const BLOCK_MS = 5000; const BATCH_SIZE = 10; const PENDING_CHECK_INTERVAL_MS = 60000; const CLAIM_MIN_IDLE_MS = 30000; // Producer async function addOrderEvent( redis: Redis, event: { orderId: string; action: string; data: Record<string, unknown> }, ): Promise<string> { return redis.xadd( STREAM_KEY, "*", "orderId", event.orderId, "action", event.action, "data", JSON.stringify(event.data), ); } // Consumer with pending message recovery async function startConsumer( redis: Redis, consumerName: string, ): Promise<void> { // Ensure group exists try { await redis.xgroup("CREATE", STREAM_KEY, GROUP_NAME, "0", "MKSTREAM"); } catch (err) { if (!(err instanceof Error) || !err.message.includes("BUSYGROUP")) { throw err; } } // Process pending messages first (messages claimed but not ACKed) await processPending(redis, consumerName); // Then process new messages while (true) { const results = await redis.xreadgroup( "GROUP", GROUP_NAME, consumerName, "COUNT", String(BATCH_SIZE), "BLOCK", String(BLOCK_MS), "STREAMS", STREAM_KEY, ">", ); if (!results) continue; for (const [, messages] of results) { for (const [id, fields] of messages) { const event = parseStreamFields(fields); try { await processOrder(event); await redis.xack(STREAM_KEY, GROUP_NAME, id); } catch (err) { console.error(`Failed to process ${id}:`, err); // Will be retried via pending recovery } } } } } // Claim and reprocess messages stuck in pending state async function processPending( redis: Redis, consumerName: string, ): Promise<void> { const pending = await redis.xpending( STREAM_KEY, GROUP_NAME, "-", "+", String(BATCH_SIZE), ); for (const [id, , idleTime] of pending) { if (Number(idleTime) > CLAIM_MIN_IDLE_MS) { const claimed = await redis.xclaim( STREAM_KEY, GROUP_NAME, consumerName, CLAIM_MIN_IDLE_MS, id, ); for (const [claimedId, fields] of claimed) { try { const event = parseStreamFields(fields); await processOrder(event); await redis.xack(STREAM_KEY, GROUP_NAME, claimedId); } catch (err) { console.error(`Failed to reprocess ${claimedId}:`, err); } } } } } function parseStreamFields(fields: string[]): { orderId: string; action: string; data: Record<string, unknown>; } { const obj: Record<string, string> = {}; for (let i = 0; i < fields.length; i += 2) { obj[fields[i]] = fields[i + 1]; } return { orderId: obj.orderId, action: obj.action, data: JSON.parse(obj.data), }; } export { addOrderEvent, startConsumer }; ``` **Why good:** MKSTREAM creates stream if it doesn't exist, processes pending messages on startup for crash recovery, XCLAIM reclaims stuck messages from dead consumers, XACK confirms processing, field parsing handles Redis stream key-value pairs -
sessions.md 3.4 KB
# Redis -- Session Storage Examples > Session storage patterns with Express connect-redis and Hono manual middleware. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [core.md](core.md) -- Connection setup, error handling - [caching.md](caching.md) -- Cache-aside and write-through patterns - [rate-limiting.md](rate-limiting.md) -- Rate limiting middleware integration --- ## Express Session with connect-redis > **Note:** connect-redis v9+ only supports node-redis (not ioredis). Use node-redis `createClient()` for session storage. ```typescript import express from "express"; import session from "express-session"; import { RedisStore } from "connect-redis"; import { createClient } from "redis"; const SESSION_SECRET = process.env.SESSION_SECRET; const SESSION_TTL_SECONDS = 86400; // 24 hours const SESSION_PREFIX = "sess:"; const COOKIE_MAX_AGE_MS = 86400000; // 24 hours if (!SESSION_SECRET) { throw new Error("SESSION_SECRET environment variable is required"); } const redisClient = createClient({ url: process.env.REDIS_URL }); redisClient.on("error", (err) => { console.error("Redis session store error:", err.message); }); await redisClient.connect(); const app = express(); app.use( session({ store: new RedisStore({ client: redisClient, prefix: SESSION_PREFIX, ttl: SESSION_TTL_SECONDS, }), secret: SESSION_SECRET, resave: false, // Don't save session if unmodified saveUninitialized: false, // Don't create session until something is stored cookie: { secure: process.env.NODE_ENV === "production", httpOnly: true, // Prevent XSS access to cookie maxAge: COOKIE_MAX_AGE_MS, sameSite: "lax", // CSRF protection }, }), ); export { app }; ``` **Why good:** `resave: false` prevents race conditions with parallel requests, `saveUninitialized: false` avoids empty sessions, secure cookie in production, httpOnly prevents XSS, sameSite prevents CSRF, uses node-redis as required by connect-redis v9+ --- ## Hono Session Middleware (Manual) ```typescript import type Redis from "ioredis"; import { createMiddleware } from "hono/factory"; import { getCookie, setCookie } from "hono/cookie"; import crypto from "node:crypto"; const SESSION_PREFIX = "session:"; const SESSION_TTL_SECONDS = 86400; const SESSION_COOKIE_NAME = "sid"; function sessionMiddleware(redis: Redis) { return createMiddleware(async (c, next) => { const sessionId = getCookie(c, SESSION_COOKIE_NAME) ?? crypto.randomUUID(); const key = `${SESSION_PREFIX}${sessionId}`; const raw = await redis.get(key); const session = raw ? JSON.parse(raw) : {}; c.set("session", session); c.set("sessionId", sessionId); await next(); // Save session after response const updatedSession = c.get("session"); await redis.set( key, JSON.stringify(updatedSession), "EX", SESSION_TTL_SECONDS, ); // Set cookie if new session if (!getCookie(c, SESSION_COOKIE_NAME)) { setCookie(c, SESSION_COOKIE_NAME, sessionId, { path: "/", httpOnly: true, sameSite: "Lax", maxAge: SESSION_TTL_SECONDS, }); } }); } export { sessionMiddleware }; ``` **Why good:** Creates session only when needed, HttpOnly and SameSite cookie flags for security, TTL auto-expires abandoned sessions, JSON serialization for flexible session data --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
setup.md 9.3 KB
# Redis -- Setup & Connection Examples > Connection setup patterns for ioredis and node-redis. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [caching.md](caching.md) -- Cache-aside, write-through, invalidation - [data-structures.md](data-structures.md) -- Strings, hashes, lists, sets, sorted sets - [queues.md](queues.md) -- BullMQ, Redis Streams, distributed locks --- ## ioredis Connection with Error Handling ```typescript import Redis from "ioredis"; const RETRY_DELAY_BASE_MS = 50; const RETRY_DELAY_MAX_MS = 2000; function createRedisClient(): Redis { const url = process.env.REDIS_URL; if (!url) { throw new Error("REDIS_URL environment variable is required"); } const client = new Redis(url, { maxRetriesPerRequest: 3, retryStrategy(times) { const delay = Math.min(times * RETRY_DELAY_BASE_MS, RETRY_DELAY_MAX_MS); return delay; }, lazyConnect: true, }); client.on("error", (err) => { console.error("Redis connection error:", err.message); }); client.on("connect", () => { console.log("Redis connected"); }); return client; } export { createRedisClient }; ``` **Why good:** Environment variable validation, named constants for retry delays, `lazyConnect` prevents connection before ready, error event handler prevents process crash, reconnection strategy with exponential backoff capped at max delay ```typescript // ❌ Bad Example - No error handling, hardcoded config import Redis from "ioredis"; const redis = new Redis("redis://localhost:6379"); // No error handler -- unhandled errors crash the process // No retry strategy -- uses default which may not suit your needs // Hardcoded connection string -- leaks in version control ``` **Why bad:** Missing error event handler crashes Node.js process on connection failure, hardcoded URL prevents environment-specific configuration, no retry strategy customization --- ## node-redis Connection (Alternative) Use node-redis only when you need Redis Stack modules (JSON, Search, TimeSeries). ```typescript import { createClient } from "redis"; async function createNodeRedisClient() { const url = process.env.REDIS_URL; if (!url) { throw new Error("REDIS_URL environment variable is required"); } const client = createClient({ url }); client.on("error", (err) => { console.error("Redis client error:", err.message); }); await client.connect(); return client; } export { createNodeRedisClient }; ``` **When to use:** Only when you need Redis Stack modules (JSON, Search, TimeSeries) -- ioredis does not support Redis Stack modules natively. --- ## Pipelining (Non-Atomic Batching) Batch multiple commands to reduce network round-trips. ```typescript import type Redis from "ioredis"; const USER_KEY_PREFIX = "user:"; const USER_TTL_SECONDS = 3600; async function cacheMultipleUsers( redis: Redis, users: Array<{ id: string; name: string; email: string }>, ): Promise<void> { const pipeline = redis.pipeline(); for (const user of users) { const key = `${USER_KEY_PREFIX}${user.id}`; pipeline.hset(key, { name: user.name, email: user.email }); pipeline.expire(key, USER_TTL_SECONDS); } const results = await pipeline.exec(); if (!results) { throw new Error("Pipeline execution returned null"); } // Check for errors in pipeline results for (const [err] of results) { if (err) { throw new Error(`Pipeline command failed: ${err.message}`); } } } export { cacheMultipleUsers }; ``` **Why good:** Single network round-trip for all commands, error checking on each result, named constants for prefix and TTL --- ## Transactions (MULTI/EXEC -- Atomic) ```typescript import type Redis from "ioredis"; const BALANCE_KEY_PREFIX = "balance:"; async function transferBalance( redis: Redis, fromUserId: string, toUserId: string, amount: number, ): Promise<boolean> { const fromKey = `${BALANCE_KEY_PREFIX}${fromUserId}`; const toKey = `${BALANCE_KEY_PREFIX}${toUserId}`; // WATCH for optimistic locking await redis.watch(fromKey); const currentBalance = await redis.get(fromKey); if (!currentBalance || parseFloat(currentBalance) < amount) { await redis.unwatch(); return false; // Insufficient balance } // MULTI/EXEC -- atomic execution const results = await redis .multi() .decrby(fromKey, amount) .incrby(toKey, amount) .exec(); // results is null if WATCH detected a change (optimistic lock failure) if (!results) { return false; // Retry needed -- another client modified the key } return true; } export { transferBalance }; ``` **Why good:** WATCH provides optimistic locking, MULTI/EXEC ensures atomicity, null check handles concurrent modification, clear return value for retry logic ```typescript // ❌ Bad Example - Non-atomic balance transfer await redis.decrby("balance:user1", 100); await redis.incrby("balance:user2", 100); // If the process crashes between these two commands, // money disappears from user1 but never reaches user2 ``` **Why bad:** Two separate commands are not atomic, crash between them causes data inconsistency, no optimistic locking for concurrent access --- ## Cluster Mode ioredis supports Redis Cluster for horizontal scaling and high availability. ```typescript import Redis from "ioredis"; const CLUSTER_RETRY_BASE_MS = 100; const CLUSTER_RETRY_MAX_MS = 2000; const MAX_REDIRECTIONS = 16; const cluster = new Redis.Cluster( [ { host: "redis-node-1", port: 6379 }, { host: "redis-node-2", port: 6379 }, { host: "redis-node-3", port: 6379 }, ], { clusterRetryStrategy(times) { return Math.min(times * CLUSTER_RETRY_BASE_MS, CLUSTER_RETRY_MAX_MS); }, maxRedirections: MAX_REDIRECTIONS, scaleReads: "slave", // Read from replicas, write to master redisOptions: { password: process.env.REDIS_PASSWORD, }, }, ); cluster.on("error", (err) => { console.error("Cluster error:", err.message); }); // Use cluster exactly like a regular Redis client await cluster.set("key", "value"); const value = await cluster.get("key"); export { cluster }; ``` **Why good:** Multiple seed nodes for discovery, `scaleReads: "slave"` offloads reads to replicas, retry strategy with backoff, password from environment variable --- ## Sentinel Setup ```typescript import Redis from "ioredis"; const sentinel = new Redis({ sentinels: [ { host: "sentinel-1", port: 26379 }, { host: "sentinel-2", port: 26379 }, { host: "sentinel-3", port: 26379 }, ], name: "mymaster", // Sentinel group name sentinelRetryStrategy(times) { return Math.min(times * 10, 1000); }, failoverDetector: true, // Detect failover and reconnect automatically }); sentinel.on("error", (err) => { console.error("Sentinel error:", err.message); }); export { sentinel }; ``` **Why good:** Multiple sentinel nodes for redundancy, `failoverDetector: true` for automatic master failover handling, retry strategy for sentinel connectivity --- ## Auto-Pipelining ioredis can automatically batch commands issued during the same event loop tick: ```typescript import Redis from "ioredis"; const redis = new Redis(process.env.REDIS_URL!, { enableAutoPipelining: true, }); // These three commands are automatically batched into one pipeline const [name, email, role] = await Promise.all([ redis.get("user:name"), redis.get("user:email"), redis.get("user:role"), ]); export { redis }; ``` **When to use:** High-throughput applications issuing many independent commands per request. Auto-pipelining reduces network round-trips without changing application code. **When NOT to use:** When commands depend on each other's results (sequential logic), or when using WATCH/MULTI for transactions. --- ## Scanning Instead of KEYS Never use `KEYS` in production -- it blocks the Redis server while scanning all keys. ```typescript import type Redis from "ioredis"; const SCAN_BATCH_SIZE = 100; async function findKeysByPattern( redis: Redis, pattern: string, ): Promise<string[]> { const allKeys: string[] = []; const stream = redis.scanStream({ match: pattern, count: SCAN_BATCH_SIZE, }); return new Promise((resolve, reject) => { stream.on("data", (keys: string[]) => { allKeys.push(...keys); }); stream.on("end", () => resolve(allKeys)); stream.on("error", (err) => reject(err)); }); } export { findKeysByPattern }; ``` **Why good:** `scanStream` iterates incrementally without blocking Redis, `count` is a hint for batch size (not a guarantee), stream-based API handles large keyspaces --- ## Key Expiration Strategies ```typescript const TTL_SHORT_SECONDS = 60; // 1 minute -- volatile data const TTL_MEDIUM_SECONDS = 300; // 5 minutes -- API response cache const TTL_LONG_SECONDS = 3600; // 1 hour -- user profiles const TTL_SESSION_SECONDS = 86400; // 24 hours -- sessions // SET with TTL (preferred -- atomic) await redis.set("key", "value", "EX", TTL_MEDIUM_SECONDS); // SET with millisecond TTL await redis.set("key", "value", "PX", 500); // SET only if key doesn't exist (distributed lock pattern) const acquired = await redis.set("lock:resource", "owner-id", "EX", 30, "NX"); // Returns "OK" if lock acquired, null if already locked // Update TTL on existing key await redis.expire("key", TTL_LONG_SECONDS); ``` --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_
-
-
reference.md 18.3 KB
# Redis Quick Reference > Decision frameworks, command reference, connection options, anti-patterns, and production checklist. See [SKILL.md](SKILL.md) for core concepts and [examples/](examples/) for code examples. --- ## Command Quick Reference ### String Commands | Command | Description | Example | | -------------------------- | --------------------------- | ----------------------------------------- | | `SET key value` | Set a string value | `redis.set("key", "value")` | | `GET key` | Get a string value | `redis.get("key")` | | `SET key value EX seconds` | Set with TTL (seconds) | `redis.set("key", "val", "EX", 300)` | | `SET key value PX ms` | Set with TTL (milliseconds) | `redis.set("key", "val", "PX", 500)` | | `SET key value NX` | Set only if not exists | `redis.set("key", "val", "EX", 30, "NX")` | | `MGET key1 key2` | Get multiple keys | `redis.mget("k1", "k2")` | | `MSET key1 val1 key2 val2` | Set multiple keys | `redis.mset("k1", "v1", "k2", "v2")` | | `INCR key` | Increment by 1 | `redis.incr("counter")` | | `INCRBY key amount` | Increment by amount | `redis.incrby("counter", 5)` | | `DECR key` | Decrement by 1 | `redis.decr("counter")` | | `APPEND key value` | Append to string | `redis.append("log", "entry\n")` | ### Hash Commands | Command | Description | Example | | ---------------------- | ------------------- | ---------------------------------------------------- | | `HSET key field value` | Set hash field | `redis.hset("user:1", "name", "Alice")` | | `HSET key obj` | Set multiple fields | `redis.hset("user:1", { name: "Alice", age: "30" })` | | `HGET key field` | Get hash field | `redis.hget("user:1", "name")` | | `HGETALL key` | Get all fields | `redis.hgetall("user:1")` | | `HMGET key f1 f2` | Get multiple fields | `redis.hmget("user:1", "name", "age")` | | `HDEL key field` | Delete field | `redis.hdel("user:1", "age")` | | `HINCRBY key field n` | Increment field | `redis.hincrby("user:1", "visits", 1)` | | `HEXISTS key field` | Check field exists | `redis.hexists("user:1", "name")` | ### List Commands | Command | Description | Example | | ----------------------- | -------------------- | ------------------------------ | | `LPUSH key value` | Push to left (head) | `redis.lpush("queue", "item")` | | `RPUSH key value` | Push to right (tail) | `redis.rpush("queue", "item")` | | `LPOP key` | Pop from left | `redis.lpop("queue")` | | `RPOP key` | Pop from right | `redis.rpop("queue")` | | `LRANGE key start stop` | Get range | `redis.lrange("queue", 0, -1)` | | `LTRIM key start stop` | Trim to range | `redis.ltrim("queue", 0, 99)` | | `LLEN key` | Get length | `redis.llen("queue")` | | `BRPOP key timeout` | Blocking pop | `redis.brpop("queue", 5)` | ### Set Commands | Command | Description | Example | | ---------------------- | ---------------- | ---------------------------------- | | `SADD key member` | Add member | `redis.sadd("tags", "redis")` | | `SREM key member` | Remove member | `redis.srem("tags", "redis")` | | `SMEMBERS key` | Get all members | `redis.smembers("tags")` | | `SISMEMBER key member` | Check membership | `redis.sismember("tags", "redis")` | | `SCARD key` | Get count | `redis.scard("tags")` | | `SINTER key1 key2` | Intersection | `redis.sinter("tags1", "tags2")` | | `SUNION key1 key2` | Union | `redis.sunion("tags1", "tags2")` | | `SDIFF key1 key2` | Difference | `redis.sdiff("tags1", "tags2")` | ### Sorted Set Commands | Command | Description | Example | | --------------------------- | ------------------ | ------------------------------------------- | | `ZADD key score member` | Add with score | `redis.zadd("lb", 100, "player1")` | | `ZSCORE key member` | Get score | `redis.zscore("lb", "player1")` | | `ZRANK key member` | Get rank (asc) | `redis.zrank("lb", "player1")` | | `ZREVRANK key member` | Get rank (desc) | `redis.zrevrank("lb", "player1")` | | `ZRANGE key start stop` | Get range (asc) | `redis.zrange("lb", 0, 9)` | | `ZREVRANGE key start stop` | Get range (desc) | `redis.zrevrange("lb", 0, 9, "WITHSCORES")` | | `ZRANGEBYSCORE key min max` | Get by score range | `redis.zrangebyscore("lb", 50, 100)` | | `ZINCRBY key incr member` | Increment score | `redis.zincrby("lb", 10, "player1")` | | `ZREM key member` | Remove member | `redis.zrem("lb", "player1")` | | `ZCARD key` | Get count | `redis.zcard("lb")` | ### Stream Commands | Command | Description | Example | | -------------------------------------------- | ------------- | ------------------------------------------------------------------------- | | `XADD key * field value` | Add entry | `redis.xadd("stream", "*", "k", "v")` | | `XREAD COUNT n STREAMS key id` | Read entries | `redis.xread("COUNT", "10", "STREAMS", "s", "0")` | | `XRANGE key start end` | Get range | `redis.xrange("stream", "-", "+")` | | `XLEN key` | Get length | `redis.xlen("stream")` | | `XGROUP CREATE key group id` | Create group | `redis.xgroup("CREATE", "s", "g", "0", "MKSTREAM")` | | `XREADGROUP GROUP g c COUNT n STREAMS key >` | Consumer read | `redis.xreadgroup("GROUP", "g", "c", "COUNT", "10", "STREAMS", "s", ">")` | | `XACK key group id` | Acknowledge | `redis.xack("stream", "group", "id")` | | `XCLAIM key group consumer min-idle id` | Claim pending | `redis.xclaim("s", "g", "c", 30000, "id")` | ### Key Management Commands | Command | Description | Example | | ----------------------------------- | --------------------- | --------------------------------------- | | `DEL key` | Delete key (blocking) | `redis.del("key")` | | `UNLINK key` | Delete key (async) | `redis.unlink("key")` | | `EXPIRE key seconds` | Set TTL | `redis.expire("key", 300)` | | `PEXPIRE key ms` | Set TTL (ms) | `redis.pexpire("key", 500)` | | `TTL key` | Get TTL (seconds) | `redis.ttl("key")` | | `PTTL key` | Get TTL (ms) | `redis.pttl("key")` | | `EXISTS key` | Check existence | `redis.exists("key")` | | `TYPE key` | Get type | `redis.type("key")` | | `RENAME key newkey` | Rename key | `redis.rename("old", "new")` | | `SCAN cursor MATCH pattern COUNT n` | Iterate keys | `redis.scanStream({ match: "user:*" })` | --- ## ioredis Connection Options | Option | Default | Description | | ------------------------ | ------------- | -------------------------------------------------------------- | | `port` | `6379` | Redis port | | `host` | `"127.0.0.1"` | Redis host | | `family` | `4` | IP version (4 = IPv4, 6 = IPv6) | | `password` | `null` | Redis AUTH password | | `db` | `0` | Database index | | `keyPrefix` | `""` | Prefix for all keys | | `retryStrategy` | built-in | Function returning delay (ms) or non-number to stop | | `maxRetriesPerRequest` | `20` | Max retries per command (null = infinite, required for BullMQ) | | `enableReadyCheck` | `true` | Check if server is ready on connect | | `enableOfflineQueue` | `true` | Queue commands while disconnected | | `connectTimeout` | `10000` | Connection timeout (ms) | | `lazyConnect` | `false` | Don't connect until first command | | `tls` | `null` | TLS options for secure connections | | `enableAutoPipelining` | `false` | Auto-batch commands in same event loop tick | | `showFriendlyErrorStack` | `false` | Better error stack traces (slower) | ### Recommended Production Configuration ```typescript const PRODUCTION_OPTIONS = { maxRetriesPerRequest: 3, connectTimeout: 10000, retryStrategy(times: number) { const MAX_RETRY_DELAY_MS = 2000; const BASE_DELAY_MS = 50; return Math.min(times * BASE_DELAY_MS, MAX_RETRY_DELAY_MS); }, enableReadyCheck: true, enableOfflineQueue: true, enableAutoPipelining: true, } as const; ``` ### Recommended BullMQ Configuration ```typescript const BULLMQ_OPTIONS = { maxRetriesPerRequest: null, // REQUIRED enableReadyCheck: false, } as const; ``` ### Recommended Test Configuration ```typescript const TEST_OPTIONS = { maxRetriesPerRequest: 1, connectTimeout: 3000, lazyConnect: true, enableOfflineQueue: false, } as const; ``` --- ## Anti-Patterns ### Using KEYS in Production ```typescript // ❌ ANTI-PATTERN: Blocks entire Redis server const keys = await redis.keys("user:*"); // Scans ALL keys -- O(N) ``` **Why it's wrong:** `KEYS` scans the entire keyspace and blocks Redis during execution. On a server with millions of keys, this can take seconds and block all other clients. **What to do instead:** Use `SCAN` (or `scanStream` in ioredis) for incremental iteration: ```typescript const stream = redis.scanStream({ match: "user:*", count: 100 }); ``` --- ### Missing TTL on Cache Keys ```typescript // ❌ ANTI-PATTERN: Cache without expiration await redis.set("cache:user:123", JSON.stringify(user)); // Key persists forever -- stale data, unbounded memory ``` **Why it's wrong:** Without TTL, cache entries accumulate until Redis runs out of memory, and stale data is served indefinitely. **What to do instead:** Always set TTL: ```typescript const CACHE_TTL_SECONDS = 300; await redis.set( "cache:user:123", JSON.stringify(user), "EX", CACHE_TTL_SECONDS, ); ``` --- ### Assuming Pipeline Atomicity ```typescript // ❌ ANTI-PATTERN: Expecting atomicity from pipeline const pipeline = redis.pipeline(); pipeline.get("balance:user1"); // Read balance pipeline.decrby("balance:user1", 100); // Subtract pipeline.incrby("balance:user2", 100); // Add await pipeline.exec(); // NOT atomic -- another client can modify balance between get and decrby ``` **Why it's wrong:** Pipelines batch commands for network efficiency but do NOT provide atomicity. Other clients can interleave commands between pipeline steps. **What to do instead:** Use Lua scripts for atomic multi-command operations, or WATCH + MULTI/EXEC for optimistic locking. --- ### Single Connection for Pub/Sub ```typescript // ❌ ANTI-PATTERN: Reusing connection for subscribe and commands const redis = new Redis(); await redis.subscribe("channel"); await redis.get("key"); // THROWS: Connection is in subscriber mode ``` **Why it's wrong:** When a Redis connection enters subscriber mode, it can ONLY execute SUBSCRIBE, UNSUBSCRIBE, PSUBSCRIBE, PUNSUBSCRIBE, PING, and QUIT. All other commands throw errors. **What to do instead:** Create separate connections for subscribing and regular commands. --- ### Shared BullMQ Connection ```typescript // ❌ ANTI-PATTERN: Sharing connection between Queue and Worker const connection = new Redis({ maxRetriesPerRequest: null }); const queue = new Queue("tasks", { connection }); const worker = new Worker("tasks", processor, { connection }); // Queue and Worker interfere with each other's connection state ``` **Why it's wrong:** BullMQ Queue and Worker manage their connections independently. Sharing a connection causes unpredictable behavior, especially during reconnection. **What to do instead:** Create a new connection for each Queue, Worker, and QueueEvents instance: ```typescript const queue = new Queue("tasks", { connection: createConnection() }); const worker = new Worker("tasks", processor, { connection: createConnection(), }); ``` --- ### Storing Large Objects ```typescript // ❌ ANTI-PATTERN: Storing large blobs await redis.set("report", JSON.stringify(largeDataset)); // 5 MB object // Wastes memory, slow serialization, blocks Redis during transfer ``` **Why it's wrong:** Redis is optimized for small, fast operations. Large values waste memory, increase serialization time, and block the single-threaded Redis server during transfer. **What to do instead:** Store a reference (URL or ID) and fetch the actual data from object storage (S3) or database. --- ## Production Checklist ### Connection Management - [ ] Error event handlers on all Redis client instances - [ ] Retry strategy configured with exponential backoff and max delay - [ ] Separate connections for Pub/Sub subscribers - [ ] Separate connections for each BullMQ Queue/Worker/QueueEvents - [ ] `maxRetriesPerRequest: null` on BullMQ connections - [ ] TLS enabled in production (`tls: {}` option or `rediss://` URL) - [ ] Connection monitoring (track `connect`, `close`, `reconnecting` events) ### Data Management - [ ] TTL set on all cache keys - [ ] Key prefixes to separate concerns (e.g., `cache:`, `session:`, `lock:`, `queue:`) - [ ] `UNLINK` instead of `DEL` for large keys (non-blocking) - [ ] `SCAN` instead of `KEYS` for pattern matching - [ ] JSON serialization error handling for cache reads ### Performance - [ ] Auto-pipelining enabled for high-throughput scenarios - [ ] Manual pipelines for batch operations - [ ] Lua scripts for atomic multi-command operations - [ ] Connection pooling via Cluster or multiple clients - [ ] `maxmemory` and eviction policy configured on Redis server ### Security - [ ] Redis password set (`requirepass` in redis.conf) - [ ] TLS for connections in production - [ ] Bind to specific interfaces (not `0.0.0.0` in production) - [ ] Rename or disable dangerous commands (`FLUSHALL`, `CONFIG`, `DEBUG`) - [ ] Separate Redis instances or databases for different environments ### Monitoring - [ ] Track Redis memory usage (`INFO memory`) - [ ] Monitor key count and expiration rates - [ ] Alert on connection failures and reconnections - [ ] Monitor BullMQ queue depths and stalled jobs - [ ] Track Pub/Sub subscriber counts --- ## Key Naming Conventions ``` {type}:{entity}:{id} -- Standard pattern cache:user:123 -- Cached user data session:abc-def-ghi -- Session data lock:order:456 -- Distributed lock ratelimit:ip:192.168.1.1 -- Rate limit counter queue:email -- BullMQ queue data stream:events:orders -- Stream key leaderboard:global -- Sorted set for rankings {entity}:{id}:{field} -- Granular keys user:123:preferences -- User preferences ``` **Best practices:** - Use colons `:` as delimiters (Redis convention) - Include the data type or purpose as the first segment - Keep keys short but descriptive - Use `keyPrefix` option in ioredis for application-level namespacing - In Cluster mode, use `{hash-tag}` to colocate related keys: `{user:123}:profile`, `{user:123}:settings` --- ## ioredis vs node-redis Comparison | Feature | ioredis (v5.x) | node-redis (v5.x) | | --------------- | ------------------------- | ------------------------------------------- | | TypeScript | Written in TypeScript | TypeScript support | | Cluster | Built-in `Redis.Cluster` | `createCluster()` | | Sentinel | Built-in | Built-in | | Auto-pipelining | `enableAutoPipelining` | Automatic (same-tick batching) | | Lua scripting | `defineCommand()` | `client.eval()` | | Pub/Sub | Event-based (`message`) | `client.subscribe()` returns async iterator | | Streams | Full support | Full support | | Redis Stack | Not supported | Full support (JSON, Search, TimeSeries) | | BullMQ | Required (native support) | Not compatible | | Connection | Lazy or eager | Must call `.connect()` | | Command style | Lowercase (`redis.set()`) | Lowercase or camelCase (`client.set()`) | | Downloads/week | ~10M | ~7M | | Maintenance | Community (best-effort) | Redis team (active) | -
SKILL.md 19.5 KB
--- name: api-database-redis description: Redis in-memory data store patterns with ioredis and node-redis -- caching, sessions, rate limiting, pub/sub, streams, queues, transactions, cluster --- # Redis Patterns > **Quick Guide:** Use Redis as an in-memory data store for caching, session management, rate limiting, pub/sub messaging, and job queues. Use **ioredis** (v5.x) as the primary client for its superior TypeScript support, Cluster/Sentinel integration, auto-pipelining, and Lua scripting. Use **node-redis** (v5.x) only when you need Redis Stack modules (JSON, Search, TimeSeries). Always set `maxRetriesPerRequest: null` for BullMQ workers, use separate connections for Pub/Sub subscribers, and define Lua scripts via `defineCommand` for atomic multi-step operations. --- <critical_requirements> ## CRITICAL: Before Using This Skill > **All code must follow project conventions in CLAUDE.md** (kebab-case, named exports, import ordering, `import type`, named constants) **(You MUST use a SEPARATE Redis connection for Pub/Sub subscribers -- a subscribed connection enters a special mode and cannot execute other commands)** **(You MUST set `maxRetriesPerRequest: null` on any ioredis connection passed to BullMQ -- BullMQ requires infinite retries and will throw if this is not set)** **(You MUST use Lua scripts (`defineCommand` or `eval`) for any operation requiring atomicity across multiple Redis commands -- separate commands are NOT atomic even in a pipeline)** **(You MUST handle the `error` event on every Redis client instance -- unhandled errors crash the Node.js process)** </critical_requirements> --- ## Examples - [Core Patterns](examples/core.md) -- ioredis/node-redis connection, error handling, reconnection, cluster, sentinel, pipelining, transactions - [Caching Patterns](examples/caching.md) -- Cache-aside, write-through, invalidation, stampede prevention, multi-key pipeline - [Data Structures](examples/data-structures.md) -- Strings, hashes, lists, sets, sorted sets with typed helpers - [Sessions](examples/sessions.md) -- Express connect-redis (node-redis required for v9+), Hono manual middleware - [Pub/Sub](examples/pub-sub.md) -- Publish/subscribe, event broadcasting, pattern subscriptions - [Rate Limiting](examples/rate-limiting.md) -- Sliding window (Lua), token bucket (Lua), middleware integration - [Queues & Locks](examples/queues.md) -- BullMQ job queues, Redis Streams with consumer groups, distributed locks **Additional resources:** - [reference.md](reference.md) -- Command cheat sheet, connection options, anti-patterns, production checklist --- **Auto-detection:** Redis, ioredis, node-redis, createClient, RedisStore, BullMQ, Queue, Worker, pub/sub, MULTI, EXEC, pipeline, Lua script, defineCommand, xadd, xread, cache-aside, rate limit, session store, connect-redis, Redis.Cluster, Sentinel **When to use:** - Caching database queries or API responses (cache-aside, write-through) - Session storage for Express/Hono/Fastify applications - Distributed rate limiting (sliding window, token bucket) - Real-time messaging with Pub/Sub - Background job processing with BullMQ queues - Leaderboards, counters, and real-time analytics with sorted sets - Distributed locks and atomic operations with Lua scripts **Key patterns covered:** - ioredis connection setup, configuration, and error handling - Data structures (strings, hashes, lists, sets, sorted sets, streams) - Cache-aside and write-through caching with TTL management - Session storage with connect-redis - Rate limiting with Lua scripts (sliding window, token bucket) - Pub/Sub messaging with separate connections - Redis Streams for persistent message queues - BullMQ for job queues with retries and scheduling - Pipelining and transactions (MULTI/EXEC) - Lua scripting for atomic operations - Cluster mode and Sentinel for high availability **When NOT to use:** - Primary database for relational data (use your relational database) - Document storage with complex queries (use a document database) - Large binary file storage (use S3/object storage) - Data that must survive total memory loss without persistence configured --- <philosophy> ## Philosophy Redis is an **in-memory data store** used as a cache, message broker, and streaming engine. The core principle: **use Redis for fast, ephemeral, or real-time data -- not as a primary database.** **Core principles:** 1. **Cache, don't store** -- Redis complements your primary database. Cache frequently accessed data, but always have a source of truth elsewhere. 2. **Atomic operations** -- Use Lua scripts or MULTI/EXEC for operations spanning multiple keys. Individual Redis commands are atomic, but sequences are not. 3. **Separate concerns** -- Use different Redis databases (or key prefixes) for caching, sessions, and queues. Use separate connections for Pub/Sub. 4. **Set TTLs on everything** -- Memory is finite. Every cached key should expire. Use `EX` (seconds) or `PX` (milliseconds) on SET commands. 5. **Fail gracefully** -- Redis is a cache, not a database. If Redis is down, the application should degrade gracefully (bypass cache, use database directly). </philosophy> --- <patterns> ## Core Patterns ### Pattern 1: ioredis Connection Setup Configure ioredis with proper error handling and reconnection strategy. See [examples/core.md](examples/core.md) for full examples including node-redis and cluster configuration. ```typescript // ✅ Good Example - Proper ioredis setup with error handling import Redis from "ioredis"; const RETRY_DELAY_BASE_MS = 50; const RETRY_DELAY_MAX_MS = 2000; function createRedisClient(): Redis { const url = process.env.REDIS_URL; if (!url) { throw new Error("REDIS_URL environment variable is required"); } const client = new Redis(url, { maxRetriesPerRequest: 3, retryStrategy(times) { return Math.min(times * RETRY_DELAY_BASE_MS, RETRY_DELAY_MAX_MS); }, lazyConnect: true, }); client.on("error", (err) => { console.error("Redis connection error:", err.message); }); return client; } export { createRedisClient }; ``` **Why good:** Environment variable validation, named constants for retry delays, `lazyConnect` prevents connection before ready, error event handler prevents process crash ```typescript // ❌ Bad Example - No error handling, hardcoded config import Redis from "ioredis"; const redis = new Redis("redis://localhost:6379"); // No error handler, no retry strategy, hardcoded URL ``` **Why bad:** Missing error event handler crashes Node.js process on connection failure, hardcoded URL prevents environment-specific configuration --- ### Pattern 2: Cache-Aside Check cache first, fall back to database on miss, populate cache. See [examples/caching.md](examples/caching.md) for write-through, invalidation, stampede prevention, and multi-key cache patterns. ```typescript // ✅ Good Example - Generic cache-aside helper const CACHE_TTL_SECONDS = 300; async function cacheAside<T>( redis: Redis, key: string, fetcher: () => Promise<T>, ttlSeconds: number = CACHE_TTL_SECONDS, ): Promise<T> { const cached = await redis.get(key); if (cached !== null) { return JSON.parse(cached) as T; } const data = await fetcher(); // Fire-and-forget cache write redis.set(key, JSON.stringify(data), "EX", ttlSeconds).catch((err) => { console.error(`Cache write failed for ${key}:`, err.message); }); return data; } ``` **Why good:** Generic `cacheAside<T>` works with any data type, fire-and-forget cache write prevents cache failure from blocking response, configurable TTL --- ### Pattern 3: Pipelining and Transactions Batch commands for network efficiency (pipeline) or atomicity (MULTI/EXEC). See [examples/core.md](examples/core.md) for full examples. ```typescript // ✅ Pipeline - batch for network efficiency (NOT atomic) const pipeline = redis.pipeline(); pipeline.hset(key, { name: user.name, email: user.email }); pipeline.expire(key, TTL_SECONDS); await pipeline.exec(); // ✅ Transaction - atomic execution with optimistic locking await redis.watch(fromKey); const results = await redis .multi() .decrby(fromKey, amount) .incrby(toKey, amount) .exec(); if (!results) { /* WATCH detected change, retry */ } ``` **Why good:** Pipeline reduces round-trips, MULTI/EXEC provides atomicity, WATCH enables optimistic locking ```typescript // ❌ Bad Example - Non-atomic balance transfer await redis.decrby("balance:user1", 100); await redis.incrby("balance:user2", 100); // Crash between commands causes data loss ``` **Why bad:** Two separate commands are not atomic, crash between them causes data inconsistency --- ### Pattern 4: Lua Scripting for Atomicity Lua scripts execute atomically on the Redis server. See [examples/rate-limiting.md](examples/rate-limiting.md) for complete Lua-based rate limiters. ```typescript // ✅ Define custom atomic command redis.defineCommand("rateLimit", { numberOfKeys: 1, lua: ` local key = KEYS[1] local limit = tonumber(ARGV[1]) local window = tonumber(ARGV[2]) local now = tonumber(ARGV[3]) redis.call('ZREMRANGEBYSCORE', key, 0, now - window) local count = redis.call('ZCARD', key) if count < limit then redis.call('ZADD', key, now, now .. '-' .. math.random(1000000)) redis.call('EXPIRE', key, window) return 1 end return 0 `, }); ``` **Why good:** `defineCommand` uses EVALSHA internally for performance, Lua script is atomic -- no race conditions --- ### Pattern 5: Pub/Sub Messaging Requires **separate connections** for subscribing and publishing. See [examples/pub-sub.md](examples/pub-sub.md) for event broadcasting system. ```typescript // ✅ Good Example - Separate connections const publisher = new Redis(url); const subscriber = new Redis(url); // MUST be separate await subscriber.subscribe("notifications"); subscriber.on("message", (channel, message) => { handleNotification(JSON.parse(message)); }); await publisher.publish("notifications", JSON.stringify(data)); ``` **Why good:** Separate connections for pub and sub (required by Redis protocol), error handlers on both ```typescript // ❌ Bad Example - Same connection const redis = new Redis(); await redis.subscribe("channel"); await redis.set("key", "value"); // THROWS: connection is in subscriber mode ``` **Why bad:** A subscribed connection cannot execute non-pub/sub commands --- ### Pattern 6: BullMQ Job Queues Robust job queuing with retries, scheduling, and priorities. See [examples/queues.md](examples/queues.md) for complete examples with workers, events, and graceful shutdown. ```typescript // ✅ Good Example - BullMQ connection factory function createBullMQConnection(): Redis { return new Redis(process.env.REDIS_URL!, { maxRetriesPerRequest: null, // REQUIRED for BullMQ }); } const emailQueue = new Queue<EmailJobData>(QUEUE_NAME, { connection: createBullMQConnection(), defaultJobOptions: { attempts: MAX_RETRY_ATTEMPTS, backoff: { type: "exponential", delay: BACKOFF_DELAY_MS }, }, }); ``` **Why good:** `maxRetriesPerRequest: null` is required for BullMQ, separate connection per Queue/Worker ```typescript // ❌ Bad Example - Missing required config const connection = new Redis(); // Missing maxRetriesPerRequest: null const queue = new Queue("emails", { connection }); // BullMQ will throw: "maxRetriesPerRequest must be null" ``` **Why bad:** BullMQ requires infinite retries -- without `null`, ioredis gives up after a set number of attempts --- ### Pattern 7: Redis Streams Persistent, ordered message logs with consumer groups. See [examples/queues.md](examples/queues.md) for full producer/consumer implementation with pending message recovery. ```typescript // Producer const entryId = await redis.xadd( STREAM_KEY, "*", "orderId", id, "action", action, ); // Consumer group setup (idempotent) try { await redis.xgroup("CREATE", STREAM_KEY, GROUP_NAME, "0", "MKSTREAM"); } catch (err) { if (!(err instanceof Error) || !err.message.includes("BUSYGROUP")) throw err; } // Consumer read + acknowledge const results = await redis.xreadgroup( "GROUP", GROUP_NAME, consumerName, "COUNT", "10", "BLOCK", "5000", "STREAMS", STREAM_KEY, ">", ); await redis.xack(STREAM_KEY, GROUP_NAME, messageId); ``` **Why good:** MKSTREAM creates stream if absent, BUSYGROUP handling for idempotent setup, XACK confirms processing </patterns> --- <performance> ## Performance Optimization - **Auto-pipelining** -- Enable `enableAutoPipelining: true` to batch commands issued during the same event loop tick. Does NOT work with WATCH/MULTI or blocking commands. See [examples/core.md](examples/core.md) for setup. - **Manual pipelining** -- Use `redis.pipeline()` for explicit command batching. See [examples/core.md](examples/core.md). - **Key expiration** -- Set TTLs on all cache keys. Use named constants (`TTL_SHORT_SECONDS = 60`, `TTL_MEDIUM_SECONDS = 300`, etc.). See [examples/core.md](examples/core.md). - **SCAN over KEYS** -- Never use `KEYS` in production (blocks Redis). Use `redis.scanStream({ match: "pattern", count: 100 })`. See [examples/core.md](examples/core.md). - **UNLINK over DEL** -- Use `unlink` for large keys (non-blocking deletion). </performance> --- <decision_framework> ## Decision Framework ### Which Redis Client? ``` Which Redis client should I use? ├─ Need Redis Stack modules (JSON, Search, TimeSeries)? -> node-redis (v5.x) ├─ Using BullMQ for job queues? -> ioredis (BullMQ requires it) ├─ Need Cluster or Sentinel support? -> ioredis (built-in, battle-tested) ├─ Need auto-pipelining? -> ioredis (enableAutoPipelining option) └─ General caching/sessions/pub-sub? -> ioredis (recommended default) ``` ### Which Caching Strategy? ``` How should I cache this data? ├─ Read-heavy, tolerates brief staleness? -> Cache-aside with TTL ├─ Needs strong consistency after writes? -> Write-through (update DB + invalidate cache) ├─ Write-heavy, can tolerate brief data loss? -> Write-behind (async cache update) └─ Data changes rarely? -> Cache-aside with long TTL + manual invalidation ``` ### Which Data Structure? ``` What Redis data structure should I use? ├─ Simple key-value (cache, sessions)? -> Strings (GET/SET) ├─ Object with multiple fields? -> Hashes (HSET/HGET) ├─ Ordered ranking/leaderboard? -> Sorted Sets (ZADD/ZRANGE) ├─ Queue (FIFO/LIFO)? -> Lists (LPUSH/RPOP) ├─ Unique collection (tags, categories)? -> Sets (SADD/SMEMBERS) ├─ Persistent message log with consumers? -> Streams (XADD/XREAD) └─ Rate limiting (sliding window)? -> Sorted Sets + Lua script ``` ### Which Messaging Pattern? ``` How should I implement real-time messaging? ├─ Fire-and-forget broadcast? -> Pub/Sub (no persistence) ├─ Need message persistence and replay? -> Streams with consumer groups ├─ Need reliable job processing with retries? -> BullMQ (built on Redis) └─ Need request-reply pattern? -> Pub/Sub with correlation IDs ``` ### Atomicity Decision ``` Do I need atomicity across multiple commands? ├─ YES -> Are the commands on the same key? │ ├─ YES -> Use a single atomic command (INCR, SETNX, etc.) │ └─ NO -> Use Lua script (defineCommand) ├─ NO, but I want batching -> Use pipeline (non-atomic, single round-trip) └─ Need optimistic locking? -> Use WATCH + MULTI/EXEC ``` </decision_framework> --- <integration> ## Integration Guide **Common integration patterns:** - **Database caching** -- Cache query results using cache-aside pattern; invalidate cache on writes - **HTTP framework sessions** -- Session storage middleware, rate limiting middleware - **BullMQ job queues** -- Built on Redis (requires ioredis with `maxRetriesPerRequest: null`) - **WebSocket scaling** -- Redis adapter for distributing WebSocket connections across servers **Replaces / Conflicts with:** - **In-memory caches** -- Redis provides distributed caching across multiple app instances - **Database-backed sessions** -- Redis sessions are faster and reduce database load - **Simple message queues** -- Redis Streams and BullMQ cover most queue use cases; use dedicated message brokers for complex routing or massive throughput </integration> --- <red_flags> ## RED FLAGS **High Priority Issues:** - Using the same connection for Pub/Sub subscribe and regular commands -- subscribed connections cannot execute non-pub/sub commands - Missing `maxRetriesPerRequest: null` on BullMQ connections -- BullMQ throws immediately without this setting - Using `KEYS` command in production -- blocks the entire Redis server while scanning all keys - No `error` event handler on Redis client -- unhandled errors crash the Node.js process - Storing large objects (> 1 MB) in Redis -- degrades performance and wastes memory; store a reference and fetch from object storage **Medium Priority Issues:** - Missing TTL on cached keys -- causes unbounded memory growth until Redis runs out of memory - Using `del` with many keys instead of `unlink` -- `del` blocks Redis; `unlink` frees memory asynchronously - Not using pipelining for batch operations -- each command is a separate network round-trip - Serializing/deserializing complex objects without error handling -- malformed JSON in cache crashes on parse - Sharing a single Redis connection across BullMQ Queue and Worker -- each needs its own connection **Common Mistakes:** - Assuming pipeline commands are atomic -- pipelines batch for network efficiency but do not provide atomicity (use MULTI/EXEC or Lua) - Forgetting that `hgetall` returns an empty object `{}` for non-existent keys (not `null`) -- check `Object.keys(result).length === 0` - Using `MULTI/EXEC` without `WATCH` for conditional updates -- transactions execute unconditionally unless you WATCH keys first - Not handling `null` returns from `GET` -- cache misses return `null`, not `undefined` - Connecting to Redis without TLS in production -- credentials sent in plaintext over the network **Gotchas & Edge Cases:** - Redis `HGETALL` returns all values as strings -- numbers stored with `HSET` come back as strings, requiring explicit parsing - `EXPIRE` resets when a key is overwritten with `SET` -- if you `SET` a key that already has a TTL, the TTL is removed unless you include `EX`/`PX` in the `SET` command - Pub/Sub messages are fire-and-forget -- if no subscriber is listening when a message is published, it is lost forever (use Streams for persistence) - Redis Cluster does not support multi-key operations across different hash slots -- use `{hash-tag}` prefix to force related keys to the same slot - `WATCH` is connection-scoped -- concurrent requests sharing a connection will interfere with each other's WATCH state - ioredis auto-pipelining does not work with `WATCH`/`MULTI` or blocking commands (`BRPOP`, `BLPOP`, `XREAD BLOCK`) </red_flags> --- <critical_reminders> ## CRITICAL REMINDERS > **All code must follow project conventions in CLAUDE.md** (kebab-case, named exports, import ordering, `import type`, named constants) **(You MUST use a SEPARATE Redis connection for Pub/Sub subscribers -- a subscribed connection enters a special mode and cannot execute other commands)** **(You MUST set `maxRetriesPerRequest: null` on any ioredis connection passed to BullMQ -- BullMQ requires infinite retries and will throw if this is not set)** **(You MUST use Lua scripts (`defineCommand` or `eval`) for any operation requiring atomicity across multiple Redis commands -- separate commands are NOT atomic even in a pipeline)** **(You MUST handle the `error` event on every Redis client instance -- unhandled errors crash the Node.js process)** **Failure to follow these rules will cause pub/sub failures, BullMQ connection errors, race conditions, and application crashes.** </critical_reminders>
Comments (0)
Sign in to join the conversation.
Reviews (0)
No reviews yet.
No comments yet.