Skip to content

fix(multifactor): make lazy calculation thread-safe - #495

Merged
fasiondog merged 9 commits into
fasiondog:masterfrom
woleigegg:fix/multifactor-lazy-calculation-thread-safe-pr
Aug 3, 2026
Merged

fix(multifactor): make lazy calculation thread-safe#495
fasiondog merged 9 commits into
fasiondog:masterfrom
woleigegg:fix/multifactor-lazy-calculation-thread-safe-pr

Conversation

@woleigegg

@woleigegg woleigegg commented Aug 3, 2026

Copy link
Copy Markdown
Contributor

RANK插件修复系列PR,后续修复还在进行中。

问题概述

MultiFactorBase 的惰性计算 calculate() 存在 C++ 数据竞争:m_calculated 为普通 bool,锁外读 + 锁内写且锁后不二次检查,多线程首次访问同一未计算实例时可能重复重建内部容器,导致持有 getter 内部引用的线程撕裂读或悬空引用。同时,计算失败时异常被吞掉仍设 m_calculated = true,向下游发布半成品。

此外,风格因子中性化路径中进程级 Eigen::setNbThreads 全局切换在并发 MF 实例间互相污染,wait_for_all_non_blocking 在子任务抛异常时提前返回会遗留未执行任务。

根因分析

缺陷 1:calculate() check-then-lock 无 double-check

// master: MultiFactorBase.cpp:816-855
bool m_calculated{false};               // 普通 bool, 非原子

void MultiFactorBase::calculate() {
    HKU_IF_RETURN(m_calculated, void());  // 锁外读, 无内存屏障
    std::lock_guard<std::mutex> lock(m_mutex);
    // 拿到锁后没有重查 m_calculated
    _checkData();
    _buildIndex();
    m_calculated = true;                 // 滞后写入
}
  • 普通 bool 锁外读与锁内写构成 C++ data race (UB)。
  • 锁后不二次检查:T1 释放锁后 T2 拿到锁会重算第二遍,在 rebuild m_stk_factor_by_date 时若有线程通过 getAllScores() 返回的引用正在读取,造成撕裂读/悬空引用。

缺陷 2:失败时发布半成品

// master: calculate() catch 块
} catch (const std::exception& e) {
    HKU_ERROR(e.what());          // 吞异常
} catch (...) {
    HKU_ERROR_UNKNOWN;            // 吞异常
}
m_calculated = true;              // 仍标记为已计算

_checkData() 或构建过程抛异常时,异常被吞掉,m_calculated 仍设为 true,下游 getter 返回空的/半成品容器,调用方无法感知失败。

缺陷 3:Eigen::setNbThreads 进程级全局竞态

// master: MultiFactorBase.cpp:680,749
Eigen::setNbThreads(1);
// ... 并行回归 ...
Eigen::setNbThreads(std::thread::hardware_concurrency());

两个 MF 实例并发 calculate() 时,一方 setNbThreads(16) 恢复全局值,另一方仍在回归中——进程级全局状态被互相污染。

缺陷 4:wait_for_all_non_blocking 提前返回遗留任务

// master: algorithm.h:228-230
if (!pool.run_available_task_once()) {
    return;                        // 一个 future 抛异常时直接返回
}

一个子任务抛异常、另一个 delayed 子任务尚未启动时 run_available_task_once() 返回 false → 直接 return,遗留未执行任务。calculate() 的 catch 块进入清理时后台仍有任务在跑,可能访问已释放资源。

修复方案

1. m_calculatedstd::atomic<bool> + acquire/release DCLP

std::atomic<bool> m_calculated{false};

void calculate() {
    if (m_calculated.load(std::memory_order_acquire))   // 快路径: lock-free
        return;
    std::lock_guard<std::mutex> lock(m_mutex);
    if (m_calculated.load(std::memory_order_relaxed))   // 锁内二次检查
        return;
    clearCalculatedData();                              // 构建前清理
    try {
        _checkData();
        getAllSrcFactors();
        _buildIndex();
    } catch (...) {
        clearCalculatedData();                           // 失败清理半成品
        m_calculated.store(false, std::memory_order_relaxed);
        throw;                                           // 原异常向上传播
    }
    m_calculated.store(true, std::memory_order_release); // 发布: 此前写入可见
}
  • acquire/release 配对:构建线程 release store 前的所有写入,对读取线程 acquire load 看到 true 后全部可见。
  • 锁内二次检查用 relaxed:mutex 已提供慢路径同步,只需确认状态。
  • 失败时清理派生数据、保持 false、原异常传播,允许下一调用者重试。

