📖 数据密集型设计

分布式缓存与多级缓存架构

深入探讨分布式缓存原理、多级缓存架构设计及缓存一致性保障

一、分布式缓存概述

分布式缓存是数据密集型应用的关键组件,通过将热点数据缓存到内存中,能够显著提升系统性能。多级缓存架构通过组合本地缓存和分布式缓存,能够进一步优化缓存命中率和减少网络开销。

二、多级缓存架构

2.1 多级缓存层次

graph TD A[客户端请求] --> B[应用层] B --> C[本地缓存 L1] C -->|命中| D[返回数据] C -->|未命中| E[分布式缓存 L2] E -->|命中| F[更新本地缓存] F --> D E -->|未命中| G[数据库] G -->|查询数据| H[更新分布式缓存] H --> I[更新本地缓存] I --> D J[数据更新] --> K[失效本地缓存] K --> L[失效分布式缓存] L --> M[更新数据库]

2.2 缓存层级对比

层级 缓存类型 特点 优点 缺点
L1 本地缓存 进程内内存 速度快、无网络开销 内存有限、数据不一致
L2 分布式缓存 独立缓存服务 容量大、数据一致 网络延迟、单点故障
L3 数据库 持久化存储 数据持久、事务支持 性能差、IO开销

三、分布式缓存实现

3.1 Redis集群缓存

public class RedisClusterCache : IDistributedCache
{
    private readonly IConnectionMultiplexer _connectionMultiplexer;
    
    public RedisClusterCache(IConnectionMultiplexer connectionMultiplexer)
    {
        _connectionMultiplexer = connectionMultiplexer;
    }
    
    public async Task GetAsync(string key)
    {
        var db = _connectionMultiplexer.GetDatabase();
        var value = await db.StringGetAsync(key);
        
        if (!value.HasValue)
        {
            return default;
        }
        
        return JsonSerializer.Deserialize(value);
    }
    
    public async Task SetAsync(string key, T value, TimeSpan? expiry = null)
    {
        var db = _connectionMultiplexer.GetDatabase();
        var json = JsonSerializer.Serialize(value);
        
        await db.StringSetAsync(key, json, expiry);
    }
    
    public async Task RemoveAsync(string key)
    {
        var db = _connectionMultiplexer.GetDatabase();
        await db.KeyDeleteAsync(key);
    }
    
    public async Task ExistsAsync(string key)
    {
        var db = _connectionMultiplexer.GetDatabase();
        return await db.KeyExistsAsync(key);
    }
    
    public async Task GetOrCreateAsync(string key, Func> factory, TimeSpan? expiry = null)
    {
        var value = await GetAsync(key);
        
        if (value != null)
        {
            return value;
        }
        
        value = await factory();
        
        if (value != null)
        {
            await SetAsync(key, value, expiry);
        }
        
        return value;
    }
    
    public async Task IncrementAsync(string key, long value = 1)
    {
        var db = _connectionMultiplexer.GetDatabase();
        return await db.StringIncrementAsync(key, value);
    }
    
    public async Task DecrementAsync(string key, long value = 1)
    {
        var db = _connectionMultiplexer.GetDatabase();
        return await db.StringDecrementAsync(key, value);
    }
}

3.2 本地缓存实现

public class LocalCache
{
    private readonly ConcurrentDictionary> _cache = new();
    private readonly TimeSpan _defaultExpiry;
    private readonly int _maxSize;
    
    public LocalCache(TimeSpan defaultExpiry, int maxSize = 1000)
    {
        _defaultExpiry = defaultExpiry;
        _maxSize = maxSize;
    }
    
    public T Get(string key)
    {
        if (_cache.TryGetValue(key, out var entry))
        {
            if (entry.ExpiryTime > DateTime.UtcNow)
            {
                entry.AccessTime = DateTime.UtcNow;
                return entry.Value;
            }
            
            _cache.TryRemove(key, out _);
        }
        
        return default;
    }
    
    public void Set(string key, T value, TimeSpan? expiry = null)
    {
        if (_cache.Count >= _maxSize)
        {
            EvictLeastRecentlyUsed();
        }
        
        _cache[key] = new CacheEntry
        {
            Value = value,
            ExpiryTime = DateTime.UtcNow + (expiry ?? _defaultExpiry),
            AccessTime = DateTime.UtcNow
        };
    }
    
    public bool Remove(string key)
    {
        return _cache.TryRemove(key, out _);
    }
    
