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

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)

问题原因分析

  1. 初始代码核心错误:全局unique_lock对象默认构造,未关联任何互斥量,调用lock()会直接触发std::system_error——这是非法操作,因为unique_lock默认构造后没有绑定互斥量,不允许执行lock/unlock操作。

  2. 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;
}

修正说明

  1. 同步逻辑:通过read_cv和process_cv两个条件变量,实现reader线程和主线程的交替执行:
    • reader线程等待主线程处理完成后才开始读取
    • 主线程等待reader线程读取完成后才开始处理
  2. 锁管理:利用unique_lock的RAII特性,在作用域内自动管理锁的生命周期,避免手动解锁错误。
  3. 退出条件:通过eof标记判断文件是否读完,主线程和reader线程都能正常退出循环。
  4. 缓冲区安全:只有在同步条件满足时才交换缓冲区,确保读写操作不会访问同一缓冲区。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 15:15:44