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

Linux共享内存中std::queue跨进程操作问题及代码排查

问题
  1. 能否将以下结构体放入Linux共享内存,并在不同进程中执行队列操作?
struct Data
{
   std::queue<int> mydata;
}
  1. 除了使用数组自定义实现队列外,还有哪些替代方法?
  2. 编写的收发测试代码中,接收端始终显示队列为空,请求解决。

发送端代码(sender.cpp):

//sender.cpp
#include <iostream>
#include <sys/mman.h>
#include <fcntl.h>
#include <unistd.h>
#include <deque>
#include <condition_variable>
#include <mutex>

struct SharedData
{
    std::deque<int> queue;
};

int main()
{
    const char *sharedMemoryName = "/my_shared_memory";
    // Create shared memory
    int sharedMemoryFd = shm_open(sharedMemoryName, O_CREAT | O_RDWR, 0666);
    if (sharedMemoryFd == -1)
    {
        std::cout << "Failed to create shared memory" << std::endl;
        return 1;
    }
    else
    {
        std::cout << "create shared memory success" << std::endl;
    }
    // Set the size of shared memory
    size_t sharedMemorySize = sizeof(SharedData);
    if (ftruncate(sharedMemoryFd, sharedMemorySize) == -1)
    {
        std::cout << "Failed to set shared memory size" << std::endl;
        shm_unlink(sharedMemoryName);
        return 1;
    }
    else
    {
        std::cout << "success to set shared memory size" << std::endl;
    }
    // Map the shared memory into the process address space
    void *sharedMemoryPtr = mmap(NULL, sharedMemorySize, PROT_READ | PROT_WRITE, MAP_SHARED, sharedMemoryFd, 0);
    if (sharedMemoryPtr == MAP_FAILED)
    {
        std::cout << "Failed to map shared memory" << std::endl;
        shm_unlink(sharedMemoryName);
        return 1;
    }
    else
    {
        std::cout << "success to map shared memory" << std::endl;
    }
    // Construct the deque in shared memory
    SharedData *sharedData = new (sharedMemoryPtr) SharedData;
    // Producer process

    // Access the shared deque
    std::deque<int> &sharedQueue = sharedData->queue;
    // Produce some data
    for (int i = 0; i < 10; ++i)
    {
        // Acquire lock or use appropriate synchronization mechanism
        // Add data to the shared deque
        sharedQueue.push_back(i);
        std::cout << "Added data : " << i << std::endl;
        // Release lock or synchronization mechanism
        sleep(1);
    }
    std::cout << "Waiting..." << std::endl;

    // Wait for child processes to complete
    std::mutex mutex_;
    std::condition_variable cv;
    std::unique_lock<std::mutex> lock(mutex_);
    cv.wait(lock);
    // Cleanup
    munmap(sharedMemoryPtr, sharedMemorySize);
    shm_unlink(sharedMemoryName);
    return 0;
}

接收端代码(receiver.cpp):

#include <iostream>
#include <sys/mman.h>
#include <fcntl.h>
#include <unistd.h>
#include <deque>
#include <condition_variable>
#include <mutex>

struct SharedData
{
    std::deque<int> queue;
};

int main()
{
    const char *sharedMemoryName = "/my_shared_memory";
    // Create shared memory
    int sharedMemoryFd = shm_open(sharedMemoryName, O_RDWR, 0666);
    if (sharedMemoryFd == -1)
    {
        std::cout << "Failed to create shared memory" << std::endl;
        return 1;
    }
    else
    {
        std::cout << "success to create shared memory" << std::endl;
    }
    // Set the size of shared memory
    size_t sharedMemorySize = sizeof(SharedData);
    if (ftruncate(sharedMemoryFd, sharedMemorySize) == -1)
    {
        std::cout << "Failed to set shared memory size" << std::endl;
        shm_unlink(sharedMemoryName);
        return 1;
    }
    else
    {
        std::cout << "success to set shared memory size" << std::endl;
    }
    // Map the shared memory into the process address space
    void *sharedMemoryPtr = mmap(NULL, sharedMemorySize, PROT_READ | PROT_WRITE, MAP_SHARED, sharedMemoryFd, 0);
    if (sharedMemoryPtr == MAP_FAILED)
    {
        std::cout << "Failed to map shared memory" << std::endl;
        shm_unlink(sharedMemoryName);
        return 1;
    }
    else
    {
        std::cout << "success to map shared memory" << std::endl;
    }
    // Construct the deque in shared memory
    SharedData *sharedData = new (sharedMemoryPtr) SharedData;

    // Consumer process

    // Access the shared deque
    std::deque<int> &sharedQueue = sharedData->queue;
    // Consume data
    for (int i = 0; i < 10; ++i)
    {
        // Acquire lock or use appropriate synchronization mechanism
        // Check if the deque is not empty
        std::cout<<"Is empty : "<<sharedQueue.empty()<<std::endl;
        if (!sharedQueue.empty())
        {
            // Process the front element
            int data = sharedQueue.front();
            sharedQueue.pop_front();
            std::cout << "Data recived  : " << data << std::endl;
        }
        // Release lock or synchronization mechanism
        sleep(1);
    }

    // Wait for child processes to complete
    std::mutex mutex_;
    std::condition_variable cv;
    std::unique_lock<std::mutex> lock(mutex_);
    cv.wait(lock);

    return 0;
}

解答

一、关于std::queue放入共享内存的可行性

不能直接将包含std::queue或std::deque的结构体放入共享内存,核心原因有两点:

  1. STL容器的元素存储在进程私有堆内存中,这部分内存不属于共享内存区域,其他进程无法访问。
  2. 容器内部的指针、迭代器等成员变量是进程地址空间内的绝对地址,在其他进程中指向的是无效内存,完全不具备跨进程有效性。

