第 2 章 · 02 Engine.Run 生命周期三阶段 本节摘要:本节精读 的 方法——Lean 单个 backtest/live job 的完整生命周期管理器。Engine 类的文档字符串自称 (引擎入口点),它是 Program.Main 与 AlgorithmManager 之间的中间层。Run 方法分三阶段:阶段 A 初始化(启动结果线程→L113 选 Synchronizer,这是回测/实盘的第一个分叉点→反射加载用户 DLL 实例化 IAlgorithm→创建券商→SecurityService/DataManager 粘合→DataFeed.Initialize→HistoryProvider 装配→执行 algorithm.
本节摘要:本节精读
Engine/Engine.cs的Run方法——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——实盘禁用并行历史请求,否则多个并行历史回调会与实时数据流竞争,破坏顺序性。
阅读完本节,你应当能够:
_marketHoursDatabaseTask 的异步预热。parallelHistoryRequestsEnabled: !_liveMode 为什么实盘要禁并行历史。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。
Run 方法(Engine.cs:87-474)很长,但分三阶段后很清晰。阶段 A 是初始化,信息量最大。
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 上限调缓存,防止回测爆内存。
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 实盘模式的真正定义。
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 才调。
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 消息能进结果流。
这是阶段 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。
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 看到不一致的状态。回测没有实时流,并行是安全的性能优化。
阶段 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 方法会跳过主循环直接进清理。
阶段 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 做两件事:
MaximumRuntime(回测默认 1 小时左右,实盘无限),超时强制中止。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(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 字段。
阶段 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 是"主角"。
Task.Run(StaticInitializations) 后台读 market-hours/symbol-properties,Run 时阻塞取结果。parallelHistoryRequestsEnabled: !_liveMode)→ Setup.Setup 执行用户 Initialize。下一节,我们钻进
algorithmManager.Run的 foreach timeSlice 循环——按固定 13 步顺序处理每个时间片,最终把数据交给用户的 OnData 回调。