第 5 章 · 02 Synchronizer 时间同步与 TimeSlice 本节摘要:本节钻取数据流子系统的"消费者"——Synchronizer。它的职责是把 SubscriptionCollection 里多个订阅(每个 symbol 独立线程在生产数据)按数据时间合并成单一时间序列,产出 流。整个接口只有一个方法: ——返回 TimeSlice 枚举,引擎 消费。回测用 (按数据时间推进,源码里还有"时间不前进"的检测),实盘用 (override 后用 墙钟等待下一秒)。这是回测/实盘的第二个分叉点( 的 )。本节还会拆解 这个引擎内部数据包的九个字段。 内容来源:原项目源码 、 (共 243 行)、 、 (共 106 行)、 。 ⚠️ 注意:本节是数据流的"心脏"。
本节摘要:本节钻取数据流子系统的"消费者"——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.cs、Engine/DataFeeds/Synchronizer.cs(共 243 行)、Engine/DataFeeds/LiveSynchronizer.cs、Engine/DataFeeds/TimeSlice.cs(共 106 行)、Engine/Engine.cs。
⚠️ 注意:本节是数据流的"心脏"。Synchronizer 既是消费者(从 SubscriptionCollection 取),又是生产者(给 AlgorithmManager 喂 TimeSlice)。理解它就理解了"为什么用户 OnData 是单线程的"——因为 Synchronizer 把多线程生产的数据串行化了。
阅读完本节,你应当能够:
ISynchronizer 接口为什么只有一个方法 StreamData。Synchronizer.StreamData 回测实现的 yield return timeSlice 与时间不前进检测。LiveSynchronizer.StreamData 实盘实现怎么用墙钟等待。Engine.cs:113。TimeSlice 的九个字段及含义。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是这种"消费者驱动节奏"的最自然表达,且不引入额外线程。
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 }
逐段解读:
SubscriptionSynchronizer.Sync(...) 才是真正做时间合并的——它遍历所有订阅,找出"下一个最早的数据时间",把那个时刻所有订阅的数据打包成一个 TimeSlice。本节不钻 SubscriptionSynchronizer 内部(它在 Engine/DataFeeds/SubscriptionSynchronizer.cs,核心是一个按时间堆排序的合并循环)。IsTimePulse 表示这是个"空脉冲"(用来推动时间但不带数据),如果算法时间已经在那个点,跳过避免重复。timeSlice.Time != previousEmitTime(时间真的前进了)才 yield return。否则可能重试一次(L150),还不行就判定"时间不前进了,数据耗尽"结束回测。⚠️ 重要说明:时间不前进检测(L137-156)是回测鲁棒性的关键。如果某个订阅的枚举器卡住一直返回同一时刻的数据(理论上不应该发生),Synchronizer 不会无限循环——重试一次后判定数据耗尽,优雅结束。这避免了"回测挂死"的灾难。
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 的墙钟等待:
MoveNext 立刻返回下一条。LiveSynchronizer 用 _newLiveDataEmitted(一个 ManualResetEventSlim)等待:要么订阅有新数据(在 OnSubscriptionNewDataAvailable 里 Set),要么墙钟跨秒(通过 _realTimeScheduleEventService 调度)。这保证了即使没数据,算法每秒至少被唤醒一次(刷新指标/检查定时任务)。实盘的 _frontierTimeProvider 是 LiveTimeProvider(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)) 消费。
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。
IEnumerable<TimeSlice> StreamData(CancellationToken),迭代器模式,AlgorithmManager 单线程拉取。SubscriptionSynchronizer.Sync 做时间合并,yield return timeSlice,带"时间不前进"检测(重试一次后判数据耗尽)。_newLiveDataEmitted(ManualResetEventSlim)等墙钟或新数据,保证每秒至少唤醒一次。Engine.cs:113:_liveMode ? new LiveSynchronizer() : new Synchronizer(),与 DataFeed 分叉配套。Time/Slice(用户可见)/SecuritiesUpdateData/ConsolidatorUpdateData/SecurityChanges/UniverseData/CustomData/Data/IsTimePulse,是引擎一个时间步的全部载荷。下一节,我们钻取用户真正能看到的 Slice 对象——ExtendedDictionary 字典,Bars/QuoteBars/Ticks/OptionChains/FuturesChains + Splits/Dividends,以及缺失数据怎么用 Fill-Forward 前向填充。