Linux下C++异步IO场景中如何正确等待condition variable?
解决Linux异步AIO读取器中条件变量notify先于wait触发的问题
我在Linux下用C++实现异步IO文件读取器,用两个缓冲区交替处理:首次阻塞读取第一块数据,之后循环里启动异步IO读取下一块,同时处理当前数据块,处理完后等待条件变量,预期AIO完成回调会通知条件变量。但实际运行中notify经常在wait之前触发,导致程序卡在wait上无法继续。
问题根源
- 条件变量的
wait未结合谓词:如果AIO回调在主循环调用wait前就执行了notify,wait会错过信号,进入无限等待。 - 共享状态管理缺失:没有明确标记异步IO是否完成的状态变量,全局变量
bytesRead的访问存在竞态条件。 - AIO读取偏移未更新:每次读取都从文件起始位置开始,导致重复读取同一数据块。
解决方案
- 添加共享状态变量
io_completed,标记异步IO是否完成,所有访问必须在互斥锁保护下进行。 - 调用
cv.wait()时传入谓词,检查io_completed是否为true,确保即使notify先触发,wait也能直接返回。 - 每次启动异步IO前重置
io_completed为false,保证等待的是当前IO的完成信号。 - 修正AIO的读取偏移,每次读取后更新偏移量,实现连续读取。
- 所有共享变量的读写操作均在互斥锁范围内执行,避免竞态。
修改后的完整代码
#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
相关产品推荐
相关产品推荐

