A production-ready distributed caching system built from first principles. Features consistent hashing, LRU eviction, TTL support, active cache invalidation, and horizontal scaling without downtime.
- Architecture
- Request Flow
- Core Features
- Redis Caching Strategy
- Queue System (Cache Invalidation)
- Distributed ID Generation
- Database Sharding
- Performance Benchmarks
- Quick Start
- Configuration
- Monitoring
- API Reference
- Deployment
- Contributing
- License
| Decision | Rationale |
|---|---|
| 16,384 hash slots | Balance between granularity (smooth rebalancing) and overhead |
| Virtual nodes per physical node | 4,096 slots/node default, reduces hotspotting |
| Async replication | Better availability; sync for critical paths only |
| Gossip protocol | Decentralized topology; no single point of failure |
| Jump consistent hash | O(1) computation, minimal key remapping on node changes |
Consistent Hashing Ring (Distributed Cache)
package consistenthash
import (
"hash/crc32"
"sort"
)
type Ring struct {
slots []uint16 // 16384 slots
nodes map[uint16]string // slot -> node_id
replicas int // virtual nodes per physical node
}
func NewRing(nodes []string) *Ring {
ring := &Ring{
slots: make([]uint16, 16384),
nodes: make(map[uint16]string),
replicas: 4096, // 16384 / 4 nodes default
}
ring.rebalance(nodes)
return ring
}
func (r *Ring) GetNode(key string) string {
hash := crc32.ChecksumIEEE([]byte(key))
slot := hash % 16384
return r.nodes[r.slots[slot]]
}
func (r *Ring) AddNode(nodeID string) {
// Rebalance: move ~4096 slots with minimal key remapping
// Only 1/N keys need to move (N = new node count)
}
func (r *Ring) RemoveNode(nodeID string) {
// Hand off slots to neighbors
// Background migration with dual-read
}Key Properties:
Monotonicity: Adding a node only moves keys to the new node, never between existing nodes
Balance: Each node owns roughly equal number of slots (±1)
Spread: Minimal key movement during topology changes (~1/N keys)
type LRUCache struct {
capacity int
ttl time.Duration
items map[string]*list.Element
evictList *list.List // container/list
mu sync.RWMutex
}
type entry struct {
key string
value []byte
expireAt time.Time
frequency int // for LFU mode
}
func (c *LRUCache) Get(key string) ([]byte, bool) {
c.mu.RLock()
elem, ok := c.items[key]
c.mu.RUnlock()
if !ok {
return nil, false
}
// Check TTL expiration
if time.Now().After(elem.Value.(*entry).expireAt) {
c.Delete(key)
return nil, false
}
c.mu.Lock()
c.evictList.MoveToFront(elem) // Mark as recently used
c.mu.Unlock()
return elem.Value.(*entry).value, true
}
func (c *LRUCache) Set(key string, value []byte, ttl time.Duration) {
c.mu.Lock()
defer c.mu.Unlock()
if elem, ok := c.items[key]; ok {
c.evictList.MoveToFront(elem)
elem.Value.(*entry).value = value
elem.Value.(*entry).expireAt = time.Now().Add(ttl)
return
}
// Evict if at capacity
if c.evictList.Len() >= c.capacity {
c.evictLRU()
}
ent := &entry{
key: key,
value: value,
expireAt: time.Now().Add(ttl),
}
elem := c.evictList.PushFront(ent)
c.items[key] = elem
}
func (c *LRUCache) evictLRU() {
back := c.evictList.Back()
if back == nil {
return
}
ent := back.Value.(*entry)
delete(c.items, ent.key)
c.evictList.Remove(back)
}type TTLManager struct {
buckets [3600]*list.List // 1-second buckets for next hour
wheel []chan string // 3600 channels
current int // current bucket index
}
func (t *TTLManager) Start() {
ticker := time.NewTicker(time.Second)
go func() {
for range ticker.C {
t.current = (t.current + 1) % 3600
// Expire all keys in current bucket
for _, key := range t.buckets[t.current] {
cache.Delete(key)
}
t.buckets[t.current] = list.New() // Reset bucket
}
}()
}
func (t *TTLManager) Schedule(key string, ttl time.Duration) {
bucket := (t.current + int(ttl.Seconds())) % 3600
t.buckets[bucket].PushBack(key)
}Pattern 1: Cache-Aside (Lazy Loading)
func (s *UserService) GetUser(ctx context.Context, userID string) (*User, error) {
cacheKey := fmt.Sprintf("user:%s", userID)
// 1. Attempt cache read
cached, err := s.cache.Get(ctx, cacheKey)
if err == nil {
var user User
if err := json.Unmarshal(cached, &user); err == nil {
metrics.CacheHits.Inc()
return &user, nil
}
}
metrics.CacheMisses.Inc()
// 2. Cache miss - load from database
user, err := s.db.QueryUser(ctx, userID)
if err != nil {
return nil, fmt.Errorf("database error: %w", err)
}
// 3. Serialize and store with TTL
data, _ := json.Marshal(user)
s.cache.Set(ctx, cacheKey, data, 5*time.Minute)
return user, nil
}When to use: Read-heavy workloads, eventual consistency acceptable.
This pattern ensures database is always the source of truth, and cache is synchronized immediately after writes.
func (s *UserService) UpdateUser(ctx context.Context, userID string, update UpdateRequest) (*User, error) {
cacheKey := fmt.Sprintf("user:%s", userID)
// 1. Update database first (source of truth)
tx, err := s.db.Begin(ctx)
if err != nil {
return nil, err
}
defer tx.Rollback()
user, err := tx.UpdateUser(ctx, userID, update)
if err != nil {
return nil, err
}
if err := tx.Commit(); err != nil {
return nil, err
}
// 2. Invalidate cache immediately
if err := s.cache.Delete(ctx, cacheKey); err != nil {
log.Warn("cache invalidation failed", "key", cacheKey, "error", err)
// Queue retry for eventual consistency recovery
s.invalidQueue.Publish(InvalidationEvent{
Key: cacheKey,
Retry: true,
})
}
// 3. Optional cache warm-up (write-through optimization)
data, _ := json.Marshal(user)
s.cache.Set(ctx, cacheKey, data, 5*time.Minute)
return user, nil
}- Write-heavy workloads
- Strong consistency required
- Critical data (wallets, inventory, user profiles)
Prevents cache avalanche / stampede problem when many requests hit expired keys.
func (s *CacheService) GetWithStampedeProtection(
ctx context.Context,
key string,
ttl time.Duration,
compute func() ([]byte, error),
) ([]byte, error) {
value, err := s.cache.Get(ctx, key)
if err != nil {
// Cache miss → compute immediately
return s.computeAndStore(ctx, key, ttl, compute)
}
// Check remaining TTL
remaining, err := s.cache.TTL(ctx, key)
if err != nil {
return value, nil // fail-safe: return stale data
}
// Probabilistic early refresh
beta := 1.0
probability := math.Exp(
float64(remaining) / (float64(ttl) / beta),
)
if rand.Float64() < probability {
lockKey := fmt.Sprintf("lock:recompute:%s", key)
acquired, _ := s.cache.SetNX(
ctx,
lockKey,
"1",
10*time.Second,
)
if acquired {
go func() {
defer s.cache.Delete(ctx, lockKey)
_ = s.computeAndStore(ctx, key, ttl, compute)
}()
}
}
return value, nil
}- Hot keys (viral content, trending products)
- Expensive recomputation
- High QPS systems
Optimizes performance when fetching multiple records.
func (s *UserService) GetUsersBatch(ctx context.Context, userIDs []string) ([]*User, error) {
// 1. Build cache keys
keys := make([]string, len(userIDs))
for i, id := range userIDs {
keys[i] = fmt.Sprintf("user:%s", id)
}
// 2. Pipeline MGET (batch cache fetch)
cached, err := s.cache.MGet(ctx, keys...)
if err != nil {
return nil, err
}
var misses []string
results := make([]*User, len(userIDs))
// 3. Separate cache hits vs misses
for i, val := range cached {
if val != nil {
var user User
_ = json.Unmarshal(val, &user)
results[i] = &user
} else {
misses = append(misses, userIDs[i])
}
}
// 4. Fetch missing data from DB in batch
if len(misses) > 0 {
users, err := s.db.QueryUsers(ctx, misses)
if err != nil {
return nil, err
}
// 5. Pipeline cache writes
pipe := s.cache.Pipeline()
for _, user := range users {
data, _ := json.Marshal(user)
pipe.Set(
ctx,
fmt.Sprintf("user:%s", user.ID),
data,
5*time.Minute,
)
}
_ = pipe.Exec(ctx)
}
return results, nil
}Invalidation Message Schema
{
"schema_version": "1.0",
"type": "INVALIDATE",
"action": "DELETE",
"keys": ["user:123", "user:123:profile"],
"pattern": null,
"version": 1699123456789,
"source": {
"service": "user-service",
"instance": "user-svc-7b9c4d5f2-xv8p9",
"timestamp": "2024-01-15T10:30:00.123Z"
},
"causality": {
"vector_clock": {
"user-service": 42,
"cache-cluster": 41
}
},
"ttl_override": 0,
"propagate": true
}
| Action | Description | Example |
|---|---|---|
DELETE |
Remove specific keys | DELETE user:123 |
PATTERN |
Remove by glob pattern | PATTERN session:* |
FLUSH |
Clear entire cache | FLUSH (emergency only) |
REFRESH |
Force recompute | REFRESH leaderboard:daily |
DECR |
Atomic decrement | DECR stock:789 |
type InvalidationConsumer struct {
cache *CacheCluster
subscriber pubsub.Subscriber
buffer chan InvalidationEvent
workers int
}
func (c *InvalidationConsumer) Start(ctx context.Context) error {
// Subscribe to invalidation channel
ch, err := c.subscriber.Subscribe(ctx, "cache-invalidation")
if err != nil {
return err
}
// Start worker pool
for i := 0; i < c.workers; i++ {
go c.worker(ctx, i)
}
// Read from pub/sub and dispatch
go func() {
for msg := range ch {
var event InvalidationEvent
if err := json.Unmarshal(msg.Payload, &event); err != nil {
log.Error("invalid invalidation message", "error", err)
continue
}
// Non-blocking send to buffer
select {
case c.buffer <- event:
default:
metrics.InvalidationBufferFull.Inc()
}
}
}()
return nil
}
func (c *InvalidationConsumer) worker(ctx context.Context, id int) {
for event := range c.buffer {
// Check causality: only process if version > local version
localVer := c.cache.GetVersion(event.Keys[0])
if event.Version <= localVer {
metrics.InvalidationStale.Inc()
continue // Already newer in cache
}
switch event.Action {
case "DELETE":
for _, key := range event.Keys {
if err := c.cache.Delete(ctx, key); err != nil {
log.Error("invalidation failed", "key", key, "error", err)
}
}
case "PATTERN":
if err := c.cache.DeletePattern(ctx, event.Pattern); err != nil {
log.Error("pattern invalidation failed", "pattern", event.Pattern)
}
case "DECR":
for _, key := range event.Keys {
c.cache.Decr(ctx, key, 1)
}
}
metrics.InvalidationsProcessed.Inc()
}
}A scalable approach for generating globally unique, time-ordered IDs in distributed systems.
0 1 2 3
0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
|0| 41-bit Timestamp (milliseconds) |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
| 10-bit Node ID | 12-bit Sequence |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+- 1 bit → Always
0(ensures positive number) - 41 bits → Timestamp in milliseconds (≈ 69 years lifespan)
- 10 bits → Node ID (supports 1024 distributed nodes)
- 12 bits → Sequence number (4096 IDs per millisecond per node)
- No central coordination required
- Time-ordered IDs (perfect for databases)
- High throughput generation
- Collision-free across distributed nodes
package snowflake
import (
"fmt"
"sync"
"time"
)
const (
epoch = 1609459200000 // 2021-01-01 00:00:00 UTC
nodeBits = 10
sequenceBits = 12
maxNodeID = (1 << nodeBits) - 1
maxSequence = (1 << sequenceBits) - 1
nodeShift = sequenceBits
timestampShift = sequenceBits + nodeBits
)
type Generator struct {
mu sync.Mutex
nodeID uint16
sequence uint16
lastTime int64
}
func NewGenerator(nodeID uint16) (*Generator, error) {
if nodeID > maxNodeID {
return nil, fmt.Errorf("node ID must be between 0 and %d", maxNodeID)
}
return &Generator{
nodeID: nodeID,
}, nil
}
func (g *Generator) Generate() (uint64, error) {
g.mu.Lock()
defer g.mu.Unlock()
now := time.Now().UnixMilli()
if now < g.lastTime {
return 0, fmt.Errorf("clock regression detected")
}
if now == g.lastTime {
g.sequence++
if g.sequence > maxSequence {
for now <= g.lastTime {
now = time.Now().UnixMilli()
}
g.sequence = 0
}
} else {
g.sequence = 0
}
g.lastTime = now
id := uint64((now-epoch)<<timestampShift) |
uint64(g.nodeID)<<nodeShift |
uint64(g.sequence)
return id, nil
}func (g *Generator) ExtractTimestamp(id uint64) time.Time {
ms := int64((id >> timestampShift) + epoch)
return time.UnixMilli(ms)
}
func (g *Generator) ExtractNodeID(id uint64) uint16 {
return uint16((id >> nodeShift) & maxNodeID)
}| Component | Bits | Range | Purpose |
|---|---|---|---|
| Timestamp | 41 | ~69 years | Ordering + time reference |
| Node ID | 10 | 0–1023 | Distributed node uniqueness |
| Sequence | 12 | 0–4095 | High throughput per ms |
IDs naturally sort by creation time.
Each node generates IDs independently.
Up to 4096 IDs/ms per node.
You can extract:
- Creation time
- Origin node
- Generation order
- URL shortener ID generation
- Distributed databases
- Event tracking systems
- Messaging systems
- Analytics pipelines
In a system like :contentReference[oaicite:0]{index=0}, Snowflake IDs ensure:
- Unique short URLs at global scale
- Fast generation without DB locks
- Natural ordering for analytics queries
- High-performance write paths
Sharding Strategy

