【问题标题】:Merge sort large file in parallel with memory limit (Linux)与内存限制并行合并排序大文件(Linux)
【发布时间】:2018-09-11 01:14:27
【问题描述】:

我需要使用t 线程对大小为M 的大型二进制文件进行排序。文件中的记录大小相同。该任务明确表示我可以分配的内存量是m,并且比M 小得多。还保证硬盘驱动器至少有2 * M 可用空间。这需要合并排序 ofc,但结果并不是那么明显。我在这里看到了三种不同的方法:

一个。将文件inputtemp1temp2 映射到内存中。执行合并排序input -> temp1 -> temp2 -> temp1 ...,直到其中一个临时被排序。线程只竞争选择下一部分工作,不竞争读/写。

B。每个 fopen 3 个文件 t 次,每个线程获得 3 个 FILE 指针,每个文件一个。同样,他们只争夺下一部分工作,读取和写入应该并行工作。

C。每次打开 3 个文件,将它们放在互斥锁下,所有线程并行工作,但为了获取更多工作或读取或写入,它们会锁定各自的互斥锁。

注意事项:

在现实生活中我肯定会选择A。但它不是破坏了缓冲有限的整个目的吗? (换句话说,这不是作弊吗?)。使用这种方法,我什至可以在没有额外缓冲区的情况下对整个文件进行基数排序。而且这个解决方案是 Linux 特有的,我认为 Linux 是从对话中暗示出来的,但在任务描述中没有明确说明。

关于 B,我认为它可以在 Linux 上运行,但不能移植,请参阅上面的 Linux 说明。

