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

82 lines
2.2 KiB
C#

using System.Collections.Concurrent;
using YLErp.Abstract;
namespace YLErp.Modules.TradeRiskCalcModule.TaskRunner
{
/// <summary>
/// [线程安全]收益互换多空组合数据源
/// </summary>
class TradeSwapDataSource : ITradeExtendUpdater, IDataUpdater
{
readonly Guid _eventGuid;
readonly ConcurrentDictionary<int, trade_swap_detail[]> _dic;
private TradeSwapDataSource()
{
_dic = new ConcurrentDictionary<int, trade_swap_detail[]>();
_eventGuid = Events.EventBus.Subscribe<OtcTradeUpdateEvent>(t =>
{
UpdateData(t.TradeIds);
});
}
~TradeSwapDataSource()
{
Events.EventBus.Unsubscribe(_eventGuid);
}
public List<trade_swap_detail> GetDatas(int tradId)
{
if (!_dic.TryGetValue(tradId, out var datas))
{
using (var db = DbContextFactory.GetYLDbContext())
{
_dic[tradId] = datas = db.trade_swap_detail.AsNoTracking()
.Where(n => n.TradeId == tradId && n.ValidState != "InValid").ToArray();
}
}
return datas.Select(t => t.Clone()).ToList();
}
/// <summary>
/// 跟随trade更新
/// </summary>
public void UpdateData(IEnumerable<int> tradeIds)
{
if (_dic.Count > 5000)
{
_dic.Clear();
return;
}
if (tradeIds == null || !tradeIds.Any())
{
return;
}
foreach (var tradId in tradeIds)
{
_dic.TryRemove(tradId, out _);
}
}
public void UpdateData(IEnumerable<string> updateKeyIds)
{
var tradeIds = DataConvert.ConvertToInt32Array(updateKeyIds);
this.UpdateData(tradeIds);
}
/// <summary>
/// 偷个懒使用这个
/// </summary>
public static readonly TradeSwapDataSource Default;
public string TableName => "trade_swap_detail";
static TradeSwapDataSource()
{
Default = new TradeSwapDataSource();
}
}
}