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

Boost async_read_some读取数据出现覆盖问题求助

问题:异步读写DiscoveryEvent数据时出现数据覆盖与乱码

我尝试实现异步读写DiscoveryEvent数据,编写了发送函数和读取用的Plugin类,但运行后读取到的数据被覆盖,输出出现乱码。

发送DiscoveryEvent的代码

/*
struct DDSEntity
{
    std::string key;
    std::string participant_key;
    std::string topic_name;
    std::string topic_type;
    bool keyless;
    dds_qos_t qos;
    std::map<std::string, RouteStatus> routes;    
};

struct DiscoveryEvent
{   
    enum DiscoveryEventType
    {
        DiscoveredPublication, 
        UndiscoveredPublication, 
        DiscoveredSubscription, 
        UndiscoveredSubscription 
    };

    std::shared_ptr<DDSEntity> entity;
    DiscoveryEventType event_type;
};

*/
boost::mutex guard; //global mutex

void send_discovery_event(const dds_entity_t dp, boost::asio::local::stream_protocol::socket *socket,
    const DiscoveryEvent& event)
{
    SPDLOG_DEBUG("Send discovery event");
    boost::async([socket, event]() {
            boost::mutex::scoped_lock scoped_lock(guard);
            auto bufs = boost::asio::buffer(&event, sizeof(rodds::dds_discovery::DiscoveryEvent));
            auto size = boost::asio::write(*socket, bufs);
        });
}

异步读取的Plugin类及主函数代码

class Plugin
{
public:
    Plugin(
        const dds_entity_t &dp, boost::asio::local::stream_protocol::socket &rx)
    : _reader(&rx), _dp(dp), _buffer(&_de, sizeof(rodds::dds_discovery::DiscoveryEvent))
    {
        SPDLOG_INFO("Plugin initialized");
        
        _reader->async_read_some(
            _buffer, 
             boost::bind(
                &Plugin::async_read_handler,
                this,
                boost::asio::placeholders::error,
                boost::asio::placeholders::bytes_transferred 
            ));
    }

    void async_read_handler(const boost::system::error_code &error, std::size_t bytes_trans)
    {
        assert(!error);
        assert(bytes_trans == sizeof(rodds::dds_discovery::DiscoveryEvent));
        
      
        if (_de.event_type == rodds::dds_discovery::DiscoveryEvent::DiscoveredPublication || 
            _de.event_type == rodds::dds_discovery::DiscoveryEvent::DiscoveredSubscription)
            SPDLOG_INFO("Catch discovery event:{0}, {1}, {2}", _de.event_type, _de.entity->topic_name, _de.entity->topic_type);
        else
            SPDLOG_INFO("Catch discovery event:{0}", _de.event_type);

        _reader->async_read_some(
            _buffer, 
             boost::bind(
                &Plugin::async_read_handler,
                this,
                boost::asio::placeholders::error,
                boost::asio::placeholders::bytes_transferred 
            ));
    }

private:
    boost::asio::local::stream_protocol::socket* _reader;
    boost::asio::mutable_buffer _buffer;
    rodds::dds_discovery::DiscoveryEvent _de;
    dds_entity_t _dp;
};

int main(int argc, char* argv[])
{
    spdlog::set_level(spdlog::level::debug);
    // programm can create reader for only one dds_topic
    if (argc == 1 || argc > 2) 
    {
        SPDLOG_ERROR("Provide topic name for forwading reader process");
        return 0;
    }
    SPDLOG_INFO("Provided topic to read: {0}", argv[1]);

    // create domain_participant, reader and writer 
    // sockets to catch rodds::dds_discovery::DiscoveryEvent`s
    SPDLOG_INFO("Generate DDS domain participant");
    const dds_entity_t dp = dds_create_participant(0, NULL, NULL);
    boost::asio::io_service io_service;
    SPDLOG_INFO("Create reader/writer sockets");
    boost::asio::local::stream_protocol::socket tx(io_service), rx(io_service);
    boost::asio::local::connect_pair(tx, rx);

    boost::asio::io_service::work work(io_service); 

    // create Plugin instance
    Plugin plugin(dp, rx);

    rodds::dds_discovery::run_discovery(dp, &tx);
    
    io_service.run();

    return 0;
}

错误输出