    public bool TryGetValue(string key, out T value)
    {
        value = Get(key);
        return value != null;
    }
    
    public void Clear()
    {
        _cache.Clear();
    }
    
    public int Count => _cache.Count;
    
    private void EvictLeastRecentlyUsed()
    {
        var oldest = _cache.OrderBy(kv => kv.Value.AccessTime).FirstOrDefault();
        
        if (oldest.Key != null)
        {
            _cache.TryRemove(oldest.Key, out _);
        }
    }
    
    private class CacheEntry
    {
        public TValue Value { get; set; }
        public DateTime ExpiryTime { get; set; }
        public DateTime AccessTime { get; set; }
    }
}

四、多级缓存服务

4.1 多级缓存管理器

public class MultiLevelCacheManager : IMultiLevelCache
{
    private readonly LocalCache _localCache;
    private readonly IDistributedCache _distributedCache;
    private readonly TimeSpan _localCacheExpiry = TimeSpan.FromMinutes(5);
    private readonly TimeSpan _distributedCacheExpiry = TimeSpan.FromHours(1);
    
    public MultiLevelCacheManager(IDistributedCache distributedCache)
    {
        _localCache = new LocalCache(_localCacheExpiry, 1000);
        _distributedCache = distributedCache;
    }
    
    public async Task GetAsync(string key)
    {
        var localValue = _localCache.Get(key);
        
        if (localValue != null)
        {
            return (T)localValue;
        }
        
        var distributedValue = await _distributedCache.GetAsync(key);
        
        if (distributedValue != null)
        {
            _localCache.Set(key, distributedValue, _localCacheExpiry);
        }
        
        return distributedValue;
    }
    
    public async Task SetAsync(string key, T value, TimeSpan? expiry = null)
    {
        _localCache.Set(key, value, expiry ?? _localCacheExpiry);
        await _distributedCache.SetAsync(key, value, expiry ?? _distributedCacheExpiry);
    }
    
    public async Task RemoveAsync(string key)
    {
        _localCache.Remove(key);
        await _distributedCache.RemoveAsync(key);
    }
    
    public async Task GetOrCreateAsync(string key, Func> factory, TimeSpan? expiry = null)
    {
        var localValue = _localCache.Get(key);
        
        if (localValue != null)
        {
            return (T)localValue;
        }
        
        var distributedValue = await _distributedCache.GetAsync(key);
        
        if (distributedValue != null)
        {
            _localCache.Set(key, distributedValue, expiry ?? _localCacheExpiry);
            return distributedValue;
        }
        
        var value = await factory();
        
        if (value != null)
        {
            _localCache.Set(key, value, expiry ?? _localCacheExpiry);
            await _distributedCache.SetAsync(key, value, expiry ?? _distributedCacheExpiry);
        }
        
        return value;
    }
    
    public async Task InvalidateCacheAsync(string pattern)
    {
        await _distributedCache.RemoveByPatternAsync(pattern);
    }
    
    public void ClearLocalCache()
    {
        _localCache.Clear();
    }
}
            
            

4.2 缓存策略实现

public class CacheStrategyService
{
    public async Task CacheAsideGetAsync(string key, Func> databaseGetter)
    {
        var cached = await _cache.GetAsync(key);
        
        if (cached != null)
        {
            return cached;
        }
        
        var data = await databaseGetter();
        
        if (data != null)
        {
            await _cache.SetAsync(key, data);
        }
        
        return data;
    }
    
    public async Task CacheAsideUpdateAsync(string key, T data, Func databaseUpdater)
    {
        await databaseUpdater();
        
        await _cache.RemoveAsync(key);
    }
    
    public async Task ReadThroughGetAsync(string key, Func> databaseGetter)
    {
        var cached = await _cache.GetAsync(key);
        
        if (cached != null)
        {
            return cached;
        }
        
        var data = await databaseGetter();
        
        if (data != null)
        {
            await _cache.SetAsync(key, data);
        }
        
        return data;
    }
    
    public async Task WriteThroughUpdateAsync(string key, T data, Func databaseUpdater)
    {
        await databaseUpdater();
        
        await _cache.SetAsync(key, data);
    }
    
    public async Task WriteBehindUpdateAsync(string key, T data, Func databaseUpdater)
    {
        await _cache.SetAsync(key, data);
        
        _backgroundQueue.QueueBackgroundWorkItem(async token =>
        {
            await databaseUpdater();
        });
    }
}

五、缓存一致性

5.1 缓存一致性策略

