多线程数据丢失问题:主线程偶读空数据(数据实际非空)
线程间数据共享库的同步异常排查与修复
核心代码
bool Broker::publish(const std::string& topic, void* data, size_t len) { int queue_pos = getQueueHead(); if (queue_pos == -1) { std::cout << "WARNING: broker buffer is full\n"; return false; } std::vector<std::string> parsed_topic = parseTopic(topic); void* data_ptr = malloc(len); memcpy(data_ptr, data, len); std::shared_ptr<void> data_share(data_ptr, free); _buffer[queue_pos] = Buffer{parsed_topic, data_share, len}; // <- 写入数据(_buffer是普通数组) sem_post(&_work[queue_pos]); // <- 通知主线程可以继续 return true; } void Broker::run() { while (_should_run) { // 队列空则休眠 if (queueSize() == 0) { sem_wait(&_work[_queue_tail]); // <- 等待信号量 continue; } // 从缓冲区取数据 const auto [topic, data, len] = std::move(_buffer[_queue_tail]); // <- 读取数据 // 释放队列空间 _queue_tail = (_queue_tail + 1) % BROKER_QUEUE_SIZE; // 后续代码... }
问题描述
主线程通过sem_wait等待信号量后,偶尔会读取到空的Buffer结构,但实际数据已经被写入。如果先打印空数据,休眠10微秒后重新读取,第二次就能获取到有效数据。
尝试过在sem_post前添加asm("" ::: "memory");内存屏障,也试过_mm_mfence,均无效。GDB查看汇编确认是先写入数组再调用信号量。
问题根源
- 队列状态的内存可见性问题:
queueSize()的实现未考虑线程间同步,主线程得到“队列空”的判断时,publish线程的写入操作可能还未对主线程可见,导致后续sem_wait唤醒后读取到旧数据。 - 队列指针的非原子操作:
_queue_head和_queue_tail如果是普通变量,其读写不具备原子性和内存可见性,可能出现指针更新与数据写入的顺序错乱。 - 信号量与内存操作的同步缺失:虽然汇编上是先写数据再调用
sem_post,但CPU的乱序执行或缓存一致性问题,可能导致主线程看到sem_post的信号却看不到之前的数据写入。
修复方案
1. 原子化队列指针
将_queue_head和_queue_tail改为std::atomic<size_t>类型,确保指针的读写操作是原子的,且跨线程可见:
std::atomic<size_t> _queue_head = 0; std::atomic<size_t> _queue_tail = 0;
2. 修复queueSize()的线程安全
基于原子指针计算队列长度,使用内存屏障保证读取的指针是最新值:
size_t Broker::queueSize() const { auto head = _queue_head.load(std::memory_order_acquire); auto tail = _queue_tail.load(std::memory_order_acquire); return head >= tail ? head - tail : BROKER_QUEUE_SIZE - (tail - head); }
3. 强化内存屏障与信号量的同步
在publish线程写入_buffer后,添加释放屏障,确保所有写入操作在sem_post前完成并对其他线程可见:
_buffer[queue_pos] = Buffer{parsed_topic, data_share, len}; // 释放屏障:确保之前的内存写入对其他线程可见 std::atomic_thread_fence(std::memory_order_release); sem_post(&_work[queue_pos]);
4. 主线程读取时的同步调整
主线程被信号量唤醒后,添加获取屏障,确保能读取到publish线程的最新写入:
void Broker::run() { while (_should_run) { auto current_tail = _queue_tail.load(std::memory_order_acquire); sem_wait(&_work[current_tail]); // 获取屏障:确保后续读取_buffer能看到最新数据 std::atomic_thread_fence(std::memory_order_acquire); const auto [topic, data, len] = std::move(_buffer[current_tail]); _queue_tail.store((current_tail + 1) % BROKER_QUEUE_SIZE, std::memory_order_release); // 后续代码... }
5. 移除冗余的队列空判断
原代码中if (queueSize() == 0)存在竞态条件,直接依赖sem_wait等待信号量即可,信号量本身已保证有数据可读取。
内容的提问来源于stack exchange,提问作者Ishay Trattner
相关产品推荐
相关产品推荐

