第 5 章 · 01 DataFeed 与 DataManager 中介者 本节摘要:本节钻取第 5 层入口——数据流子系统。Lean 的 目录是整个项目最大的子模块(120+ 文件),职责是把多品种、多频率的原始数据按时间对齐,统一喂给算法。核心三件套:DataFeed(生产者,多线程生产数据)、DataManager(中介者,统一管理订阅生命周期)、Subscription(单个订阅对象)。回测用 (注释原话"incredibly fast"——从磁盘读 zip+CSV),实盘用 (从券商实时拉)。DataManager 同时给 DataFeed 和 SubscriptionManager 提供订阅管理,是两者之间的中介者。
本节摘要:本节钻取第 5 层入口——数据流子系统。Lean 的
Engine/DataFeeds目录是整个项目最大的子模块(120+ 文件),职责是把多品种、多频率的原始数据按时间对齐,统一喂给算法。核心三件套:DataFeed(生产者,多线程生产数据)、DataManager(中介者,统一管理订阅生命周期)、Subscription(单个订阅对象)。回测用FileSystemDataFeed(注释原话"incredibly fast"——从磁盘读 zip+CSV),实盘用LiveTradingDataFeed(从券商实时拉)。DataManager 同时给 DataFeed 和 SubscriptionManager 提供订阅管理,是两者之间的中介者。本节先讲这三者的职责分工,02 节讲 Synchronizer 怎么消费。
内容来源:原项目源码
Engine/DataFeeds/FileSystemDataFeed.cs、Engine/DataFeeds/DataManager.cs、Engine/DataFeeds/Subscription.cs、Engine/DataFeeds/SubscriptionCollection.cs。
⚠️ 注意:
Engine/DataFeeds子目录巨多(Enumerators/Queues/Transport/WorkScheduling/Factories 等),本节只钻取顶层骨架,具体枚举器(Fill/Refresh/Frontfill)在 03 节、Subscription 内部状态机在第 7 章。
阅读完本节,你应当能够:
IDataFeed 接口与两个实现(回测 FileSystemDataFeed/实盘 LiveTradingDataFeed)的分工。Subscription/SubscriptionCollection/SubscriptionDataReader 三个核心对象。UniverseManager.CollectionChanged 监听 universe 增删。IDataFeed 定义数据源契约(方法 Initialize/IsActive 等,在 Engine/DataFeeds/IDataFeed.cs)。两个生产实现:
FileSystemDataFeed(Engine/DataFeeds/FileSystemDataFeed.cs:40):
36 /// <summary> 37 /// Historical datafeed stream reader for processing files on a local disk. 38 /// </summary> 39 /// <remarks>Filesystem datafeeds are incredibly fast</remarks> 40 public class FileSystemDataFeed : IDataFeed 41 { 42 private IAlgorithm _algorithm; 43 private ITimeProvider _timeProvider; ... 49 private SubscriptionCollection _subscriptions; ... 51 private SubscriptionDataReaderSubscriptionEnumeratorFactory _subscriptionFactory;
注释里的"incredibly fast"是源码原话。回测模式下,所有历史数据都是本地 zip+CSV(第 6 章详讲格式),FileSystemDataFeed 给每个订阅起一个 SubscriptionDataReader 枚举器,从磁盘流式读取,几乎没有 IO 瓶颈——Lean 回测能跑到每秒几百万数据点的速度就靠这个。
LiveTradingDataFeed(Engine/DataFeeds/LiveTradingDataFeed.cs):实盘实现,数据来自券商 API 实时推送(通过 IDataQueue 适配器,见 03 节 Queues 子目录)。它内部要做"实时数据 → 时间对齐 → 喂 Synchronizer"的转换,比 FileSystemDataFeed 复杂得多。
💡 钻取要点:回测与实盘在 DataFeed 层就已经分叉——这是 Lean 的第一个分叉点(回测 FileSystemDataFeed vs 实盘 LiveTradingDataFeed)。第二个分叉点在 Synchronizer 层(02 节)。这种"同一接口两个实现"的设计让回测策略几乎可以无缝切到实盘。
DataManager 是数据流子系统的"中介者"(Engine/DataFeeds/DataManager.cs:34):
31 /// <summary> 32 /// DataManager will manage the subscriptions for both the DataFeeds and the SubscriptionManager 33 /// </summary> 34 public class DataManager : IAlgorithmSubscriptionManager, IDataFeedSubscriptionManager, IDataManager 35 { 36 private readonly IDataFeed _dataFeed; 37 private readonly MarketHoursDatabase _marketHoursDatabase; 38 private readonly ITimeKeeper _timeKeeper; 39 private readonly bool _liveMode; ... 44 private readonly IAlgorithm _algorithm;
注释直说:"为 DataFeed 和 SubscriptionManager 共同管理订阅"。它同时实现三个接口:IAlgorithmSubscriptionManager(给算法侧用)、IDataFeedSubscriptionManager(给 DataFeed 侧用)、IDataManager(给 Synchronizer 用)。这样 DataFeed 和 SubscriptionManager 不直接耦合,都通过 DataManager 协调订阅的增删。
核心事件(L58-63):
57 /// <summary> 58 /// Event fired when a new subscription is added 59 /// </summary> 60 public event EventHandler<Subscription> SubscriptionAdded; 61 62 /// <summary> 63 /// Event fired when an existing subscription is removed 64 /// </summary> 65 public event EventHandler<Subscription> SubscriptionRemoved;
Synchronizer、LiveTradingDataFeed 等都会订阅这两个事件,DataManager 在增删订阅时通知它们。
挂 UniverseManager.CollectionChanged(L90):
89 // wire ourselves up to receive notifications when universes are added/removed 90 algorithm.UniverseManager.CollectionChanged += (sender, args) => 91 { 92 var universe = args.Value; 93 switch (args.Action) 94 { 95 case NotifyCollectionChangedAction.Replace: 96 case NotifyCollectionChangedAction.Add: 97 var config = universe.Configuration; 98 var start = algorithm.UtcTime; 99 if (algorithm.GetLocked() && args.Action == NotifyCollectionChangedAction.Add && universe is UserDefinedUniverse) 100 { 101 // 如果在 Initialize 后(运行时)动态加,把时间往后推 1 tick 102 start = start.AddTicks(1); 103 } // ... 创建 SubscriptionRequest,触发 DataFeed 建订阅 ... };
这是"运行时动态加证券"的关键链路:用户在 OnData 里调 AddEquity("AAPL") → UniverseManager 加一个 UserDefinedUniverse → 触发 CollectionChanged → DataManager 收到 → 给 DataFeed 下 SubscriptionRequest → DataFeed 起一个新的 Subscription 枚举器 → 数据进来。整个链路是事件驱动的,不需要用户手动管订阅生命周期。
💡 钻取要点:中介者模式(GoF)的标准应用。没有 DataManager 的话,DataFeed 要直接监听 UniverseManager、又要直接调 SubscriptionManager,耦合乱成一团。DataManager 把"订阅生命周期"统一封装,三个角色(算法/DataFeed/SubscriptionManager)都只跟它对话。
Subscription(Engine/DataFeeds/Subscription.cs:32):单个订阅的运行时对象,实现 IEnumerator<SubscriptionData>:
32 public class Subscription : IEnumerator<SubscriptionData> { public SubscriptionDataConfig Configuration { get; init; } ... public virtual bool MoveNext() { ... } // 拉下一条数据 public SubscriptionData Current { get; private set; }
它包装一个底层枚举器(从磁盘读/从实时队列读),对外提供 MoveNext/Current。每个订阅(每个 symbol+resolution+datatype 组合)一个 Subscription 实例。
SubscriptionCollection(Engine/DataFeeds/SubscriptionCollection.cs):所有活跃 Subscription 的集合,线程安全的并发字典。DataFeed 多线程往里塞数据,Synchronizer 单线程从里取。这是生产者消费者链的缓冲区。
SubscriptionDataReader(Engine/DataFeeds/SubscriptionDataReader.cs):回测专用的磁盘读取器,实现 IEnumerator<BaseData>。它打开 zip 文件,逐行解析 CSV,做 fill-forward/enumerator 链组装(03 节详讲),产出 BaseData 流。FileSystemDataFeed 内部就是用它做实际 IO。
三者关系:SubscriptionCollection 持有多个 Subscription,每个 Subscription 内部的枚举器链最底层是 SubscriptionDataReader(回测)或 IDataQueue(实盘)。
整个数据流是经典的生产者消费者模式:
DataFeed (多线程生产) │ 每个 Subscription 内部一条线程,从磁盘/实时源拉数据 ▼ SubscriptionCollection (缓冲) │ 线程安全的并发集合,生产者写,消费者读 ▼ Synchronizer (单线程消费) │ SubscriptionSynchronizer.Sync() 合并所有订阅,按数据时间推进 │ 产 IEnumerable<TimeSlice> ▼ AlgorithmManager (单线程喂算法) │ 每个 TimeSlice → 算法.OnData(timeSlice.Slice) ▼ 用户 OnData
关键点:
💡 钻取要点:这个模式跟 vnpy 的"事件引擎 Queue+Thread"同构——都是"多生产者 → 队列 → 单消费者"的本质。区别是 Lean 的数据量级大得多(回测动辄几亿条),所以生产端用真正的多线程(每个 symbol 独立 IO),而不是 vnpy 那种单线程事件循环。
DataManager 的 SubscriptionAdded 事件有几个订阅者,完整链路:
SubscriptionRemoved 同理,触发 DataFeed 停止读、Synchronizer 移出合并、Subscription 释放枚举器。
FileSystemDataFeed(读 zip+CSV,源码注释"incredibly fast")/ 实盘 LiveTradingDataFeed(券商 API 实时)。这是回测/实盘第一个分叉点。IAlgorithmSubscriptionManager/IDataFeedSubscriptionManager/IDataManager 三个接口,统一管理订阅生命周期,解耦 DataFeed 与 SubscriptionManager。UniverseManager.CollectionChanged,用户 AddEquity → universe 增 → DataManager 建 SubscriptionRequest → DataFeed 起枚举器。Subscription(单订阅,IEnumerator<SubscriptionData>)、SubscriptionCollection(并发集合,缓冲)、SubscriptionDataReader(磁盘读取器,回测用)。下一节,我们钻取 Synchronizer 怎么把 SubscriptionCollection 里多个订阅的数据按时间合并成 TimeSlice 流,以及回测 Synchronizer 与实盘 LiveSynchronizer 的第二个分叉点。