【问题标题】:One upstream stream feeding multiple downstream streams一个上游流向多个下游流输送
【发布时间】:2015-07-13 02:53:29
【问题描述】:

我有一个一般的 Streams API 问题,我想“有效地”解决。假设我有一个(可能非常大,可能是无限的)流。我想以某种方式对其进行预处理,例如,过滤掉一些项目,并对一些项目进行变异。假设这种预处理是复杂的、时间和计算密集型的,所以我不想做两次。

接下来,我想对项目序列执行两组不同的操作,并使用不同的流类型构造处理每个不同序列的远端。对于无限流,这将是一个 forEach,对于一个有限流,它可能是一个收集器或其他什么。

显然,我可能会将中间结果收集到一个列表中,然后从该列表中拖出两个单独的流,分别处理每个流。这适用于有限流,但 a) 它看起来“丑陋”,b) 对于非常大的流来说可能不切实际,而且完全不适用于无限流。

我想我可以将 peek 用作一种“三通”。然后,我可以对 peek 下游的结果执行一系列处理,并以某种方式强制消费者在 peek 中执行“其他”工作,但现在第二条路径不再是流。

我发现我可以创建一个 BlockingQueue,使用 peek 将内容推送到该队列中,然后从队列中获取一个流。这似乎是一个好主意,实际上效果很好,尽管我不明白流是如何关闭的(它实际上是这样,但我看不到如何关闭)。这是说明这一点的示例代码:

List<Student> ls = Arrays.asList(
  new Student("Fred", 2.3F)
  // more students (and Student definition) elided ...
);

BlockingQueue<Student> pipe = new LinkedBlockingQueue<>();

ls.stream()
  .peek(s -> {
     try {
       pipe.put(s);
     } catch (InterruptedException ioe) {
       ioe.printStackTrace();
     }
   })
   .forEach(System.out::println);

   new Thread(
     new Runnable() {
       public void run() {
         Map<String, Double> map = 
           pipe.stream()
             .collect(Collectors.groupingBy(s->s.getName(),
                      Collectors.averagingDouble(s->s.getGpa())));
         map.forEach(
           (k,v)->
             System.out.println(
               "Students called " + k 
               + " average " + v));

       }
     }).start();

那么,第一个问题是:有没有“更好”的方法来做到这一点?

第二个问题,BlockingQueue 上的流到底是如何关闭的?

干杯, 托比

【问题讨论】:

  • 我认为,您所有的流都没有关闭。当流关闭时,将调用 onClose() 处理程序。注册这样的处理程序并检查它是否被调用。 Stream 不会自动关闭:您应该使用 try-with-resources 语句或手动关闭它。可能您所说的“流已关闭”是指不同的意思?

标签: java-8 java-stream


【解决方案1】:

有趣的问题。我先回答第二个问题,因为这是一个更简单的问题。

第二个问题,BlockingQueue 上的流到底是如何关闭的?

“关闭”我认为您的意思是,流具有一定数量的元素,然后它完成,不考虑将来可能添加到队列中的任何元素。原因是队列中的流仅代表创建流时队列的当前内容。它不代表任何未来元素,即其他线程可能在未来添加的元素。

如果您想要一个表示队列当前和未来内容的流,那么您可以使用other answer 中描述的技术。基本上使用Stream.generate() 调用queue.take()。不过,我不认为这是你想要做的,所以我不会在这里进一步讨论。

现在是您的更大问题。

您有一个对象来源,您希望对其进行一些处理,包括过滤。然后,您想要获取结果并通过两个不同的下游处理步骤发送它们。基本上你有一个生产者和两个消费者。

您必须处理的一个基本问题是如何处理不同处理步骤以不同速率发生的情况。假设我们已经解决了如何从队列中获取流而不使流过早终止的问题。如果生产者生产元素的速度快于消费者处理此队列中元素的速度,则队列将累积元素,直到填满所有可用内存。

您还必须决定如何以不同的速率处理不同的消费者处理元素。如果一个消费者明显比另一个消费者慢,则可能需要缓冲任意数量的元素(这可能会填满内存),或者必须放慢速度较快的消费者以匹配较慢消费者的平均速率。

让我简要介绍一下您将如何进行。不过,我不知道你的实际要求,所以我不知道这是否会令人满意。需要注意的一点是,在这种应用程序中使用并行流可能会出现问题,因为并行流不能很好地处理阻塞和负载平衡。

首先,我将从生产者的流处理元素开始,并将它们累积到ArrayBlockingQueue

BlockingQueue<T> queue = new ArrayBlockingQueue<>(capacity);
producer.map(...)
        .filter(...)
        .forEach(queue::put);

(注意put会抛出InterruptedException,所以你不能只把queue::put放在这里。你必须在这里放一个try-catch块,或者写一个帮助方法。但是要做什么并不明显如果InterruptedException 被抓到就做。)

如果队列已满,这将阻塞管道。要么在它自己的线程中按顺序运行,要么在并行的情况下在专用线程池中运行,以避免阻塞公共池。

接下来是消费者:

while (true) {
    // wait until the queue is full, or a timeout has expired,
    // depending upon how frequently you want to continue
    // processing elements emitted by the producer
    List<T> list = new ArrayList<>();
    queue.drainTo(list);
    downstream1 = list.stream().filter(...).map(...).collect(...);
    downstream2 = list.stream().filter(...).map(...).collect(...);
    // deal with results downstream1 and downstream2
}

这里的关键是从生产者到消费者的切换是使用drainTo 方法分批完成的,该方法将队列的元素添加到目标并原子地清空队列。这样,消费者不必等待生产者完成其处理(如果它是无限的,则不会发生)。此外,消费者正在对已知数量的数据进行操作,并且不会在处理过程中阻塞。因此,如果有帮助的话,每个消费者流都可以并行运行。

在这里,我让消费者​​同步运行。如果您希望消费者以不同的速率运行,则必须构建额外的队列(或其他东西)来独立缓冲他们的工作负载。

如果消费者总体上比生产者慢,队列最终将被填满并被阻塞,将生产者减慢到消费者可以接受的速度。如果消费者平均比生产者快,那么也许你不需要担心消费者的相对处理速度。您可以让它们循环并拾取生产者设法放入队列的任何内容,甚至可以让它们阻塞直到有可用的东西。

我应该说,我所概述的是一种非常简单的多阶段流水线方法。如果您的应用程序对性能至关重要,您可能会发现自己在调整内存消耗、负载平衡、增加吞吐量和减少延迟方面做了大量工作。还有其他框架可能更适合您的应用程序。例如,您可以看看LMAX Disruptor,虽然我自己没有任何经验。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-09-07
    • 2019-04-06
    • 2019-01-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多