using BaseOUDAL; using System.Collections; using YieldChain.Commons; using YLErp.Abstract; using YLErp.Model; namespace YLErp.Modules.DataCacheModule { /// /// 数据缓存管理静态类(提供基础数据缓存及跟踪数据更新事件) /// 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) }; } /// /// 当前最大traceID /// public static long MaxId => _yLDBDataUpdater.MaxId; /// /// 轮询间隔 /// public static int Interval { get; private set; } = 2000; /// /// 调试日志是否输出 /// public static bool DebugOut { get; set; } /// /// 这个方法在WEB项目IRegisteredObject接口中调用非常关键, /// 可以有效避免定时更新被dispose掉 /// public static void StartUpdate(int interval = 2000) { if (interval > 2000) { Interval = interval; } _updateTimer.Change(0, Interval); } /// /// /// public static void StopUpdate() { _updateTimer?.Change(Timeout.Infinite, 0); } /// /// 更新一次 /// public static void UpdateOnce() { TimerTask(null); } /// /// 注册数据跟踪监听器(暂时只支持ylcms数据库) /// public static void RegisterDataTraceSink(IDataTraceSink sink) { _yLDBDataUpdater.RegisterDataTraceSink(sink); } /// /// 取消数据跟踪监听器注册(暂时只支持ylcms数据库) /// public static void UnRegisterDataTraceSink(IDataTraceSink sink) { _yLDBDataUpdater.UnRegisterDataTraceSink(sink); } public static void ResetDataCache() { _yLDBDataUpdater.ResetDataCache(); _clientDBDataUpdater.ResetDataCache(); } /// /// 定时调用缓存更新操作 /// 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-----缓存的基础数据源------ /// /// 交易品种 /// public static IDataSourceKey2 GetVarietyDataSource() { return VarietyDataSource.Default; } /// /// 期权标的 /// public static IUnderlyingDataSource GetUnderlyingDataSource() { return new UnderlyingDbDataSource(); } /// /// 场内期权基础数据 /// public static IDataSourceKey2 GetExchangeListOptionDataSource() { return ExchangeListOptionDataSource.Default; } /// /// 标的相关性 /// public static IDataSource GetCorrelationDataSource() { return CorrelationDataSource.Default; } /// /// 个股黑白名单 /// public static IDataSource GetStockBlackWhiteDataSource() { return StockBlackWhiteDataSource.Default; } /// /// 股票手续费设定 /// public static IDataSource GetStockCommissionDataSource() { return StockCommissionDataSource.Default; } /// /// 对冲账户 /// public static IDataSource GetExchangeAccountDataSource() { return ExchangeAccountDataSource.Default; } /// /// 簿记账户 /// public static IDataSource GetAssetUnitDataSource() { return AssetUnitDataSource.Default; } /// /// 簿记账户组 /// public static IDataSource GetAssetUnitGroupDataSource() { return AssetUnitDataGroupSource.Default; } /// /// 客户黑名单 /// public static IDataSource GetClientBlackDataSource() { return ClientBlackDataSource.Default; } /// /// 客户信息 /// public static IDataSource GetClientDataSource() { return ClientDataSource.Default; } /// /// 客户等级信息 /// public static IDataSource GetClientLevelDataSource() { return ClientLevelDataSource.Default; } /// /// 市场信息 /// public static IDataSource GetMarketDataSource() { return MarketDataSource.Default; } #endregion #region----用于缓存数据调试---- public static List GetSinkTables() { var tableNames = new List(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 SinkList { get; set; } /// /// 更新跟踪数据 /// 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(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缓存更新出错"); } } /// /// 注册数据跟踪监听器 /// public void RegisterDataTraceSink(IDataTraceSink sink) { if (sink == null) { throw new ArgumentNullException(nameof(sink)); } lock (SinkList) { if (!SinkList.Contains(sink)) { SinkList.Add(sink); } } } /// /// 取消数据跟踪监听器注册 /// public void UnRegisterDataTraceSink(IDataTraceSink sink) { if (sink != null) { lock (SinkList) { SinkList.Remove(sink); } } } /// /// 重置数据缓存 /// public void ResetDataCache() { MaxId = -1; foreach (var sink in SinkList) { sink.Reset(); } } public void GetSinkTables(List 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)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() { baseDataTrace }; //设置基础数据更新时间 var handler = new EventHandler((object sender, EventArgs e) => { Events.EventBus.Publish(); }); 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)GetClientDataSource(), (IDataUpdater)GetClientLevelDataSource(), (IDataUpdater)GetClientBlackDataSource() }; //基础数据6秒刷新 var baseDataTrace = new DataTraceSinkGroup(6, basicDataUpdaters.ToArray()); SinkList = new List() { baseDataTrace }; //设置基础数据更新时间 var handler = new EventHandler((object sender, EventArgs e) => { Events.EventBus.Publish(); }); 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 } /// /// 数据跟踪监听接口 /// public interface IDataTraceSink { void Reset(); void UpdateCache(); void OnDataTrace(DataTraceInfo dataTrace); } /// /// 数据跟踪监听类库 /// public class DataTraceSink : IDataTraceSink { readonly List _list; readonly ThrottleAction _throttle; /// /// 构造函数--不使用throttle更新方式 /// public DataTraceSink(IDataUpdater dataUpdater) { _list = new List { "0" }; Updater = dataUpdater ?? throw new ArgumentNullException(nameof(dataUpdater)); } /// /// 构造函数--使用throttle更新方式 /// 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); } } } /// /// 更新数据 /// 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; } /// /// /// public string TableName => Updater.TableName; } /// /// /// public class DataTraceSinkGroup : IDataTraceSink, IEnumerable { readonly ThrottleAction _throttle; readonly IEnumerable _sinks; /// /// /// /// 更新器集合 /// 每次更新等待的时间(秒) 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(); } } /// /// Throttle机制的更新行为调用方法 /// private void ThrottleUpdate() { foreach (var sink in _sinks) { sink.UpdateCache(); } } /// /// 更新缓存数据(采用了Throttle机制) /// public void UpdateCache() { _throttle.Execute(); } /// /// /// public void OnDataTrace(DataTraceInfo dataTrace) { foreach (var sink in _sinks) { sink.OnDataTrace(dataTrace); } } public IEnumerator GetEnumerator() { return _sinks.GetEnumerator(); } IEnumerator IEnumerable.GetEnumerator() { return _sinks.GetEnumerator(); } } }