基于Concurrent Dictionary锁的钱包存款并发处理异常问题
钱包系统并发存款处理重复入账问题排查与修复
问题概述
开发钱包系统时,处理支付服务商的存款成功Webhook通知,通过Hangfire入队后台任务,逻辑是匹配交易凭证与金额后为用户钱包入账。原本用C# ConcurrentDictionary 结合lock,针对交易凭证和钱包ID加锁以支持并行处理,但测试同一交易凭证的5次并发调用时,所有线程都执行成功,导致钱包余额变为预期的5倍,预期仅第一次执行成功。
核心问题分析
- 锁对象被过早移除:核心处理代码的
finally块中调用_transactionReferenceLocks.TryRemove(transactionReference, out _),第一个线程释放锁后立即移除锁对象,后续线程的GetOrAdd会创建新的锁对象,多个线程拿到不同锁,无法实现互斥,导致所有请求同时进入业务逻辑。 - 内存锁的局限性:如果系统是多实例集群部署,内存中的
ConcurrentDictionary锁完全无效,每个实例有独立的锁集合,无法跨实例互斥。 - 业务逻辑缺乏幂等性保障:即使锁生效,查询交易状态到更新状态的过程中仍存在竞态条件,仅依赖内存锁无法从根本上防止重复处理。
修复方案
1. 修正内存锁的生命周期
移除finally块中的锁对象移除逻辑,确保同一交易凭证的所有并发请求共享同一个锁对象。为避免内存泄漏,添加定期清理机制,移除不再使用的锁对象。
修改核心处理代码:
lock (_transactionReferenceLocks.GetOrAdd(transactionReference, () => new object())) { try { // 原有业务逻辑:查询Webhook事件、验证交易、更新状态、入账等 unitOfWork.BeginTransaction(); // ... 省略其他代码 unitOfWork.Commit(); } finally { // 移除这行代码:_transactionReferenceLocks.TryRemove(transactionReference, out _); } }
添加锁对象定期清理逻辑(在服务类中):
private readonly Timer _lockCleanupTimer; public TransactionService(ILogger<TransactionService> logger, IUnitOfWork unitOfWork) { // 其他初始化逻辑 _lockCleanupTimer = new Timer(CleanupUnusedLocks, null, TimeSpan.FromMinutes(5), TimeSpan.FromMinutes(5)); } private void CleanupUnusedLocks(object state) { foreach (var key in _transactionReferenceLocks.Keys.ToList()) { // 尝试移除锁对象,成功则说明无线程持有该锁 if (_transactionReferenceLocks.TryRemove(key, out _)) { logger.LogDebug("清理未使用的交易锁:{TransactionReference}", key); } } } // 实现IDisposable接口释放定时器 public void Dispose() { _lockCleanupTimer?.Dispose(); }
2. 数据库层面添加幂等性保障
这是最可靠的防重复处理方案,通过数据库原子更新确保只有第一个请求能将交易从Pending改为Successful。
替换原有直接修改交易状态的代码,改用数据库批量更新:
// 原子更新交易状态,仅当当前状态为Pending时生效 var updatedCreditCount = unitOfWork.WalletTransactions.UpdateTransactionStatus( creditTransaction.Id, WalletTransactionStatus.Pending, WalletTransactionStatus.Successful, transactionCompletionTimestampInUtc ); var updatedDebitCount = unitOfWork.WalletTransactions.UpdateTransactionStatus( debitTransaction.Id, WalletTransactionStatus.Pending, WalletTransactionStatus.Successful, transactionCompletionTimestampInUtc ); // 如果更新行数为0,说明交易已被处理,直接回滚并返回 if (updatedCreditCount == 0 || updatedDebitCount == 0) { logger.LogWarning("交易已处理,跳过重复请求:{TransactionReference}", transactionReference); unitOfWork.Rollback(); return; }
对应的仓储层方法(以EF Core为例):
public int UpdateTransactionStatus(Guid transactionId, WalletTransactionStatus oldStatus, WalletTransactionStatus newStatus, DateTime timestamp) { return _dbContext.WalletTransactions .Where(t => t.Id == transactionId && t.TransactionStatus == oldStatus) .ExecuteUpdate(t => t .SetProperty(s => s.TransactionStatus, newStatus) .SetProperty(s => s.TransactionTimestamp, timestamp) ); }
3. 利用Hangfire内置去重机制
Hangfire支持任务去重,避免同一交易凭证的任务被重复入队或执行:
- 使用
DisableConcurrentExecution特性,禁止同一任务并发执行:
[DisableConcurrentExecution(60)] // 60秒内禁止并发执行同一任务 public async Task CompleteWalletDepositAsync(string transactionReference) { // 业务逻辑 }
- 入队时指定唯一JobId,确保同一交易凭证只创建一个任务:
BackgroundJob.Enqueue<ITransactionService>( x => x.CompleteWalletDepositAsync(transactionReference), new EnqueuedState { JobId = transactionReference } );
额外建议
- 多实例部署需用分布式锁:如果系统是集群部署,内存锁无法跨实例生效,需使用Redis分布式锁或数据库分布式锁,确保同一交易凭证的请求在所有实例中互斥。
- 钱包入账的锁优化:目前
LockAndCreditWallet和LockAndDebitWallet中在修改余额后立即移除钱包锁对象,这会导致同一钱包的并发请求可能拿到不同锁,建议保留钱包锁对象并定期清理,或者直接在数据库层面用乐观锁更新余额:
// 数据库原子更新余额 var updatedCount = unitOfWork.Wallets.UpdateBalance( walletToCredit.Id, walletToCredit.Balance, walletToCredit.Balance + transactionAmount ); if (updatedCount == 0) { logger.LogWarning("钱包余额已被修改,处理失败:{WalletId}", walletToCredit.Id); unitOfWork.Rollback(); return; }
内容的提问来源于stack exchange,提问作者hiddenhenry
相关产品推荐
相关产品推荐

