use*_*988 5 java java-8 java-stream
我正在尝试了解 Java 的 Stream API 的内部调用。
我有以下代码,它有两个过滤器(中间)操作和一个终端操作。
IntStream.of(1,2,3)
.filter(e->e%2==0)
.filter(e->e==2)
.forEach(e->System.out.println(e));
Run Code Online (Sandbox Code Playgroud)
Stream -> 返回带有覆盖过滤器的 Stream -> 返回带有覆盖过滤器的 Stream -> 终端
我看到对于每个中间操作,都会返回一个带有重写filter方法的新流。一旦命中终端方法,流就会执行filter. 我看到filter()如果有两个filter 操作而不是一次,它会运行两次。
我想了解一次流遍历如何能够两次调用过滤器。
粘贴下面的 IntPipeline 代码,该代码在 Stream 中被过滤器方法命中。
@Override
public final IntStream filter(IntPredicate predicate) {
Objects.requireNonNull(predicate);
return new StatelessOp<Integer>(this, StreamShape.INT_VALUE,
StreamOpFlag.NOT_SIZED) {
@Override
Sink<Integer> opWrapSink(int flags, Sink<Integer> sink) {
return new Sink.ChainedInt<Integer>(sink) {
@Override
public void begin(long size) {
downstream.begin(-1);
}
@Override
public void accept(int t) {
if (predicate.test(t)) ///line 11
downstream.accept(t);
}
};
}
};
}
Run Code Online (Sandbox Code Playgroud)
在filter()返回一个新的流,其谓语被设置为e%2==0,然后再返回一个新的流,其谓语e==2。一旦终端操作被击中,对于每次遍历,谓词代码就会在第 11 行执行。
编辑:我看到它downstream用于将中间操作链接为 LinkedList。那么所有实现都添加到linkedlist前一阶段并在遍历开始时调用?
我认为您对 Streams 的理解有一些负担和混合概念。
您的困惑与过滤器(或任何其他操作)实现无关;也不存在覆盖中间(或任何其他)流操作的概念;
覆盖是一个完全不同的概念,它与继承有关
每个流都是一个单独的管道,它有 (1)开始,(2) 可选的中间部分,以及 (3)结束/结论;
中间操作不被覆盖;相反,它们形成一个连续的操作链,流的每个元素都必须按照相同的顺序进行(除非某个元素在某些中间操作中被丢弃)。
将 K 个对象的 Stream 视为一个管道,其中 K 的实例将通过。该管道有一个开始/源(对象进入管道的地方)和结束/目的地;然而,在这两者之间,可能有几个中间操作将(或者)过滤、变换、装箱等这些对象。流是惰性的,这意味着在调用终端操作之前不会执行中间操作(这是流的一个很大的特性),因此,当调用终端操作时,每一项,一次一个,将通过这个管道.
此外,请阅读此代码段:
为了执行计算,流操作被组合成一个流管道。流管道由一个源(可能是数组、集合、生成器函数、I/O 通道等)、零个或多个中间操作(将流转换为另一个流,例如
filter(Predicate))和一个终端操作(产生结果或副作用,例如count()或forEach(Consumer))。流是懒惰的;对源数据的计算只在终端操作启动时进行,源元素只在需要时被消费。
请记住范式,该 Stream 包括:
如果您仍然感到困惑(现在不应该是这种情况),您还可以参考一些重要的观点,它们可能会更清楚地了解您的困惑:
中间操作返回一个新的流。他们总是很懒惰;执行诸如 filter() 之类的中间操作实际上并不执行任何过滤,而是创建一个新流,该流在遍历时包含与给定谓词匹配的初始流的元素。管道源的遍历直到管道的终端操作执行完毕才开始;
延迟处理流可以显着提高效率;在诸如上面的 filter-map-sum 示例之类的管道中,过滤、映射和求和可以融合为数据的单次传递,中间状态最少。懒惰还可以避免在不必要时检查所有数据;对于诸如“查找第一个长度超过 1000 个字符的字符串”之类的操作,只需检查足够多的字符串即可找到具有所需特征的字符串,而无需检查源中所有可用的字符串。(当输入流是无限的而不仅仅是大时,这种行为变得更加重要。)
我将尽力解释 Stream API 幕后发生的事情,首先你应该改变你对到目前为止的编程方式的看法,尝试获得这个新想法。
因此,举一个现实世界想象工厂的例子(我的意思是现实世界中的真实工厂,而不是工厂设计模式),在工厂中,我们有一些原材料和不同阶段的一些连续流程,将原材料转化为成品。要掌握这个概念,请参见下图:
(stage1)原材料 -> (stage2)处理输入并将输出传递到下一个阶段 -> (stage3)处理输入并将输出传递到下一个阶段 -> .....
因此,第一阶段的输出是原材料,所有后续阶段对其输入进行一些处理并将其传递(例如,它可以将阶段的输入转换为其他内容,或者由于其质量低而完全拒绝该输入)然后它将把输出交给它前面的另一个阶段。从现在开始,我们将这个连续的阶段统称为管道。
什么是处理?,它可以是任何东西,例如一个阶段可以决定将输入转换为完全不同的东西并将其传递(这正是
mapStream API 中为我们提供的),另一个阶段可能允许输入基于传递在某些条件下(这正是filterStream API 中所做的)。
Java Stream API 的作用类似于工厂。每个 Stream 都是一个管道,您可以向每个管道添加另一个阶段并创建一个新管道,因此当您编写时,IntStream.of(1,2,3)您已经创建了一个管道,因此IntStream让我们分解您的代码:
IntStream intStream = IntStream.of(1,2,3)
Run Code Online (Sandbox Code Playgroud)
这相当于我们工厂里的原材料,因此它是一条只有一级的管道。然而,仅通过原材料的管道没有任何好处。让我们在之前的管道中添加另一个阶段并创建一个新的管道:
IntStream evenNumbePipeline = intStream.filter(e -> e%2==0);
Run Code Online (Sandbox Code Playgroud)
请注意,在这里您创建了新的管道,并且该管道正是前一个管道加上另一个仅允许偶数通过并拒绝其他管道的阶段。当您调用过滤器方法时,以下部分代码将创建一个新的管道:
@Override
public final IntStream filter(IntPredicate predicate) {
Objects.requireNonNull(predicate);
return new StatelessOp<Integer>(this, StreamShape.INT_VALUE,
StreamOpFlag.NOT_SIZED) {
@Override
Sink<Integer> opWrapSink(int flags, Sink<Integer> sink) {
return new Sink.ChainedInt<Integer>(sink) {
@Override
public void begin(long size) {
downstream.begin(-1);
}
@Override
public void accept(int t) {
if (predicate.test(t)) ///line 11
downstream.accept(t);
}
};
}
};
}
Run Code Online (Sandbox Code Playgroud)
您可以看到过滤器返回扩展 IntPipeline 的新实例,StatelessOp<Integer>如下所示:
abstract static class StatelessOp<E_IN> extends IntPipeline<E_IN>
Run Code Online (Sandbox Code Playgroud)
让我们暂时停一下,问一个问题:到目前为止有进行过任何操作吗?答案是否定的,当你创建一个工厂或工厂的管道时,还没有生产任何产品,你应该向管道提供原材料以通过工厂获得成品,但到目前为止我们还没有这样做。因此,当我们在流中调用过滤器和其他操作时,我们只是在设计管道程序,我们并没有真正处理任何内容,我们只是在管道中添加另一个阶段,然后说嘿,当您收到输入时,您应该执行此程序在上面,
在我们的例子中,我们将 stage2 添加到我们的工厂中,并告诉它当您收到输入时检查它是否为偶数,然后如果它是偶数则允许它通过。我们现在正在回答您的问题,让我们在管道中添加另一个阶段:
IntStream onlyTwo = evenNumbePipeline.filter(e -> e==2);
Run Code Online (Sandbox Code Playgroud)
在这里,您创建新的管道,它获取先前的管道(即evenNumbePipeline)并向该管道添加另一个阶段(evenNumbePipeline 没有更改,我们创建其中包含 EvenNumbePipeline 的新管道)。让我们看一下到目前为止我们的管道:
raw material(stage1) -> filter even number(stage2) -> filter only 2(stage3)
Run Code Online (Sandbox Code Playgroud)
将其视为我们管道中各个阶段的定义,而不是操作,也许我们还没有原材料,但我们可以设计我们的工厂,以便以后可以为其提供原材料。您可以看到该管道具有三个阶段,每个阶段都会对前一个阶段的输出执行一些操作。管道将由原材料一一提供(现在忘记并行流),因此当您向该管道提供 1 作为原材料时,它会经历这些阶段。这个阶段的每个阶段都是java中的new Object。
那么让我们来谈谈你在问题中所说的话
我想了解一次流遍历如何能够调用过滤器两次。
从我们到目前为止的调查来看,您认为 Stream 的遍历调用了两次过滤器方法,还是当我们创建管道时我们调用了两次过滤器方法?
我们调用此filter方法两次,因为我们希望管道中有两个不同的阶段。考虑一个工厂,我们两次调用过滤器方法,因为我们在设计工厂时希望有两个不同的过滤器阶段。我们还没有进入工厂阶段,也还没有生产出任何成品。
让我们玩得开心并产生一些输出:
onlyTwo.forEach(e -> System.out.println(e));
Run Code Online (Sandbox Code Playgroud)
作为foreach终端运营,它启动了我们的工厂并为我们工厂的管道提供原材料。因此,例如1先经过stage2,然后经过stage3,然后传递到foreach语句。
但还有一个问题:当我们设计管道时,如何定义每个阶段的作用?
在我们的示例中,当我们设计管道时,我们调用过滤器方法来创建新阶段并将过滤器阶段应该执行的操作作为参数传递给它,该参数的名称是谓词。当(每个阶段)收到输入时确切执行的操作由opWrapSink每个阶段的方法定义。因此,我们在创建阶段时必须实现此方法,所以让我们回到过滤器方法,其中 Stream 类为我们创建新的阶段:
@Override
public final IntStream filter(IntPredicate predicate) {
Objects.requireNonNull(predicate);
return new StatelessOp<Integer>(this, StreamShape.INT_VALUE,
StreamOpFlag.NOT_SIZED) {
@Override
Sink<Integer> opWrapSink(int flags, Sink<Integer> sink) {
return new Sink.ChainedInt<Integer>(sink) {
@Override
public void begin(long size) {
downstream.begin(-1);
}
@Override
public void accept(int t) {
if (predicate.test(t)) ///line 11
downstream.accept(t);
}
};
}
};
}
Run Code Online (Sandbox Code Playgroud)
您可以看到每个阶段的opWrapSink方法都返回 aSink但什么是Sink?
为了消除此接口中的大量复杂性,它是一个消费者,并具有如下的接受方法(它还有许多其他用于原始类型的接受方法,以避免不必要的装箱和拆箱):
void accept(T t);
Run Code Online (Sandbox Code Playgroud)
当您实现此接口时,您应该定义要如何处理将在阶段中作为输入传递的输入值。您不需要在程序中实现此接口,因为Stream实现中的方法已经为您完成了繁重的工作。让我们看看它在过滤器案例中是如何实现的:
@Override
Sink<Integer> opWrapSink(int flags, Sink<Integer> sink) {
return new Sink.ChainedInt<Integer>(sink) {
@Override
public void begin(long size) {
downstream.begin(-1);
}
@Override
public void accept(int t) {
if (predicate.test(t)) ///line 11
downstream.accept(t);
}
};
}
Run Code Online (Sandbox Code Playgroud)
Stream框架提供了opWrapSink下一阶段的方法Sink(作为调用该方法时的第二个参数),这意味着我们知道Pipeline中的下一阶段如何完成他们的工作(在this的帮助下Sink),但我们应该为他们提供一个输入,它很明显,下一级的输入是当前阶段的输出。产生当前阶段的输出所需的另一个参数是当前阶段的输入。
输入到当前阶段 -> 在输入上执行当前阶段的操作 -> 将输出传递到下一个阶段(管道中的后续阶段)
因此,在accept方法中,我们将当前阶段的输入作为参数,t我们应该对此输入执行一些操作(作为当前阶段对输入的操作),然后将其传递到下一个阶段。在我们的过滤阶段,我们需要检查阶段的输入是否通过ta predicate(在我们的例子中是 e%2==0),然后我们应该将其传递到下一个阶段Sink。这正是我们的接受方法所做的(下游正是Sink管道中以下阶段的):
@Override
public void accept(int t) {
if (predicate.test(t)) ///line 11
downstream.accept(t);
}
Run Code Online (Sandbox Code Playgroud)
在这个方法实现中你应该注意到的accept是,如果它传递了一个谓词(在我们的例子中是 e%2==0),并且如果它确实传递了当前阶段的输入(即 t),它只会将当前阶段的输入(即 t)传递到下一个阶段。 not pass the predicate it does not pass it through(这正是我们期望过滤阶段要做的事情);
| 归档时间: |
|
| 查看次数: |
458 次 |
| 最近记录: |