基于.NET Core MongoDB的通用仓储如何实现工作单元模式?
MongoDB 工作单元模式实现方案(适配现有通用仓储)
前置说明
MongoDB 4.0及以上版本支持多文档事务,这是实现工作单元的核心基础。工作单元的目标是将多个集合操作纳入同一事务,保证原子性——要么全部成功提交,要么全部回滚。
1. 定义工作单元接口
创建IUnitOfWork接口,规范事务操作与仓储获取逻辑:
/// <summary> /// MongoDB工作单元接口 /// </summary> public interface IUnitOfWork : IDisposable { /// <summary> /// 获取指定文档类型的仓储实例 /// </summary> /// <typeparam name="TDocument">文档类型,需实现IAudiableEntity</typeparam> /// <returns>对应的仓储实例</returns> IMongoRepository<TDocument> GetRepository<TDocument>() where TDocument : IAudiableEntity; /// <summary> /// 开启事务 /// </summary> Task StartTransactionAsync(); /// <summary> /// 提交事务 /// </summary> Task CommitAsync(); /// <summary> /// 回滚事务 /// </summary> Task RollbackAsync(); }
2. 实现工作单元类
基于MongoDB会话机制实现工作单元,同时维护仓储实例缓存避免重复创建:
/// <summary> /// MongoDB工作单元实现类 /// </summary> public class MongoUnitOfWork : IUnitOfWork { private readonly IMongoDatabase _database; private IClientSessionHandle _session; private readonly Dictionary<Type, object> _repositories = new Dictionary<Type, object>(); private bool _disposed = false; public MongoUnitOfWork(IMongoDbSettings settings) { var client = new MongoClient(settings.ConnectionString); _database = client.GetDatabase(settings.DatabaseName); } /// <summary> /// 获取仓储实例,复用已创建的实例 /// </summary> public IMongoRepository<TDocument> GetRepository<TDocument>() where TDocument : IAudiableEntity { var documentType = typeof(TDocument); if (!_repositories.ContainsKey(documentType)) { // 传入当前会话给仓储,确保操作在事务内执行 var repository = new MongoRepository<TDocument>(_session, settings); _repositories.Add(documentType, repository); } return (IMongoRepository<TDocument>)_repositories[documentType]; } /// <summary> /// 开启事务会话 /// </summary> public async Task StartTransactionAsync() { _session = await _database.Client.StartSessionAsync(); _session.StartTransaction(); } /// <summary> /// 提交事务 /// </summary> public async Task CommitAsync() { if (_session == null) throw new InvalidOperationException("事务未开启,请先调用StartTransactionAsync"); await _session.CommitTransactionAsync(); } /// <summary> /// 回滚事务 /// </summary> public async Task RollbackAsync() { if (_session == null) throw new InvalidOperationException("事务未开启,请先调用StartTransactionAsync"); await _session.AbortTransactionAsync(); } /// <summary> /// 释放资源 /// </summary> public void Dispose() { Dispose(true); GC.SuppressFinalize(this); } protected virtual void Dispose(bool disposing) { if (_disposed) return; if (disposing) { _session?.Dispose(); _repositories.Clear(); } _disposed = true; } ~MongoUnitOfWork() { Dispose(false); } }
3. 修改现有通用仓储,支持事务会话
调整MongoRepository构造函数,增加事务会话参数,让所有操作绑定到工作单元的事务上下文:
调整后的仓储接口(保持原有方法,仅补充中文注释)
/// <summary> /// MongoDB通用仓储接口 /// </summary> /// <typeparam name="TDocument">文档类型,需实现IAudiableEntity</typeparam> public interface IMongoRepository<TDocument> where TDocument : IAudiableEntity { /// <summary> /// 插入单条文档 /// </summary> /// <param name="document">待插入文档</param> void InsertOne(TDocument document); /// <summary> /// 异步插入单条文档 /// </summary> /// <param name="document">待插入文档</param> /// <returns>插入后的文档</returns> Task<TDocument> InsertOneAsync(TDocument document); /// <summary> /// 插入多条文档 /// </summary> /// <param name="documents">待插入文档集合</param> void InsertMany(ICollection<TDocument> documents); /// <summary> /// 异步插入多条文档 /// </summary> /// <param name="documents">待插入文档集合</param> Task InsertManyAsync(ICollection<TDocument> documents); /// <summary> /// 替换单条文档(带Slides脏标记处理) /// </summary> /// <param name="document">待替换文档</param> void ReplaceOne(TDocument document); /// <summary> /// 异步替换单条文档(带Slides脏标记处理) /// </summary> /// <param name="document">待替换文档</param> Task ReplaceOneAsync(TDocument document); /// <summary> /// 根据条件删除单条文档 /// </summary> /// <param name="filterExpression">过滤条件</param> void DeleteOne(Expression<Func<TDocument, bool>> filterExpression); /// <summary> /// 异步根据条件删除单条文档 /// </summary> /// <param name="filterExpression">过滤条件</param> Task DeleteOneAsync(Expression<Func<TDocument, bool>> filterExpression); /// <summary> /// 根据ID删除文档 /// </summary> /// <param name="id">文档ID</param> void DeleteById(string id); /// <summary> /// 异步根据ID删除文档 /// </summary> /// <param name="id">文档ID</param> Task DeleteByIdAsync(string id); /// <summary> /// 根据条件删除多条文档 /// </summary> /// <param name="filterExpression">过滤条件</param> void DeleteMany(Expression<Func<TDocument, bool>> filterExpression); /// <summary> /// 异步根据条件删除多条文档 /// </summary> /// <param name="filterExpression">过滤条件</param> Task DeleteManyAsync(Expression<Func<TDocument, bool>> filterExpression); }
调整后的仓储实现类
/// <summary> /// MongoDB通用仓储实现类 /// </summary> /// <typeparam name="TDocument">文档类型,需实现IAudiableEntity</typeparam> public class MongoRepository<TDocument> : IMongoRepository<TDocument> where TDocument : IAudiableEntity { private readonly IMongoCollection<TDocument> _collection; private readonly IClientSessionHandle _session; // 事务会话 private readonly IMongoDbSettings _settings; // 非事务场景构造函数 public MongoRepository(IMongoDbSettings settings) { _settings = settings; var database = new MongoClient(settings.ConnectionString).GetDatabase(settings.DatabaseName); _collection = database.GetCollection<TDocument>(GetCollectionName(typeof(TDocument))); } // 事务场景构造函数 public MongoRepository(IClientSessionHandle session, IMongoDbSettings settings) { _session = session; _settings = settings; var database = session.Client.GetDatabase(settings.DatabaseName); _collection = database.GetCollection<TDocument>(GetCollectionName(typeof(TDocument))); } /// <summary> /// 获取文档对应的集合名称(从BsonCollectionAttribute读取) /// </summary> /// <param name="documentType">文档类型</param> /// <returns>集合名称</returns> private protected string GetCollectionName(Type documentType) { return ((BsonCollectionAttribute)documentType.GetCustomAttributes( typeof(BsonCollectionAttribute), true) .FirstOrDefault())?.CollectionName; } public virtual void InsertOne(TDocument document) { if (_session != null) _collection.InsertOne(_session, document); else _collection.InsertOne(document); } public virtual async Task<TDocument> InsertOneAsync(TDocument document) { if (_session != null) await _collection.InsertOneAsync(_session, document); else await _collection.InsertOneAsync(document); return document; } public void InsertMany(ICollection<TDocument> documents) { if (_session != null) _collection.InsertMany(_session, documents); else _collection.InsertMany(documents); } public virtual async Task InsertManyAsync(ICollection<TDocument> documents) { if (_session != null) await _collection.InsertManyAsync(_session, documents); else await _collection.InsertManyAsync(documents); } public void ReplaceOne(TDocument document) { var filter = Builders<TDocument>.Filter.Eq(doc => doc.Id, document.Id); if (document is Slides slide) { slide.IsDirty = true; } if (_session != null) _collection.FindOneAndReplace(_session, filter, document); else _collection.FindOneAndReplace(filter, document); } public void ReplaceWODirtyOne(TDocument document) { var filter = Builders<TDocument>.Filter.Eq(doc => doc.Id, document.Id); if (_session != null) _collection.FindOneAndReplace(_session, filter, document); else _collection.FindOneAndReplace(filter, document); } public virtual async Task ReplaceOneAsync(TDocument document) { var filter = Builders<TDocument>.Filter.Eq(doc => doc.Id, document.Id); if (document is Slides slide) { slide.IsDirty = true; } if (_session != null) await _collection.FindOneAndReplaceAsync(_session, filter, document); else await _collection.FindOneAndReplaceAsync(filter, document); } public virtual async Task ReplaceWODirtyOneAsync(TDocument document) { var filter = Builders<TDocument>.Filter.Eq(doc => doc.Id, document.Id); if (_session != null) await _collection.FindOneAndReplaceAsync(_session, filter, document); else await _collection.FindOneAndReplaceAsync(filter, document); } public void DeleteOne(Expression<Func<TDocument, bool>> filterExpression) { if (_session != null) _collection.FindOneAndDelete(_session, filterExpression); else _collection.FindOneAndDelete(filterExpression); } public Task DeleteOneAsync(Expression<Func<TDocument, bool>> filterExpression) { if (_session != null) return _collection.FindOneAndDeleteAsync(_session, filterExpression); else return _collection.FindOneAndDeleteAsync(filterExpression); } public void DeleteById(string id) { var objectId = new ObjectId(id); var filter = Builders<TDocument>.Filter.Eq(doc => doc.Id, objectId); if (_session != null) _collection.FindOneAndDelete(_session, filter); else _collection.FindOneAndDelete(filter); } public Task DeleteByIdAsync(string id) { var objectId = new ObjectId(id); var filter = Builders<TDocument>.Filter.Eq(doc => doc.Id, objectId); if (_session != null) return _collection.FindOneAndDeleteAsync(_session, filter); else return _collection.FindOneAndDeleteAsync(filter); } public void DeleteMany(Expression<Func<TDocument, bool>> filterExpression) { if (_session != null) _collection.DeleteMany(_session, filterExpression); else _collection.DeleteMany(filterExpression); } public Task DeleteManyAsync(Expression<Func<TDocument, bool>> filterExpression) { if (_session != null) return _collection.DeleteManyAsync(_session, filterExpression); else return _collection.DeleteManyAsync(filterExpression); } }
4. 依赖注入配置(ASP.NET Core)
在Program.cs中注册工作单元与仓储:
builder.Services.AddSingleton<IMongoDbSettings>(sp => { return new MongoDbSettings { ConnectionString = builder.Configuration.GetConnectionString("MongoDB"), DatabaseName = builder.Configuration["MongoDB:DatabaseName"] }; }); // 工作单元用Scoped生命周期(与请求周期一致) builder.Services.AddScoped<IUnitOfWork, MongoUnitOfWork>(); // 非事务场景可直接注入仓储 builder.Services.AddScoped(typeof(IMongoRepository<>), typeof(MongoRepository<>));
5. 使用示例(服务层)
public class OrderService { private readonly IUnitOfWork _unitOfWork; public OrderService(IUnitOfWork unitOfWork) { _unitOfWork = unitOfWork; } public async Task CreateOrderWithUser() { try { // 开启事务 await _unitOfWork.StartTransactionAsync(); // 获取两个仓储实例 var userRepo = _unitOfWork.GetRepository<User>(); var orderRepo = _unitOfWork.GetRepository<Order>(); // 执行跨集合操作 var newUser = new User { Name = "新用户" }; await userRepo.InsertOneAsync(newUser); await orderRepo.InsertOneAsync(new Order { UserId = newUser.Id.ToString() }); // 提交事务 await _unitOfWork.CommitAsync(); } catch (Exception) { // 异常回滚 await _unitOfWork.RollbackAsync(); throw; } } }
注意事项
- 确保MongoDB版本≥4.0且部署为副本集(单节点不支持事务)。
- 工作单元生命周期建议用Scoped,与ASP.NET Core请求周期对齐。
- 非事务场景仍可直接注入
IMongoRepository<>使用,无需通过工作单元。
内容的提问来源于stack exchange,提问作者ramakrishnan r
相关产品推荐
相关产品推荐

