将MemoryCache内容写入数据库及竞态条件问题咨询
解决方案
先梳理现有代码的核心问题
- 非线程安全集合风险:用
HashSet<string>存储值,多线程并发添加时会出现数据丢失或异常,因为HashSet不是线程安全的。 - 定时逻辑竞态:
_lastUpdate是普通DateTime,多线程读写非原子,会导致多个线程同时触发数据库写入操作。 - 缓存操作原子性缺失:修改缓存后调用
Set覆盖旧值,若多个线程同时修改同一个key的集合,会覆盖彼此的修改。 - 枚举集合的潜在异常:遍历
_memoryCacheKeys时,若有线程同时增删键,会抛出枚举修改异常。
单实例并发安全改造方案
1. 替换为线程安全集合
用ConcurrentDictionary<string, byte>模拟线程安全的去重集合(键存value,值无意义),替代非线程安全的HashSet。
2. 改用Timer触发定时写入
把定时写入逻辑从消息处理流程中剥离,用Timer单独调度,避免每个消息都触发判断逻辑,减少竞态。
3. 加锁确保单次写入
用Monitor锁保证同一时间只有一个写入任务在执行,避免重复写入。
4. 原子操作处理缓存读写
通过GetOrCreate获取或创建缓存项,写入DB前先原子移除缓存,避免后续消息修改已待写入的数据。
改造后代码示例
public class MyMessageConsumer : BaseMessageConsumer, IDisposable { private readonly IMemoryCache _memoryCache; private readonly Timer _flushTimer; private readonly object _flushLock = new object(); private readonly ConcurrentDictionary<string, bool> _cacheKeys = new ConcurrentDictionary<string, bool>(); public MyMessageConsumer(IMemoryCache memoryCache) { _memoryCache = memoryCache; // 初始化Timer:30分钟后首次执行,之后每30分钟执行一次 _flushTimer = new Timer(FlushToDbAsync, null, TimeSpan.FromMinutes(30), TimeSpan.FromMinutes(30)); } internal async Task<MessageConsumeResult> UpdateAsync(string cacheKey, string value) { // 获取或创建线程安全的去重集合 var valuesSet = _memoryCache.GetOrCreate(cacheKey, entry => { entry.AbsoluteExpirationRelativeToNow = TimeSpan.FromHours(48); return new ConcurrentDictionary<string, byte>(); }); // 线程安全添加值(自动去重) valuesSet.TryAdd(value, byte.MinValue); // 跟踪缓存键 _cacheKeys.TryAdd(cacheKey, true); return MessageConsumeResult.Completed(); } private async void FlushToDbAsync(object state) { // 避免并发执行写入任务 if (!Monitor.TryEnter(_flushLock)) return; try { // 遍历键的副本,避免枚举时集合被修改 var keysToFlush = _cacheKeys.Keys.ToList(); foreach (var key in keysToFlush) { // 原子取出并移除缓存项 if (_memoryCache.TryGetValue(key, out ConcurrentDictionary<string, byte> valuesSet)) { _memoryCache.Remove(key); _cacheKeys.TryRemove(key, out _); var values = valuesSet.Keys.ToList(); if (values.Any()) { await WriteToSomeDbAsync(key, values); } } } } finally { Monitor.Exit(_flushLock); } } private async Task WriteToSomeDbAsync(string key, List<string> values) { // 数据库写入需用原子追加逻辑,避免多实例/线程冲突 // 示例(SQL Server):用MERGE结合JSON操作实现原子追加 using (var conn = new SqlConnection("your_connection_string")) { await conn.OpenAsync(); var sql = @" MERGE INTO KeyValueStore AS Target USING (SELECT @Key AS CacheKey, @Values AS NewValues) AS Source ON Target.CacheKey = Source.CacheKey WHEN MATCHED THEN UPDATE SET Values = JSON_MODIFY(Target.Values, 'append $', @NewValues) WHEN NOT MATCHED THEN INSERT (CacheKey, Values) VALUES (Source.CacheKey, JSON_QUERY(@NewValues));"; await conn.ExecuteAsync(sql, new { CacheKey = key, NewValues = JsonSerializer.Serialize(values) }); } } public void Dispose() { _flushTimer?.Dispose(); GC.SuppressFinalize(this); } }
多实例部署额外注意事项
- 消息队列消费隔离:配置消息队列的消费者组(如Kafka)或独占队列(如RabbitMQ),确保同一条消息仅被一个实例消费,避免重复处理。
- 数据库原子写入:必须使用数据库的原子操作(如
MERGE、INSERT ... ON DUPLICATE KEY UPDATE),保证多个实例同时写入同一个key时,值不会丢失或覆盖。 - 可选:改用分布式缓存:若需要多实例共享缓存状态,可替换
MemoryCache为Redis,利用Redis的原子操作(如HSET、SMEMBERS)管理值列表,再定期从Redis同步到DB。
内容的提问来源于stack exchange,提问作者fricadelle
相关产品推荐
相关产品推荐

