You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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. 端口1传入x1:cache1_更新为x1,输出{x1, nullptr}
  2. 端口1传入x2:cache1_覆盖为x2,输出{x2, nullptr}
  3. 端口2传入y1:cache2_更新为y1,输出{x2, y1}
  4. 端口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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 05:17:52