2. 提取 clearCalculatedData()

清除所有计算后派生数据 (m_ref_dates/m_stk_map/m_all_factors/m_date_index/m_stk_factor_by_date/m_ic),保留配置成员。calculate() 构建前和异常后均调用,确保重试基于干净状态。

3. getters 统一委托 calculate()

所有 getter (getDatetimeList/getFactor/getAllFactors/getScores/getAllScores/getIC) 改为直接调 calculate(),由其 lock-free 快路径承担已就绪检查。getIC()m_mutex 锁在 calculate() 之后获取,严格顺序不嵌套,无自死锁。

4. setters 用 relaxed store false

setQuery/setRefStock/setStockList/setRefFactorSet/setNormalize/addSpecialNormalize/paramChanged 改为 m_calculated.store(false, relaxed)。setter 不与 getter 并发执行(见已知局限),只需让下次 calculate() 知道该重算。

5. reset() 全程持锁

void reset() {
    std::lock_guard<std::mutex> lock(m_mutex);
    _reset();                    // 虚函数
    clearCalculatedData();
    m_calculated.store(false, std::memory_order_release);
}

6. 提取 StyleRegression.cpp 串行内核

calculate_residualsMultiFactorBase.cpp 静态函数提取为 StyleRegression.cppcalculate_style_residuals,删除进程级 Eigen::setNbThreads 调用。回归矩阵均为栈局部对象,外层按日并行天然可重入,不依赖全局线程配置。

7. wait_for_all_non_blocking drain 修正

删除 run_available_task_once() 返回 false 时的提前 return,改为继续等待所有 future 就绪,确保调用方 get() 抛异常时无遗留任务。

验证

编译

  • 环境:MSVC 2022 (v17.14.19) + xmake,release,Windows x86_64
  • 结果:build ok, spent 319.5s,0 error(2 个既有 C4018 warning,非本 PR 引入)

功能验证

PR1 新增 13 个测试用例,逐个单独运行确认通过:

测试用例 验证点 结果
test_MF_thread_safe_concurrent_first_access 32 线程混合 getter 并发首次访问,build 恰好一次,结果一致 1 passed
test_MF_failed_first_call_clean_retry 首次计算注入异常 → 清理 + 异常传播 + 第二次成功重试 1 passed
test_MF_reset_recalculates_clean reset 后重算无陈旧结果 1 passed
test_MF_nested_calculate_no_deadlock 嵌套 MF 惰性触发无死锁 1 passed
test_MF_clone_independent_state ready 原件 + clone 并发访问,只 clone 重算,结果一致 1 passed
test_MF_serialization_load_recalculates 序列化 load 不发布陈旧状态,触发重算 1 passed
test_style_regression_perfect_fit 一元完美拟合残差全 0 1 passed
test_style_regression_golden_values OLS 解析解核对 [0.4, -1.2, 1.2, -0.4] 1 passed
test_style_regression_insufficient_samples 样本不足返回全 NaN 不崩溃 1 passed
test_style_regression_nan_row 自变量含 NaN 行残差置 NaN 1 passed
test_style_regression_rank_deficient 秩亏共线不崩溃 1 passed
test_style_regression_eigen_threads_unchanged 并发调用前后 Eigen::nbThreads() 不变(防回归) 1 passed
test_global_wait_drains_all_futures_on_exception 一个任务抛异常时另一 delayed 任务也被 drain 1 passed

回归测试

完整 unit-test 套件(排除 benchmark):

test cases:    791 |    791 passed | 0 failed | 0 skipped
assertions: 208669 | 208669 passed | 0 failed
Status: SUCCESS!
total time:   8.585 s

0 回归。

改动范围

