第 5 章 · 02 Synchronizer 时间同步与 TimeSlice


文档摘要

第 5 章 · 02 Synchronizer 时间同步与 TimeSlice 本节摘要:本节钻取数据流子系统的"消费者"——Synchronizer。它的职责是把 SubscriptionCollection 里多个订阅(每个 symbol 独立线程在生产数据)按数据时间合并成单一时间序列,产出 流。整个接口只有一个方法: ——返回 TimeSlice 枚举,引擎 消费。回测用 (按数据时间推进,源码里还有"时间不前进"的检测),实盘用 (override 后用 墙钟等待下一秒)。这是回测/实盘的第二个分叉点( 的 )。本节还会拆解 这个引擎内部数据包的九个字段。 内容来源:原项目源码 、 (共 243 行)、 、 (共 106 行)、 。 ⚠️ 注意:本节是数据流的"心脏"。

第 5 章 · 02 Synchronizer 时间同步与 TimeSlice

本节摘要:本节钻取数据流子系统的"消费者"——Synchronizer。它的职责是把 SubscriptionCollection 里多个订阅(每个 symbol 独立线程在生产数据)按数据时间合并成单一时间序列,产出 TimeSlice 流。整个接口只有一个方法:IEnumerable<TimeSlice> StreamData(CancellationToken)——返回 TimeSlice 枚举,引擎 foreach 消费。回测用 Synchronizer(按数据时间推进,源码里还有"时间不前进"的检测),实盘用 LiveSynchronizer(override 后用 ITimeProvider 墙钟等待下一秒)。这是回测/实盘的第二个分叉点(Engine.cs:113_liveMode ? new LiveSynchronizer() : new Synchronizer())。本节还会拆解 TimeSlice 这个引擎内部数据包的九个字段。

内容来源:原项目源码 Engine/DataFeeds/ISynchronizer.csEngine/DataFeeds/Synchronizer.cs(共 243 行)、Engine/DataFeeds/LiveSynchronizer.csEngine/DataFeeds/TimeSlice.cs(共 106 行)、Engine/Engine.cs

⚠️ 注意:本节是数据流的"心脏"。Synchronizer 既是消费者(从 SubscriptionCollection 取),又是生产者(给 AlgorithmManager 喂 TimeSlice)。理解它就理解了"为什么用户 OnData 是单线程的"——因为 Synchronizer 把多线程生产的数据串行化了。

学习目标

阅读完本节,你应当能够:

  1. 说清 ISynchronizer 接口为什么只有一个方法 StreamData
  2. 读懂 Synchronizer.StreamData 回测实现的 yield return timeSlice 与时间不前进检测。
  3. 读懂 LiveSynchronizer.StreamData 实盘实现怎么用墙钟等待。
  4. 找到回测/实盘第二个分叉点 Engine.cs:113
  5. 列出 TimeSlice 的九个字段及含义。

一、ISynchronizer 接口:只有一个方法

Engine/DataFeeds/ISynchronizer.cs:25-32:

