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

多线程数据丢失问题:主线程偶读空数据(数据实际非空)

线程间数据共享库的同步异常排查与修复

核心代码

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查看汇编确认是先写入数组再调用信号量。

问题根源

  1. 队列状态的内存可见性问题:queueSize()的实现未考虑线程间同步,主线程得到“队列空”的判断时,publish线程的写入操作可能还未对主线程可见,导致后续sem_wait唤醒后读取到旧数据。
  2. 队列指针的非原子操作:_queue_head和_queue_tail如果是普通变量,其读写不具备原子性和内存可见性,可能出现指针更新与数据写入的顺序错乱。
  3. 信号量与内存操作的同步缺失:虽然汇编上是先写数据再调用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 16:48:09