ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

RxJava 数学与聚合算子详解:count、reduce、collect、toList、toMap 全解析

RxJava 数学与聚合算子详解:count、reduce、collect、toList、toMap 全解析 RxJava 数学与聚合算子详解count、reduce、collect、toList、toMap 全解析【免费下载链接】RxJavaRxJava – Reactive Extensions for the JVM – a library for composing asynchronous and event-based programs using observable sequences for the Java VM.项目地址: https://gitcode.com/gh_mirrors/rx/RxJava本文以 docs/Mathematical-and-Aggregate-Operators.md 为骨架系统讲解 RxJava 中作用于“整条序列”的数学与聚合算子包括来自rxjava2-extensions扩展模块的 8 个数学算子averageDouble、averageFloat、max、min、sumDouble、sumFloat、sumInt、sumLong以及 10 个标准聚合算子count、reduce、reduceWith、collect、collectInto、toList、toSortedList、toMap、toMultimap。读完后你将理解这些算子为何必须缓冲整个序列、返回Single/Maybe与返回Observable/Flowable的设计差异并能结合 Observable.java 的源码与内部实现类正确选型和使用这些算子。聚合算子的共同前提必须等上游完整结束这类算子对Observable或Flowable发出的整个物品序列执行数学或其他操作。由于它们必须等待源序列发出所有物品甚至完成信号之后才能构造自己的发射值而且通常必须缓冲这些物品因此在可能拥有非常长甚至是无限序列的流上使用这类算子是有危险的。这一限制在所有聚合算子上都成立使用时请注意对无限流如Observable.never()、interval、持续推送的 WebSocket 数据调用toList()或reduce()缓冲区会持续增长直至内存耗尽对超长流应优先考虑分段聚合配合buffer、window先切块再聚合或流式近似方案。另一个贯穿全文的关键设计差异标准聚合算子count、reduce等返回Single或Maybe因为输出物品的数量恒定为至多一个数学算子MathObservable/MathFlowable提供的一组则返回Observable和Flowable这是两者签名上最根本的不同。数学算子来自 rxjava2-extensions 模块本节的算子属于RxJava2Extensions项目并非 RxJava 核心库自带。需要将rxjava2-extensions模块作为依赖加入项目groupId 为com.github.akarnokd。以下示例假设已引入MathObservable和MathFlowable类import hu.akarnokd.rxjava2.math.MathObservable; import hu.akarnokd.rxjava2.math.MathFlowable;这组算子均可用于Flowable与Observable但不适用于Maybe、Single、Completable。averageDouble计算Observable发出的Number的平均值并以Double形式发出该平均值。ObservableInteger numbers Observable.just(1, 2, 3); MathObservable.averageDouble(numbers).subscribe((Double avg) - System.out.println(avg)); // prints 2.0averageFloat计算Observable发出的Number的平均值并以Float形式发出该平均值。ObservableInteger numbers Observable.just(1, 2, 3); MathObservable.averageFloat(numbers).subscribe((Float avg) - System.out.println(avg)); // prints 2.0max发出源Observable发出的最大值。可以指定一个Comparator用于比较源发出的元素。ObservableInteger numbers Observable.just(4, 9, 5); MathObservable.max(numbers).subscribe(System.out::println); // prints 9指定Comparator可以在源中查找最长的Stringfinal ObservableString names Observable.just(Kirk, Spock, Chekov, Sulu); MathObservable.max(names, Comparator.comparingInt(String::length)) .subscribe(System.out::println); // prints Chekovmin发出源Observable发出的最小值。同样可以指定Comparator自定义比较逻辑。ObservableInteger numbers Observable.just(4, 9, 5); MathObservable.min(numbers).subscribe(System.out::println); // prints 4sumDouble将Observable发出的Double相加并发出这个和。ObservableDouble numbers Observable.just(1.0, 2.0, 3.0); MathObservable.sumDouble(numbers).subscribe((Double sum) - System.out.println(sum)); // prints 6.0sumFloat将Observable发出的Float相加并发出这个和。ObservableFloat numbers Observable.just(1.0F, 2.0F, 3.0F); MathObservable.sumFloat(numbers).subscribe((Float sum) - System.out.println(sum)); // prints 6.0sumInt将Observable发出的Integer相加并发出这个和。ObservableInteger numbers Observable.range(1, 100); MathObservable.sumInt(numbers).subscribe((Integer sum) - System.out.println(sum)); // prints 5050sumLong将Observable发出的Long相加并发出这个和。ObservableLong numbers Observable.rangeLong(1L, 100L); MathObservable.sumLong(numbers).subscribe((Long sum) - System.out.println(sum)); // prints 5050注意从源码结构看这些sum*算子在MathObservable/MathFlowable中是按数值类型分别提供的本仓库中搜索不到MathObservable、averageDouble等相关定义证实它们来自外部的rxjava2-extensions模块选择时应与源元素类型匹配以避免不必要的装箱/拆箱与类型转换。标准聚合算子注意这些标准聚合算子返回Single或Maybe因为输出物品的数量恒定为至多一个。count计算Observable发出的物品数量并以Long形式发出该计数。Maybe同样提供该算子Single与Completable不提供。Observable.just(1, 2, 3).count().subscribe(System.out::println); // prints 3在 Observable.java 中count()是一个final方法直接构造内部实现类。从源码结构看核心库为计数场景准备了两个实现ObservableCount.java返回Observable的形态与 ObservableCountSingle.java返回SingleLong并实现FuseToObservable支持融合优化。计数本身无需缓冲任何物品只维护一个Long计数器因此它是聚合算子中内存开销最小的一类。reduce按顺序将函数应用于每个发出的物品只发出最终累积值。reduce(seed, reducer)以一个初始种子值启动累积过程。Observable.range(1, 5) .reduce((product, x) - product * x) .subscribe(System.out::println); // prints 120在 Observable.java 中签名形如public final NonNull R SingleR reduce(R seed, NonNull BiFunctionR, ? super T, R reducer)其内部只需要在内存中保存一个“当前累积值”并不缓冲整个序列——这是reduce相对toList的重要优势对超长序列做纯累积运算求和、求最大、拼接字符串等时reduce的内存占用是 O(1) 的而先toList再处理则是 O(n) 的。reduceWith与reduce相同——按顺序应用函数、只发出最终累积值——区别在于初始值由一个Supplier提供Observable.java#L8985public final NonNull R SingleR reduceWith(NonNull SupplierR seedSupplier, NonNull BiFunctionR, ? super T, R reducer)这使得每个订阅者都可以拥有独立创建的累积容器。Observable.just(1, 2, 2, 3, 4, 4, 4, 5) .reduceWith(TreeSet::new, (set, x) - { set.add(x); return set; }) .subscribe(System.out::println); // prints [1, 2, 3, 4, 5]collect将源Observable发出的物品收集进一个由你创建的可变数据结构返回发出该结构的SingleObservable.java#L5455。Observable.just(Kirk, Spock, Chekov, Sulu) .collect(() - new StringJoiner( \uD83D\uDD96 ), StringJoiner::add) .map(StringJoiner::toString) .subscribe(System.out::println); // prints Kirk Spock Chekov Sulu签名中第一个参数是Supplier? extends U initialItemSupplier第二个参数是BiConsumer? super U, ? super T collector。与reduceWith不同collect的回调不需要返回值——容器就地修改即可代码更简洁。Java 8 还可用Collector风格的collectObservable.java#L14495其实现位于 ObservableCollectWithCollector.java 与 ObservableCollectWithCollectorSingle.java可以复用 JDK 标准库中的Collector实例。collectInto与collect类似但初始容器由调用方直接传入而非由Supplier创建Observable.java#L5490。注意将用于收集物品的可变值如下例的StringBuilder会在多个订阅者之间共享。Observable.just(R, x, J, a, v, a) .collectInto(new StringBuilder(), StringBuilder::append) .map(StringBuilder::toString) .subscribe(System.out::println); // prints RxJava这一“共享可变状态”的特性意味着collectInto更适合一次性消费的冷场景若被多次订阅多个订阅者会并发写入同一容器。从源码结构看collect的返回类型是SingleU每个订阅触发新建容器而collectInto同样返回Single但对热源或重复订阅需要自行处理容器的可见性与线程安全。toList收集Observable的所有物品作为单个List发出Observable.java#L12758。Observable.just(2, 1, 3) .toList() .subscribe(System.out::println); // prints [2, 1, 3]除无参版本外还有两个实用重载toList(int capacityHint)Observable.java#L12791提示内部ArrayList的初始容量在已知数据规模时可减少扩容开销toList(SupplierU collectionSupplier)Observable.java#L12826用自定义集合实现如LinkedList、TreeSet等替代默认ArrayList。toSortedList收集所有物品并作为已排序的List发出。默认使用自然排序也可指定Comparator与capacityHinttoSortedList()Observable.java#L13168内部使用Functions.naturalComparator()toSortedList(Comparator)Observable.java#L13196toSortedList(Comparator, int capacityHint)Observable.java#L13229toSortedList(int capacityHint)Observable.java#L13264。Observable.just(2, 1, 3) .toSortedList(Comparator.reverseOrder()) .subscribe(System.out::println); // prints [3, 2, 1]toMap将源发出的序列转换为以指定 key 函数为键的Map。提供三个重载toMap(keySelector)Observable.java#L12858物品本身作为 valuetoMap(keySelector, valueSelector)Observable.java#L12893分别指定 key 与 value 选择器toMap(keySelector, valueSelector, mergeFunction)Observable.java#L12930遇到重复 key 时用 merge 函数合并而不是抛异常。Observable.just(1, 2, 3, 4) .toMap((x) - { // defines the key in the Map return x; }, (x) - { // defines the value that is mapped to the key return (x % 2 0) ? even : odd; }) .subscribe(System.out::println); // prints {1odd, 2even, 3odd, 4even}使用提示如果源中可能出现重复 key 且你不关心合并规则应选择带mergeFunction的重载避免下游收到IllegalStateException。toMultimap将源发出的序列转换为一个既是Collection又是Map的结构——相同 key 对应多个 value天然规避了toMap的重复 key 问题。同样提供单/双选择器及带容器工厂的重载Observable.java#L12964 至 Observable.java#L13080。Observable.just(1, 2, 3, 4) .toMultimap((x) - { // defines the key in the Map return (x % 2 0) ? even : odd; }, (x) - { // defines the value that is mapped to the key return x; }) .subscribe(System.out::println); // prints {even[2, 4], odd[1, 3]}源码级实现与选型小结从源码结构看上述标准聚合算子在核心库Observable与Flowable中是对称提供的Flowable对应实现位于src/main/java/io/reactivex/rxjava4/internal/operators/flowable/下同名类且聚合路径普遍实现了FuseToObservable如 ObservableCountSingle.java、ObservableCollectSingle.java即在与相邻算子组合时可以跳过中间队列、直接同步传递结果降低订阅延迟。选型时可以按以下决策链快速判断需求推荐算子内存特征返回类型只统计条数count()O(1)不缓冲物品SingleLong纯累积运算和、积、最大、拼接reduce/reduceWithO(1) 累积值SingleR收集到自定义可变结构StringJoiner等collect/collectIntoO(n)SingleU收集为标准List可指定容量/集合类型toListO(n)SingleListT需要排序结果toSortedListO(n)SingleListT一对一按键索引toMap重 key 用 merge 重载O(n)SingleMapK, V一对多按键分组toMultimapO(n)SingleMapK, CollectionV数值型平均/求和核心库未覆盖的类型MathObservable/MathFlowable的average*、sum*O(1)~O(n)Observable/Flowable最后重申使用边界所有聚合算子都要等上游终止才能产出结果并通常缓冲整个序列count、reduce这类仅需累积值的除外因此对长流或无限流应谨慎使用需要外部依赖的数学算子rxjava2-extensions请确认项目中已引入该模块而count、reduce、toList、toMap、toMultimap、toSortedList、collect、collectInto均为 RxJava 核心库自带Maybe额外提供count。【免费下载链接】RxJavaRxJava – Reactive Extensions for the JVM – a library for composing asynchronous and event-based programs using observable sequences for the Java VM.项目地址: https://gitcode.com/gh_mirrors/rx/RxJava创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表