
RxJava 实战问题求解用 Observable 操作符链解决 Euler 求和与斐波那契数列生成【免费下载链接】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本篇文章围绕 RxJava 仓库中 Problem-Solving-Examples-in-RxJava.md 提供的两道经典编程谜题展开用 RxJava 操作符求解 Project Euler 第 1 题3 与 5 的倍数求和以及用纯操作符组合生成斐波那契数列。通过逐步拆解range、filter、map、takeWhile、merge、distinct、reduce、create、repeat、scan、take等操作符的调用链读者将掌握用声明式数据流思维替代命令式循环的 RxJava 解题范式并理解这些操作符在当前仓库源码中的真实实现行为。从两道谜题理解 RxJava 的思考方式RxJava 的核心思想是把异步与事件驱动的程序组织成可组合的可观察序列observable sequences。传统的命令式写法用循环和中间变量描述怎么做而 RxJava 用操作符链描述要什么——数据源从上游流出经过一个个操作符的变换、过滤、合并、聚合最终被订阅者消费。下面的两个例子恰好覆盖了 RxJava 操作符的几大典型类别创建型range、create、fromArray、repeat变换型map、scan过滤型filter、takeWhile、take、distinct合并型merge聚合型reduceProject Euler 第 1 题3 或 5 的倍数求和原题是这样的热身题如果我们列出 10 以下所有是 3 或 5 倍数的自然数会得到 3、5、6 和 9。这些数的和是 23。请找出 1000 以下所有是 3 或 5 的倍数的自然数之和。答案是 233168。下面用两种 RxJava 思路分别求解。思路一rangefilter从全集里筛最直接的做法是枚举 1000 以下的所有自然数然后用filter保留那些是 3 或 5 的倍数的项ObservableInteger threesAndFives Observable.range(1, 999).filter(e - e % 3 0 || e % 5 0);在 Groovy 中写法几乎一样条件取反来兼容 Groovy 的 truthinessdef threesAndFives Observable.range(1,999).filter({ !((it % 3) (it % 5)) });其中range(start, count)是 RxJava 的静态创建操作符range(1, 999)会按顺序发射 1 到 999 共 999 个整数。从当前仓库源码看Observable.range 会对参数做严格校验count 0直接抛出IllegalArgumentException(count 0 required but it was ...)count 0返回empty()count 1退化为just(start)并且检测start count - 1是否会溢出Integer.MAX_VALUE最终包装为内部的ObservableRange并经过RxJavaPlugins.onAssembly装配后返回。filter则接收一个Predicate? super TObservable.filter 对每个通过谓词判断的元素向下游发射其余元素被丢弃。思路二maptakeWhilemergedistinct从倍数集合合并另一种做法是构造两个序列一个专门生成 3 的倍数一个专门生成 5 的倍数。方法是先用map把每个自然数映射成它的倍数再用takeWhile在倍数小于 1000 时截断最后用merge合并两条流ObservableInteger threes Observable.range(1, 999).map(e - e * 3).takeWhile(e - e 1000); ObservableInteger fives Observable.range(1, 999).map(e - e * 5).takeWhile(e - e 1000); ObservableInteger threesAndFives Observable.merge(threes, fives).distinct();def threes Observable.range(1,999).map({it*3}).takeWhile({it1000}); def fives Observable.range(1,999).map({it*5}).takeWhile({it1000}); def threesAndFives Observable.merge(threes, fives).distinct();两个细节值得注意takeWhile在这里是边界守卫map会把 1..999 映射成 3..2997但我们只关心 1000 以下的倍数takeWhile(e - e 1000)一旦遇到第一个不满足条件的值就停止发射并完成。这也是 Observable.takeWhile 的标准语义——逐项检查谓词遇到首个为false的项即终止序列。不能漏掉distinct15、30 这类数同时是 3 和 5 的倍数会同时出现在两条流里merge只是把两条流的时间线交错合并并不会去重。distinct()见 Observable.distinct内部维护一个HashSet保证每个值只向下游发射一次。merge的实现也值得看一下静态重载 Observable.merge(ObservableSource? extends ObservableSource? extends T) 在底层复用了ObservableFlatMap以Functions.identity()作为映射函数、Integer.MAX_VALUE作为最大并发订阅数——也就是说merge本质上是不做变换的无限并发 flatMap上游多条流交错地向下游推数据。求和reduce把序列折叠成一个数现在得到了去重后的倍数序列接下来要把它求和。如果你安装了可选的rxjava-math模块可以直接用sumInteger/sumLong之类的数学操作符见 Mathematical-and-Aggregate-Operators。但如果不想依赖这个模块标准 RxJava 里最合适的聚合操作符就是reduceSingleInteger summer threesAndFives.reduce(0, (a, b) - a b);def summer threesAndFives.reduce(0, { a, b - ab });注意这里的返回类型当前仓库中 Observable.reduce(R seed, BiFunction) 返回的是SingleR不是Observable——这正是 RxJava 3/4 时代的类型设计一个确定会恰好发射一个结果并完成的流用Single表达语义更准确。其底层被包装为ObservableReduceSeedSingle。reduce的折叠过程如下以0作为种子值seed每当上游threesAndFives发射一个元素b就调用累加函数(a, b) - a b把当前种子a和发射值相加返回的新值覆盖种子当上游完成时把最终的种子值作为唯一一次发射推给下游。下表演示了折叠过程的前几步与最后一步整个序列共有 466 个去重后的倍数iterationseedemissionreduce1033235838614…………466232169999233168订阅输出到目前为止我们只是搭建了一条管道还没有人消费它。RxJava 是惰性的只有调用subscribe才会真正触发数据流动summer.subscribe(System.out::print);summer.subscribe({println(it);});Single.subscribe(Consumer)在成功时把唯一的结果233168交给回调因此上面这行会直接打印出答案。这正体现了构建数据流与订阅数据流两个阶段的分工。生成斐波那契数列第二个谜题是如何创建一个发射斐波那契数列的Observable方法一create 循环从头造一个 Observable最直接的方式是用create操作符从零开始制造一个Observable在传给它的回调里用传统循环生成序列ObservableInteger fibonacci Observable.create(emitter - { int f1 0, f2 1, f 1; while (!emitter.isDisposed()) { emitter.onNext(f); f f1 f2; f1 f2; f2 f; } });def fibonacci Observable.create({ emitter - def f10, f21, f1; while(!emitter.isDisposed()) { emitter.onNext(f); f f1f2; f1 f2; f2 f; }; });这种写法能工作但它太像普通线性编程了——本质上还是命令式的循环只是把输出换成了emitter.onNext(...)。而且注意这个循环只有在emitter.isDisposed()即订阅者取消订阅时才会停下所以在实际使用中一定要配合take(n)限制发射数量否则它会无限发射下去。从源码看Observable.create 接收一个ObservableOnSubscribeT这里用 lambda 实现通过Objects.requireNonNull拒绝空值并包装为内部类ObservableCreate后经插件装配返回。emitter.isDisposed()对应ObservableEmitter的取消订阅状态检查——这正是响应式取消在自定义序列里的典型用法。方法二纯操作符组合用repeatscanmap驱动状态机有没有更Rx 风格的做法文档给出的方案是用一个只发射 0 的无限序列作为磨盘驱动scan这个状态机不断推进ObservableInteger fibonacci Observable.fromArray(0) .repeat() .scan(new int[]{0, 1}, (a, b) - new int[]{a[1], a[0] a[1]}) .map(a - a[1]);def fibonacci Observable.from(0).repeat().scan([0,1], { a,b - [a[1], a[0]a[1]] }).map({it[1]});这个方法略显取巧原文称之为 janky但拆开看每一步都很有教学价值fromArray(0).repeat()先发射一个0然后repeat()反复重放这个序列Observable.repeat 在源完成后自动重新订阅于是得到一条无限长的全 0 流。这些 0 本身毫无意义纯粹是给scan提供事件节拍。scan(new int[]{0, 1}, (a, b) - new int[]{a[1], a[0] a[1]})scan与reduce类似都带种子值并逐项累积但区别在于scan把每一步的中间结果都发射出去而reduce只在最后发射一次。这里的累加器完全忽略上游传来的b反正都是 0只把寄存器从[a, b]变换为[b, ab]。带初值的重载对应源码 Observable.scan(R initialValue, BiFunction)。map(a - a[1])寄存器按如下顺序演变[0,1] → [1,1] → [1,2] → [2,3] → [3,5] → [5,8] → …数组中第二个元素恰好就是斐波那契数列 1, 1, 2, 3, 5, 8, …用map把它单独取出来即可。也就是说scan在这里扮演了一个随每次输入变换内部状态并对外发布状态快照的有状态变换器这正是用操作符组合替代显式循环变量的关键技巧。打印前 15 项无论用哪种方法生成输出前 15 项都一样简单fibonacci.take(15).subscribe(System.out::println);fibonnaci.take(15).subscribe({println(it)})];take(15)在发射了 15 个元素后立即让上游完成并断开订阅见 Observable.take这对于上述无限序列至关重要——它既是只取一部分的需求表达也是防止无限循环泄漏资源的机制。有没有更优雅的方案create版本太命令式repeatscan版本又有点绕它专门造了一条毫无信息量的全 0 流只是为了让scan的曲柄转起来。原文档提到一个名为generate的操作符本可以更优雅地解决这个问题——它允许用一个状态 生成函数直接描述迭代式序列无需借助无意义的占位流但这个操作符在当时尚未并入 RxJava。作为练习你可以尝试思考更简洁的组合方式例如用concatMap配合状态对象或在create回调里维护状态但保持声明式风格。这道开放题本身就是理解 RxJava 操作符组合力的最好训练。从两道题看 RxJava 的通用解题套路把上面的两个例子放在一起可以提炼出 RxJava 处理数学计算 事件流问题的通用模式选择数据源有限范围用range/fromArray无限序列用create或repeat 占位流。变换与筛选map做一对一映射filter按谓词剔除takeWhile/take控制序列边界distinct去重。合并多源merge交错多条流底层是flatMap identity 映射默认无限并发。聚合收口要每一步都看中间结果用scan只要最终一个结果用reduce注意当前仓库中带种子值的reduce返回Single。订阅触发所有链式构建都是惰性的subscribe才是数据流动的起点对无限序列务必配合take之类的终止操作符。这些操作符的语义细节参数校验、返回值类型、惰性机制都可以在 src/main/java/io/reactivex/rxjava4/core/Observable.java 中找到对应实现配套的测试用例如 flowable/FlowableReduceTests.java、flowable/FlowableConcatTests.java则从行为层面验证了这些语义。如果你想继续深入仓库的 Observable.md 汇总了订阅协议Creating-Observables、Transforming-Observables、Filtering-Observables、Combining-Observables 与 Mathematical-and-Aggregate-Operators 则分别覆盖了本文用到的各类操作符。【免费下载链接】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),仅供参考