[2022-09-11 13:23:19.226] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0, rt/rosout,|msg::dds_::Log_
[2022-09-11 13:23:19.226] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0,|rametersReply,|srv::dds_::GetParameters_Response_
[2022-09-11 13:23:19.226] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0, rameter_typesReply,srv::dds_::GetParameterTypes_Response_
[2022-09-11 13:23:19.226] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0, CrametersReply,srv::dds_::SetParameters_Response_
[2022-09-11 13:23:19.227] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0,|rameters_atomicallyReply,|srv::dds_::SetParametersAtomically_Response_
[2022-09-11 13:23:19.227] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0, |be_parametersReply,|srv::dds_::DescribeParameters_Response_
[2022-09-11 13:23:19.227] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0, |arametersReply,|srv::dds_::ListParameters_Response_
[2022-09-11 13:23:19.227] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0,  |nts,|msg::dds_::ParameterEvent_
[2022-09-11 13:23:19.228] [info] [forwading_dds_reader.cpp:61] Catch discovery event:0, rt/chatter,|ds_::String_

问题根源与修复方案

核心问题

  1. 复杂类型无法直接内存拷贝传输
    DiscoveryEvent包含std::shared_ptr<DDSEntity>,而DDSEntity内有std::string、std::map这类非POD(普通旧数据)类型。直接用boost::asio::buffer(&event, sizeof(...))做内存拷贝,只能传输结构体的固定大小部分(比如指针、枚举值),但字符串、map的实际内容存在堆上,接收端拿到的是无效指针,访问时就会出现乱码或内存错误。

  2. async_read_some无法保证完整读取
    async_read_some只保证读取至少一个字节,不能确保一次读取完整的DiscoveryEvent结构体。如果某次只读到部分字节,后续的读取会覆盖_de里的不完整数据,导致结构体解析混乱。

  3. 发送端数据生命周期风险
    send_discovery_event的lambda捕获了event的拷贝,但如果event.entity指向的DDSEntity在序列化完成前被销毁,接收端的shared_ptr会指向无效内存。

修复方案

  1. 实现序列化/反序列化逻辑
    必须把DiscoveryEvent和DDSEntity的所有成员(包括字符串、map)序列化为可传输的字节流,比如用Boost.Serialization、Protobuf,或者自定义序列化函数。示例思路:
// 序列化DDSEntity到输出流
void serialize(const DDSEntity& entity, std::ostream& os) {
    // 序列化字符串:先写长度,再写内容
    size_t len = entity.key.size();
    os.write(reinterpret_cast<const char*>(&len), sizeof(len));
    os.write(entity.key.data(), len);
    // 同理处理其他std::string成员
    os.write(reinterpret_cast<const char*>(&entity.keyless), sizeof(entity.keyless));
    // 序列化std::map:先写元素个数,再逐个序列化键值对
    len = entity.routes.size();
    os.write(reinterpret_cast<const char*>(&len), sizeof(len));
    for (const auto& pair : entity.routes) {
        // 序列化键和RouteStatus(需给RouteStatus实现序列化)
        size_t key_len = pair.first.size();
        os.write(reinterpret_cast<const char*>(&key_len), sizeof(key_len));
        os.write(pair.first.data(), key_len);
        serialize(pair.second, os);
    }
}
// 对应实现DiscoveryEvent的序列化,以及接收端的反序列化函数
  1. 用async_read替代async_read_some
    要保证读取完整的序列化数据,建议先读取数据长度,再读取对应长度的字节:
// 修改Plugin类,添加成员变量:size_t _data_len; std::vector<char> _read_buffer;
void Plugin::start_read() {
    // 先读取数据长度
    boost::asio::async_read(*_reader, boost::asio::buffer(&_data_len, sizeof(_data_len)),
        boost::bind(&Plugin::handle_read_length, this, boost::asio::placeholders::error));
}

void Plugin::handle_read_length(const boost::system::error_code& error) {
    if (!error) {
        _read_buffer.resize(_data_len);
        // 读取完整数据
        boost::asio::async_read(*_reader, boost::asio::buffer(_read_buffer),
            boost::bind(&Plugin::handle_read_data, this, boost::asio::placeholders::error));
    }
}

void Plugin::handle_read_data(const boost::system::error_code& error) {
    if (!error) {
        // 从缓冲区反序列化为DiscoveryEvent
        std::istringstream is(std::string(_read_buffer.begin(), _read_buffer.end()));
        DiscoveryEvent event;
        deserialize(event, is);
        // 处理事件...
        start_read(); // 继续读取下一个事件
    }
}
  1. 确保发送端数据生命周期
    在send_discovery_event中,确保event及其关联的DDSEntity在序列化完成前不会被销毁,比如依赖std::shared_ptr的引用计数维持数据存活,直到发送完成。

内容的提问来源于stack exchange,提问作者Splinter1984

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 16:00:26