BsonDocument到结果的动态表达或内存缓存

Nod*_*.JS 6 c# mongodb mongodb-.net-driver

我正在尝试为其编写代理类,IMongoCollection以便可以将内存中的高速缓存用于某些方法的实现。但是,问题在于,几乎所有过滤器都是类型FilterDefinition<T>,这意味着我们可以调用Render它们以获得a BsonDocument。我想知道是否有一种方法可以将过滤器转换BsonDocument为动态Expression,以便可以在内存中运行它List<T>。或者,也许还有一种我不知道的更好的内存缓存方法。谢谢。

更新:

我很想按照@ simon-mourier的建议编写解决方案,但这个棘手的解决方案的问题是C#mongo驱动程序返回IAsyncCursor<T>查找操作,该操作基本上是BsonDocuments 的流,并且在每次读取后,它都指向最后一个索引并自行处理。而且无法将流重置为其初始位置。这意味着下面的代码是第一次工作,但是之后,我们得到一个例外,那就是游标在流的末尾并且已经被处理掉了。

using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using DAL.Extensions;
using MongoDB.Bson;
using MongoDB.Bson.Serialization;
using MongoDB.Driver;

namespace DAL.Proxies
{
    public static class MongoCollectionProxy
    {
        private static readonly Dictionary<Type, object> _instances = new Dictionary<Type, object>();

        public static IMongoCollection<T> New<T>(IMongoCollection<T> proxy)
        {
            return ((IMongoCollection<T>)_instances.AddOrUpdate(typeof(T), () => new MongoCollectionBaseProxyImpl<T>(proxy)));
        }
    }

    public class MongoCollectionBaseProxyImpl<T> : MongoCollectionBase<T>
    {
        private readonly IMongoCollection<T> _proxy;

        private readonly ConcurrentDictionary<string, object> _cache = new ConcurrentDictionary<string, object>();

        public MongoCollectionBaseProxyImpl(IMongoCollection<T> proxy)
        {
            _proxy = proxy;
        }

        public override Task<IAsyncCursor<TResult>> AggregateAsync<TResult>(PipelineDefinition<T, TResult> pipeline,
            AggregateOptions options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            return _proxy.AggregateAsync(pipeline, options, cancellationToken);
        }

        public override Task<BulkWriteResult<T>> BulkWriteAsync(IEnumerable<WriteModel<T>> requests,
            BulkWriteOptions options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            return _proxy.BulkWriteAsync(requests, options, cancellationToken);
        }

        [Obsolete("Use CountDocumentsAsync or EstimatedDocumentCountAsync instead.")]
        public override Task<long> CountAsync(FilterDefinition<T> filter, CountOptions options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            return _proxy.CountAsync(filter, options, cancellationToken);
        }

        public override Task<IAsyncCursor<TField>> DistinctAsync<TField>(FieldDefinition<T, TField> field,
            FilterDefinition<T> filter, DistinctOptions options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            return _proxy.DistinctAsync(field, filter, options, cancellationToken);
        }

        public override async Task<IAsyncCursor<TProjection>> FindAsync<TProjection>(FilterDefinition<T> filter,
            FindOptions<T, TProjection> options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            // ReSharper disable once SpecifyACultureInStringConversionExplicitly
            return await CacheResult(filter.Render().ToString(), () => _proxy.FindAsync(filter, options, cancellationToken));
        }

        public override async Task<TProjection> FindOneAndDeleteAsync<TProjection>(FilterDefinition<T> filter,
            FindOneAndDeleteOptions<T, TProjection> options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            return await InvalidateCache(_proxy.FindOneAndDeleteAsync(filter, options, cancellationToken));
        }

        public override async Task<TProjection> FindOneAndReplaceAsync<TProjection>(FilterDefinition<T> filter,
            T replacement,
            FindOneAndReplaceOptions<T, TProjection> options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            return await InvalidateCache(_proxy.FindOneAndReplaceAsync(filter, replacement, options,
                cancellationToken));
        }

        public override async Task<TProjection> FindOneAndUpdateAsync<TProjection>(FilterDefinition<T> filter,
            UpdateDefinition<T> update,
            FindOneAndUpdateOptions<T, TProjection> options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            return await InvalidateCache(_proxy.FindOneAndUpdateAsync(filter, update, options, cancellationToken));
        }

        public override Task<IAsyncCursor<TResult>> MapReduceAsync<TResult>(BsonJavaScript map, BsonJavaScript reduce,
            MapReduceOptions<T, TResult> options = null,
            CancellationToken cancellationToken = new CancellationToken())
        {
            return _proxy.MapReduceAsync(map, reduce, options, cancellationToken);
        }

        public override IFilteredMongoCollection<TDerivedDocument> OfType<TDerivedDocument>()
        {
            return _proxy.OfType<TDerivedDocument>();
        }

        public override IMongoCollection<T> WithReadPreference(ReadPreference readPreference)
        {
            return _proxy.WithReadPreference(readPreference);
        }

        public override IMongoCollection<T> WithWriteConcern(WriteConcern writeConcern)
        {
            return _proxy.WithWriteConcern(writeConcern);
        }

        public override CollectionNamespace CollectionNamespace => _proxy.CollectionNamespace;

        public override IMongoDatabase Database => _proxy.Database;

        public override IBsonSerializer<T> DocumentSerializer => _proxy.DocumentSerializer;

        public override IMongoIndexManager<T> Indexes => _proxy.Indexes;

        public override MongoCollectionSettings Settings => _proxy.Settings;

        private async Task<TResult> CacheResult<TResult>(string key, Func<Task<TResult>> result)
        {
            return _cache.ContainsKey(key) ? (TResult) _cache[key] : (TResult) _cache.AddOrUpdate(key, await result());
        }

        private TResult InvalidateCache<TResult>(TResult result)
        {
            _cache.Clear();

            return result;
        }
    }
}
Run Code Online (Sandbox Code Playgroud)