使用Boost.Interprocess发送序列化结构体的正确方式及重连问题解决
Boost.Interprocess消息队列重启接收端后无法恢复通信的解决方案
问题背景
使用Boost.Interprocess消息队列发送序列化后的自定义结构体二进制数据,代码编译运行正常,但关闭并重启接收端后,数据无法正常收发,需要实现接收端重启后自动恢复通信的能力。
问题分析
- 反序列化包含无效数据:接收端中
serialized_string.resize(MAX_SIZE)后将整个缓冲区内容写入流,但mq.receive返回的recvd_size才是实际消息长度,多余的缓冲区垃圾数据会破坏二进制归档结构,导致反序列化失败。 - 旧消息残留:接收端异常关闭时,队列中未读取的旧消息会残留,重启后优先读取这些旧消息,若旧消息无法正常反序列化,会导致接收端卡住或无法处理新消息。
- 异常处理缺失:接收循环未处理反序列化异常,单个坏消息会导致整个接收端崩溃。
解决方案与修正代码
1. 结构体定义文件(datastruct.hpp)
添加序列化版本号,避免结构体版本变化导致的兼容问题:
#ifndef DATASTRUCT_HPP #define DATASTRUCT_HPP #include <string> #include <vector> #include <boost/serialization/vector.hpp> #include <boost/serialization/version.hpp> #define MAX_SIZE 150000 namespace dataStruct { struct VUserPoint { float PositionX; float PositionY; float PositionZ; template <typename Archive> void serialize(Archive& ar, unsigned int const version) { ar & PositionX; ar & PositionY; ar & PositionZ; } }; struct FvpData { std::vector<VUserPoint> UserPoints; int frameNumber; template <typename Archive> void serialize(Archive& ar, unsigned int const version) { ar & UserPoints; ar & frameNumber; } }; } // namespace dataStruct // 为结构体添加版本号,确保序列化兼容性 BOOST_CLASS_VERSION(dataStruct::VUserPoint, 1) BOOST_CLASS_VERSION(dataStruct::FvpData, 1) #endif // DATASTRUCT_HPP
2. 发送端代码(sender.cpp)
优化消息队列创建逻辑,避免每次发送都重新初始化队列:
#include "DataSender.h" #include <boost/archive/binary_oarchive.hpp> #include <boost/interprocess/ipc/message_queue.hpp> #include <sstream> using namespace boost::interprocess; int num = 1; // 全局消息队列实例,避免重复创建开销 static message_queue mq(open_or_create, "mq", 100, MAX_SIZE); void DataSender::Grab(std::vector<Eigen::Vector3d> points) { dataStruct::FvpData data; data.UserPoints.reserve(points.size()); // 预分配内存提升效率 for (auto& p : points) { dataStruct::VUserPoint pnt; pnt.PositionX = static_cast<float>(p.x()); pnt.PositionY = static_cast<float>(p.y()); pnt.PositionZ = static_cast<float>(p.z()); data.UserPoints.push_back(pnt); } data.frameNumber = ++num; try { std::stringstream oss; boost::archive::binary_oarchive oa(oss); oa << data; std::string serialized_string = oss.str(); mq.send(serialized_string.data(), serialized_string.size(), 1); std::cout << "发送帧号: " << data.frameNumber << std::endl; } catch (interprocess_exception& ex) { std::cerr << "发送错误: " << ex.what() << std::endl; } }
3. 接收端代码(receiver.cpp)
修复反序列化逻辑,添加旧消息清理,完善异常处理:
#include <iostream> #include <string> #include <sstream> #include <boost/archive/binary_iarchive.hpp> #include <boost/interprocess/ipc/message_queue.hpp> #include <boost/archive/archive_exception.hpp> #include "DataStruct.hpp" using namespace boost::interprocess; int main() { try { message_queue mq(open_or_create, "mq", 100, MAX_SIZE); // 启动时清理队列中残留的旧消息(适合不需要保留旧消息的场景) message_queue::size_type recvd_size; unsigned int priority; char temp_buf[MAX_SIZE]; while (!mq.empty()) { mq.receive(temp_buf, MAX_SIZE, recvd_size, priority); std::cout << "清理残留旧消息,长度: " << recvd_size << std::endl; } std::cout << "接收端启动,开始监听消息..." << std::endl; while (true) { dataStruct::FvpData me; std::string serialized_string; serialized_string.resize(MAX_SIZE); mq.receive(&serialized_string[0], MAX_SIZE, recvd_size, priority); // 只保留实际接收到的有效数据 serialized_string.resize(recvd_size); try { std::stringstream iss(serialized_string); boost::archive::binary_iarchive ia(iss); ia >> me; std::cout << "收到帧号: " << me.frameNumber << ",点数量: " << me.UserPoints.size() << std::endl; } catch (boost::archive::archive_exception& ex) { std::cerr << "反序列化错误: " << ex.what() << ",跳过该消息" << std::endl; continue; // 跳过坏消息,继续接收后续消息 } } } catch (interprocess_exception& ex) { std::cerr << "消息队列错误: " << ex.what() << std::endl; return 1; } // 程序正常退出时可选择清理消息队列(需确保所有使用进程已退出) // message_queue::remove("mq"); return 0; }
补充说明
- 消息队列生命周期:若需要在所有进程退出后清理队列,可在发送端或接收端的退出逻辑中调用
message_queue::remove("mq"),但需确保所有关联进程已终止。 - 旧消息处理可选方案:若需要保留部分旧消息,可在结构体中添加同步标识字段,接收端重启后读取消息直到收到最新的同步标识,再开始处理有效数据。
- 异常容错:接收循环中捕获反序列化异常,确保单个损坏消息不会导致接收端停止工作。
内容的提问来源于stack exchange,提问作者anti
相关产品推荐
相关产品推荐

