【问题标题】:Multithreaded file processing and reporting多线程文件处理和报告
【发布时间】:2012-07-24 21:38:06
【问题描述】:

我有一个应用程序,它处理存储在来自输入目录的多个文件中的数据,然后根据该数据生成一些输出。

到目前为止,应用程序是按顺序工作的,即它启动一个“管理器”线程,该线程

  • 将输入目录的内容读入File[]数组
  • 按顺序处理每个文件并存储结果
  • 处理完所有文件后终止

我想把它转换成一个多线程应用程序,其中的“管理器”线程

  • 将输入目录的内容读入File[]数组
  • 启动多个“处理器”线程,每个线程处理单个文件、存储结果并将该文件的摘要报告返回给“管理器”线程
  • 在处理完所有文件后终止

“处理器”线程的数量最多等于文件的数量,因为它们将通过ThreadPoolExecutor 回收。

任何避免使用join() 或wait()/notify() 的解决方案都是可取的。

基于以上场景:

  1. 让那些“处理器”线程向“管理器”线程报告的最佳方式是什么?基于Callable 和Future 的实现在这里有意义吗?
  2. “管理器”线程如何知道所有“处理器”线程何时完成,即所有文件何时处理完毕?
  3. 是否有一种方法可以“定时”处理器线程并在它花费“太长时间”时终止它(即,尽管经过了预先配置的时间量,它仍未返回结果)?

任何指向(伪)源代码的指针或示例将不胜感激。

【问题讨论】:

  • 我不是线程专家,但这里是我的快速镜头给你一个方向: 1. 你可以使用等待/通知机制来控制报告行为。 2. 主管理器线程可以在读取文件时处理数据。是否需要等到所有文件都处理完? 3. TimerTask 在你的情况下可以派上用场。
  • 感谢您的快速回答。管理器线程不需要等待所有处理器线程完成,但它确实需要知道它们何时完成。我宁愿不使用 wait()/notify(),因此引用 Callable 接口。无论如何,请将您的 cmets 放入答案中,以便我投票。 :-)
  • 您是否需要同步对公共资源的访问?就像您在其中收集结果的文件或缓存一样?
  • 是的,每个单独的处理器线程的结果都会有一个共同的“接收器”,但这会自行处理同步。

标签: java multithreading file callable


【解决方案1】:

您绝对可以在不使用join() 或wait()/notify() 自己的情况下做到这一点。

你应该先看看java.util.concurrent.ExecutorCompletionService。

我认为你应该编写以下类:

  • FileSummary - 保存单个文件处理结果的简单值对象
  • FileProcessor implements Callable<FileSummary> - 将文件转换为 FileSummary 结果的策略
  • File Manager - 创建 FileProcessor 实例、将它们提交到工作队列然后聚合结果的高级管理器。

FileManager 看起来像这样:

class FileManager {
   private CompletionService<FileSummary> cs; // Initialize this in constructor

   public FinalResult processDir(File dir) {
      int fileCount = 0;
      for(File f : dir.listFiles()) {
         cs.submit(new FileProcessor(f));
         fileCount++;
      }

      for(int i = 0; i < fileCount; i++) {
         FileSummary summary = cs.take().get();
         // aggregate summary into final result;
      }
   }

如果你想实现超时,你可以在 CompletionService 上使用poll() 方法而不是take()。

【讨论】:

  • 这与我对使用 Callable/Future 对象的想法非常接近。 ExecutorCompletionService 是一个很酷的建议! poll() 的替代方法是使用带有超时的 get()。请编辑 get() 调用,它应该是 cs.take().get(),而不是 cs.take.get()。谢谢! :-)
  • @PNS 实际上我不确定在 get() 上使用超时是否会起作用,因为 take() 会阻塞直到有可用的结果,即 get() 永远不会阻塞。
  • 如果使用ExecutorCompletionService则正确,否则可以使用超时。
【解决方案2】:

wait()/notify() 是非常低级的原语,你想避免它们是对的。

最简单的解决方案是使用线程安全的队列(或堆栈等——在这种情况下并不重要)。在启动工作线程之前,您的主线程可以将所有Files 添加到线程安全队列/堆栈中。然后启动工作线程,让它们全部拉取Files 并处理它们,直到没有剩余为止。

工作线程可以将结果添加到另一个线程安全队列/堆栈,主线程可以从中获取结果。主线程知道有多少Files,所以当它检索到相同数量的结果时,它就会知道工作已经完成。

java.util.concurrent.BlockingQueue 之类的东西可以使用,java.util.concurrent 中还有其他线程安全的集合也可以。

您还询问了终止耗时过长的工作线程的问题。我会提前告诉你:如果你能让在工作线程上运行的代码足够健壮,你可以放心地把这个特性排除在外,你会让事情变得简单得多。

如果您确实需要此功能,最简单和最可靠的解决方案是为每个线程设置一个“终止”标志,并让工作任务代码经常检查该标志并在设置时退出。为工人创建一个自定义类,并为此目的包含一个volatile boolean 字段。还包括一个setter方法(因为volatile,它不需要是synchronized)。

如果工作人员发现其“终止”标志已设置,它可以将其File 对象推回工作队列/堆栈,以便另一个线程可以处理它。当然,如果出现问题,意味着File无法被成功处理,这将导致无限循环。

最好的办法是让工作代码非常简单和健壮,这样你就不用担心它“不会终止”。

【讨论】:

    【解决方案3】:
    1. 他们无需报告。只需计算剩余要完成的作业数量,并在完成时减少线程计数。

    2. 当剩余待完成作业的计数达到零时,所有“处理器”线程都已完成。

    3. 当然,只需将该代码添加到线程即可。当它开始工作时,检查时间并计算停止时间。定期(比如当您从文件中读取更多内容时),检查它是否超过了停止时间,如果是,则停止。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2012-03-17
      • 1970-01-01
      • 1970-01-01
      • 2013-09-11
      • 2017-02-14
      • 1970-01-01
      相关资源
      最近更新 更多