【问题标题】:Parallelize search in a Java set在 Java 集中并行搜索
【发布时间】:2013-01-11 05:12:03
【问题描述】:

我有一个名为 linesList<String> 和一个名为 voc 的巨大 (~3G) Set<String>。我需要从lines 中找到voc 中的所有行。这种多线程方式可以吗?

目前我有这个简单的代码:

for(String line: lines) {
  if (voc.contains(line)) {
    // Great!!
  }
}

有没有办法同时搜索几行?可能有现成的解决方案吗?

PS:我正在使用javolution.util.FastMap,因为它在填充时表现更好。

【问题讨论】:

  • 目前大约 50000 行在 1 秒内被搜索到。 vox 是一个集合,而不是一个列表。
  • 我相信@threadswarm 提供了一个很好的答案。你应该接受他的回答。

标签: java multithreading collections


【解决方案1】:

绝对有可能使用多个线程并行化。您可以执行以下操作:

  1. 将列表分解为不同的“块”,每个执行搜索的线程一个。
  2. 让每个线程查看其块,检查每个字符串是否在集合中,如果是,则将字符串添加到结果集合中。

例如,您可能有以下线程例程:

public void scanAndAdd(List<String> allStrings, Set<String> toCheck,
                       ConcurrentSet<String> matches, int start, int end) {
    for (int i = start; i < end; i++) {
        if (toCheck.contains(allStrings.get(i))) {
            matches.add(allStrings.get(i));
        }
    }
}

然后,您可以根据需要生成尽可能多的线程来运行上述方法并等待所有线程完成。然后将结果匹配存储在matches

为简单起见,我将输出设置为ConcurrentSet,它会自动消除由于写入而导致的竞争条件。由于您只对要检查的字符串列表和字符串集进行读取,因此从allStrings 读取或在toCheck 中执行查找时不需要同步。

希望这会有所帮助!

【讨论】:

  • “由于您只对要检查的字符串列表和字符串集进行读取,因此不需要同步”。错误的!您仍然需要正确发布它。
  • @unbeli- 抱歉 - 我的意思是说在阅读过程中不需要同步。我假设该列表没有被另一个线程同时填充。我还明确使用了ConcurrentMap 来生成输出列表,该列表隐式处理并发。还是我在这里错过了更深层次的东西?
【解决方案2】:

如果您正在寻找这个,只需在不同线程之间拆分 就会(至少在 Oracle JVM 中)将工作分散到所有 CPU 中。 我喜欢使用 CyclicBarrier,让这些线程更容易控制。

http://javarevisited.blogspot.cz/2012/07/cyclicbarrier-example-java-5-concurrency-tutorial.html

【讨论】:

  • 您可以根据需要启动尽可能多的线程,行数除以 CPU 数加上一个线程可能是最好的。 CyclicBarrier 将允许您按顺序使用它们,等到它们全部完成。
  • 如果不清楚:启动尽可能多的线程,可用 CPU 的数量是最好的 np =Runtime.getRuntime().availableProcessors(),取 numThreadsNeeded = np +1。将循环障碍传递给其构造函数中的所有线程,并让它们在完成部分工作后对其调用 await()
【解决方案3】:

这是一个可能的实现。请注意,错误/中断处理已被省略,但这可能会给您一个起点。我包含了一个 main 方法,因此您可以将其复制并粘贴到您的 IDE 中以进行快速演示。

编辑:稍微清理一下以提高可读性和列表分区

import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.CompletionService;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorCompletionService;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class ParallelizeListSearch {

    public static void main(String[] args) throws InterruptedException, ExecutionException {
        List<String> searchList = new ArrayList<String>(7);
        searchList.add("hello");
        searchList.add("world");
        searchList.add("java");
        searchList.add("debian");
        searchList.add("linux");
        searchList.add("jsr-166");
        searchList.add("stack");

        Set<String> targetSet = new HashSet<String>(searchList);

        Set<String> matchSet = findMatches(searchList, targetSet);
        System.out.println("Found " + matchSet.size() + " matches");
        for(String match : matchSet){
            System.out.println("match:  " + match);
        }
    }

    public static Set<String> findMatches(List<String> searchList, Set<String> targetSet) throws InterruptedException, ExecutionException {
        Set<String> locatedMatchSet = new HashSet<String>();

        int threadCount = Runtime.getRuntime().availableProcessors();   

        List<List<String>> partitionList = getChunkList(searchList, threadCount);

        if(partitionList.size() == 1){
            //if we only have one "chunk" then don't bother with a thread-pool
            locatedMatchSet = new ListSearcher(searchList, targetSet).call();
        }else{  
            ExecutorService executor = Executors.newFixedThreadPool(threadCount);
            CompletionService<Set<String>> completionService = new ExecutorCompletionService<Set<String>>(executor);

            for(List<String> chunkList : partitionList)
                completionService.submit(new ListSearcher(chunkList, targetSet));

            for(int x = 0; x < partitionList.size(); x++){
                Set<String> threadMatchSet = completionService.take().get();
                locatedMatchSet.addAll(threadMatchSet);
            }

            executor.shutdown();
        }


        return locatedMatchSet;
    }

    private static class ListSearcher implements Callable<Set<String>> {

        private final List<String> searchList;
        private final Set<String> targetSet;
        private final Set<String> matchSet = new HashSet<String>();

        public ListSearcher(List<String> searchList, Set<String> targetSet) {
            this.searchList = searchList;
            this.targetSet = targetSet;
        }

        @Override
        public Set<String> call() {
            for(String searchValue : searchList){
                if(targetSet.contains(searchValue))
                    matchSet.add(searchValue);
            }

            return matchSet;
        }

    }

    private static <T> List<List<T>> getChunkList(List<T> unpartitionedList, int splitCount) {
        int totalProblemSize = unpartitionedList.size();
        int chunkSize = (int) Math.ceil((double) totalProblemSize / splitCount);

        List<List<T>> chunkList = new ArrayList<List<T>>(splitCount);

        int offset = 0;
        int limit = 0;
        for(int x = 0; x < splitCount; x++){
            limit = offset + chunkSize;
            if(limit > totalProblemSize)
                limit = totalProblemSize;
            List<T> subList = unpartitionedList.subList(offset, limit);
            chunkList.add(subList);
            offset = limit;
        }

        return chunkList;
    }

}

【讨论】:

  • +1 看起来不错,我喜欢使用完成服务的想法。
【解决方案4】:

另一种选择是使用Akka,它可以非常简单地完成这些事情。

实际上,在使用 Akka 完成了一些搜索工作后,我可以告诉您的其中一件事是它支持两种并行化此类事物的方式:通过可组合期货或代理。对于您想要的,可组合期货就足够了。然后,Akka 实际上并没有增加那么多:Netty 提供了大规模并行的 io 基础设施,Futures 是 jdk 的一部分,但 Akka 确实让将这两者放在一起并在需要时/如果需要扩展它们变得超级简单。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-02-26
    • 1970-01-01
    • 2021-03-12
    • 1970-01-01
    • 2012-11-20
    • 1970-01-01
    • 2021-01-10
    • 2014-10-03
    相关资源
    最近更新 更多