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

778 lines
23 KiB
C#

using BaseOUDAL;
using System.Collections;
using YieldChain.Commons;
using YLErp.Abstract;
using YLErp.Model;
namespace YLErp.Modules.DataCacheModule
{
/// <summary>
/// 数据缓存管理静态类(提供基础数据缓存及跟踪数据更新事件)
/// </summary>
public static partial class DataCacheManager
{
const int MaxTimeOut = 2 * 60 * 1000;
static DateTime _lastTime;
static readonly IYcLogger _logger;
static readonly Timer _updateTimer;
static readonly ThrottleAction _throttleValueDate;
static readonly TimeSpan[] MarketPriceTimes;
static readonly YLDBDataCacheUpdater _yLDBDataUpdater;
static readonly ClientDBDataCacheUpdater _clientDBDataUpdater;
static DataCacheManager()
{
_updateTimer = new Timer(TimerTask, null, 1000, Interval);
_throttleValueDate = new ThrottleAction(() => BLL.valuedateBLL.ResetValueDate(), 30);
_yLDBDataUpdater = new YLDBDataCacheUpdater();
_clientDBDataUpdater = new ClientDBDataCacheUpdater();
_logger = LogFactory.GetLogger(nameof(DataCacheManager));
DebugOut = AppContext.TryGetSwitch("AppSettings:DataCacheDebug", out _);
MarketPriceTimes = new TimeSpan[] {
new TimeSpan(3,0,0),new TimeSpan(8,30,0),
new TimeSpan(15,30,0),new TimeSpan(20,30,0)
};
}
/// <summary>
/// 当前最大traceID
/// </summary>
public static long MaxId => _yLDBDataUpdater.MaxId;
/// <summary>
/// 轮询间隔
/// </summary>
public static int Interval { get; private set; } = 2000;
/// <summary>
/// 调试日志是否输出
/// </summary>
public static bool DebugOut { get; set; }
/// <summary>
/// 这个方法在WEB项目IRegisteredObject接口中调用非常关键,
/// 可以有效避免定时更新被dispose掉
/// </summary>
public static void StartUpdate(int interval = 2000)
{
if (interval > 2000)
{
Interval = interval;
}
_updateTimer.Change(0, Interval);
}
/// <summary>
///
/// </summary>
public static void StopUpdate()
{
_updateTimer?.Change(Timeout.Infinite, 0);
}
/// <summary>
/// 更新一次
/// </summary>
public static void UpdateOnce()
{
TimerTask(null);
}
/// <summary>
/// 注册数据跟踪监听器(暂时只支持ylcms数据库)
/// </summary>
public static void RegisterDataTraceSink(IDataTraceSink sink)
{
_yLDBDataUpdater.RegisterDataTraceSink(sink);
}
/// <summary>
/// 取消数据跟踪监听器注册(暂时只支持ylcms数据库)
/// </summary>
public static void UnRegisterDataTraceSink(IDataTraceSink sink)
{
_yLDBDataUpdater.UnRegisterDataTraceSink(sink);
}
public static void ResetDataCache()
{
_yLDBDataUpdater.ResetDataCache();
_clientDBDataUpdater.ResetDataCache();
}
/// <summary>
/// 定时调用缓存更新操作
/// </summary>
private static void TimerTask(object state)
{
try
{
_updateTimer.Change(MaxTimeOut, MaxTimeOut);
//控制基础数据更新在10s的频率上
if (_lastTime.AddSeconds(10) < DateTime.Now)
{
_lastTime = DateTime.Now;
_yLDBDataUpdater.Update();
_clientDBDataUpdater.Update();
}
//更新标的价格
try
{
var delay = false;
var time = DateTime.Now.TimeOfDay;
for (var i = 0; i < MarketPriceTimes.Length; i += 2)
{
delay = time > MarketPriceTimes[i] && time < MarketPriceTimes[i + 1];
if (delay)
{
break;
}
}
if (GetUnderlyingDataSource().UpdatePrices(delay) && DebugOut)
{
_logger.Debug("更新标的价格");
}
}
catch (Exception ex)
{
_logger.Error(ex, "更新标的价格出错");
}
_throttleValueDate.Execute();
}
catch (Exception ex)
{
_logger.Error(ex, "TimerTask");
}
finally
{
_updateTimer.Change(Interval, Interval);
}
}
#region-----缓存的基础数据源------
/// <summary>
/// 交易品种
/// </summary>
public static IDataSourceKey2<Variety> GetVarietyDataSource()
{
return VarietyDataSource.Default;
}
/// <summary>
/// 期权标的
/// </summary>
public static IUnderlyingDataSource GetUnderlyingDataSource()
{
return new UnderlyingDbDataSource();
}
/// <summary>
/// 场内期权基础数据
/// </summary>
public static IDataSourceKey2<ExchangeListOption> GetExchangeListOptionDataSource()
{
return ExchangeListOptionDataSource.Default;
}
/// <summary>
/// 标的相关性
/// </summary>
public static IDataSource<CorrelationTable> GetCorrelationDataSource()
{
return CorrelationDataSource.Default;
}
/// <summary>
/// 个股黑白名单
/// </summary>
public static IDataSource<StockBlackWhite> GetStockBlackWhiteDataSource()
{
return StockBlackWhiteDataSource.Default;
}
/// <summary>
/// 股票手续费设定
/// </summary>
public static IDataSource<StockCommissionConfig> GetStockCommissionDataSource()
{
return StockCommissionDataSource.Default;
}
/// <summary>
/// 对冲账户
/// </summary>
public static IDataSource<ExchangeAccount> GetExchangeAccountDataSource()
{
return ExchangeAccountDataSource.Default;
}
/// <summary>
/// 簿记账户
/// </summary>
public static IDataSource<AssetUnit> GetAssetUnitDataSource()
{
return AssetUnitDataSource.Default;
}
/// <summary>
/// 簿记账户组
/// </summary>
public static IDataSource<AssetUnitGroup> GetAssetUnitGroupDataSource()
{
return AssetUnitDataGroupSource.Default;
}
/// <summary>
/// 客户黑名单
/// </summary>
public static IDataSource<client_black> GetClientBlackDataSource()
{
return ClientBlackDataSource.Default;
}
/// <summary>
/// 客户信息
/// </summary>
public static IDataSource<Client> GetClientDataSource()
{
return ClientDataSource.Default;
}
/// <summary>
/// 客户等级信息
/// </summary>
public static IDataSource<ClientLevel> GetClientLevelDataSource()
{
return ClientLevelDataSource.Default;
}
/// <summary>
/// 市场信息
/// </summary>
public static IDataSource<Market> GetMarketDataSource()
{
return MarketDataSource.Default;
}
#endregion
#region----用于缓存数据调试----
public static List<string> GetSinkTables()
{
var tableNames = new List<string>(30);
_yLDBDataUpdater.GetSinkTables(tableNames);
_clientDBDataUpdater.GetSinkTables(tableNames);
return tableNames;
}
public static string GetSinkData(string tableName)
{
if (string.IsNullOrEmpty(tableName))
{
return "tableName参数不能为空";
}
var str = _yLDBDataUpdater.GetSinkData(tableName);
if (str.StartsWith("未找到数据"))
{
str = _clientDBDataUpdater.GetSinkData(tableName);
}
return str;
}
#endregion
#region----DataUpdater----
abstract class DataCacheUpdaterBase
{
//初始值设置为-1支持datatrace表数据空的情况
public long MaxId { get; private set; } = -1;
protected List<IDataTraceSink> SinkList { get; set; }
/// <summary>
/// 更新跟踪数据
/// </summary>
protected void InnerUpdate(DataTraceDbContext dbContext, string dbName)
{
var maxid = MaxId;
try
{
if (MaxId < 0)
{
MaxId = dbContext.DataTrace.Max(n => (long?)n.id) ?? 0;
}
else
{
var query = from a in dbContext.DataTrace
where a.id > MaxId
orderby a.id
select new DataTraceInfo
{
id = a.id,
TableName = a.TableName,
DataKeyId = a.DataKeyId
};
var traceDatas = query.ToArray();
if (traceDatas.Length > 0)
{
MaxId = traceDatas[traceDatas.Length - 1].id;
var set = new HashSet<string>(traceDatas.Length, StringComparer.OrdinalIgnoreCase);
foreach (var data in traceDatas)
{
var key = string.Concat(data.TableName, "[|]", data.DataKeyId);
if (set.Add(key))
{
try
{
lock (SinkList)
{
foreach (var sink in SinkList)
{
sink.OnDataTrace(data);
}
}
}
catch (Exception ex)
{
_logger.Error(ex, "UpdateTraceData");
}
}
}
}
}
if (MaxId == maxid)
{
return;
}
if (DebugOut)
{
_logger.Debug($"[{dbName}.Datatrace MaxId]{MaxId}");
}
}
catch (Exception ex)
{
_logger.Error(ex, dbName + ".Datatrace更新出错");
}
try
{
lock (SinkList)
{
foreach (var sink in SinkList)
{
sink.UpdateCache();
}
}
}
catch (Exception ex)
{
MaxId = -1;
_logger.Error(ex, dbName + ".Datatrace缓存更新出错");
}
}
/// <summary>
/// 注册数据跟踪监听器
/// </summary>
public void RegisterDataTraceSink(IDataTraceSink sink)
{
if (sink == null)
{
throw new ArgumentNullException(nameof(sink));
}
lock (SinkList)
{
if (!SinkList.Contains(sink))
{
SinkList.Add(sink);
}
}
}
/// <summary>
/// 取消数据跟踪监听器注册
/// </summary>
public void UnRegisterDataTraceSink(IDataTraceSink sink)
{
if (sink != null)
{
lock (SinkList)
{
SinkList.Remove(sink);
}
}
}
/// <summary>
/// 重置数据缓存
/// </summary>
public void ResetDataCache()
{
MaxId = -1;
foreach (var sink in SinkList)
{
sink.Reset();
}
}
public void GetSinkTables(List<string> resultList)
{
foreach (var sink in SinkList)
{
if (sink is DataTraceSink ds)
{
resultList.Add(ds.TableName);
}
else if (sink is DataTraceSinkGroup dsg)
{
resultList.AddRange(dsg.AsEnumerable().Select(n => n.TableName));
}
else
{
resultList.Add("unknown table sink");
}
}
}
public string GetSinkData(string tableName)
{
if (string.IsNullOrEmpty(tableName))
{
return "tableName参数不能为空";
}
foreach (var sink in SinkList)
{
if (sink is DataTraceSink ds)
{
if (ds.TableName == tableName)
{
return (ds.Updater as IJsonSerializable)?.ToJson() ?? "未实现JSON序列化接口:" + tableName;
}
}
else if (sink is DataTraceSinkGroup dsg)
{
foreach (var sink2 in dsg)
{
if (sink2.TableName == tableName)
{
return (sink2.Updater as IJsonSerializable)?.ToJson() ?? "未实现JSON序列化接口:" + tableName;
}
}
}
}
return "未找到数据更新容器:" + tableName;
}
}
class YLDBDataCacheUpdater : DataCacheUpdaterBase
{
public YLDBDataCacheUpdater()
{
var basicDataUpdaters = new List<IDataUpdater>() {
(IDataUpdater)GetMarketDataSource(),
(IDataUpdater)GetVarietyDataSource(),
(IDataUpdater)GetCorrelationDataSource(),
(IDataUpdater)GetUnderlyingDataSource(),
(IDataUpdater)GetExchangeListOptionDataSource(),
(IDataUpdater)GetStockBlackWhiteDataSource(),
(IDataUpdater)GetClientBlackDataSource(),
(IDataUpdater)GetStockCommissionDataSource(),
(IDataUpdater)GetExchangeAccountDataSource(),
(IDataUpdater)GetAssetUnitDataSource(),
(IDataUpdater)GetAssetUnitGroupDataSource()
};
if (AppManager.SubSystem != Enums.SubSystemName.OtcWeb)
{
basicDataUpdaters.Insert(0, AppConfigUpdater.Default);
}
//基础数据6秒刷新
var baseDataTrace = new DataTraceSinkGroup(6, basicDataUpdaters.ToArray());
SinkList = new List<IDataTraceSink>() { baseDataTrace };
//设置基础数据更新时间
var handler = new EventHandler((object sender, EventArgs e) =>
{
Events.EventBus.Publish<Events.DataCacheUpdateEvent>();
});
foreach (var item in basicDataUpdaters)
{
if (item is IDataSourceEvent dataSourceEvent)
{
dataSourceEvent.DataSourceUpdated += handler;
}
}
}
public void Update()
{
using (var dbtrace = DbContextFactory.GetYLDbContext().GetDataTraceDbContext())
{
InnerUpdate(dbtrace, "ylcmsDB");
}
}
}
class ClientDBDataCacheUpdater : DataCacheUpdaterBase
{
public ClientDBDataCacheUpdater()
{
var basicDataUpdaters = new List<IDataUpdater>() {
(IDataUpdater)GetClientDataSource(),
(IDataUpdater)GetClientLevelDataSource(),
(IDataUpdater)GetClientBlackDataSource()
};
//基础数据6秒刷新
var baseDataTrace = new DataTraceSinkGroup(6, basicDataUpdaters.ToArray());
SinkList = new List<IDataTraceSink>() { baseDataTrace };
//设置基础数据更新时间
var handler = new EventHandler((object sender, EventArgs e) =>
{
Events.EventBus.Publish<Events.DataCacheUpdateEvent>();
});
foreach (var item in basicDataUpdaters)
{
if (item is IDataSourceEvent dataSourceEvent)
{
dataSourceEvent.DataSourceUpdated += handler;
}
}
}
public void Update()
{
using (var dbtrace = DbContextFactory.GetClientDbContext(null).GetDataTraceDbContext())
{
InnerUpdate(dbtrace, "clientDB");
}
}
}
#endregion
}
/// <summary>
/// 数据跟踪监听接口
/// </summary>
public interface IDataTraceSink
{
void Reset();
void UpdateCache();
void OnDataTrace(DataTraceInfo dataTrace);
}
/// <summary>
/// 数据跟踪监听类库
/// </summary>
public class DataTraceSink : IDataTraceSink
{
readonly List<string> _list;
readonly ThrottleAction _throttle;
/// <summary>
/// 构造函数--不使用throttle更新方式
/// </summary>
public DataTraceSink(IDataUpdater dataUpdater)
{
_list = new List<string> { "0" };
Updater = dataUpdater ?? throw new ArgumentNullException(nameof(dataUpdater));
}
/// <summary>
/// 构造函数--使用throttle更新方式
/// </summary>
public DataTraceSink(int throttleUpdateWaitSeconds, IDataUpdater dataUpdater) : this(dataUpdater)
{
_throttle = new ThrottleAction(DoUpdateCache, throttleUpdateWaitSeconds < 1 ? 2 : throttleUpdateWaitSeconds, 60);
}
public void Reset()
{
if (Updater is IDataSource dataSource)
{
dataSource.ResetDataSource();
}
_list.Add("0");
}
public void OnDataTrace(DataTraceInfo dataTrace)
{
if (dataTrace != null && Updater.TableName.Equals(dataTrace.TableName, StringComparison.OrdinalIgnoreCase))
{
lock (_list)
{
_list.Add(dataTrace.DataKeyId);
}
}
}
/// <summary>
/// 更新数据
/// </summary>
public void UpdateCache()
{
if (_throttle != null)
{
_throttle.Execute();
}
else
{
DoUpdateCache();
}
}
private void DoUpdateCache()
{
string[] idArr = null;
try
{
lock (_list)
{
if (_list.Count < 1)
{
return;
}
idArr = _list.ToArray();
_list.Clear();
}
AppManager.SetSysInfo("[DataTrace]" + Updater.TableName, "执行,top10-Id:" + string.Join(",", idArr.Take(10)));
Updater.UpdateData(idArr);
}
catch (Exception ex)
{
if (idArr != null && idArr.Length > 0)
{
lock (_list)
{
_list.AddRange(idArr);
}
}
LogFactory.GetLogger("DataTrace_" + Updater.TableName).Error(ex, "更新缓存失败");
}
}
public IDataUpdater Updater { get; }
/// <summary>
///
/// </summary>
public string TableName => Updater.TableName;
}
/// <summary>
///
/// </summary>
public class DataTraceSinkGroup : IDataTraceSink, IEnumerable<DataTraceSink>
{
readonly ThrottleAction _throttle;
readonly IEnumerable<DataTraceSink> _sinks;
/// <summary>
///
/// </summary>
/// <param name="updaters">更新器集合</param>
/// <param name="throttleUpdateWaitSeconds">每次更新等待的时间(秒)</param>
public DataTraceSinkGroup(int throttleUpdateWaitSeconds, params IDataUpdater[] updaters)
{
if (updaters == null)
{
throw new ArgumentNullException(nameof(updaters));
}
_throttle = new ThrottleAction(ThrottleUpdate, throttleUpdateWaitSeconds < 1 ? 2 : throttleUpdateWaitSeconds, 60);
_sinks = updaters.Where(n => !string.IsNullOrEmpty(n?.TableName)).Select(n => new DataTraceSink(n)).ToArray();
}
public void Reset()
{
foreach (var sink in _sinks)
{
sink.Reset();
}
}
/// <summary>
/// Throttle机制的更新行为调用方法
/// </summary>
private void ThrottleUpdate()
{
foreach (var sink in _sinks)
{
sink.UpdateCache();
}
}
/// <summary>
/// 更新缓存数据(采用了Throttle机制)
/// </summary>
public void UpdateCache()
{
_throttle.Execute();
}
/// <summary>
///
/// </summary>
public void OnDataTrace(DataTraceInfo dataTrace)
{
foreach (var sink in _sinks)
{
sink.OnDataTrace(dataTrace);
}
}
public IEnumerator<DataTraceSink> GetEnumerator()
{
return _sinks.GetEnumerator();
}
IEnumerator IEnumerable.GetEnumerator()
{
return _sinks.GetEnumerator();
}
}
}