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

在单生产者单消费者模型中实现子循环缓冲区的无锁定位

在单生产者单消费者模型中实现子循环缓冲区的无锁定位

我完全理解你现在的困境——在单生产者单消费者的异步缓冲阅读器里,想实现已缓冲数据内的无锁快速定位,但卡在了m_readLeft的安全更新上,还担心去掉m_readLeft后没法维护m_count。咱们一步步拆解这个问题,利用单生产者单消费者模型的特性来简化逻辑。

首先明确核心优势:因为只有两个职责明确的线程(生产者只写缓冲区、处理跨缓冲的seek;消费者只读/定位),我们不需要复杂的锁,用原子变量+合适的内存顺序就能保证线程安全,甚至消费者独用的变量连原子都不用。

问题分析与优化方向

  1. 缺失缓冲区与文件位置的映射:你没法快速判断seek目标是否在已缓冲范围内,需要新增变量关联缓冲区位置和文件偏移。
  2. 冗余的m_count变量:已缓冲字节数完全可以通过m_head和m_tail的差值推导出来,不需要单独维护,这直接解决了去掉m_readLeft后的m_count维护问题。
  3. m_readLeft的线程安全误解:因为只有消费者线程(主线程)会读写m_readLeft,它不需要声明为原子变量,普通变量直接修改就安全。

具体解决方案

第一步:调整核心变量

修改类内的成员变量,适配无锁逻辑:

// 将m_head、m_tail改为原子变量,保证跨线程可见性
std::atomic<size_t> m_head = 0;
std::atomic<size_t> m_tail = 0;

// 新增:记录缓冲区头部对应的文件偏移,建立缓冲区与文件位置的映射
std::atomic<qint64> m_bufferStartPos = 0;

// m_readPos需要被生产者读取,所以保留原子性;m_readLeft改为普通变量(消费者独用)
std::atomic<size_t> m_readPos = 0;
size_t m_readLeft = 0;

// 移除单独的m_count变量,后续用m_head和m_tail推导已缓冲大小

第二步:实现无锁Fast Path Seek

在seek函数中先判断目标位置是否在已缓冲范围内,是的话直接调整消费者读取状态,不需要锁:

bool AsyncBufferedReader::seek(qint64 pos)
{
    if (pos < 0 || pos > size()) {
        return false;
    }

    // Fast Path:目标位置在已缓冲数据内,无锁处理
    qint64 bufferedStart = m_bufferStartPos.load(std::memory_order_acquire);
    size_t bufferedSize = (m_tail.load(std::memory_order_acquire) - m_head.load(std::memory_order_acquire) + m_capacity) % m_capacity;
    
    if (pos >= bufferedStart && pos < bufferedStart + bufferedSize) {
        // 计算目标位置在缓冲区中的偏移
        size_t offsetInBuffer = static_cast<size_t>(pos - bufferedStart);
        size_t newReadPos = (m_head.load(std::memory_order_acquire) + offsetInBuffer) % m_capacity;

        // 更新消费者读取状态(单线程操作,完全安全)
        m_readPos.store(newReadPos, std::memory_order_release);
        m_readLeft = bufferedSize - offsetInBuffer;

        // 更新QIODevice的当前位置
        QIODevice::seek(pos);
        return true;
    }

    // Slow Path:目标位置不在已缓冲范围内,需要通知生产者处理
    QMutexLocker locker(&m_mutex);
    if (m_aborted || !m_workerRunning) {
        return false;
    }

    m_seekRequested = true;
    m_seekPos = pos;
    m_seekSuccess = false;

    // 唤醒生产者处理seek请求
    m_bufferSpaceWait.wakeOne();
    // 等待seek完成
    while (m_seekRequested && !m_aborted) {
        m_seekFinishedWait.wait(&m_mutex);
    }

    if (m_aborted) {
        return false;
    }

    bool success = m_seekSuccess;
    if (success) {
        QIODevice::seek(pos);
        // 重置消费者读取状态
        m_readPos.store(m_head.load(std::memory_order_acquire), std::memory_order_release);
        m_readLeft = (m_tail.load(std::memory_order_acquire) - m_head.load(std::memory_order_acquire) + m_capacity) % m_capacity;
    }

    m_seekRequested = false;
    return success;
}

