【问题标题】:Get results of scheduled non-blocking operations in Java获取 Java 中计划的非阻塞操作的结果
【发布时间】:2019-10-30 20:57:09
【问题描述】:

我正在尝试以计划和非阻塞方式执行一些阻塞操作(例如 HTTP 请求)。假设我有 10 个请求,一个请求需要 3 秒,但我不想等待 3 秒,而是等待 1 秒并发送下一个请求。在所有执行完成后,我想将所有结果收集到一个列表中并返回给用户。

下面是我的场景原型(线程睡眠用作阻塞操作而不是HTTP请求)

    public static List<Integer> getResults(List<Integer> inputs) throws InterruptedException, ExecutionException {
        List<Integer> results = new LinkedList<Integer>();
        Queue<Callable<Integer>> tasks = new LinkedList<Callable<Integer>>();
        List<Future<Integer>> futures = new LinkedList<Future<Integer>>();
        for (Integer input : inputs) {
            Callable<Integer> task = new Callable<Integer>() {
                public Integer call() throws InterruptedException {
                    Thread.sleep(3000);
                    return input + 1000;
                }
            };
            tasks.add(task);
        }

        ExecutorService es = Executors.newCachedThreadPool();
        ScheduledExecutorService ses = Executors.newScheduledThreadPool(1);
        ses.scheduleAtFixedRate(new Runnable() {
            @Override
            public void run() {
                Callable<Integer> task = tasks.poll();
                if (task == null) {
                    ses.shutdown();
                    es.shutdown();
                    return;
                }
                futures.add(es.submit(task));
            }
        }, 0, 1000, TimeUnit.MILLISECONDS);

        while(true) {
            if(futures.size() == inputs.size()) {
                for (Future<Integer> future : futures) {
                    Integer result = future.get();
                    results.add(result);
                }
                return results;
            }
        }
    }

    public static void main(String[] args) throws InterruptedException, ExecutionException {
        List<Integer> results = getResults(new LinkedList<Integer>(Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)));
        System.out.println(Arrays.toString(results.toArray()));
    }

我正在等待一个 while 循环,直到所有任务都返回正确的结果。但它永远不会进入中断条件,它会无限循环。每当我放置一个像记录器甚至断点这样的 I/O 操作时,它都会中断 while 循环,一切都会好起来的。

我对 Java 并发性比较陌生,并试图了解正在发生的事情以及这是否是正确的做法。我猜 I/O 操作会触发线程调度程序上的某些内容并使其检查集合的大小。

【问题讨论】:

    标签: java concurrency scheduling


    【解决方案1】:

    您需要同步您的线程。您有两个不同的线程(主线程和执行服务线程)访问futures 列表,由于LinkedList 不同步,这两个线程看到futures 的两个不同值。

    while(true) {
      synchronized(futures) {
        if(futures.size() == inputs.size()) {
          ...
        }
      }
    }
    

    这是因为 java 中的线程使用 cpu 缓存来提高性能。因此,在同步之前,每个线程都可以有不同的变量值。 这个 SO question 有更多关于这方面的信息。

    同样来自this 回答:

    一切都与记忆有关。线程通过共享内存进行通信,但是当一个系统中有多个 CPU 都试图访问同一个内存系统时,内存系统就会成为瓶颈。因此,典型的多 CPU 计算机中的 CPU 可以延迟、重新排序和缓存内存操作以加快速度。

    当线程不相互交互时,这很有效,但当它们真正想要交互时会导致问题:如果线程 A 将值存储到普通变量中,Java 不保证何时(或什至是否)线程B 会看到值的变化。

    为了在重要的时候解决这个问题,Java 为您提供了某些同步线程的方法。也就是说,让线程就程序内存的状态达成一致。 volatile关键字和synchronized关键字是线程间建立同步的两种方式。

    最后,futures 列表不会在您的代码中更新,因为主线程一直被占用,因为无限 while 块。在您的 while 循环中执行任何 I/O 操作都会为 cpu 提供足够的喘息空间来更新其本地缓存。

    无限 while 循环通常不是一个好主意,因为它非常耗费资源。在下一次迭代之前添加一个小延迟可以使它更好一点(尽管仍然效率低下)。

    【讨论】:

    • 非常感谢。我刚刚注意到你的意思。另一种解决方案可能是将futures 创建为Vector,它也是线程安全的,而不是LinkedList。我同意你关于无限while循环的观点。您知道查看列表并返回结果的更好方法吗?我应该使用Observable 吗?有什么推荐吗?
    • 我刚刚发现 ExecutorService 有函数awaitTermination 等待所有线程完成。它只是阻塞主线程,直到该 ExecutorService 管理的所有线程都完成为止。这似乎是正确的,而不是无限循环。谢谢。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多