25 public interface ISynchronizer 26 { 27 /// <summary> 28 /// Returns an enumerable which provides the data to stream to the algorithm 29 /// </summary> 30 IEnumerable<TimeSlice> StreamData(CancellationToken cancellationToken); 31 }

整个接口就一个方法,返回 IEnumerable<TimeSlice>。这是 C# 的迭代器模式——实现里用 yield return 逐个产出 TimeSlice,调用方用 foreach 消费。CancellationToken 用于优雅停机。

💡 钻取要点:为什么用 IEnumerable 而非"事件"或"回调"?因为 Lean 的 AlgorithmManager 是单线程拉取模型——它在一个 while 循环里 foreach (var timeSlice in synchronizer.StreamData(token)),每次拿到一个 TimeSlice 就同步执行算法的 OnData,执行完再拉下一个。IEnumerable + yield 是这种"消费者驱动节奏"的最自然表达,且不引入额外线程。

二、Synchronizer.StreamData:回测实现

Engine/DataFeeds/Synchronizer.cs:86-161(精简):

86 public virtual IEnumerable<TimeSlice> StreamData(CancellationToken cancellationToken) 87 { 88 PostInitialize(); 89 SubscriptionSynchronizer.SetTimeProvider(GetTimeProvider()); 94 var previousEmitTime = DateTime.MaxValue; 96 var enumerator = SubscriptionSynchronizer 97 .Sync(SubscriptionManager.DataFeedSubscriptions, cancellationToken) 98 .GetEnumerator(); 99 var previousWasTimePulse = false; 101 var retried = false; 102 while (!cancellationToken.IsCancellationRequested) 103 { 104 TimeSlice timeSlice; 105 try 106 { 107 if (!enumerator.MoveNext()) 108 { 109 break; // 枚举器结束,所有订阅数据耗尽 111 } 112 timeSlice = enumerator.Current; 114 catch (Exception err) 115 { 116 Algorithm.SetRuntimeError(err, "Synchronizer"); 117 break; 120 } 122 if (timeSlice == null || cancellationToken.IsCancellationRequested) break; 124 if (timeSlice.IsTimePulse && Algorithm.UtcTime == timeSlice.Time) 125 { 126 previousWasTimePulse = timeSlice.IsTimePulse; 127 continue; // 跳过算法已在当前时间的时间脉冲 130 } 137 if (timeSlice.Time != previousEmitTime || previousWasTimePulse || timeSlice.UniverseData.Count != 0) 138 { 139 previousEmitTime = timeSlice.Time; 140 previousWasTimePulse = timeSlice.IsTimePulse; 141 retried = false; 143 yield return timeSlice; // 产 TimeSlice 给 AlgorithmManager 145 } 146 else 147 { 148 // 时间没前进,且 slice 有数据则重试一次,否则结束 150 if (!timeSlice.Slice.HasData || retried) break; 151 retried = true; 155 } 157 } 159 enumerator.DisposeSafely(); 160 Log.Trace("Synchronizer.GetEnumerator(): Exited thread."); 161 }

逐段解读:

  1. L96-98 委托给 SubscriptionSynchronizer:SubscriptionSynchronizer.Sync(...) 才是真正做时间合并的——它遍历所有订阅,找出"下一个最早的数据时间",把那个时刻所有订阅的数据打包成一个 TimeSlice。本节不钻 SubscriptionSynchronizer 内部(它在 Engine/DataFeeds/SubscriptionSynchronizer.cs,核心是一个按时间堆排序的合并循环)。
  2. L102 while 循环:单线程,从 enumerator 拉 TimeSlice。
  3. L107 MoveNext:拉下一个 TimeSlice,失败(所有订阅数据耗尽)就 break 结束回测。
  4. L124 时间脉冲跳过:IsTimePulse 表示这是个"空脉冲"(用来推动时间但不带数据),如果算法时间已经在那个点,跳过避免重复。
  5. L137 时间前进检测:核心——只有当 timeSlice.Time != previousEmitTime(时间真的前进了)才 yield return。否则可能重试一次(L150),还不行就判定"时间不前进了,数据耗尽"结束回测。

⚠️ 重要说明:时间不前进检测(L137-156)是回测鲁棒性的关键。如果某个订阅的枚举器卡住一直返回同一时刻的数据(理论上不应该发生),Synchronizer 不会无限循环——重试一次后判定数据耗尽,优雅结束。这避免了"回测挂死"的灾难。

三、LiveSynchronizer.StreamData:实盘实现

Engine/DataFeeds/LiveSynchronizer.cs:85-130(精简):

85 public override IEnumerable<TimeSlice> StreamData(CancellationToken cancellationToken) 86 { 87 PostInitialize(); 89 var shouldSendExtraEmptyPacket = false; 90 var nextEmit = DateTime.MinValue; 91 var lastLoopStart = DateTime.UtcNow; 93 var enumerator = SubscriptionSynchronizer 94 .Sync(SubscriptionManager.DataFeedSubscriptions, cancellationToken) 95 .GetEnumerator(); 97 var previousWasTimePulse = false; 98 while (!cancellationToken.IsCancellationRequested) 99 { 100 var now = DateTime.UtcNow; 101 if (!previousWasTimePulse) 102 { 103 if (!_newLiveDataEmitted.IsSet && !Algorithm.IsWarmingUp) 104 { 105 // 没有新数据,且非预热期:看是否同一秒 109 if (lastLoopStart.Second == now.Second) 110 { 111 _realTimeScheduleEventService.ScheduleEvent( TimeSpan.FromMilliseconds(GetPulseDueTime(now)), now); 112 _newLiveDataEmitted.Wait(); // 等新数据或下一秒脉冲 114 } 116 } 117 _newLiveDataEmitted.Reset(); 119 } 120 lastLoopStart = now; 122 TimeSlice timeSlice; 123 try { if (!enumerator.MoveNext()) break; timeSlice = enumerator.Current; } catch ... 129 // ... yield return ... 130 }

实盘与回测的核心差异在 L98-119 的墙钟等待:

  • 回测:磁盘数据无穷尽,Synchronizer 全速拉,MoveNext 立刻返回下一条。
  • 实盘:数据是券商实时推的,可能这一秒没数据。LiveSynchronizer_newLiveDataEmitted(一个 ManualResetEventSlim)等待:要么订阅有新数据(在 OnSubscriptionNewDataAvailable 里 Set),要么墙钟跨秒(通过 _realTimeScheduleEventService 调度)。这保证了即使没数据,算法每秒至少被唤醒一次(刷新指标/检查定时任务)。

实盘的 _frontierTimeProviderLiveTimeProvider(L60),它包装真实墙钟。回测用的是 SubscriptionFrontierTimeProvider(数据驱动)。

💡 钻取要点:这是"墙钟驱动 vs 数据驱动"的本质区别。回测时间是数据时间(数据说了算),实盘时间是真实时间(墙钟说了算)。LiveSynchronizer 通过 override 把这个差异封装在 StreamData 里,对外仍是 IEnumerable<TimeSlice>,AlgorithmManager 完全无感。

四、回测/实盘第二个分叉点

Engine/Engine.cs:113:

113 var synchronizer = _liveMode ? new LiveSynchronizer() : new Synchronizer();

_liveMode 由 Engine 构造时传入(L74)。这是引擎装配阶段的选择——第一个分叉点是 DataFeed(FileSystemDataFeed vs LiveTradingDataFeed,01 节),第二个分叉点就是 Synchronizer。两个分叉点必须配套:回测 FileSystemDataFeed + Synchronizer,实盘 LiveTradingDataFeed + LiveSynchronizer。

后续 L113 的 synchronizer 被传给 AlgorithmManager(L173 附近),AlgorithmManager 在主循环里 foreach (var timeSlice in synchronizer.StreamData(token)) 消费。

五、TimeSlice:引擎内部数据包

Engine/DataFeeds/TimeSlice.cs:28-104 定义 TimeSlice 类。它是 Synchronizer 产出的、AlgorithmManager 一个时间步用到的所有载荷:

28 public class TimeSlice 29 { 33 public int DataPointCount { get; } // 数据点计数 38 public DateTime Time { get; } // UTC 时间 43 public List<DataFeedPacket> Data { get; } // 原始数据包列表 48 public Slice Slice { get; } // 用户可见部分(OnData 收到的) 53 public List<UpdateData<ISecurityPrice>> SecuritiesUpdateData { get; } // 更新 Security 价格 58 public List<UpdateData<SubscriptionDataConfig>> ConsolidatorUpdateData { get; } // 更新 Consolidator 63 public List<UpdateData<ISecurityPrice>> CustomData { get; } // 自定义数据 68 public SecurityChanges SecurityChanges { get; } // 证券增删(宇宙选择结果) 73 public Dictionary<Universe, BaseDataCollection> UniverseData { get; } // 宇宙数据 78 public bool IsTimePulse { get; } // 是否时间脉冲 104 }

九个字段分工:

字段 用途
Time 这个时间步的 UTC 时间
DataPointCount 这一步总共多少数据点(性能统计)
Data 原始 DataFeedPacket 列表(引擎内部)
Slice 用户可见部分,传给 OnData(03 节详讲)
SecuritiesUpdateData 用来更新 Security.Price/Cache 的数据
ConsolidatorUpdateData 用来喂数据给 Consolidator(合成更大周期)
CustomData 自定义数据流
SecurityChanges 这一时间步宇宙选择增加/删除了哪些证券
UniverseData 宇宙选择用的原始数据
IsTimePulse 时间脉冲标志(无数据只推时间)

💡 钻取要点:TimeSlice 是"引擎视角",Slice 是"用户视角"。TimeSlice 含 SecuritiesUpdateData/ConsolidatorUpdateData 这些引擎内部要处理的数据,而 Slice 只是用户 OnData 看到的 Bars/Ticks 等。AlgorithmManager 拿到 TimeSlice 后:1) 用 SecuritiesUpdateData 更新 Security 价格;2) 用 ConsolidatorUpdateData 喂 Consolidator;3) 调用户 OnData(timeSlice.Slice)。用户永远看不到 TimeSlice 全貌,只看到 Slice。

本节要点回顾

  1. ISynchronizer 只一个方法 IEnumerable<TimeSlice> StreamData(CancellationToken),迭代器模式,AlgorithmManager 单线程拉取。
  2. Synchronizer.StreamData 回测:委托 SubscriptionSynchronizer.Sync 做时间合并,yield return timeSlice,带"时间不前进"检测(重试一次后判数据耗尽)。
  3. LiveSynchronizer.StreamData 实盘:override 后用 _newLiveDataEmitted(ManualResetEventSlim)等墙钟或新数据,保证每秒至少唤醒一次。
  4. 第二个分叉点 Engine.cs:113:_liveMode ? new LiveSynchronizer() : new Synchronizer(),与 DataFeed 分叉配套。
  5. TimeSlice 九字段:Time/Slice(用户可见)/SecuritiesUpdateData/ConsolidatorUpdateData/SecurityChanges/UniverseData/CustomData/Data/IsTimePulse,是引擎一个时间步的全部载荷。

下一节,我们钻取用户真正能看到的 Slice 对象——ExtendedDictionary 字典,Bars/QuoteBars/Ticks/OptionChains/FuturesChains + Splits/Dividends,以及缺失数据怎么用 Fill-Forward 前向填充。


作者与出处
原作者: 灏天文库
整理: 灏天文库整理
本站整理收录,版权归原作者/开源协议所有;欢迎通过原文链接访问源仓库。
发布者: 作者: 灏天文库 转发
评论区 (0)
U