graph TD A[数据更新] --> B{一致性策略} B -->|Cache-Aside| C[更新数据库] C --> D[删除缓存] D --> E[下次读取时回填] B -->|Write-Through| F[更新数据库] F --> G[同步更新缓存] B -->|Write-Behind| H[更新缓存] H --> I[异步更新数据库] J[缓存失效] --> K{失效策略} K -->|TTL| L[自动过期] K -->|主动删除| M[更新时删除] K -->|消息通知| N[订阅变更消息]

5.2 分布式缓存一致性

public class CacheConsistencyService
{
    public async Task UpdateDataWithCacheAsync(string key, T data, Func updateDatabase)
    {
        await updateDatabase();
        
        await RemoveCacheAsync(key);
    }
    
    public async Task RemoveCacheAsync(string key)
    {
        await _distributedCache.RemoveAsync(key);
        
        await _cacheInvalidationService.BroadcastInvalidationAsync(key);
    }
    
    public async Task HandleCacheInvalidationAsync(string key)
    {
        _localCache.Remove(key);
    }
    
    public async Task GetDataWithRetryAsync(string key, Func> getter, int retryCount = 3)
    {
        for (int i = 0; i < retryCount; i++)
        {
            try
            {
                return await _cache.GetAsync(key);
            }
            catch (CacheException)
            {
                if (i == retryCount - 1)
                {
                    throw;
                }
                
                await Task.Delay(100 * (i + 1));
            }
        }
        
        return default;
    }
    
    public async Task GetDataWithFallbackAsync(string key, Func> cacheGetter, Func> databaseGetter)
    {
        try
        {
            var cached = await cacheGetter();
            
            if (cached != null)
            {
                return cached;
            }
        }
        catch (Exception)
        {
        }
        
        return await databaseGetter();
    }
}

public class CacheInvalidationService
{
    private readonly IMessageBus _messageBus;
    private const string InvalidationChannel = "cache-invalidation";
    
    public CacheInvalidationService(IMessageBus messageBus)
    {
        _messageBus = messageBus;
    }
    
    public async Task BroadcastInvalidationAsync(string key)
    {
        var message = new CacheInvalidationMessage { Key = key, Timestamp = DateTime.UtcNow };
        
        await _messageBus.PublishAsync(InvalidationChannel, message);
    }
    
    public async Task SubscribeToInvalidationAsync(Func handler)
    {
        await _messageBus.SubscribeAsync(InvalidationChannel, async (message) =>
        {
            var invalidationMessage = JsonSerializer.Deserialize(message);
            
            if (invalidationMessage != null)
            {
                await handler(invalidationMessage.Key);
            }
        });
    }
}

public class CacheInvalidationMessage
{
    public string Key { get; set; }
    public DateTime Timestamp { get; set; }
}

六、缓存问题解决方案

6.1 缓存穿透解决方案

public class CachePenetrationGuard
{
    private readonly ISet _bloomFilter;
    private readonly int _expectedInsertions = 1000000;
    private readonly double _falsePositiveRate = 0.01;
    
    public CachePenetrationGuard()
    {
        _bloomFilter = new BloomFilter(_expectedInsertions, _falsePositiveRate);
    }
    
    public async Task GetWithPenetrationGuardAsync(string key, Func> databaseGetter)
    {
        if (!_bloomFilter.Contains(key))
        {
            return default;
        }
        
        var cached = await _cache.GetAsync(key);
        
        if (cached != null)
        {
            return cached;
        }
        
        var data = await databaseGetter();
        
        if (data != null)
        {
            await _cache.SetAsync(key, data);
        }
        else
        {
            await _cache.SetAsync(key, default(T), TimeSpan.FromMinutes(5));
            
            _bloomFilter.Add(key);
        }
        
        return data;
    }
    
    public void AddKeyToFilter(string key)
    {
        _bloomFilter.Add(key);
    }
    
    public void RemoveKeyFromFilter(string key)
    {
        _bloomFilter.Remove(key);
    }
}

public class BloomFilter
{
    private readonly BitArray _bitArray;
    private readonly int _hashCount;
    private readonly int _size;
    
    public BloomFilter(int expectedInsertions, double falsePositiveRate)
    {
        _size = CalculateSize(expectedInsertions, falsePositiveRate);
        _hashCount = CalculateHashCount(_size, expectedInsertions);
        _bitArray = new BitArray(_size);
    }
    
    private int CalculateSize(int n, double p)
    {
        return (int)Math.Ceiling(-n * Math.Log(p) / Math.Pow(Math.Log(2), 2));
    }
    
