第 2 章 · 02 Engine.Run 生命周期三阶段


文档摘要

第 2 章 · 02 Engine.Run 生命周期三阶段 本节摘要:本节精读 的 方法——Lean 单个 backtest/live job 的完整生命周期管理器。Engine 类的文档字符串自称 (引擎入口点),它是 Program.Main 与 AlgorithmManager 之间的中间层。Run 方法分三阶段:阶段 A 初始化(启动结果线程→L113 选 Synchronizer,这是回测/实盘的第一个分叉点→反射加载用户 DLL 实例化 IAlgorithm→创建券商→SecurityService/DataManager 粘合→DataFeed.Initialize→HistoryProvider 装配→执行 algorithm.

第 2 章 · 02 Engine.Run 生命周期三阶段

本节摘要:本节精读 Engine/Engine.csRun 方法——Lean 单个 backtest/live job 的完整生命周期管理器。Engine 类的文档字符串自称 LEAN ALGORITHMIC TRADING ENGINE: ENTRY POINT(引擎入口点),它是 Program.Main 与 AlgorithmManager 之间的中间层。Run 方法分三阶段:阶段 A 初始化(启动结果线程→L113 选 Synchronizer,这是回测/实盘的第一个分叉点→反射加载用户 DLL 实例化 IAlgorithm→创建券商→SecurityService/DataManager 粘合→DataFeed.Initialize→HistoryProvider 装配→执行 algorithm.Initialize 用户代码);阶段 B 主循环(Isolator 限制下的 algorithmManager.Run);阶段 C 清理(Dispose 4 个 handler + 断开券商)。读完本节,你看清了 Engine 怎么把"config + job + 用户 DLL"组装成一个跑得起来的算法实例。

内容来源:原项目源码 Engine/Engine.cs(共 608 行),聚焦 Engine 构造器(L72-78)与 Run 方法(L87-474)。

⚠️ 注意:L113 的三元 _liveMode ? new LiveSynchronizer() : new Synchronizer() 是整个引擎第一个回测/实盘分叉点——它决定了数据怎么"流"给算法:回测按历史时间片推(LiveSynchronizer 不需要),实盘按实时墙钟推。Python 对照:这一段没有任何用户代码介入,纯引擎内部;Python 算法的 Initialize 在阶段 A 末尾的 Setup.Setup 里被反射调用,与 C# 算法走同一套代码路径。配置坑:parallelHistoryRequestsEnabled: !_liveMode——实盘禁用并行历史请求,否则多个并行历史回调会与实时数据流竞争,破坏顺序性。

学习目标

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

  1. 解释 Engine 类文档字符串"ENTRY POINT"的含义与它的真实层级位置。
  2. 读懂 Engine 构造器里 _marketHoursDatabaseTask 的异步预热。
  3. 逐段讲清 Run 方法三阶段(初始化 / 主循环 / 清理)各自的关键动作。
  4. 说出 L113 Synchronizer vs LiveSynchronizer 是回测实盘的第一个分叉点。
  5. 解释 Setup.CreateAlgorithmInstance 怎么反射加载用户 DLL。
  6. 说清 SecurityService + DataManager + UniverseSelection 怎么粘合成数据通路。
  7. 解释 parallelHistoryRequestsEnabled: !_liveMode 为什么实盘要禁并行历史。
  8. 读懂阶段 C 的"等 4 个 handler IsActive=false,最多 30 秒"超时机制。

一、Engine 类头:ENTRY POINT 的自我定位

Engine 类的文档注释(Engine/Engine.cs:43-49)写得很有意思:

43 /// <summary> 44 /// LEAN ALGORITHMIC TRADING ENGINE: ENTRY POINT. 45 /// 46 /// The engine loads new tasks, create the algorithms and threads, and sends them 47 /// to Algorithm Manager to be executed. It is the primary operating loop. 48 /// </summary> 49 public class Engine

自称 ENTRY POINT(入口点)——但实际上 Program.Main 才是进程入口,Engine 是 Main 与 AlgorithmManager 之间的中间层。这段注释的语境是:对算法执行流程而言,Engine 是入口(Program.Main 只是进程壳)。Engine 负责"加载 task → 创建算法与线程 → 交给 AlgorithmManager 执行",这正是注释 L46-47 的话。

Engine 类的字段(Engine.cs:51-64):

51 private bool _historyStartDateLimitedWarningEmitted; 52 private bool _historyNumericalPrecisionLimitedWarningEmitted; 53 private readonly bool _liveMode; 54 private readonly Task<MarketHoursDatabase> _marketHoursDatabaseTask; 59 public LeanEngineSystemHandlers SystemHandlers { get; } 64 public LeanEngineAlgorithmHandlers AlgorithmHandlers { get; }

四个字段:两个一次性告警开关(避免刷屏)、_liveMode、后台预热的 _marketHoursDatabaseTask、两套 handler 的只读属性。前两个字段是"只发一次警告"的常见模式——历史数据起止日期警告、数值精度警告各发一次就够了。

二、构造器:异步预热静态数据库

Engine 构造器(Engine.cs:72-78)在第 2 章第 01 节已展开过,这里只回顾要点:

72 public Engine(LeanEngineSystemHandlers systemHandlers, LeanEngineAlgorithmHandlers algorithmHandlers, bool liveMode) 73 { 74 _liveMode = liveMode; 75 SystemHandlers = systemHandlers; 76 AlgorithmHandlers = algorithmHandlers; 77 _marketHoursDatabaseTask = Task.Run(StaticInitializations); 78 }

构造器不做任何"重活",只起一个后台预热 Task。真正的初始化都在 Run 方法里。这种"构造廉价、Run 贵"的设计,是为了让 Program.Main 能先 new Engine 拿到对象引用(在 try/finally 里更安全),再调用 Run——即使 Run 抛异常,finally 也能 dispose 已构造的 Engine。

三、阶段 A 初始化:Run 方法的上半场

Run 方法(Engine.cs:87-474)很长,但分三阶段后很清晰。阶段 A 是初始化,信息量最大

1. 启动结果线程(L95-109)

95 Log.Initialize(job.UserId, job.ProjectId, job.AlgorithmId); 97 Messages.SetAlgorithmLanguage(job.Language); 98 Log.Trace($"Engine.Run(): Resource limits '{job.Controls.CpuAllocation}' CPUs. {job.Controls.RamAllocation} MB RAM."); 99 TextSubscriptionDataSourceReader.SetCacheSize((int)(job.RamAllocation * 0.4)); ... 108 AlgorithmHandlers.Results.Initialize(new(job, SystemHandlers.Notify, SystemHandlers.Api, AlgorithmHandlers.Transactions, AlgorithmHandlers.MapFileProvider));

Results.Initialize 启动结果处理线程(BacktestingResultHandler / LiveTradingResultHandler 各自的实现)。TextSubscriptionDataSourceReader.SetCacheSize 把数据源读取缓存设成 RAM 配额的 40%——这是按 job 给的 RAM 上限调缓存,防止回测爆内存。

2. L113 选 Synchronizer:第一个分叉点(核心!)

110 IBrokerage brokerage = null; 111 DataManager dataManager = null; 112 var performanceTrackingTool = new PerformanceTrackingTool(); 113 var synchronizer = _liveMode ? new LiveSynchronizer() : new Synchronizer();

这是回测与实盘的第一个分叉点。两个 Synchronizer 都实现 ISynchronizer,但行为完全不同:

数据流方式 用途
Synchronizer 按历史时间片推(StreamData 吐 TimeSlice) 回测
LiveSynchronizer 按实时墙钟推,等待真实数据到来 实盘

后续 algorithmManager.Run(..., synchronizer, ...) 把这个 synchronizer 传进去,主循环 foreach (var timeSlice in Stream(algorithm, synchronizer, ...)) 直接消费它吐的 TimeSlice。L113 的这个选择,决定了主循环看到的数据是"快进的历史"还是"实时的当下"。

💡 钻取要点:这个分叉点比第 1 章的"换 6-7 个 handler"更本质——handler 是"用什么实现",Synchronizer 是"数据怎么流"。即使你把 handler 都换成 live 版,只要 synchronizer 还是 Synchronizer,引擎还是会按历史节奏推(这就是 paper trading 的本质:live handler + 历史节奏,实盘行情模拟成回测)。理解这一点,你就理解了 Lean 实盘模式的真正定义。

3. 反射加载用户 DLL(L118-131)

118 var marketHoursDatabase = _marketHoursDatabaseTask.Result; // 阻塞等异步预热完成 120 AlgorithmHandlers.Setup.WorkerThread = workerThread; 122 // Save algorithm to cache, load algorithm instance: 123 algorithm = AlgorithmHandlers.Setup.CreateAlgorithmInstance(job, assemblyPath); 125 algorithm.ProjectId = job.ProjectId; 128 SystemHandlers.LeanManager.SetAlgorithm(algorithm); 131 AlgorithmHandlers.ObjectStore.Initialize(job.UserId, job.ProjectId, job.UserToken, job.Controls, algorithm.AlgorithmMode);

L118 阻塞取构造器启动的预热结果——至此异步预热的并行收益兑现。L123 是用户算法实例化的关键:Setup.CreateAlgorithmInstance(job, assemblyPath) 反射加载 assemblyPath 指向的 DLL(或 .py),按 config 的 algorithm-type-name 找类,new 出 IAlgorithm 实例。C# 算法用 Assembly.LoadFrom(dll) + Activator.CreateInstance,Python 算法用 AlgorithmPythonWrapper 包装 .py。

注意此时还没有执行用户代码——只是构造了对象。用户的 Initialize() 要等到阶段 A 最后的 Setup.Setup 才调。

4. 创建券商(L143-147)

143 IBrokerageFactory factory; 144 brokerage = AlgorithmHandlers.Setup.CreateBrokerage(job, algorithm, out factory); 147 brokerage.Message += (_, e) => AlgorithmHandlers.Results.BrokerageMessage(e);

回测时 CreateBrokerage 返回内置的 BacktestingBrokerage(本地成交模拟);实盘时按 environment.*.live-mode-brokerage 实例化真实券商(InteractiveBrokersBrokerage 等)。L147 把券商的 Message 事件转发给 ResultHandler——所以券商的 disconnect/reconnect 消息能进结果流。

5. SecurityService + DataManager 粘合(L149-177)

这是阶段 A 里最"胶水"的一段:

149 var symbolPropertiesDatabase = SymbolPropertiesDatabase.FromDataFolder(); 150 var mapFilePrimaryExchangeProvider = new MapFilePrimaryExchangeProvider(AlgorithmHandlers.MapFileProvider); 151 var registeredTypesProvider = new RegisteredSecurityDataTypesProvider(); 152 var securityService = new SecurityService(algorithm.Portfolio.CashBook, 153 marketHoursDatabase, symbolPropertiesDatabase, algorithm, registeredTypesProvider, 154 new SecurityCacheProvider(algorithm.Portfolio), mapFilePrimaryExchangeProvider, algorithm, 155 new IndicatorBasedOptionPriceModelProvider(algorithm.Securities)); 162 algorithm.Securities.SetSecurityService(securityService); 164 dataManager = new DataManager(AlgorithmHandlers.DataFeed, 165 new UniverseSelection(algorithm, securityService, AlgorithmHandlers.DataPermissionsManager, AlgorithmHandlers.DataProvider), 166 algorithm, algorithm.TimeKeeper, marketHoursDatabase, _liveMode, registeredTypesProvider, AlgorithmHandlers.DataPermissionsManager); 177 algorithm.SubscriptionManager.SetDataManager(dataManager);

构造了一条数据通路:SecurityService(单只证券的工厂)→ UniverseSelection(决定哪些证券进订阅)→ DataManager(数据总管)→ SubscriptionManager(算法侧的订阅 API)。用户的 AddEquity("SPY") 最终走这条路:调用 SubscriptionManager → 触发 UniverseSelection → 用 SecurityService 造一个 SPY 的 Security → 注册到 DataManager。

6. DataFeed.Initialize + HistoryProvider 装配(L185-222)

179 synchronizer.Initialize(algorithm, dataManager, performanceTrackingTool); 185 AlgorithmHandlers.DataFeed.Initialize( 186 algorithm, job, AlgorithmHandlers.Results, AlgorithmHandlers.MapFileProvider, 190 AlgorithmHandlers.FactorFileProvider, AlgorithmHandlers.DataProvider, dataManager, 194 (IDataFeedTimeProvider)synchronizer, AlgorithmHandlers.DataPermissionsManager.DataChannelProvider); 197 var historyProvider = GetHistoryProvider(); 198 historyProvider.SetBrokerage(brokerage); 199 historyProvider.Initialize( 200 new HistoryProviderInitializeParameters(job, SystemHandlers.Api, AlgorithmHandlers.DataProvider, 204 AlgorithmHandlers.DataCacheProvider, AlgorithmHandlers.MapFileProvider, AlgorithmHandlers.FactorFileProvider, ... 217 // disable parallel history requests for live trading 218 parallelHistoryRequestsEnabled: !_liveMode,

DataFeed.Initialize 注入全部依赖(算法、map/factor provider、dataManager、作为时间提供者的 synchronizer)。HistoryProvider 装配时,L218 parallelHistoryRequestsEnabled: !_liveMode 是关键——实盘禁用并行历史请求。原因:实盘里多个历史回调并发完成,会与实时数据流的顺序竞争,可能让算法的 OnData 看到不一致的状态。回测没有实时流,并行是安全的性能优化。

7. Setup.Setup 执行用户 Initialize(L263-265)

阶段 A 最后一步,执行用户的 Initialize 代码:

263 initializeComplete = AlgorithmHandlers.Setup.Setup(new SetupHandlerParameters(dataManager.UniverseSelection, algorithm, 264 brokerage, job, AlgorithmHandlers.Results, AlgorithmHandlers.Transactions, AlgorithmHandlers.RealTime, 265 AlgorithmHandlers.DataCacheProvider, AlgorithmHandlers.MapFileProvider));

到这一行,用户在 Initialize() 里写的 SetStartDate / AddEquity / SetCash / SetWarmUp 全部执行完毕。initializeComplete 是布尔结果——任何用户 Initialize 抛异常,这里都会被捕获并标记为 false,Run 方法会跳过主循环直接进清理。

四、阶段 B 主循环:Isolator 限制下的 algorithmManager.Run

阶段 A 成功后(initializeComplete == true),进阶段 B(Engine.cs:322-417)。核心两行:

346 var isolator = new Isolator(); 349 var complete = isolator.ExecuteWithTimeLimit(AlgorithmHandlers.Setup.MaximumRuntime, algorithmManager.TimeLimit.IsWithinLimit, () => 350 { 351 try 352 { 357 algorithmManager.Run(job, algorithm, synchronizer, AlgorithmHandlers.Transactions, AlgorithmHandlers.Results, AlgorithmHandlers.RealTime, SystemHandlers.LeanManager, isolator.CancellationTokenSource, performanceTrackingTool); 358 } ... 366 }, job.Controls.RamAllocation, workerThread: workerThread, sleepIntervalMillis: algorithm.LiveMode ? 10000 : 1000);

Isolator.ExecuteWithTimeLimit 做两件事:

  1. 时间预算:给算法一个最大运行时长 MaximumRuntime(回测默认 1 小时左右,实盘无限),超时强制中止。
  2. RAM 预算:job.Controls.RamAllocation 监控内存,超 RAM 抛 OutOfMemory。

主循环本身在 algorithmManager.Run 里(下一节详讲 13 步)。Engine 这层只负责"圈住"它,任何异常或超时都能被外部捕获。

回测结束后,L398-405 打印著名的性能横幅:

404 AlgorithmHandlers.Results.DebugMessage($"Algorithm Id:({job.AlgorithmId}) completed in {totalSeconds:F2} seconds at {kps:F0}k data points per second. Processing total of {dataPoints:N0} data points.");

这就是你跑完 Lean 回测看到的那条 "completed in 12.34 seconds at 5k data points per second"。

五、阶段 C 清理:Dispose + 等 30 秒

阶段 C(Engine.cs:419-455)的清理顺序很重要:

419 synchronizer.DisposeSafely(); 421 AlgorithmHandlers.DataFeed.Exit(); 424 AlgorithmHandlers.Results.Exit(); 427 var millisecondInterval = 10; 428 var millisecondTotalWait = 0; 429 while ((AlgorithmHandlers.Results.IsActive 430 || (AlgorithmHandlers.Transactions != null && AlgorithmHandlers.Transactions.IsActive) 431 || (AlgorithmHandlers.DataFeed != null && AlgorithmHandlers.DataFeed.IsActive) 432 || (AlgorithmHandlers.RealTime != null && AlgorithmHandlers.RealTime.IsActive)) 433 && millisecondTotalWait < 30 * 1000) 434 { 435 Thread.Sleep(millisecondInterval); 436 if (millisecondTotalWait % (millisecondInterval * 10) == 0) 437 { 438 Log.Trace("Waiting for threads to exit..."); 439 } 440 millisecondTotalWait += millisecondInterval; 441 }

L429-433 是关键的优雅退出等待:依次检查 4 个 handler 的 IsActive 标志,最多等 30 秒。每 10ms 检查一次,每 100ms 打一条 "Waiting for threads to exit..." 日志。30 秒还没退完就放弃等待,继续往下走(此时可能有线程泄漏,但比卡死整个进程强)。

接着断开券商、释放 setup(Engine.cs:443-453):

443 if (brokerage != null) 444 { 445 Log.Trace("Engine.Run(): Disconnecting from brokerage..."); 446 brokerage.Disconnect(); 447 brokerage.Dispose(); 448 } 449 if (AlgorithmHandlers.Setup != null) 450 { 451 Log.Trace("Engine.Run(): Disposing of setup handler..."); 452 AlgorithmHandlers.Setup.Dispose(); 453 }

finally 块(L462-473)兜底再 Exit 一遍 4 个 handler + DataMonitor,保证任何异常路径下 handler 都被通知退出:

464 if (_liveMode && algorithmManager.State != AlgorithmStatus.Running && algorithmManager.State != AlgorithmStatus.RuntimeError) 465 SystemHandlers.Api.SetAlgorithmStatus(job.AlgorithmId, algorithmManager.State); 467 AlgorithmHandlers.Results.Exit(); 468 AlgorithmHandlers.DataFeed.Exit(); 469 AlgorithmHandlers.Transactions.Exit(); 470 AlgorithmHandlers.RealTime.Exit(); 471 AlgorithmHandlers.DataMonitor.Exit();

L464 是实盘专属——保证实盘异常退出时,云端 API 一定知道这个算法不再运行(避免云端以为算法还在跑)。

⚠️ 注意:IsActive 是个只读属性,各 handler 内部用自己的退出标志维护它。ResultHandler / TransactionHandler / DataFeed / RealTime 都有独立的后台线程,Exit() 是"请求退出",IsActive=false 是"真的退完了"。这种"请求 + 等待确认"的优雅退出模式,与第 1 章 vnpy 的 _active 标志位 + join 思路完全一致——只是 Lean 用 IsActive 属性而非 _active 字段。

六、JOB HANDLERS 日志:阶段 A 装配的可观测性

阶段 A 末尾 L310-319 有一条信息密度极高的日志,把所有装配的 handler 类型全打出来:

310 Log.Trace($"JOB HANDLERS:{Environment.NewLine}" + 311 $" DataFeed: {AlgorithmHandlers.DataFeed.GetType().FullName}{Environment.NewLine}" + 312 $" Setup: {AlgorithmHandlers.Setup.GetType().FullName}{Environment.NewLine}" + 313 $" RealTime: {AlgorithmHandlers.RealTime.GetType().FullName}{Environment.NewLine}" + 314 $" Results: {AlgorithmHandlers.Results.GetType().FullName}{Environment.NewLine}" + 315 $" Transactions: {AlgorithmHandlers.Transactions.GetType().FullName}{Environment.NewLine}" + 316 $" Object Store: {AlgorithmHandlers.ObjectStore.GetType().FullName}{Environment.NewLine}" + 317 $" History Provider: {historyProviderName}{Environment.NewLine}" + 318 $" Brokerage: {brokerage?.GetType().FullName}{Environment.NewLine}" + 319 $" Data Provider: {AlgorithmHandlers.DataProvider.GetType().FullName}{Environment.NewLine}");

跑 Lean 时看到的那块 JOB HANDLERS: 列表就是它。这是排查"我到底在跑哪个模式"的第一手信息——回测会看到 FileSystemDataFeed / BacktestingBrokerage,实盘会看到 LiveTradingDataFeed / InteractiveBrokersBrokerage。出问题时先看这段日志,就能确认 config.json 的 environment 有没有正确叠加。

💡 钻取要点:Engine 这一层(约 470 行 Run)其实不处理任何一根数据——它只是把所有零件装起来,然后交给 AlgorithmManager 跑,跑完清理。真正的"逐根数据"在下一节的 AlgorithmManager 13 步。Engine 是"导演",AlgorithmManager 是"主角"。

本节要点回顾

  1. Engine 自称 ENTRY POINT:对算法执行流程而言是入口,实际位于 Program.Main 与 AlgorithmManager 之间。
  2. 构造器异步预热:Task.Run(StaticInitializations) 后台读 market-hours/symbol-properties,Run 时阻塞取结果。
  3. 阶段 A 初始化:启动结果线程 → L113 选 Synchronizer(回测实盘第一个分叉点) → 反射加载用户 DLL → 创建券商 → SecurityService/DataManager/UniverseSelection 粘合 → DataFeed.Initialize → HistoryProvider 装配(parallelHistoryRequestsEnabled: !_liveMode)→ Setup.Setup 执行用户 Initialize。
  4. 阶段 B 主循环:Isolator 限制(RAM + 时间预算)下的 algorithmManager.Run;回测结束打 "completed in X seconds at Yk data points per second"。
  5. 阶段 C 清理:DisposeSafely(synchronizer)→ Exit(DataFeed/Results)→ 等 4 个 handler IsActive=false 最多 30 秒 → Disconnect 券商 → Dispose Setup。
  6. finally 兜底:实盘专属 SetAlgorithmStatus,再 Exit 一遍 5 个 handler,保证异常路径下 handler 一定退出。
  7. JOB HANDLERS 日志:L310-319 打印全部装配的 handler 类型,排查模式问题的第一手信息。

下一节,我们钻进 algorithmManager.Run 的 foreach timeSlice 循环——按固定 13 步顺序处理每个时间片,最终把数据交给用户的 OnData 回调。


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