You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 02:49:55