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

OCaml多线程异常求助:Sys_error("Bad file descriptor")及性能问题

OCaml多线程文件读取问题修复

问题1:Thread killed on uncaught exception Sys_error ("Bad file descriptor")

原因

  • 多个线程共用同一个文件句柄pipe,当某个线程遇到End_of_file后关闭了文件,其他线程仍尝试调用input_line读取已关闭的句柄,触发文件描述符错误。
  • 异常处理仅捕获End_of_file,未处理其他IO异常(如Sys_error),导致线程因未捕获异常崩溃。

修复方案

  • 添加全局标志位标记文件是否已读取完成,避免重复关闭文件或后续线程读取已关闭句柄。
  • 扩展异常处理范围,捕获所有IO相关异常,确保线程正常退出。
  • 仅让第一个遇到End_of_file的线程执行文件关闭操作,其他线程直接释放资源退出。

问题2:rock方法控制台输出速度极慢

原因

  • read_from_file创建4个工作线程后立即返回,rock中仅等待read_from_file线程结束,未等待4个工作线程完成读取就开始处理db,此时db仍在被并发写入,引发数据不一致和竞态开销。
  • 输出时每次调用flush stdout,频繁刷新终端会大幅降低输出效率。

修复方案

  • 在read_from_file中保存所有工作线程句柄,创建完成后等待所有线程执行完毕再返回。
  • 减少flush调用频率,改为批量输出后统一刷新终端。

修改后的完整代码

type person = {
  firstn : string; 
  lastn : string; 
  age : int; 
  id : int; 
  salary : float option
}

let db = ref []
let readlock = Mutex.create ()
let writelock = Mutex.create ()
let semaphore = Semaphore.Counting.make 4
let file_closed = ref false
let file_closed_mutex = Mutex.create ()

let read_from_file () =
  let pipe = open_in "/home/gc/ocamldir/ocamladv/lib/bigdata.csv" in

  let helper () =
    let rec loop () =
      Semaphore.Counting.acquire semaphore;
      try
        Mutex.lock readlock;
        (* 检查文件是否已关闭,避免无效读取 *)
        Mutex.lock file_closed_mutex;
        let is_closed = !file_closed in
        Mutex.unlock file_closed_mutex;
        if is_closed then (
          Mutex.unlock readlock;
          Semaphore.Counting.release semaphore;
          ()
        ) else (
          let input = input_line pipe in
          Mutex.unlock readlock;
          Mutex.lock writelock;
          let detail = String.split_on_char ',' input in
          db := {
            firstn = List.nth detail 0;
            lastn = List.nth detail 1;
            age = int_of_string @@ List.nth detail 2;
            id = int_of_string @@ List.nth detail 3;
            salary = None
          } :: !db;
          Mutex.unlock writelock;
          Semaphore.Counting.release semaphore;
          loop ()
        )
      with
      | End_of_file ->
          Semaphore.Counting.release semaphore;
          Mutex.unlock readlock;
          (* 仅第一个触发EOF的线程关闭文件 *)
          Mutex.lock file_closed_mutex;
          if not !file_closed then (
            file_closed := true;
            close_in pipe
          );
          Mutex.unlock file_closed_mutex
      | Sys_error msg ->
          Semaphore.Counting.release semaphore;
          Mutex.unlock readlock;
          Printf.eprintf "Thread error: %s\n" msg
    in
    loop ()
  in
  let threads = ref [] in
  let rec create_threads count =
    if count > 0 then (
      let t = Thread.create helper () in
      threads := t :: !threads;
      create_threads (count - 1)
    )
  in
  create_threads 4;
  (* 等待所有工作线程完成读取 *)
  List.iter Thread.join !threads

let rock () =
  read_from_file ();
  !db
  |> List.filter (fun {age; _} -> age < 18)
  |> List.filter (fun {firstn; _} -> String.length firstn < 6)
  |> List.iteri (fun counter {firstn; lastn; id; _} ->
         ANSITerminal.print_string
           [ANSITerminal.Background ANSITerminal.Cyan;
            ANSITerminal.Foreground ANSITerminal.Red]
           (Printf.sprintf "\nIllegal account: {%d} " counter);
         ANSITerminal.print_string
           [ANSITerminal.Background ANSITerminal.Cyan;
            ANSITerminal.Foreground ANSITerminal.Blue]
           (Printf.sprintf "full name: %10s %10s id: %5d" firstn lastn id)
       );
  flush stdout (* 最后统一刷新终端 *)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 11:37:05