fix(worker): 消除回测子进程结果消息丢失竞态

现象: 整年区间回测偶发 'backtest worker exited without result (exitcode=0)'。
根因: 子进程 event_queue.put 只是把消息交给后台 feeder 线程, 主线程随即退出,
feeder 随进程销毁, 大结果/高负载下消息尾部未刷入管道; 父进程 0.1s 轮询
恰在 Empty+进程已死时跳出循环, 未读到的 result 被丢弃。

修复 (双侧):
- 子进程 finally 中 close+join_thread 队列, 保证退出前消息完整刷入管道
- 父进程在判定无结果前做一次兜底排空 (join 后 1s×2 轮 get)
This commit is contained in:
shy3130
2026-08-30 19:05:19 +08:00
parent d4aa9ed6ad
commit b8fa08f027
+21
View File
@@ -254,6 +254,12 @@ def _worker_entry(task: dict[str, Any], event_queue, cancel_event) -> None:
if store is not None:
with suppress(Exception):
store.db.close()
# 保证结果消息在进程退出前完整刷入管道: put 只是入队,
# 实际写管道的是后台 feeder 线程; 不 join 的话主线程先退出,
# feeder 随进程销毁, 消息尾部丢失 → 父进程误判 "exited without result"。
with suppress(Exception):
event_queue.close()
event_queue.join_thread()
def run_worker_task(
@@ -310,6 +316,21 @@ def run_worker_task(
elif message_type == "error":
failure = message
# 子进程退出后, 队列读线程可能尚未把管道尾部的 result/error 搬进本地缓冲
# (0.1s 轮询在系统高负载下会先看到 Empty+进程已死)。join 后做一次兜底排空,
# 只要消息完整刷入过管道就一定能取到。
if result is None and failure is None:
for _ in range(2):
try:
message = events.get(timeout=1.0)
except queue.Empty:
break
message_type = message.get("type")
if message_type == "result":
result = message["payload"]
elif message_type == "error":
failure = message
process.join(timeout=10.0)
if process.is_alive():
process.terminate()