Files
zszq-trs/YLErpDAL/Modules/TradeRiskCalcModule/TaskRunner/TradeExtendDataSource.cs
T
2024-05-09 14:06:26 +08:00

85 lines
2.1 KiB
C#

using System.Collections.Concurrent;
namespace YLErp.Modules.TradeRiskCalcModule.TaskRunner
{
/// <summary>
/// [线程安全]交易扩展数据源
/// </summary>
class TradeExtendDataSource<TData> : ITradeExtendUpdater where TData : TradeExtendBase
{
readonly Guid _eventGuid;
readonly ConcurrentDictionary<int, TData> _dic;
private TradeExtendDataSource()
{
_dic = new ConcurrentDictionary<int, TData>();
_eventGuid = Events.EventBus.Subscribe<OtcTradeUpdateEvent>(t =>
{
UpdateData(t.TradeIds);
});
}
~TradeExtendDataSource()
{
Events.EventBus.Unsubscribe(_eventGuid);
}
/// <summary>
/// 清空缓存
/// </summary>
public void Clear()
{
_dic.Clear();
}
/// <summary>
/// 获取数据
/// </summary>
public TData GetData(int tradId)
{
if (!_dic.TryGetValue(tradId, out var data))
{
var table = DbContextFactory.GetYLDbContext().Set<TData>();
data = table.FirstOrDefault(n => n.TradeId == tradId);
if (data != null)
{
_dic.AddOrUpdate(tradId, data, (k, t) => data);
}
}
return data;
}
/// <summary>
/// 跟随trade更新
/// </summary>
public void UpdateData(IEnumerable<int> tradeIds)
{
if (_dic.Count > 1000)
{
_dic.Clear();
return;
}
if (tradeIds == null || !tradeIds.Any())
{
return;
}
foreach (var tradId in tradeIds)
{
_dic.TryRemove(tradId, out _);
}
}
/// <summary>
/// 使用这个主要是偷个懒
/// </summary>
public static readonly TradeExtendDataSource<TData> Default;
static TradeExtendDataSource()
{
Default = new TradeExtendDataSource<TData>();
}
}
}