【问题标题】:How to check the size() or isEmpty() for ConcurrentLinkedQueue如何检查 ConcurrentLinkedQueue 的 size() 或 isEmpty()
【发布时间】:2016-05-28 20:34:42
【问题描述】:

我正在尝试用 Java 为 Web 爬虫创建一个简单结构的原型。到目前为止,原型只是尝试执行以下操作:

  • 使用起始 URL 列表初始化队列
  • 从队列中取出一个 URL 并提交到一个新线程
  • 做一些工作,然后将该 URL 添加到一组已访问的 URL 中

对于起始 URL 的队列,我使用 ConcurrentLinkedQueue 进行同步。 为了产生新的线程,我使用ExecutorService

但是在创建新线程时,应用程序需要检查ConcurrentLinkedQueue 是否为空。我尝试使用:

  • .size()
  • .isEmpty()

但两者似乎都没有返回 ConcurrentLinkedQueue 的真实状态。

问题出在下面的块中:

while (!crawler.getUrl_horizon().isEmpty()) {
                workers.submitNewWorkerThread(crawler);
            }

因此,ExecutorService 会在其限制范围内创建所有线程,即使输入只有 2 个 URL。

这里实现多线程的方式有问题吗?如果没有,检查 ConcurrentLinkedQueue 状态的更好方法是什么?

应用程序的起始类:

public class CrawlerApp {

    private static Crawler crawler;

    public static void main(String[] args) {
        crawler = = new Crawler();
        initializeApp();
        startCrawling();

    }

    private static void startCrawling() {
        crawler.setUrl_visited(new HashSet<URL>());
        WorkerManager workers = WorkerManager.getInstance();
        while (!crawler.getUrl_horizon().isEmpty()) {
            workers.submitNewWorkerThread(crawler);
        }
        try {
            workers.getExecutor().shutdown();
            workers.getExecutor().awaitTermination(10, TimeUnit.MINUTES);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }

    private static void initializeApp() {

        Properties config = new Properties();
        try {
            config.load(CrawlerApp.class.getClassLoader().getResourceAsStream("url-horizon.properties"));
            String[] horizon = config.getProperty("urls").split(",");
            ConcurrentLinkedQueue<URL> url_horizon = new ConcurrentLinkedQueue<>();
            for (String link : horizon) {
                URL url = new URL();
                url.setURL(link);
                url_horizon.add(url);
            }
            crawler.setUrl_horizon(url_horizon);
        } catch (IOException e) {
            e.printStackTrace();
        }

    }

}

Crawler.java 维护 URL 队列和已访问 URL 集。

public class Crawler implements Runnable {
    private ConcurrentLinkedQueue<URL> url_horizon;

    public void setUrl_horizon(ConcurrentLinkedQueue<URL> url_horizon) {
        this.url_horizon = url_horizon;
    }

    public ConcurrentLinkedQueue<URL> getUrl_horizon() {
        return url_horizon;
    }

    private Set<URL> url_visited;

    public void setUrl_visited(Set<URL> url_visited) {
        this.url_visited = url_visited;
    }

    public Set<URL> getUrl_visited() {
        return Collections.synchronizedSet(url_visited);
    }

    @Override
    public void run() {
        URL url = nextURLFromHorizon();
        scrap(url);
        addURLToVisited(url);

    }

    private URL nextURLFromHorizon() {
        if (!getUrl_horizon().isEmpty()) {
            URL url = url_horizon.poll();
            if (getUrl_visited().contains(url)) {
                return nextURLFromHorizon();
            }
            System.out.println("Horizon URL:" + url.getURL());
            return url;

        }
        return null;

    }

    private void scrap(URL url) {
        new Scrapper().scrap(url);
    }

    private void addURLToVisited(URL url) {
        System.out.println("Adding to visited set:" + url.getURL());
        getUrl_visited().add(url);
    }

}

URL.java 只是一个带有private String url 并覆盖hashCode()equals() 的类。

另外,Scrapper.scrap() 到目前为止只有虚拟实现:

public void scrap(URL url){
        System.out.println("Done scrapping:"+url.getURL());
    }

WorkerManager 创建线程:

public class WorkerManager {
    private static final Integer WORKER_LIMIT = 10;
    private final ExecutorService executor = Executors.newFixedThreadPool(WORKER_LIMIT);

    public ExecutorService getExecutor() {
        return executor;
    }

    private static volatile WorkerManager instance = null;

    private WorkerManager() {
    }

    public static WorkerManager getInstance() {
        if (instance == null) {
            synchronized (WorkerManager.class) {
                if (instance == null) {
                    instance = new WorkerManager();
                }
            }
        }

        return instance;
    }

    public Future submitNewWorkerThread(Runnable run) {
        return executor.submit(run);
    }

}

【问题讨论】:

    标签: java multithreading concurrency web-crawler executorservice


    【解决方案1】:

    问题

    您最终创建的线程数多于队列中的 URL 的原因是,有可能(实际上很可能)执行程序的任何线程都不会启动,直到您经历了很多 while 循环次。

    无论何时使用线程,您都应始终牢记线程是独立调度的,并按照自己的节奏运行,除非您明确同步它们。在这种情况下,线程可以在 submit() 调用之后的任何时间启动,即使您似乎希望每个线程在您的 while 循环中的下一次迭代之前启动并经过 nextURLFromHorizon

    解决方案

    在将Runnable 提交给执行程序之前,请考虑将 URL 从队列中取出。我还建议定义一个CrawlerTask 提交给Executor 一次,而不是一个Crawler 重复提交。在这样的设计中,您甚至不需要为要抓取的 URL 提供线程安全的容器。

    class CrawlerTask extends Runnable {
       URL url;
    
       CrawlerTask(URL url) { this.url = url; }
    
       @Override
       public void run() {
         scrape(url);
         // add url to visited?
       }
    }
    
    class Crawler {
      ExecutorService executor;
      Queue urlHorizon;
    
      //...
    
      private static void startCrawling() {
        while (!urlHorizon.isEmpty()) {
          executor.submit(new CrawlerTask(urlHorizon.poll());
        }
        // ...
      }
    }
    

    【讨论】:

    • 你的设计显然更有意义。但是如果我需要显式地同步线程,我应该做一个类级别的synchronized 方法还是创建一个Object lock 并使用synchronized 块?
    • 另外,即使我创建了一个CrawlerTask,我仍然需要传递对Crawler 对象的引用,该对象具有一组已访问过的URL。
    • synchronized 不能单独确保线程按照您想要的顺序执行。它只是防止两个线程同时进入块/方法的主体。如果你想让一个线程等到另一个线程做某事,你需要使用Object.waitObject.notify。在生产者-消费者模式中使用并发 Queue 在内部使用这些机制,但随后一个线程 .add()s 和另一个线程 .remove()s 来自队列。
    • 要回答你的第二个问题,是的,你可以传递Crawler 或只是将CrawlerTask 设为Crawler 的内部(非静态)类,以便它可以访问Crawlers范围。然后为以前访问过的 URL 使用线程安全集。
    猜你喜欢
    • 2016-03-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-06-26
    • 1970-01-01
    • 2010-10-11
    • 1970-01-01
    • 2014-08-08
    相关资源
    最近更新 更多