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

使用Boost.Interprocess发送序列化结构体的正确方式及重连问题解决

Boost.Interprocess消息队列重启接收端后无法恢复通信的解决方案

问题背景

使用Boost.Interprocess消息队列发送序列化后的自定义结构体二进制数据,代码编译运行正常,但关闭并重启接收端后,数据无法正常收发,需要实现接收端重启后自动恢复通信的能力。

问题分析

  1. 反序列化包含无效数据:接收端中serialized_string.resize(MAX_SIZE)后将整个缓冲区内容写入流,但mq.receive返回的recvd_size才是实际消息长度,多余的缓冲区垃圾数据会破坏二进制归档结构,导致反序列化失败。
  2. 旧消息残留:接收端异常关闭时,队列中未读取的旧消息会残留,重启后优先读取这些旧消息,若旧消息无法正常反序列化,会导致接收端卡住或无法处理新消息。
  3. 异常处理缺失:接收循环未处理反序列化异常,单个坏消息会导致整个接收端崩溃。

解决方案与修正代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 11:30:55