第 5 章 · 01 DataFeed 与 DataManager 中介者


文档摘要

第 5 章 · 01 DataFeed 与 DataManager 中介者 本节摘要:本节钻取第 5 层入口——数据流子系统。Lean 的 目录是整个项目最大的子模块(120+ 文件),职责是把多品种、多频率的原始数据按时间对齐,统一喂给算法。核心三件套:DataFeed(生产者,多线程生产数据)、DataManager(中介者,统一管理订阅生命周期)、Subscription(单个订阅对象)。回测用 (注释原话"incredibly fast"——从磁盘读 zip+CSV),实盘用 (从券商实时拉)。DataManager 同时给 DataFeed 和 SubscriptionManager 提供订阅管理,是两者之间的中介者。

第 5 章 · 01 DataFeed 与 DataManager 中介者

本节摘要:本节钻取第 5 层入口——数据流子系统。Lean 的 Engine/DataFeeds 目录是整个项目最大的子模块(120+ 文件),职责是把多品种、多频率的原始数据按时间对齐,统一喂给算法。核心三件套:DataFeed(生产者,多线程生产数据)、DataManager(中介者,统一管理订阅生命周期)、Subscription(单个订阅对象)。回测用 FileSystemDataFeed(注释原话"incredibly fast"——从磁盘读 zip+CSV),实盘用 LiveTradingDataFeed(从券商实时拉)。DataManager 同时给 DataFeed 和 SubscriptionManager 提供订阅管理,是两者之间的中介者。本节先讲这三者的职责分工,02 节讲 Synchronizer 怎么消费。

内容来源:原项目源码 Engine/DataFeeds/FileSystemDataFeed.csEngine/DataFeeds/DataManager.csEngine/DataFeeds/Subscription.csEngine/DataFeeds/SubscriptionCollection.cs

⚠️ 注意:Engine/DataFeeds 子目录巨多(Enumerators/Queues/Transport/WorkScheduling/Factories 等),本节只钻取顶层骨架,具体枚举器(Fill/Refresh/Frontfill)在 03 节、Subscription 内部状态机在第 7 章。

学习目标

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

  1. 说清 IDataFeed 接口与两个实现(回测 FileSystemDataFeed/实盘 LiveTradingDataFeed)的分工。
  2. 理解 DataManager 作为"中介者"的角色——为什么需要它同时管 DataFeed 和 SubscriptionManager。
  3. 认识 Subscription/SubscriptionCollection/SubscriptionDataReader 三个核心对象。
  4. 描述"生产者消费者"链:DataFeed 多线程生产 → SubscriptionCollection 缓冲 → Synchronizer 单线程消费 → AlgorithmManager 单线程喂算法。
  5. 知道 DataManager 怎么挂 UniverseManager.CollectionChanged 监听 universe 增删。

一、IDataFeed 接口与两个实现

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:中介者

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 / SubscriptionCollection / SubscriptionDataReader

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

关键点:

  • 生产者多线程:FileSystemDataFeed 给每个 Subscription 起独立任务,并行读多个 symbol 的磁盘文件。
  • 消费者单线程:Synchronizer 和 AlgorithmManager 都单线程,保证用户策略代码不需要加锁。
  • 缓冲解耦:SubscriptionCollection 用并发集合,生产快了不会丢(消费者跟不上会反压),生产慢了消费者会等。
  • 时间对齐:Synchronizer 的核心职责是把不同 symbol 不同频率的数据按时间戳合并成 TimeSlice(02 节详讲)。

💡 钻取要点:这个模式跟 vnpy 的"事件引擎 Queue+Thread"同构——都是"多生产者 → 队列 → 单消费者"的本质。区别是 Lean 的数据量级大得多(回测动辄几亿条),所以生产端用真正的多线程(每个 symbol 独立 IO),而不是 vnpy 那种单线程事件循环。

五、SubscriptionAdded 事件下游

DataManager 的 SubscriptionAdded 事件有几个订阅者,完整链路:

  1. Synchronizer(02 节):收到新订阅,把它纳入时间合并逻辑。
  2. LiveTradingDataFeed(实盘):收到新订阅,告诉券商 API 开始推这个 symbol 的实时数据。
  3. AggregationManager:负责把多个订阅的更新聚合成给指标用。

SubscriptionRemoved 同理,触发 DataFeed 停止读、Synchronizer 移出合并、Subscription 释放枚举器。

本节要点回顾

  1. IDataFeed 两实现:回测 FileSystemDataFeed(读 zip+CSV,源码注释"incredibly fast")/ 实盘 LiveTradingDataFeed(券商 API 实时)。这是回测/实盘第一个分叉点。
  2. DataManager 中介者:同时实现 IAlgorithmSubscriptionManager/IDataFeedSubscriptionManager/IDataManager 三个接口,统一管理订阅生命周期,解耦 DataFeed 与 SubscriptionManager。
  3. 挂 CollectionChanged:DataManager 监听 UniverseManager.CollectionChanged,用户 AddEquity → universe 增 → DataManager 建 SubscriptionRequest → DataFeed 起枚举器。
  4. 三个核心对象:Subscription(单订阅,IEnumerator<SubscriptionData>)、SubscriptionCollection(并发集合,缓冲)、SubscriptionDataReader(磁盘读取器,回测用)。
  5. 生产者消费者链:DataFeed 多线程生产 → SubscriptionCollection 缓冲 → Synchronizer 单线程消费 → AlgorithmManager 单线程喂算法 → 用户 OnData。

下一节,我们钻取 Synchronizer 怎么把 SubscriptionCollection 里多个订阅的数据按时间合并成 TimeSlice 流,以及回测 Synchronizer 与实盘 LiveSynchronizer 的第二个分叉点。


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