文件 改动
hikyuu_cpp/hikyuu/trade_sys/multifactor/MultiFactorBase.h m_calculatedatomic<bool>clearCalculatedData() 声明,并发语义注释,序列化 load 改 atomic store
hikyuu_cpp/hikyuu/trade_sys/multifactor/MultiFactorBase.cpp DCLP calculate,clearCalculatedData,getters 委托,setters relaxed,reset 持锁,删除 setNbThreads,删除内联 calculate_residuals
hikyuu_cpp/hikyuu/trade_sys/multifactor/StyleRegression.h calculate_style_residuals 声明(HKU_API 导出)
hikyuu_cpp/hikyuu/trade_sys/multifactor/StyleRegression.cpp 串行 QR 回归内核实现
hikyuu_cpp/hikyuu/utilities/thread/algorithm.h wait_for_all_non_blocking drain 修正
hikyuu_cpp/hikyuu/xmake.lua StyleRegression.cpp 排除 unity build 组
hikyuu_cpp/unit_test/hikyuu/trade_sys/multifactor/test_MF_ThreadSafe.cpp 6 个并发测试用例
hikyuu_cpp/unit_test/hikyuu/trade_sys/multifactor/test_style_regression.cpp 6 个风格回归测试用例
hikyuu_cpp/unit_test/hikyuu/utilities/thread/test_algorithm.cpp 1 个 drain 测试用例
hikyuu_cpp/unit_test/xmake.lua unit-test target 添加 eigen package

总计 10 文件,+883/-147 行。

说明:style_regression.* 重命名为 StyleRegression.*(PascalCase,与 MultiFactorBase 等编译单元命名一致),同时更新 MultiFactorBase.cpp 的 include、xmake.lua 的 unity build 排除规则。

性能影响

  • 快路径m_calculated.load(acquire) 为 lock-free 原子读,与旧版 if (m_calculated) 普通读相比可忽略。
  • 慢路径:锁内多一次 relaxed 原子读,可忽略。
  • 风格回归:保持按日并行,仅删除全局 setNbThreads 调用,回归矩阵为栈局部对象不依赖全局配置。
  • 高并发首次 miss 场景下的锁争用优化见后续 PR6(single-flight),本 PR 优先保证正确性。

已知局限与后续工作

1. 并发边界:不支持 getter 与配置 mutation 并发

现状:本 PR 实现多线程同时首次触发 calculate/getter 的安全性及计算完成后的并发只读。头文件注释 (MultiFactorBase.h:175-184) 明确划界。

局限:不支持 getter 与 reset/setQuery/setStockList/setRefFactorSet/setParam 同时执行,不支持调用方持有 getter 返回的内部引用期间另一线程修改实例。getter 返回内部容器引用,若要支持读写完全并发需改为不可变快照或复制返回,属更大 API 重构。

后续:独立 issue 评估 getter 返回值改为 shared_ptr<const ...> 的 API 影响。

2. Python 子类 _reset() 自死锁风险

现状reset() 改为全程持锁并调用虚函数 _reset()。所有内置派生类 (EqualWeight/IC/ICIR/Weight) 均未覆盖 _reset(),用基类空实现,当前无死锁。头文件注释 (MultiFactorBase.cpp:179) 已声明约束。

局限MultiFactorBase 通过 pybind trampoline 导出,Python 子类可覆盖 _reset()。若 Python _reset() 内调用 self.calculate()/getIC()(均需 m_mutex),会自死锁于非递归 std::mutex。文档约束已声明,但语言层面未阻止。

后续:可选方案(独立 issue)— reset() 拆分为锁外 _resetUnsafe() + 锁内清理,或改用 std::recursive_mutex

3. DCLP 非 single-flight,高并发首次 miss 锁争用

现状:采用 acquire/release DCLP + 锁后二次检查。首个拿到锁的线程构建,后续线程通过二次检查跳过重算。

局限:在首个构建线程释放锁前,其他等待锁的线程排队进入慢路径,存在锁串行化延迟。

后续:独立 PR 引入 Pending/Ready 分离 + promise/future single-flight,需满足嵌套 RANK 和线程池兼容验收标准后方可合并。

4. ThreadSanitizer 未验证

现状:新增 13 个并发测试用例覆盖 PLAN §17.1 验收矩阵的 4/5 项。

局限:开发环境为 Windows MSVC,无 ThreadSanitizer 支持,未做动态竞态检测。正确性依靠内存序推理(release/acquire happens-before)和代码审查。

后续:在 Linux + clang/gcc + TSan 环境补跑 unit-test,可作为独立 CI 步骤。

5. ABI 变更

现状m_calculatedbool 改为 std::atomic<bool>,改变 MultiFactorBase 对象布局。