    private int CalculateHashCount(int m, int n)
    {
        return (int)Math.Ceiling((m / n) * Math.Log(2));
    }
    
    public void Add(string item)
    {
        for (int i = 0; i < _hashCount; i++)
        {
            var hash = ComputeHash(item, i);
            _bitArray[hash % _size] = true;
        }
    }
    
    public bool Contains(string item)
    {
        for (int i = 0; i < _hashCount; i++)
        {
            var hash = ComputeHash(item, i);
            if (!_bitArray[hash % _size])
            {
                return false;
            }
        }
        
        return true;
    }
    
    private int ComputeHash(string item, int seed)
    {
        return item.GetHashCode() ^ seed;
    }
}

6.2 缓存击穿解决方案

public class CacheBreakdownGuard
{
    private readonly ConcurrentDictionary _lockDictionary = new();
    
    public async Task GetWithBreakdownGuardAsync(string key, Func> databaseGetter, TimeSpan expiry)
    {
        var cached = await _cache.GetAsync(key);
        
        if (cached != null)
        {
            return cached;
        }
        
        var semaphore = _lockDictionary.GetOrAdd(key, _ => new SemaphoreSlim(1, 1));
        
        try
        {
            await semaphore.WaitAsync();
            
            cached = await _cache.GetAsync(key);
            
            if (cached != null)
            {
                return cached;
            }
            
            var data = await databaseGetter();
            
            if (data != null)
            {
                await _cache.SetAsync(key, data, expiry);
            }
            
            return data;
        }
        finally
        {
            semaphore.Release();
            
            _lockDictionary.TryRemove(key, out _);
        }
    }
    
    public async Task GetWithRandomExpiryAsync(string key, Func> databaseGetter, TimeSpan baseExpiry)
    {
        var cached = await _cache.GetAsync(key);
        
        if (cached != null)
        {
            return cached;
        }
        
        var data = await databaseGetter();
        
        if (data != null)
        {
            var randomOffset = TimeSpan.FromSeconds(new Random().Next(60, 300));
            await _cache.SetAsync(key, data, baseExpiry + randomOffset);
        }
        
        return data;
    }
}

6.3 缓存雪崩解决方案

public class CacheAvalancheGuard
{
    private readonly TimeSpan _defaultExpiry = TimeSpan.FromHours(1);
    
    public async Task GetWithAvalancheGuardAsync(string key, Func> databaseGetter)
    {
        var cached = await _cache.GetAsync(key);
        
        if (cached != null)
        {
            await ExtendExpiryAsync(key);
            
            return cached;
        }
        
        var data = await databaseGetter();
        
        if (data != null)
        {
            await _cache.SetAsync(key, data, GetRandomExpiry());
        }
        
        return data;
    }
    
    private TimeSpan GetRandomExpiry()
    {
        var random = new Random();
        var offset = TimeSpan.FromMinutes(random.Next(30, 90));
        
        return _defaultExpiry + offset;
    }
    
    private async Task ExtendExpiryAsync(string key)
    {
        await _cache.SetExpiryAsync(key, GetRandomExpiry());
    }
    
    public async Task SetWithAvalancheGuardAsync(string key, T value)
    {
        await _cache.SetAsync(key, value, GetRandomExpiry());
    }
    
    public async Task WarmUpCacheAsync(List keys, Func> dataProvider)
    {
        foreach (var key in keys)
        {
            var data = await dataProvider(key);
            
            if (data != null)
            {
                await _cache.SetAsync(key, data, GetRandomExpiry());
            }
        }
    }
}

七、分布式缓存最佳实践

7.1 缓存策略选择

策略 适用场景 优点 缺点
Cache-Aside 读多写少 实现简单 可能不一致
Read-Through 读密集 透明缓存 需要封装
Write-Through 写密集 强一致性 写延迟
Write-Behind 极高写入 性能最好 数据丢失风险

7.2 多级缓存最佳实践

  • 合理设置缓存层级
  • 设置合理的过期时间
  • 实现缓存失效通知
  • 处理缓存穿透/击穿/雪崩
  • 监控缓存命中率

八、总结

分布式缓存与多级缓存架构是数据密集型应用的关键技术。通过组合本地缓存和分布式缓存,能够兼顾性能和一致性。合理选择缓存策略,处理缓存穿透、击穿、雪崩等问题,能够构建稳定可靠的缓存系统。缓存一致性保障和监控是生产环境中必须考虑的问题。