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_
问题根源与修复方案
核心问题
复杂类型无法直接内存拷贝传输
DiscoveryEvent包含std::shared_ptr<DDSEntity>,而DDSEntity内有std::string、std::map这类非POD(普通旧数据)类型。直接用boost::asio::buffer(&event, sizeof(...))做内存拷贝,只能传输结构体的固定大小部分(比如指针、枚举值),但字符串、map的实际内容存在堆上,接收端拿到的是无效指针,访问时就会出现乱码或内存错误。async_read_some无法保证完整读取
async_read_some只保证读取至少一个字节,不能确保一次读取完整的DiscoveryEvent结构体。如果某次只读到部分字节,后续的读取会覆盖_de里的不完整数据,导致结构体解析混乱。发送端数据生命周期风险
send_discovery_event的lambda捕获了event的拷贝,但如果event.entity指向的DDSEntity在序列化完成前被销毁,接收端的shared_ptr会指向无效内存。
修复方案
- 实现序列化/反序列化逻辑
必须把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的序列化,以及接收端的反序列化函数
- 用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(); // 继续读取下一个事件 } }
- 确保发送端数据生命周期
在send_discovery_event中,确保event及其关联的DDSEntity在序列化完成前不会被销毁,比如依赖std::shared_ptr的引用计数维持数据存活,直到发送完成。
内容的提问来源于stack exchange,提问作者Splinter1984
相关产品推荐
相关产品推荐

