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

如何保证共享内存数据可用性?mmap跨进程共享+semaphore通知场景验证

mmap共享文件+条件变量通知的数据一致性疑问

我通过mmap映射文件实现进程间数据共享,在完成所有文件数据写入后,立即通过共享内存中的条件变量发送就绪通知。现在有个疑问:读进程接收通知后,是否可能无法读取到部分文件的最新数据?

我编写了模拟代码(如下),测试中数据始终可用,但不确定是否因测试时长不足导致(同事反馈类似场景偶现数据未更新问题)。请问是否有理论依据或官方文档可验证该代码的正确性?

#include <boost/interprocess/sync/interprocess_mutex.hpp>
#include <boost/interprocess/sync/interprocess_condition.hpp>
#include <boost/interprocess/shared_memory_object.hpp>
#include <boost/interprocess/file_mapping.hpp>
#include <boost/interprocess/mapped_region.hpp>
#include <boost/interprocess/sync/scoped_lock.hpp>
#include <boost/date_time/posix_time/posix_time.hpp>
#include <optional>
#include <fstream>
#include <ios>
#include <thread>
// data to notify if update is done
template <class T> struct SharedData {
    SharedData(const T &init_value) : data_(init_value) {}
    // Mutex to protect access to the data
    boost::interprocess::interprocess_mutex mutex_;
    // Condition to wait before the data is set
    boost::interprocess::interprocess_condition cond_;
    T data_;
};

template <class T> class SharedDataWriter {
public:
    SharedDataWriter(const std::string& name, const T &init_value)
    : shm_(boost::interprocess::open_or_create, name.data(), boost::interprocess::read_write)
    , name_(name) {
        shm_.truncate(sizeof(SharedData<T>));
        region_.emplace(shm_, boost::interprocess::read_write);
        void *addr = region_->get_address();
        shared_data_ = new (addr) SharedData<T>(init_value);
    }
    ~SharedDataWriter() {
        boost::interprocess::shared_memory_object::remove(name_.data());
    }

    void set_value(const T& value) {
        boost::interprocess::scoped_lock<boost::interprocess::interprocess_mutex> lock(shared_data_->mutex_);
        shared_data_->data_ = value;
        shared_data_->cond_.notify_all();
    }
private:
    boost::interprocess::shared_memory_object shm_;
    std::optional<boost::interprocess::mapped_region> region_;
    std::string name_;
    SharedData<T>* shared_data_ = nullptr;
};

template <class T> class SharedDataReader {
public:
    SharedDataReader(const std::string& name)
    : shm_(boost::interprocess::open_only, name.data(), boost::interprocess::read_write)
    , region_(shm_, boost::interprocess::read_write) {
        void * addr = region_.get_address();
        shared_data_ = (SharedData<T> *)addr;
    }

    T wait_value() {
        boost::interprocess::scoped_lock<boost::interprocess::interprocess_mutex> lock(
            shared_data_->mutex_);
        shared_data_->cond_.wait(lock);
        return shared_data_->data_;
    }
private:
    boost::interprocess::shared_memory_object shm_;
    boost::interprocess::mapped_region region_;
    SharedData<T> *shared_data_ = nullptr;
};
// file content is basically a table
struct FixedHeader {
    int32_t col_size_;
    int32_t row_size_;

    void * get_value_start_addr(char * begin) {
        return begin + sizeof(FixedHeader);
    }
};

struct SingleMappedFile {
    SingleMappedFile(const std::string& filename)
    : mapped_file_(filename.data(), boost::interprocess::read_write)
    , region_(mapped_file_, boost::interprocess::read_write) {}
    FixedHeader& get_header() {
        return *((FixedHeader*)get_start_addr());
    }
    void write(int32_t row_index, int32_t value) {
        int32_t * value_start = (int32_t *)get_value_start_addr();
        for (int32_t i = 0; i < get_header().col_size_; ++i) {
            value_start[row_index * get_header().col_size_ + i] = value;
        }
    }
    int32_t* get_row(int32_t row_index) {
        int32_t * value_start = (int32_t *)get_value_start_addr();
        return value_start + row_index * get_header().col_size_;
    }
    void * get_start_addr() const { return region_.get_address(); }
    const char * get_file_name() const { return mapped_file_.get_name(); }
private:
    char * get_value_start_addr() {
        return (char*)get_header().get_value_start_addr((char*)get_start_addr());
    }
    boost::interprocess::file_mapping mapped_file_;
    boost::interprocess::mapped_region region_;
};

