尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Akka Streams `Source.range` 操作符完全指南:生成整数序列、步长控制与底层实现原理

Akka Streams `Source.range` 操作符完全指南:生成整数序列、步长控制与底层实现原理 后端并发编程异步编程【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址https://gitcode.com/gh_mirrors/ak/akka-core点击查看免费下载Source.range是 Akka Streams 中用于快速创建整数数据流的核心工厂方法它将一个闭区间内的整数逐个发射为流元素并支持自定义步长包括负步长实现倒序。本文基于当前仓库的官方文档、Java/Scala 双 DSL 源码及测试用例完整讲解Source.range的依赖配置、三种典型调用形态、Reactive Streams 语义、底层Range.inclusive实现机制与边界行为帮助你直接将其运用于分页、批处理、基准测试等实战场景。操作符概览发射闭区间内的整数Source.range属于 Source 工厂类Source操作符集合其核心语义可以概括为一句话按顺序发射区间[start, end]含两端内的每一个整数且可选步长大于 1。官方文档range.md给出的 Reactive Streams 语义定义如下emits发射当下游有需求demand时发射下一个值completes完成到达区间末端时流正常完成。这意味着Source.range是一个完全被动的惰性数据源它不会一次性生成全部整数而是严格遵循下游背压backpressure逐个产出元素区间穷尽后以onComplete正常结束不会抛出异常。这一点让它可以安全地用于产生海量序列如1到1_000_000内存占用恒定。依赖配置引入 akka-stream 模块在使用Source.range前需要引入akka-stream模块。官方文档推荐通过 Akka 的 BOMBill of Materials统一管理版本// sbt libraryDependencies com.typesafe.akka %% akka-stream % AkkaVersion!-- Maven -- dependency groupIdcom.typesafe.akka/groupId artifactIdakka-stream_2.13/artifactId version${akka.version}/version /dependency// Gradle dependencies { implementation platform(com.typesafe.akka:akka-bom_2.13:$akkaVersion) implementation com.typesafe.akka:akka-stream_2.13 }其中akka-bom是官方推荐的版本管理方式bomGroupcom.typesafe.akka、bomArtifactakka-bom_$scala.binary.version$AkkaVersion指向当前仓库使用的 Akka 版本。另外需要注意官方文档特别说明Akka 依赖托管在 Akka 的安全库仓库secure library repository中需要按 https://account.akka.io/token 的说明配置带 token 的仓库地址才能拉取。Java 用法三种调用形态Java DSL 的官方示例位于 SourceDocExamples.java。先准备好所需的 importimport akka.NotUsed; import akka.actor.ActorSystem; import akka.stream.javadsl.Source;创建ActorSystem后即可构造区间数据源ActorSystem system ActorSystem.create(Source); // 形态一默认步长 1发射 1, 2, 3, ..., 100 SourceInteger, NotUsed source Source.range(1, 100); // 形态二指定步长 5发射 1, 6, 11, ..., 96 SourceInteger, NotUsed sourceStepFive Source.range(1, 100, 5); // 形态三负步长实现倒序发射 100, 99, 98, ..., 1 SourceInteger, NotUsed sourceStepNegative Source.range(100, 1, -1);三种形态对应的签名分别为签名说明示例Source.range(start, end)步长固定为 1 的完整闭区间Source.range(1, 100)→ 1 到 100Source.range(start, end, step)指定正整数步长Source.range(1, 100, 5)→ 1, 6, 11, ...Source.range(start, end, step)负步长实现倒序Source.range(100, 1, -1)→ 100, 99, ..., 1发射出的元素经runForeach交由ActorSystem驱动执行并打印// 将流中的整数逐个打印到控制台 source.runForeach(i - System.out.println(i), system);这里runForeach(elem - ..., system)将流物化materialize并通过传入的ActorSystem运行。整个调用链印证了官方文档“用apply方法生成整数序列”的描述Java 侧的Source.range(...)在 Scala 侧等价于Source(1 to N)的写法详见下文实现原理。Scala 用法直接使用apply在 Scala DSL 中官方文档指出使用apply方法即可生成整数序列。由于 Scala 的1 to 100本身就是scala.collection.immutable.Range闭区间且Source.apply接受任意immutable.Iterable见 Source.scala 的def applyT因此最简洁的写法是import akka.actor.ActorSystem import akka.stream.scaladsl.{Sink, Source} implicit val system: ActorSystem ActorSystem(Source) // 等价于 Java 的 Source.range(1, 100) val source: Source[Int, NotUsed] Source(1 to 100) // 带步长 val stepped: Source[Int, NotUsed] Source(1 to 100 by 5) // 倒序 val reversed: Source[Int, NotUsed] Source(100 to 1 by -1) source.runWith(Sink.foreach(println))Source(1 to 100)与 Java 的Source.range(1, 100)在底层走的是同一条实现路径见下节这也是官方在 Java DSL 文档注释中特别说明“allows to createSourceout of range as simply as on ScalaSource(1 to N)”的原因。底层实现原理基于Range.inclusive的惰性包装Source.range的实现位于 Java DSL 源文件 javadsl/Source.scala/** * Creates [[Source]] that represents integer values in range [start;end], step equals to 1. * It allows to create Source out of range as simply as on Scala Source(1 to N) * * Uses [[scala.collection.immutable.Range.inclusive(Int, Int)]] internally */ def range(start: Int, end: Int): javadsl.Source[Integer, NotUsed] range(start, end, 1) /** * Creates [[Source]] that represents integer values in range [start;end], with the given step. * ... * Uses [[scala.collection.immutable.Range.inclusive(Int, Int, Int)]] internally */ def range(start: Int, end: Int, step: Int): javadsl.Source[Integer, NotUsed] new Source(scaladsl.Source(Range.inclusive(start, end, step).asInstanceOf[immutable.Iterable[Integer]]))从源码可以提炼出三个关键事实两参重载是语法糖range(start, end)直接委托给range(start, end, 1)即步长默认值为1。内部借助 Scala 标准库的Range.inclusive(start, end, step)区间是包含两端inclusive的闭区间这正是Source.range(0, 10)会产出0,1,...,10共 11 个元素、而不是 10 个元素的原因。Range本身就是惰性迭代器Range.inclusive不会预先分配一个包含全部整数的集合而是按需逐个产出元素。把它交给Source.apply后流的发射节奏完全由下游需求驱动完美契合上文“emits when there is demand”的 Reactive Streams 语义。因此即使Source.range(1, Int.MaxValue)也不会造成内存问题。此外从step参数可以推断步长既可为正升序也可为负降序如range(100, 1, -1)这与 ScalaRange对by的语义完全一致。测试用例验证闭区间与步长行为仓库中的单元测试对上述行为给出了直接印证见 javadsl/SourceTest.javaTest public void mustWorkFromRange() throws Exception { CompletionStageListInteger f Source.range(0, 10).grouped(20).runWith(Sink.head(), system); final ListInteger result f.toCompletableFuture().get(3, TimeUnit.SECONDS); assertEquals(11, result.size()); // 闭区间0..10 共 11 个元素 Integer counter 0; for (Integer i : result) assertEquals(i, counter); // 严格升序 } Test public void mustWorkFromRangeWithStep() throws Exception { CompletionStageListInteger f Source.range(0, 10, 2).grouped(20).runWith(Sink.head(), system); final ListInteger result f.toCompletableFuture().get(3, TimeUnit.SECONDS); assertEquals(6, result.size()); // 0, 2, 4, 6, 8, 10 Integer counter 0; for (Integer i : result) { assertEquals(i, counter); counter 2; } }这两条测试从实证角度确认了文档语义Source.range(0, 10)产出11个元素0到10含两端顺序严格递增Source.range(0, 10, 2)产出6个元素0, 2, 4, 6, 8, 10步长生效且末尾未超过end。同时Source.range在仓库其他测试中被广泛用作构造测试数据流的工具例如 FlowTest.java 中的Source.range(0, 2)、SinkTest.java 中的Source.range(0, 2)等可见它是测试与基准场景中最常用的数据源工厂之一。实战建议与边界注意结合文档语义与源码实现使用Source.range时有几点值得注意闭区间是约定Source.range(1, 100)包含100若需要“前 100 个整数”应写Source.range(1, 100)或Source.range(0, 99)。惰性安全由于底层是惰性Range可放心构造超大区间如Source.range(1, 1_000_000)配合grouped、take、throttle等操作符做分页拉取或限速处理而无需担心一次性占用内存。步长与方向step为正数时升序、为负数时降序但需保证方向与start/end的相对位置一致否则区间为空例如Source.range(1, 10, -1)不产生任何元素并立即完成。类型提示Java 侧返回SourceInteger, NotUsedNotUsed表示该 Source 物化时不产生有价值的材料化值materialized value适用于runForeach、runWith(Sink.xxx)等消费式用法。更多整数类数据源若需要无限递增序列可改用Source.unfold或Source.repeat组合Source.range专用于有界闭区间。小结Source.range是 Akka Streams 中“最小但最常用”的 Source 工厂之一它以闭区间整数序列为输入以背压友好的惰性发射为行为以emits when there is demand / completes when the end of range is reached为契约内部由 Scala 标准库Range.inclusive支撑Java 与 Scala 两套 DSL 共用同一实现。通过本文的示例与源码、测试佐证你可以放心地将它用于构造测试数据、实现分页逻辑或任何需要整数序列的流式场景。赞分享后端并发编程异步编程【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址https://gitcode.com/gh_mirrors/ak/akka-core点击查看免费下载相关推荐Akka Streams sliding 操作符完全指南滑动窗口、步长控制与移动平均实现Akka Streams sliding 操作符完全指南滑动窗口、步长控制与移动平均实现 本指南围绕 Akka Streams 的 sliding 操作符展开后端并发编程异步编程Akka Streams FileIO.fromPath 文件读取 Source 操作符完整指南从 API 签名到底层实现Akka Streams FileIO.fromPath 文件读取 Source 操作符完整指南从 API 签名到底层实现 本篇技术指南以 Akka 官方文档后端并发编程异步编程Akka Streams collect 操作符完全指南用偏函数一步完成过滤与转换Akka Streams collect 操作符完全指南用偏函数一步完成过滤与转换 collect 是 Akka Streams 中一个简单而强大的流式处理操后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表