局限:新 hikyuu core + 旧 extind.dll(未经重编)= 不受支持配置。仓库内所有消费者一起重编,无 in-repo 破坏。风险限于 out-of-tree C++ 插件按旧头文件编译而链接新库。

后续:插件包应声明最低核心 build/version,发布说明中明确二进制兼容性要求。

6. selector 异常传播行为变化

现状:旧 calculate() 捕获异常后仍设 m_calculated = true,下游 selector 在空数据上静默运行。新实现清理半成品、保持 false、原异常传播。MultiFactorSelector/MultiFactorSelector2_calculate()_getSelected() 中四处调用 m_mf->calculate()/getScores() 无 try/catch。

局限:这不是局限,是有意的行为修正——旧实现让 _checkData() 失败被吞掉,selector 选出 0 只股票而调用方看不到异常。新实现让异常传播至上层框架捕获,符合"失败不应伪装成功"原则。Python 端由 pybind11 自动转换异常。

后续:无需后续 issue。若上游反馈某些场景需容错,应在 selector 层加 try/catch + 降级策略,而非回退 MF 异常语义。

不在本 PR 范围的后续工作

RANK 完整修复按设计文档拆分为六步 PR。本 PR 仅含 PR1(核心 MF 并发修复)。确定性截面排序(tie 不确定)、进程级数据 revision、插件结构化缓存 key(碰撞/先发布后计算)、不可变 RankPanel、并发构建合并分别由后续 PR 独立交付,依赖 PR1 完成后推进。

- m_calculated: bool -> std::atomic<bool>, acquire/release DCLP
- Extract clearCalculatedData(), called before build and on failure
- calculate(): _checkData under try, failure cleans and rethrows
- reset(): hold mutex for whole reset, _reset then clear then store false
- getters: delegate to calculate() which has lock-free fast path
- setters: relaxed store false (mutation not concurrent with reads)
- Extract style regression QR to style_regression.cpp serial kernel
- Remove process-global Eigen::setNbThreads during neutralization
- wait_for_all_non_blocking: drain all futures before returning

ABI: MultiFactorBase layout changes, all C++ extensions must recompile
- style_regression.cpp excluded from unity group to prevent Eigen
  header leakage into sibling sources
- wait_for_all_non_blocking: yield first then bounded 1ms backoff
- drain: one task throws while another delayed task runs, exception
  returns only after all futures drained (spin-synchronized start)
- concurrent first access: 32 threads mixed getters, build exactly once
- failed first call: dirty derived state cleaned, original exception
  propagates, second call succeeds on clean state
- reset recalculates without stale results
- nested MF calculation triggers another MF lazily, no deadlock
- serialization load does not publish stale state
- style regression golden values, NaN row, rank deficient, Eigen
  thread configuration unchanged under concurrency
- style regression golden residuals corrected to OLS solution
  [0.4, -1.2, 1.2, -0.4] for y=[2,1,4,3], x=[0,1,2,3]
- Eigen thread config check unconditional: nbThreads() always exists,
  assert unchanged after concurrent calls regardless of OpenMP
- add clone independence test: ready original + clone accessed
  concurrently, only clone recalculates, results identical
test_style_regression.cpp includes Eigen/Core; the unit-test target
was missing the eigen package, unlike the core target
MULTIFACTOR_IMP declares _calculate in the class body; defining it
inline caused C2535 duplicate declaration. Move to out-of-class
definitions matching the built-in subclass pattern.
unit-test links against hikyuu.dll; without HKU_API the symbol was
not exported and the test binary failed to link (LNK2019)
@woleigegg
woleigegg marked this pull request as draft August 3, 2026 09:54
Match the multifactor directory naming convention where
compiled translation units use PascalCase (MultiFactorBase,
NormalizeBase, ScoreRecord) rather than snake_case.
@fasiondog
fasiondog marked this pull request as ready for review August 3, 2026 17:59
@fasiondog
fasiondog merged commit 779b893 into fasiondog:master Aug 3, 2026
8 checks passed

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

其实,这个wait_for_all_non_blocking 改动意义不大,这个就是微小的权衡影响,类似的很早就尝试过。后面就属于AI总在这没事微调了。

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

应该说是,不同场景下的微调,平衡所有常用场景。本身影响不大,但AI在变得不同场景下使用时,总会喜欢在这里改来改去。

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

嗯嗯,谢谢老师给的建议,我后续pr尽可能小一些,做好约束和审查。

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants