RxCpp中Subject派生Observable无法在独立线程运行问题求助
问题原因解释
- Subject的本质特性:Subject(包括BehaviorSubject)是事件转发器,本身不生成事件,所有事件由外部调用
on_next()注入。事件的初始线程完全由调用on_next()的线程决定,和订阅时的调度器配置无关。 - subscribe_on的作用边界:
subscribe_on()仅控制Observable订阅逻辑的执行线程(也就是注册观察者的过程),但Subject的订阅逻辑只是完成观察者注册,不涉及事件分发,因此无法改变事件最终到达订阅者的线程。 - observe_on的正确用法:
observe_on()才是用来指定订阅者回调(on_next/on_error/on_completed)执行线程的操作符,它会把事件切换到指定调度器的线程中分发给订阅者,这是控制Subject事件线程的关键。
符合预期的实现方案
以下是完整示例代码,实现普通Subject在主线程运行、BehaviorSubject在后台线程运行的需求:
#include <rxcpp/rx.hpp> #include <iostream> #include <thread> #include <chrono> using namespace rxcpp; using namespace rxcpp::sources; using namespace rxcpp::operators; using namespace rxcpp::schedulers; int main() { // 普通Subject:事件发送和订阅者回调均在主线程 auto normal_subject = subjects::subject<int>(); auto normal_sub = normal_subject.get_subscriber(); // 直接订阅,默认使用当前线程(主线程)执行回调 normal_subject.get_observable().subscribe( make_subscriber<int>( [](int v) { std::cout << "[主线程] 普通Subject 接收事件: " << v << " 线程ID: " << std::this_thread::get_id() << std::endl; }, [](std::exception_ptr) {}, []() {} ) ); // 主线程发送事件 normal_sub.on_next(1); normal_sub.on_next(2); normal_sub.on_completed(); // BehaviorSubject:订阅者回调在后台线程执行 auto behavior_subject = subjects::behavior_subject<int>(0); // 初始值0 auto behavior_sub = behavior_subject.get_subscriber(); // 使用observe_on指定后台线程,让订阅者回调在新线程执行 behavior_subject.get_observable() .observe_on(new_thread()) // 核心:切换回调执行线程到后台 .subscribe( make_subscriber<int>( [](int v) { std::cout << "[后台线程] BehaviorSubject 接收事件: " << v << " 线程ID: " << std::this_thread::get_id() << std::endl; }, [](std::exception_ptr) {}, []() {} ) ); // 可选:如果需要事件发送也在后台,将on_next放到后台线程调用 std::thread behavior_thread([&behavior_sub]() { behavior_sub.on_next(10); behavior_sub.on_next(20); behavior_sub.on_completed(); }); behavior_thread.join(); // 避免主线程提前退出,等待后台线程完成 std::this_thread::sleep_for(std::chrono::milliseconds(500)); return 0; }
关键说明
- 普通Subject:因为事件发送在主线程,且未指定
observe_on,所以订阅者回调默认在主线程执行,符合需求。 - BehaviorSubject:通过
observe_on(new_thread())强制将订阅者回调切换到后台线程;如果需要事件发送逻辑也在后台,直接在后台线程调用on_next()即可。 - 注意:使用后台线程时,要确保主线程不会提前退出,可通过
join()或sleep_for()等待后台任务完成。
内容的提问来源于stack exchange,提问作者Linus Deng
相关产品推荐
相关产品推荐