第三步:修改makeSpaceForMoreReading函数

现在不需要维护m_count,直接通过m_head和m_readPos计算已消费的字节数:

bool AsyncBufferedReader::makeSpaceForMoreReading()
{
    size_t currentReadPos = m_readPos.load(std::memory_order_acquire);
    size_t currentHead = m_head.load(std::memory_order_relaxed);
    size_t bytesConsumed = 0;

    // 计算消费者已消费的字节数(处理循环缓冲区的两种情况)
    if (currentHead <= currentReadPos) {
        bytesConsumed = currentReadPos - currentHead;
    } else {
        bytesConsumed = m_capacity - currentHead + currentReadPos;
    }

    if (bytesConsumed == 0) {
        return false; // 没有可释放的空间
    }

    // 更新缓冲区头部和对应的文件偏移
    m_head.store(currentReadPos, std::memory_order_release);
    m_bufferStartPos.fetch_add(static_cast<qint64>(bytesConsumed), std::memory_order_release);

    return true;
}

第四步:调整readData函数

用推导的bufferedSize代替原m_count,直接修改m_readLeft(单线程操作安全):

qint64 AsyncBufferedReader::readData(char *data, qint64 maxlen)
{
    if (maxlen <= 0) return 0;

    size_t head = m_head.load(std::memory_order_acquire);
    size_t tail = m_tail.load(std::memory_order_acquire);
    size_t bufferedSize = (tail - head + m_capacity) % m_capacity;

    // 无数据时等待生产者
    if (bufferedSize == 0) {
        QMutexLocker locker(&m_mutex);
        while (bufferedSize == 0 && !m_aborted && !m_sourceEof) {
            m_dataWait.wait(&m_mutex);
            tail = m_tail.load(std::memory_order_acquire);
            bufferedSize = (tail - head + m_capacity) % m_capacity;
        }
        if (m_aborted || (bufferedSize == 0 && m_sourceEof)) {
            return -1;
        }
    }

    // 确定实际可读字节数
    qint64 toRead = std::min<qint64>(maxlen, std::min<size_t>(m_readLeft, bufferedSize));
    if (toRead <= 0) return 0;

    size_t readPos = m_readPos.load(std::memory_order_relaxed);
    size_t bytesRead = 0;

    // 处理循环缓冲区的两种读取情况
    if (readPos + toRead <= m_capacity) {
        std::memcpy(data, &m_buffer[readPos], toRead);
        bytesRead = toRead;
        m_readPos.store(readPos + bytesRead, std::memory_order_release);
    } else {
        size_t firstChunk = m_capacity - readPos;
        std::memcpy(data, &m_buffer[readPos], firstChunk);
        std::memcpy(data + firstChunk, &m_buffer[0], toRead - firstChunk);
        bytesRead = toRead;
        m_readPos.store(toRead - firstChunk, std::memory_order_release);
    }

    // 更新消费者剩余可读字节数
    m_readLeft -= bytesRead;
    // 更新设备当前位置
    QIODevice::seek(pos() + bytesRead);

    // 通知生产者有新空间可用
    m_bufferSpaceWait.wakeOne();

    return bytesRead;
}

关键要点总结

  1. 消费者独用变量无需原子/锁:m_readLeft只有主线程读写,直接修改完全安全。
  2. 移除冗余的m_count:用m_head和m_tail的差值推导已缓冲大小,简化维护逻辑。
  3. 内存顺序要正确:生产者修改原子变量用memory_order_release,消费者读取用memory_order_acquire,保证跨线程的修改可见性。
  4. 缓冲区与文件位置映射:m_bufferStartPos是实现无锁fast path的核心,能快速判断seek目标是否在已缓冲范围内。

备注:内容来源于stack exchange,提问作者unknown.prince

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 18:07:59