using System;
using System.Collections.Generic;
using System.Linq;
using System.Linq.Expressions;
using System.Runtime.CompilerServices;
using System.Threading;
using System.Threading.Tasks;
using Dpz.Core.Entity.Base;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using MongoDB.Driver;
namespace Dpz.Core.MongodbAccess;
public sealed class Repository<T> : IRepository<T>
where T : IBaseEntity
{
private readonly MongodbAccess<T> _access;
public Repository(string? connectionString)
{
_access = new MongodbAccess<T>(connectionString);
}
[ActivatorUtilitiesConstructor]
public Repository(IConfiguration configuration)
{
var connectionString = configuration.GetConnectionString("mongodb");
_access = new MongodbAccess<T>(connectionString);
}
public Repository(IConfiguration configuration, string collectionName)
{
var connectionString = configuration.GetConnectionString("mongodb");
_access = string.IsNullOrEmpty(collectionName)
? new MongodbAccess<T>(connectionString)
: new MongodbAccess<T>(connectionString, collectionName);
}
public Repository(string? connectionString, string collectionName)
{
_access = string.IsNullOrEmpty(collectionName)
? new MongodbAccess<T>(connectionString)
: new MongodbAccess<T>(connectionString, collectionName);
}
public IQueryable<T> MongodbQueryable => _access.MongoQueryable;
public IMongoCollection<T> Collection => _access.Collection;
public IQueryable<T> SearchFor(Expression<Func<T, bool>> predicate)
{
return _access.MongoQueryable.Where(predicate);
}
public IFindFluent<T, T> SearchFor(FilterDefinition<T> filter)
{
return Collection.Find(filter);
}
public IFindFluent<T, T> SearchFor(FilterDefinition<T> filter, FindOptions options)
{
return Collection.Find(filter, options);
}
public async IAsyncEnumerable<T> SearchForAsync(
FilterDefinition<T> filter,
[EnumeratorCancellation] CancellationToken cancellationToken = default
)
{
var result = await Collection.FindAsync(filter, cancellationToken: cancellationToken);
while (await result.MoveNextAsync(cancellationToken))
{
foreach (var item in result.Current)
{
yield return item;
}
}
}
public async IAsyncEnumerable<T> SearchForAsync(
FilterDefinition<T> filter,
FindOptions<T> options,
[EnumeratorCancellation] CancellationToken cancellationToken = default
)
{
var result = await Collection.FindAsync(filter, options, cancellationToken);
while (await result.MoveNextAsync(cancellationToken))
{
foreach (var item in result.Current)
{
yield return item;
}
}
}
public async Task<T?> FindAsync(object id, CancellationToken cancellationToken = default)
{
var filter = MongodbExtensions.GetIdPropertyFilter<T>(id);
return await (
await Collection.FindAsync(filter, cancellationToken: cancellationToken)
).SingleOrDefaultAsync(cancellationToken: cancellationToken);
}
public async Task InsertAsync(T entity, CancellationToken cancellationToken = default)
{
await Collection.InsertOneAsync(entity, cancellationToken: cancellationToken);
}
public async Task InsertAsync(
IReadOnlyCollection<T> source,
CancellationToken cancellationToken = default
)
{
if (source.Count > 0)
{
await Collection.InsertManyAsync(source, cancellationToken: cancellationToken);
}
}
public async Task<DeleteResult> DeleteAsync(
Expression<Func<T, bool>> filter,
CancellationToken cancellationToken = default
)
{
var result = await Collection.DeleteManyAsync(filter, cancellationToken);
return result;
}
public async Task<DeleteResult> DeleteAsync(
FilterDefinition<T> filter,
CancellationToken cancellationToken = default
)
{
var result = await Collection.DeleteManyAsync(filter, cancellationToken);
return result;
}
public async Task<DeleteResult> DeleteAsync(
object id,
CancellationToken cancellationToken = default
)
{
var filter = MongodbExtensions.GetIdPropertyFilter<T>(id);
return await Collection.DeleteOneAsync(filter, cancellationToken);
}
public async Task<UpdateResult> UpdateAsync(
Expression<Func<T, bool>> predicate,
UpdateDefinition<T> update,
CancellationToken cancellationToken = default
)
{
var result = await Collection.UpdateManyAsync(
predicate,
update,
cancellationToken: cancellationToken
);
return result;
}
public async Task<ReplaceOneResult> UpdateAsync(
T entity,
CancellationToken cancellationToken = default
)
{
var filter = MongodbExtensions.GetIdPropertyFilter(entity);
var result = await Collection.ReplaceOneAsync(
filter,
entity,
cancellationToken: cancellationToken
);
return result;
}
public async Task<BulkWriteResult<T>> UpdateAsync(
IEnumerable<T> entities,
CancellationToken cancellationToken = default
)
{
var writes = entities
.Select(x => new ReplaceOneModel<T>(MongodbExtensions.GetIdPropertyFilter(x), x))
.ToList();
return await Collection.BulkWriteAsync(writes, cancellationToken: cancellationToken);
}
public IMongoDatabase Database => _access.Database;
}
⚠⚠ 以下内容为AI分析的结果,请根据实际情况进行判断。
下面按模块逐项解释这段 Repository
总体说明
- Repository
是一个针对 MongoDB 的泛型仓储实现,封装了对 MongodbAccess (源码未显示,推测是对 MongoDB Driver 的封装)的访问,提供常用的增删改查操作接口。 - 依赖项包括 MongoDB.Driver、Microsoft.Extensions.Configuration/DependencyInjection 等。
- 通过多种构造函数支持直接传入连接字符串、从 IConfiguration 中读取连接字符串,或同时指定 collection 名称。
- 带有 [ActivatorUtilitiesConstructor] 的构造函数方便在 ASP.NET Core DI 容器中被选为注入构造函数。
构造函数与连接/集合选择
- Repository(string? connectionString):使用传入的 connectionString 创建 MongodbAccess
。 - Repository(IConfiguration configuration):从 configuration.GetConnectionString("mongodb") 获取连接字符串并创建 MongodbAccess
(带有 ActivatorUtilitiesConstructor,供 DI 使用)。 - Repository(IConfiguration configuration, string collectionName):从 IConfiguration 取连接字符串;如果 collectionName 非空则将其传入 MongodbAccess
,否则使用默认集合名字。 - Repository(string? connectionString, string collectionName):同上但直接用传入的 connectionString 和 collectionName。
- MongodbAccess
负责实际的 MongoDB 连接/集合初始化(本代码未给出)。
公开属性
- MongodbQueryable:返回 _access.MongoQueryable,允许通过 LINQ 对集合进行查询(基于 MongoDB LINQ 支持)。
- Collection:返回 IMongoCollection
,直接暴露底层集合以便运行原生 MongoDB.Driver 操作。 - Database:返回 IMongoDatabase(来自 _access),便于执行数据库级操作。
查询方法
- SearchFor(Expression<Func<T, bool>> predicate):基于 LINQ 表达式返回 IQueryable
(延迟查询,交给驱动执行)。 - SearchFor(FilterDefinition
filter):返回 Find 的 IFindFluent<T,T>,可继续链式构造查询(例如 Sort/Limit)。 - SearchFor(FilterDefinition
filter, FindOptions options):带 FindOptions 的 Find,返回 IFindFluent。 - SearchForAsync(FilterDefinition
filter, CancellationToken):异步遍历 FindAsync 结果并通过 IAsyncEnumerable yield 返回所有文档(每批 MoveNextAsync 后 yield 当前批次元素)。 - SearchForAsync(FilterDefinition
filter, FindOptions options, CancellationToken):带 FindOptions 的异步枚举。
单条查找
- FindAsync(object id, CancellationToken):使用 MongodbExtensions.GetIdPropertyFilter
(id) 构建基于实体 id 的过滤器,然后异步查找并返回 SingleOrDefault(找到 1 条返回实体,0 条返回 null;若多条则抛异常)。
插入
- InsertAsync(T entity):InsertOneAsync 插入单个文档。
- InsertAsync(IReadOnlyCollection
source):如果 source.Count>0 则 InsertManyAsync 批量插入(注意:如果 source 为 null,会抛 NRE;只有在 Count>0 才调用 InsertMany)。
删除
- DeleteAsync(Expression<Func<T, bool>> filter):DeleteManyAsync,接受表达式筛选器。
- DeleteAsync(FilterDefinition
filter):DeleteManyAsync,接受 FilterDefinition。 - DeleteAsync(object id):用 id 构建过滤器并 DeleteOneAsync 删除单个文档。
更新 / 替换
- UpdateAsync(Expression<Func<T, bool>> predicate, UpdateDefinition
update):UpdateManyAsync,基于 UpdateDefinition 做部分更新($set 等)。 - UpdateAsync(T entity):通过 MongodbExtensions.GetIdPropertyFilter(entity) 找到文档并用 ReplaceOneAsync 完整替换该文档(整文档替换)。
- UpdateAsync(IEnumerable
entities):对传入实体集合创建 ReplaceOneModel 列表并 BulkWriteAsync 批量替换(按每个实体的 id 作为筛选条件)。
返回类型(重要)
- DeleteResult / UpdateResult / ReplaceOneResult / BulkWriteResult
:来自 MongoDB.Driver,包含操作影响的结果(已删除/已更新/匹配计数等)。 - IFindFluent<T,T>:可继续链式操作然后 ToList/First 等。
- IAsyncEnumerable
:支持异步流式消费。
依赖的辅助项
- MongodbExtensions.GetIdPropertyFilter
(...):用于基于实体或 id 构建 id 过滤器(具体实现未给出,但关键用于按照主键查找/删除/替换)。 - IConfiguration.GetConnectionString("mongodb"):约定从配置中读取名为 "mongodb" 的连接字符串。
注意事项与改进建议(可选)
- 插入的 IList/集合重载没有对 null 做保护(如果可能传入 null,可先检查以免 NRE)。
- FindAsync 返回 SingleOrDefaultAsync,若存在多条匹配会抛出异常,调用方应保证 id 唯一或使用 FirstOrDefault 视需求。
- 当前代码未处理异常或重试策略,生产环境可能需要增加错误处理与日志、索引检查、超时配置等。
- 对 UpdateAsync(IEnumerable
) 的 BulkWrite 没有使用 upsert 标记(若需要可配置 ReplaceOneModel 的 Upsert = true)。 - 并发/线程安全:MongoDB .NET Driver 的 IMongoCollection/客户端是线程安全的,通常可以复用 Repository 实例。
总结
- 这是一个围绕 MongoDB 的通用仓储类,封装了常用的 CRUD、批量操作以及基于 LINQ 的查询入口,支持通过配置或直接传参初始化连接与集合,方便在应用中以统一接口操作 MongoDB 文档实体。
AI 正在分析代码…
评论加载中...