C++异步IO锁与条件变量实现崩溃问题排查
异步文件IO实现崩溃问题排查与解决
问题背景
要实现异步文件IO逻辑:读取一块文件数据的同时,启动线程异步读取下一块数据并处理当前块。过程中多次遇到崩溃:
- 最初使用
condition_variable时程序崩溃 - 改用直接加锁方式后,首次调用
readLock.lock()触发terminate - 采用RAII修正代码后仍崩溃,第三次读取时抛出
std::system_error,错误信息为"Operation not permitted"并核心转储
错误代码
初始代码
#include <fcntl.h> #include <unistd.h> #include <condition_variable> #include <iostream> #include <thread> using namespace std; using namespace std::chrono_literals; constexpr uint32_t blockSize = 32768; char* buf1; char* buf2; char* processbuf = buf1; char* readbuf = buf2; mutex readMutex; mutex procMutex; //condition_variable readCV; //condition_variable procCV; int fh; int bytesRead; std::unique_lock<std::mutex> readLock; std::unique_lock<std::mutex> procLock; void nextRead() { while (true) { readLock.lock(); bytesRead = read(fh, readbuf, blockSize); for (int i = 0; i < 5; i++) { cerr << "reading..." << endl; usleep(100000); readLock.unlock(); } cerr << "notifying..." << endl; readLock.unlock(); if (bytesRead != blockSize) // last time, end here! return; procLock.lock(); procLock.unlock(); } } void process(char* buf, uint32_t bytesRead) { cout << "hit it!" << endl; } int thready() { buf1 = new char[blockSize]; buf2 = new char[blockSize]; fh = open("bigfile.dat", O_RDONLY); if (fh < 0) { throw std::runtime_error("cannot open file!"); } int currentBytesRead = read(fh, buf1, blockSize); thread reader(nextRead); process(processbuf, currentBytesRead); while (true) { readLock.lock(); cout << "back from wait" << endl; swap(processbuf, readbuf); // switch to other buffer for next time currentBytesRead = bytesRead; // copy locally so thread can do the other one // TODO: is the above a problem? what if readLock.unlock(); procLock.lock(); process(processbuf, currentBytesRead); procLock.unlock(); } reader.join(); delete[] buf1; delete[] buf2; } int main() { try { thready(); }catch(std::exception& e) { cerr << e.what() << '\n'; } return 0; }
RAII修正代码
#include <fcntl.h> #include <unistd.h> #include <condition_variable> #include <iostream> #include <thread> using namespace std; using namespace std::chrono_literals; constexpr uint32_t blockSize = 32768; char* buf1; char* buf2; char* processbuf = buf1; char* readbuf = buf2; mutex readMutex; mutex procMutex; //condition_variable readCV; //condition_variable procCV; int fh; int bytesRead; void nextRead() { while (true) { { unique_lock<mutex> readLock(readMutex); bytesRead = read(fh, readbuf, blockSize); for (int i = 0; i < 5; i++) { cerr << "reading..." << endl; usleep(100000); readLock.unlock(); } cerr << "notifying..." << endl; } if (bytesRead != blockSize) // last time, end here! return; // wait for process to complete unique_lock<mutex> procLock(procMutex); } } void process(char* buf, uint32_t bytesRead) { cout << "processing..." << endl; usleep(100000); } int thready() { buf1 = new char[blockSize]; buf2 = new char[blockSize]; fh = open("bigfile.dat", O_RDONLY); if (fh < 0) { throw std::runtime_error("cannot open file!"); } int currentBytesRead = read(fh, buf1, blockSize); thread reader(nextRead); process(processbuf, currentBytesRead); while (true) { { unique_lock<mutex> readLock(readMutex); cout << "back from wait" << endl; swap(processbuf, readbuf); // switch to other buffer for next time currentBytesRead = bytesRead; // create local copy // TODO: is the above a problem? what if } unique_lock<mutex> procLock(procMutex); process(processbuf, currentBytesRead); } reader.join(); delete[] buf1; delete[] buf2; } int main() { try { thready(); } catch(std::exception& e) { cerr << e.what() << '\n'; } return 0; }
程序输出
processing... reading... reading...back from wait processing... back from wait processing... terminate called after throwing an instance of 'std::system_error' what(): Operation not permitted Aborted (core dumped)
问题原因分析
初始代码核心错误:全局
unique_lock对象默认构造,未关联任何互斥量,调用lock()会直接触发std::system_error——这是非法操作,因为unique_lock默认构造后没有绑定互斥量,不允许执行lock/unlock操作。RAII修正后的错误:
- 重复解锁互斥量:
nextRead函数中,readLock在作用域内被lock一次,但在for循环中执行了5次unlock,解锁次数超过锁定次数会触发未定义行为,最终抛出"Operation not permitted"错误。 - 同步逻辑完全失效:没有使用条件变量实现正确的线程同步,主线程和reader线程的执行顺序完全混乱,缓冲区交换、读写操作的时机毫无约束,导致数据竞争和非法操作。
- 无效的互斥量使用:
procMutex仅被加锁后立即解锁,没有起到同步主线程处理完成与reader线程启动下一次读取的作用。 - 无限循环无退出条件:主线程的
while(true)永远不会终止,即使reader线程读完最后一块数据返回,主线程仍会继续尝试获取锁,引发后续错误。 - 全局变量初始化风险:
processbuf和readbuf初始指向未分配内存的buf1/buf2,虽然后续主线程先分配内存再启动线程,但全局变量本身存在线程安全隐患。
- 重复解锁互斥量:
解决办法
修正要点
- 使用条件变量实现主线程与reader线程的同步:reader读完一块后通知主线程处理,主线程处理完后通知reader读下一块。
- 严格控制
unique_lock的lock/unlock次数,利用RAII作用域自动管理锁的生命周期,避免手动重复解锁。 - 添加循环退出条件,当处理完最后一块数据后终止主线程循环,确保线程正常join。
- 避免不必要的全局变量,改用结构体封装状态,或通过参数传递(示例中仍保留全局变量但修正同步逻辑)。
修正后的完整代码
#include <fcntl.h> #include <unistd.h> #include <condition_variable> #include <iostream> #include <thread> #include <mutex> using namespace std; constexpr uint32_t blockSize = 32768; char* buf1; char* buf2; char* process_buf = nullptr; char* read_buf = nullptr; mutex mtx; condition_variable read_cv; condition_variable process_cv; int fh; int bytes_read = 0; bool read_done = false; bool process_done = true; bool eof = false; void next_read() { while (true) { unique_lock<mutex> lock(mtx); // 等待主线程处理完成,才能开始下一次读取 process_cv.wait(lock, []{ return process_done; }); process_done = false; if (eof) break; bytes_read = read(fh, read_buf, blockSize); cerr << "reading completed, read " << bytes_read << " bytes" << endl; read_done = true; read_cv.notify_one(); if (bytes_read != blockSize) { eof = true; break; } } } void process_data(char* buf, uint32_t bytes) { cout << "processing " << bytes << " bytes..." << endl; usleep(100000); } int main() { try { buf1 = new char[blockSize]; buf2 = new char[blockSize]; process_buf = buf1; read_buf = buf2; fh = open("bigfile.dat", O_RDONLY); if (fh < 0) { throw runtime_error("cannot open file!"); } // 首次读取 int current_bytes = read(fh, process_buf, blockSize); thread reader(next_read); while (true) { unique_lock<mutex> lock(mtx); // 等待读取完成 read_cv.wait(lock, []{ return read_done; }); read_done = false; process_data(process_buf, current_bytes); if (eof) break; // 交换缓冲区 swap(process_buf, read_buf); current_bytes = bytes_read; process_done = true; process_cv.notify_one(); } reader.join(); close(fh); delete[] buf1; delete[] buf2; } catch (exception& e) { cerr << e.what() << '\n'; if (fh >= 0) close(fh); delete[] buf1; delete[] buf2; return 1; } return 0; }
修正说明
- 同步逻辑:通过
read_cv和process_cv两个条件变量,实现reader线程和主线程的交替执行:- reader线程等待主线程处理完成后才开始读取
- 主线程等待reader线程读取完成后才开始处理
- 锁管理:利用
unique_lock的RAII特性,在作用域内自动管理锁的生命周期,避免手动解锁错误。 - 退出条件:通过
eof标记判断文件是否读完,主线程和reader线程都能正常退出循环。 - 缓冲区安全:只有在同步条件满足时才交换缓冲区,确保读写操作不会访问同一缓冲区。
内容的提问来源于stack exchange,提问作者Dov
相关产品推荐
相关产品推荐