void writer_main() {
    // create 50 files and set init value
    for (size_t i = 0; i < 50; ++i) {
        std::ofstream file(std::string("file_") + std::to_string(i), std::ios_base::out | std::ios_base::binary);
        int32_t col = 5000;
        int32_t row = 300;
        file.write((char*)&col, sizeof(col));
        file.write((char*)&row, sizeof(row));
        for (int32_t k = 0; k < col * row; ++k) file.write((char*)&k, sizeof(k));
    }

    // mmap all these 50 files
    std::vector<SingleMappedFile> mapped_files;
    for (size_t i = 0; i < 50; ++i) {
        mapped_files.emplace_back(std::string("file_") + std::to_string(i));
    }

    SharedDataWriter<int32_t> writer("notify", -1);

    int32_t value = 1;
    int32_t row_index = 0;
    while(true) {
        char c;
        std::cout << "input anything to write data." << std::endl;
        std::cin >> c;
        // write data into mapped files and notify
        for (auto& mfile : mapped_files) {
            mfile.write(row_index, value);
        }
        writer.set_value(row_index);
        row_index++;
    }
}

void thread_func(SharedDataReader<int32_t>& reader, SingleMappedFile& mfile) {
    // wait for notification and check if data is updated
    while (true) {
        int32_t row_index = reader.wait_value();
        int32_t* row = mfile.get_row(row_index);
        for (int32_t i = mfile.get_header().col_size_; i > 0; --i) {
            if (row[i - 1] != 1) {
                throw std::runtime_error("data not updated.");
            }
        }
        std::cout << "row_index = " << row_index << std::endl;
    }
}

void reader_main() {
    std::vector<SingleMappedFile> mapped_files;
    for (size_t i = 0; i < 50; ++i) {
        mapped_files.emplace_back(std::string("file_") + std::to_string(i));
    }
    SharedDataReader<int32_t> reader("notify");
    std::vector<std::thread> threads;
    // create 50 threads, each check one file.
    for (auto& mfile : mapped_files) {
        threads.emplace_back([&reader, &mfile]() { thread_func(reader, mfile); });
    }
    for (auto& t : threads) t.join();
}

int main(int argc, char **) {
    if (argc > 1) {
        reader_main();
    } else {
        writer_main();
    }
}

核心结论:你的代码存在数据不一致风险,偶现问题符合系统特性

1. 问题根源:异步刷盘与内存可见性缺失

  • mmap写入的是用户态缓存页,操作系统会异步将脏页刷入磁盘,这个过程没有固定时序。读进程的mmap映射页可能仍指向旧的磁盘数据,或者写进程的脏页尚未完成刷盘,导致读进程读取到旧数据。
  • 你使用的条件变量仅保证了通知信号的顺序,但没有建立写文件操作到读文件操作的内存可见性屏障,也未强制写进程的脏页同步到磁盘——条件变量的内存语义只覆盖其保护的shared_data_->data_,不延伸到mmap映射的文件内存区域。

2. 测试未复现的原因

测试场景数据量小、系统负载低,操作系统的脏页自动刷盘策略(如空闲时触发)刚好让数据在通知后完成同步,但在高负载、大数据量场景下,异步刷盘的延迟会暴露数据不一致问题。

3. 修复方案:强制同步+内存屏障

要确保数据一致性,需在写进程发送通知前执行以下操作:

  • 对每个SingleMappedFile,写完数据后调用msync(region_.get_address(), region_.get_size(), MS_SYNC),强制将脏页同步到磁盘,确保读进程的映射能获取最新数据。
  • 配合内存屏障:写进程在调用notify_all前添加std::atomic_thread_fence(std::memory_order_release),读进程在读取mmap数据前添加std::atomic_thread_fence(std::memory_order_acquire),确保写操作的内存可见性传递到读进程。

4. 理论依据

  • POSIX标准明确:mmap写入不保证立即同步到磁盘,msync是强制同步的唯一标准手段。
  • C++内存模型规定:条件变量的notify_all具备release语义,wait返回时具备acquire语义,但该语义仅作用于受互斥锁保护的变量,无法覆盖独立的mmap内存区域的操作顺序与可见性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 01:48:14