基于OneTBB库:带更新策略的joiner/indexer或composite节点实现问询
实现带自定义更新策略的Joiner节点
完全可以实现符合你需求的节点,结合joiner的多端口聚合能力和overwrite_node的缓存覆盖特性,针对你的shared_ptr类型(用nullptr作为未就绪标记),以下是具体实现方案:
核心逻辑梳理
节点需要维护每个输入端口的最新值缓存(初始为nullptr),严格匹配你的三个规则:
- 只要至少一个端口的缓存非
nullptr,就具备输出能力; - 任何端口收到新的有效值(非
nullptr)时,立即更新对应端口的缓存,并触发输出当前所有端口的缓存状态; - 未收到新值的端口,缓存保留之前的有效值(或
nullptr),不会被清空。
通用模板实现(C++)
以下是双端口场景的通用模板类,可直接适配你的shared_ptr类型:
#include <memory> #include <functional> #include <utility> template<typename T1, typename T2> class UpdateJoiner { public: // 输出类型为两个端口的shared_ptr组合 using OutputPair = std::pair<std::shared_ptr<T1>, std::shared_ptr<T2>>; // 输出回调函数类型 using OutputCallback = std::function<void(const OutputPair&)>; explicit UpdateJoiner(OutputCallback callback) : output_callback_(std::move(callback)) {} // 端口1的输入处理函数 void process_input1(std::shared_ptr<T1> new_value) { if (new_value) { // 仅处理有效非nullptr值 cache1_ = std::move(new_value); trigger_output(); } } // 端口2的输入处理函数 void process_input2(std::shared_ptr<T2> new_value) { if (new_value) { cache2_ = std::move(new_value); trigger_output(); } } private: // 触发输出逻辑:只要至少一个缓存非nullptr就输出 void trigger_output() { if (cache1_ || cache2_) { output_callback_(std::make_pair(cache1_, cache2_)); } } std::shared_ptr<T1> cache1_; // 端口1的缓存,初始为nullptr std::shared_ptr<T2> cache2_; // 端口2的缓存,初始为nullptr OutputCallback output_callback_; };
匹配你的示例场景验证
用你给出的测试流程验证:
- 端口1传入
x1:cache1_更新为x1,输出{x1, nullptr} - 端口1传入
x2:cache1_覆盖为x2,输出{x2, nullptr} - 端口2传入
y1:cache2_更新为y1,输出{x2, y1} - 端口1传入
x3:cache1_覆盖为x3,输出{x3, y1}
完全符合你的预期行为。
框架适配示例(以ROS2为例)
如果是在ROS2环境中使用,只需将上述模板类与ROS2的订阅/发布机制结合:
#include "rclcpp/rclcpp.hpp" #include "std_msgs/msg/string.hpp" #include "std_msgs/msg/int32.hpp" // 自定义复合消息类型(需提前定义) #include "my_msgs/msg/composite_data.hpp" class UpdateJoinerNode : public rclcpp::Node { public: UpdateJoinerNode() : Node("update_joiner_node") { // 初始化UpdateJoiner,绑定输出回调 joiner_ = std::make_unique<UpdateJoiner<std::string, int32_t>>( [this](const auto& output) { auto msg = std::make_unique<my_msgs::msg::CompositeData>(); if (output.first) msg->data1 = *output.first; if (output.second) msg->data2 = *output.second; publisher_->publish(std::move(msg)); } ); // 创建端口1的订阅器 sub1_ = this->create_subscription<std_msgs::msg::String>( "topic1", 10, [this](std::shared_ptr<std_msgs::msg::String> msg) { joiner_->process_input1(std::make_shared<std::string>(msg->data)); } ); // 创建端口2的订阅器 sub2_ = this->create_subscription<std_msgs::msg::Int32>( "topic2", 10, [this](std::shared_ptr<std_msgs::msg::Int32> msg) { joiner_->process_input2(std::make_shared<int32_t>(msg->data)); } ); // 创建输出发布器 publisher_ = this->create_publisher<my_msgs::msg::CompositeData>("output_topic", 10); } private: std::unique_ptr<UpdateJoiner<std::string, int32_t>> joiner_; rclcpp::Subscription<std_msgs::msg::String>::SharedPtr sub1_; rclcpp::Subscription<std_msgs::msg::Int32>::SharedPtr sub2_; rclcpp::Publisher<my_msgs::msg::CompositeData>::SharedPtr publisher_; }; int main(int argc, char* argv[]) { rclcpp::init(argc, argv); rclcpp::spin(std::make_shared<UpdateJoinerNode>()); rclcpp::shutdown(); return 0; }
扩展到多端口场景
如果需要支持更多端口,可以将模板改为可变参数类型,用std::tuple存储缓存和输出,核心逻辑保持一致:每次有端口收到新值,更新对应缓存并触发输出当前所有缓存状态。
内容的提问来源于stack exchange,提问作者SpeakX
相关产品推荐
相关产品推荐