关于C,它是可移植的,但我不知道如何优化它(例如,8 个线程足够小m 只会在队列中等待轮到他们,然后读/写一小部分数据,然后立即对其进行排序并再次相互碰撞。IMO 的工作速度不太可能超过 1 个线程。

问题:

  1. 哪种解决方案更适合该任务?
  2. 哪种解决方案在现实生活中是更好的设计(假设是 Linux)?
  3. B 有效吗?换句话说,多次打开文件并并行写入(到文件的不同部分)合法吗?
  4. 任何替代方法?

【问题讨论】:

    标签: c linux sorting


    【解决方案1】:

    你的问题有很多方面,所以我会试着把它分解一下,同时尝试回答你几乎所有的问题:

    • 您在存储设备上获得了一个大文件,该文件可能在 blocks 上运行,即您可以同时加载和存储许多条目。如果您从存储中访问单个条目,则必须处理相当大的访问延迟延迟,您只能通过同时加载许多元素来尝试隐藏它,从而在所有元素加载时间中分摊延迟。

    • 与存储相比,您的主存相当快(尤其是对于随机访问),因此您希望在主存中保留尽可能多的数据,并且只在存储上读取和写入顺序块。这也是 A 没有真正作弊的原因,因为如果您尝试使用存储进行随机访问,那么您将比使用主内存慢得多。

    结合这些结果,您可以得出以下方法,该方法基本上是A,但带有一些通常用于外部算法的工程细节。

    • 仅使用一个专用线程来读取和写入存储。 这样,每个文件只需要一个文件描述符,理论上甚至可以在很短的时间内收集和重新排序来自所有线程的读取和写入请求,以获得几乎顺序的访问模式。此外,您的线程可以将写入请求排队并继续下一个块,而无需等待 IO 完成。

    • t 个块(来自 input)加载到最大大小的主内存中,以便您可以在每个块上并行运行 mergesort。将块排序后,将它们写入存储为temp1
      重复此操作,直到文件中的所有块都已排序。

    • 现在对已排序的块进行所谓的多路合并: 每个线程从temp1 加载一定数量的 k 个连续块到内存中,并使用优先级队列或锦标赛树将它们合并,以找到下一个要插入到结果块中的最小值。一旦您的块已满,您就将其写入temp2 的存储空间,以便为下一个块释放内存。在这一步之后,在概念上交换temp1temp2

    • 您仍然需要执行多个合并步骤,但与您在 A 中可能表示的常规双向合并相比,此数字减少了 log k 倍强>。在最初的几个合并步骤之后,您的块可能太大而无法放入主内存,因此您将它们拆分为较小的块,并且从第一个小块开始,仅当所有先前的元素都已经被提取时才获取下一个块合并。在这里,您甚至可以进行一些预取,因为块访问的顺序是由块最小值预先确定的,但这可能超出了这个问题的范围。 请注意,k 的值通常仅受可用内存的限制。

    • 最后,t 大块需要合并在一起。我真的不知道是否有一个很好的并行方法,可能需要按顺序合并它们,所以你可以再次使用上面的 t-way 合并来生成单个排序文件。

    【讨论】:

    • 如果您正在寻找更详细的描述和分析,我可以通过快速搜索找到最合适的一个是Sanders and Dementiev的以下论文,尽管它使用了更通用的硬件模型(并行磁盘)。
    【解决方案2】:

    Gnu 排序是文本文件的多线程合并排序,但这里可以使用它的基本功能。将“块”定义为可以在大小为 m 的内存中排序的记录数。

    排序阶段:对于每个记录“块”,读取记录“块”,对“块”使用多线程排序,然后将记录“块”写入临时文件,以上限结束(M / m) 临时文件。 Gnu sort 对指向记录的指针数组进行排序,部分原因是记录是可变长度的。对于固定大小的记录,在我的测试中,由于缓存问题,直接对记录进行排序而不是对指向记录的指针数组进行排序(这会导致对记录的缓存不友好的随机访问)更快,除非记录大小大于 128 之间的某个值和 256 个字节。

    合并阶段:对临时文件执行单线程k-way合并(例如优先级队列),直到生成单个文件。多线程在这里没有帮助,因为它假设 k 路合并阶段是 I/O 绑定的,而不是 cpu 绑定的。对于 Gnu 排序,k 的默认值为 16(它对临时文件进行 16 路合并)。

    为避免超过 2 x M 空间,文件在读取后需要删除。

    【讨论】:

      【解决方案3】:

      如果您的文件比您的 RAM 大得多,那么这就是解决方案。 https://stackoverflow.com/a/49839773/1647320

      如果您的文件大小是 RAM 大小的 70-80%,那么以下是解决方案。它是内存中的并行归并排序。

      根据您的系统更改此行。 fpath 是您的一个大输入文件。 shared 路径是存储执行日志的位置。fdir 是存储和合并中间文件的位置。根据您的机器更改这些路径。

      public static final String fdir = "/tmp/";
          public static final String shared = "/exports/home/schatterjee/cs553-pa2a/";
          public static final String fPath = "/input/data-20GB.in";
          public static final String opLog = shared+"Mysort20GB.log";
      

      然后运行以下程序。您的最终排序文件将在 fdir 路径中使用名称 op2GB 创建。最后一行 Runtime.getRuntime().exec("valsort" + fdir + "op" + (treeHeight*100)+1 + " > " + opLog);检查输出是否排序。如果您的机器上没有安装 valsort,或者输入文件不是使用 gensort(http://www.ordinal.com/gensort.html) 生成的,请删除此行。

      另外,不要忘记更改 int totalLines = 20000000;到文件中的总行数。并且线程数 (int threadCount = 8) 应该始终是 2 的幂。

      import java.io.*;
      import java.nio.file.Files;
      import java.nio.file.Paths;
      import java.util.LinkedList;
      import java.util.Comparator;
      import java.util.HashMap;
      import java.util.stream.Stream;
      
      
      class SplitJob extends Thread {
          LinkedList<String> chunkName;
          int startLine, endLine;
      
          SplitJob(LinkedList<String> chunkName, int startLine, int endLine) {
              this.chunkName = chunkName;
              this.startLine = startLine;
              this.endLine = endLine;
          }
      
          public void run() {
              try {
                  int totalLines = endLine + 1 - startLine;
                  Stream<String> chunks =
                          Files.lines(Paths.get(Mysort2GB.fPath))
                                  .skip(startLine - 1)
                                  .limit(totalLines)
                                  .sorted(Comparator.naturalOrder());
                  chunks.forEach(line -> {
                      chunkName.add(line);
                  });
                  System.out.println(" Done Writing " + Thread.currentThread().getName());
      
              } catch (Exception e) {
                  System.out.println(e);
              }
          }
      }
      
      class MergeJob extends Thread {
          int list1, list2, oplist;
          MergeJob(int list1, int list2, int oplist) {
              this.list1 = list1;
              this.list2 = list2;
              this.oplist = oplist;
          }
      
          public void run() {
              try {
                  System.out.println(list1 + " Started Merging " + list2 );
                  LinkedList<String> merged = new LinkedList<>();
                  LinkedList<String> ilist1 = Mysort2GB.sortedChunks.get(list1);
                  LinkedList<String> ilist2 = Mysort2GB.sortedChunks.get(list2);
      
                  //Merge 2 files based on which string is greater.
                  while (ilist1.size() != 0 || ilist2.size() != 0) {
                      if (ilist1.size() == 0 ||
                              (ilist2.size() != 0 && ilist1.get(0).compareTo(ilist2.get(0)) > 0)) {
                          merged.add(ilist2.remove(0));
                      } else {
                          merged.add(ilist1.remove(0));
                      }
                  }
                  System.out.println(list1 + " Done Merging " + list2 );
                  Mysort2GB.sortedChunks.remove(list1);
                  Mysort2GB.sortedChunks.remove(list2);
                  Mysort2GB.sortedChunks.put(oplist, merged);
              } catch (Exception e) {
                  System.out.println(e);
              }
          }
      }
      
      public class Mysort2GB {
          //public static final String fdir = "/Users/diesel/Desktop/";
          public static final String fdir = "/tmp/";
          public static final String shared = "/exports/home/schatterjee/cs553-pa2a/";
          public static final String fPath = "/input/data-2GB.in";
          public static HashMap<Integer, LinkedList<String>> sortedChunks = new HashMap();
          public static final String opfile = fdir+"op2GB";
          public static final String opLog = shared + "mysort2GB.log";
      
      
          public static void main(String[] args) throws Exception{
              long startTime = System.nanoTime();
              int threadCount = 8; // Number of threads
              int totalLines = 20000000;
              int linesPerFile = totalLines / threadCount;
              LinkedList<Thread> activeThreads = new LinkedList<Thread>();
      
      
              for (int i = 1; i <= threadCount; i++) {
                  int startLine = i == 1 ? i : (i - 1) * linesPerFile + 1;
                  int endLine = i * linesPerFile;
                  LinkedList<String> thisChunk = new LinkedList<>();
                  SplitJob mapThreads = new SplitJob(thisChunk, startLine, endLine);
                  sortedChunks.put(i,thisChunk);
                  activeThreads.add(mapThreads);
                  mapThreads.start();
              }
              activeThreads.stream().forEach(t -> {
                  try {
                      t.join();
                  } catch (Exception e) {
                  }
              });
      
              int treeHeight = (int) (Math.log(threadCount) / Math.log(2));
      
              for (int i = 0; i < treeHeight; i++) {
                  LinkedList<Thread> actvThreads = new LinkedList<Thread>();
                  for (int j = 1, itr = 1; j <= threadCount / (i + 1); j += 2, itr++) {
                      int offset = i * 100;
                      int list1 = j + offset;
                      int list2 = (j + 1) + offset;
                      int opList = itr + ((i + 1) * 100);
                      MergeJob reduceThreads =
                              new MergeJob(list1,list2,opList);
                      actvThreads.add(reduceThreads);
                      reduceThreads.start();
                  }
                  actvThreads.stream().forEach(t -> {
                      try {
                          t.join();
                      } catch (Exception e) {
                      }
                  });
              }
              BufferedWriter writer = Files.newBufferedWriter(Paths.get(opfile));
              sortedChunks.get(treeHeight*100+1).forEach(line -> {
                  try {
                      writer.write(line+"\r\n");
                  }catch (Exception e){
      
                  }
              });
              writer.close();
              long endTime = System.nanoTime();
              double timeTaken = (endTime - startTime)/1e9;
              System.out.println(timeTaken);
              BufferedWriter logFile = new BufferedWriter(new FileWriter(opLog, true));
              logFile.write("Time Taken in seconds:" + timeTaken);
              Runtime.getRuntime().exec("valsort  " + opfile + " > " + opLog);
              logFile.close();
          }
      }
      
      
        [1]: https://i.stack.imgur.com/5feNb.png
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2011-09-01
        • 2015-09-15
        • 2018-10-15
        • 1970-01-01
        • 2018-03-08
        • 1970-01-01
        相关资源
        最近更新 更多