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

Linux下C++异步IO场景中如何正确等待condition variable?

解决Linux异步AIO读取器中条件变量notify先于wait触发的问题

我在Linux下用C++实现异步IO文件读取器,用两个缓冲区交替处理:首次阻塞读取第一块数据,之后循环里启动异步IO读取下一块,同时处理当前数据块,处理完后等待条件变量,预期AIO完成回调会通知条件变量。但实际运行中notify经常在wait之前触发,导致程序卡在wait上无法继续。

问题根源

  1. 条件变量的wait未结合谓词:如果AIO回调在主循环调用wait前就执行了notify,wait会错过信号,进入无限等待。
  2. 共享状态管理缺失:没有明确标记异步IO是否完成的状态变量,全局变量bytesRead的访问存在竞态条件。
  3. AIO读取偏移未更新:每次读取都从文件起始位置开始,导致重复读取同一数据块。

解决方案

  1. 添加共享状态变量io_completed,标记异步IO是否完成,所有访问必须在互斥锁保护下进行。
  2. 调用cv.wait()时传入谓词,检查io_completed是否为true,确保即使notify先触发,wait也能直接返回。
  3. 每次启动异步IO前重置io_completed为false,保证等待的是当前IO的完成信号。
  4. 修正AIO的读取偏移,每次读取后更新偏移量,实现连续读取。
  5. 所有共享变量的读写操作均在互斥锁范围内执行,避免竞态。

修改后的完整代码

#include <aio.h>
#include <fcntl.h>
#include <signal.h>
#include <unistd.h>

#include <condition_variable>
#include <cstring>
#include <iostream>
#include <thread>

using namespace std;
using namespace std::chrono_literals;

constexpr uint32_t blockSize = 512;

mutex readMutex;
condition_variable cv;

int fh;
int bytesRead = 0;
bool io_completed = false; // 标记异步IO是否完成

void process(char* buf, uint32_t bytesRead) {
    cout << "processing..." << endl;
    usleep(100000);
}

void aio_completion_handler(sigval_t sigval) {
    struct aiocb* req = (struct aiocb*)sigval.sival_ptr;

    if (aio_error(req) == 0) {
        int ret = aio_return(req);
        unique_lock<mutex> readLock(readMutex);
        bytesRead = ret; // 使用实际读取字节数而非预设值
        io_completed = true; // 标记IO完成
        cout << "ret == " << ret << endl;
        cout << (char*)req->aio_buf << endl;
    } else {
        // 处理IO错误
        unique_lock<mutex> readLock(readMutex);
        bytesRead = -1;
        io_completed = true;
        cerr << "AIO error occurred" << endl;
    }
    cv.notify_one(); // 通知主循环
}

void thready() {
    char* buf1 = new char[blockSize];
    char* buf2 = new char[blockSize];
    aiocb cb;
    char* processbuf = buf1;
    char* readbuf = buf2;
    off_t current_offset = 0; // 跟踪文件读取偏移量

    fh = open("smallfile.dat", O_RDONLY);
    if (fh < 0) {
        throw std::runtime_error("cannot open file!");
    }

    memset(&cb, 0, sizeof(aiocb));
    cb.aio_fildes = fh;
    cb.aio_nbytes = blockSize;
    cb.aio_offset = current_offset;

    // 设置AIO回调参数
    cb.aio_sigevent.sigev_notify_attributes = nullptr;
    cb.aio_sigevent.sigev_notify = SIGEV_THREAD;
    cb.aio_sigevent.sigev_notify_function = aio_completion_handler;
    cb.aio_sigevent.sigev_value.sival_ptr = &cb;

    // 首次阻塞读取
    int currentBytesRead = read(fh, buf1, blockSize);
    if (currentBytesRead <= 0) {
        close(fh);
        throw std::runtime_error("initial read failed");
    }
    current_offset += currentBytesRead; // 更新偏移量

    while (true) {
        {
            unique_lock<mutex> readLock(readMutex);
            io_completed = false; // 重置IO完成状态,准备新的异步读取
        }

        // 设置本次异步读取的缓冲区和偏移量
        cb.aio_buf = readbuf;
        cb.aio_offset = current_offset;
        int ret = aio_read(&cb);
        if (ret < 0) {
            cerr << "aio_read failed" << endl;
            break;
        }

        // 处理当前数据块
        process(processbuf, currentBytesRead);

        // 等待异步IO完成,带谓词检查避免虚假唤醒和信号丢失
        {
            unique_lock<mutex> readLock(readMutex);
            cv.wait(readLock, []{ return io_completed; });
            currentBytesRead = bytesRead; // 在锁内安全读取共享变量
        }

        // 检查读取结果,结束循环条件
        if (currentBytesRead <= 0 || currentBytesRead < blockSize) {
            break;
        }

        current_offset += currentBytesRead; // 更新偏移量
        cout << "back from wait" << endl;
        swap(processbuf, readbuf); // 交换缓冲区,准备下一轮处理
    }

    // 处理最后一块未处理的数据
    if (currentBytesRead > 0) {
        process(readbuf, currentBytesRead);
    }

    close(fh);
    delete[] buf1;
    delete[] buf2;
}

int main() {
    try {
        thready();
    } catch (std::exception& e) {
        cerr << e.what() << '\n';
    }
    return 0;
}

关键修改说明

  • 状态变量io_completed:明确标记异步IO的完成状态,主循环通过谓词检查该变量,彻底解决notify先于wait触发的问题。
  • 带谓词的wait调用:cv.wait(readLock, []{ return io_completed; })会先检查状态,若已完成则直接返回,否则进入等待,同时自动处理虚假唤醒。
  • 偏移量管理:新增current_offset跟踪读取位置,确保文件连续读取,避免重复读取同一块。
  • 共享变量安全访问:所有对bytesRead和io_completed的操作均在互斥锁保护下,消除竞态条件。
  • 错误处理:在AIO回调中处理错误场景,设置错误状态并通知主循环,避免程序无响应。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 20:18:25