二、替代方法(除数组自定义队列外)

  • 使用boost::interprocess库:该库提供了专门为进程间通信设计的容器(如boost::interprocess::queue),所有内存分配都在共享内存区域,天然支持跨进程访问。
  • 实现环形缓冲区(Ring Buffer):在共享内存中开辟固定大小的连续内存块,通过头尾指针控制读写,适合高吞吐量场景,可自行实现或使用开源成熟实现。
  • 使用POSIX消息队列:虽不属于共享内存队列,但也是进程间数据传递的常用方案,无需手动管理共享内存与同步,由系统负责消息存储和投递。
  • 自定义共享内存链表队列:在共享内存中构建链表结构,用内存偏移量替代绝对指针维护节点关系,确保跨进程地址有效性。

三、代码问题分析与修正

你的代码存在三个核心问题,直接导致接收端看不到数据:

  1. STL容器内存分配错误:std::deque的元素存在发送端私有堆中,共享内存里只有容器的空控制结构,接收端无法访问私有堆数据。
  2. 重复构造覆盖状态:接收端调用new (sharedMemoryPtr) SharedData会重新构造容器,直接覆盖发送端写入的控制结构。
  3. 缺失跨进程同步:没有进程间可用的互斥锁和条件变量,读写操作存在竞态条件,即使内存问题解决也会出现数据混乱。

修正后的代码示例(基于boost::interprocess)

发送端(sender.cpp)

#include <iostream>
#include <boost/interprocess/shared_memory_object.hpp>
#include <boost/interprocess/mapped_region.hpp>
#include <boost/interprocess/allocators/allocator.hpp>
#include <boost/interprocess/containers/deque.hpp>
#include <boost/interprocess/sync/interprocess_mutex.hpp>
#include <boost/interprocess/sync/interprocess_condition.hpp>
#include <unistd.h>

namespace bip = boost::interprocess;

// 共享内存中的数据结构,包含同步机制和共享队列
struct SharedData {
    bip::interprocess_mutex mutex;
    bip::interprocess_condition cond;
    // 使用共享内存分配器的deque
    typedef bip::allocator<int, bip::managed_shared_memory::segment_manager> ShmemAllocator;
    typedef bip::deque<int, ShmemAllocator> SharedDeque;
    SharedDeque queue;

    // 构造函数,传入共享内存分配器初始化队列
    SharedData(const ShmemAllocator& alloc) : queue(alloc) {}
};

int main() {
    const char* shm_name = "/my_shared_memory";

    // 删除残留的共享内存(如果存在)
    bip::shared_memory_object::remove(shm_name);

    // 创建共享内存段,大小足够容纳结构和数据
    bip::managed_shared_memory segment(bip::create_only, shm_name, 65536);

    // 获取共享内存分配器
    SharedData::ShmemAllocator alloc(segment.get_segment_manager());

    // 在共享内存中构造SharedData对象
    SharedData* shared_data = segment.construct<SharedData>("SharedData")(alloc);

    // 生产数据
    for (int i = 0; i < 10; ++i) {
        bip::scoped_lock<bip::interprocess_mutex> lock(shared_data->mutex);
        shared_data->queue.push_back(i);
        std::cout << "Added data: " << i << std::endl;
        // 通知消费者有新数据
        shared_data->cond.notify_one();
        lock.unlock();
        sleep(1);
    }

    // 等待消费者处理完所有数据
    bip::scoped_lock<bip::interprocess_mutex> lock(shared_data->mutex);
    while (!shared_data->queue.empty()) {
        shared_data->cond.wait(lock);
    }
    lock.unlock();

    // 清理资源
    segment.destroy<SharedData>("SharedData");
    bip::shared_memory_object::remove(shm_name);

    return 0;
}

接收端(receiver.cpp)

#include <iostream>
#include <boost/interprocess/shared_memory_object.hpp>
#include <boost/interprocess/mapped_region.hpp>
#include <boost/interprocess/allocators/allocator.hpp>
#include <boost/interprocess/containers/deque.hpp>
#include <boost/interprocess/sync/interprocess_mutex.hpp>
#include <boost/interprocess/sync/interprocess_condition.hpp>
#include <unistd.h>

namespace bip = boost::interprocess;

struct SharedData {
    bip::interprocess_mutex mutex;
    bip::interprocess_condition cond;
    typedef bip::allocator<int, bip::managed_shared_memory::segment_manager> ShmemAllocator;
    typedef bip::deque<int, ShmemAllocator> SharedDeque;
    SharedDeque queue;

    SharedData(const ShmemAllocator& alloc) : queue(alloc) {}
};

int main() {
    const char* shm_name = "/my_shared_memory";

    // 打开已存在的共享内存段
    bip::managed_shared_memory segment(bip::open_only, shm_name);

    // 查找共享内存中的SharedData对象
    SharedData* shared_data = segment.find<SharedData>("SharedData").first;

    // 消费数据
    int count = 0;
    while (count < 10) {
        bip::scoped_lock<bip::interprocess_mutex> lock(shared_data->mutex);
        // 等待队列非空
        while (shared_data->queue.empty()) {
            shared_data->cond.wait(lock);
        }
        // 取出并处理数据
        int data = shared_data->queue.front();
        shared_data->queue.pop_front();
        std::cout << "Data received: " << data << std::endl;
        count++;
        // 通知生产者数据已处理
        shared_data->cond.notify_one();
        lock.unlock();
        sleep(1);
    }

    return 0;
}

编译与运行

编译时需要链接boost库:

g++ sender.cpp -o sender -lboost_interprocess
g++ receiver.cpp -o receiver -lboost_interprocess

先启动发送端,再启动接收端,即可看到数据正常传递。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 07:32:04