RxCPP如何订阅被持续修改的STL容器生成的observable实现实时回调?
问题原因
- 你调用
rxcpp::observable<>::iterate(vec)生成的是冷Observable,它只会在创建时遍历容器当前的快照内容,你创建时vec是空的,自然不会触发任何onNext回调,后续对vec的修改它完全感知不到。 - 原生
std::vector没有内置变更通知机制,直接在其他线程修改容器本身不会主动对外广播变更事件。
正确实现思路
不要直接监听std::vector的变更,而是用RxCpp的subject作为事件中转:每次你往容器中推入新元素时,同时向subject发射对应元素的事件,所有订阅了该subject的观察者都会收到onNext回调。另外多线程操作容器需要加锁避免数据竞争。
完整可运行代码
#include <rxcpp/rx.hpp> #include <vector> #include <thread> #include <mutex> #include <cstdio> #include <chrono> #include <iostream> std::mutex vec_mtx; void updateVec(std::vector<int>& v, rxcpp::subjects::subject<int>& sub) { // 模拟持续推入元素 for (int i = 1; i <= 10; ++i) { std::lock_guard<std::mutex> lock(vec_mtx); v.push_back(i); sub.get_subscriber().on_next(i); // 推入元素的同时发射事件 std::this_thread::sleep_for(std::chrono::milliseconds(500)); // 模拟间隔推送 } sub.get_subscriber().on_completed(); // 所有元素推送完成后触发结束 } int main() { std::vector<int> vec{}; rxcpp::subjects::subject<int> sub; // 订阅事件流 sub.get_observable().subscribe( [](int v) { std::printf("OnNext-> value: %d \n", v); }, []() { std::cout << "OnCompleted" << std::endl; } ); std::thread t1(updateVec, std::ref(vec), std::ref(sub)); t1.join(); return 0; }
如果你的业务不需要保留所有历史元素,也可以省略std::vector的维护,直接通过subject发射事件即可。如果需要后续读取容器内的所有元素,读取时也需要加同一把锁保证线程安全。
内容的提问来源于stack exchange,提问作者somewhat_clueless_developer
相关产品推荐
相关产品推荐

