如何保证共享内存数据可用性?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
相关产品推荐
相关产品推荐

