Akka Streams `Source.range` 操作符完全指南:生成整数序列、步长控制与底层实现原理
2026/9/23 21:46:42 网站建设 项目流程
  • 后端
  • 并发编程
  • 异步编程

【免费下载链接】akka-core

A 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正常结束,不会抛出异常。这一点让它可以安全地用于产生海量序列(如11_000_000),内存占用恒定。

依赖配置:引入 akka-stream 模块

在使用Source.range前,需要引入akka-stream模块。官方文档推荐通过 Akka 的 BOM(Bill of Materials)统一管理版本:

// sbt libraryDependencies += "com.typesafe.akka" %% "akka-stream" % AkkaVersion
<!-- Maven --> <dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-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是官方推荐的版本管理方式(bomGroup=com.typesafe.akkabomArtifact=akka-bom_$scala.binary.version$),AkkaVersion指向当前仓库使用的 Akka 版本。另外需要注意,官方文档特别说明:Akka 依赖托管在 Akka 的安全库仓库(secure library repository)中,需要按 https://account.akka.io/token 的说明配置带 token 的仓库地址才能拉取。

Java 用法:三种调用形态

Java DSL 的官方示例位于 SourceDocExamples.java。先准备好所需的 import:

import akka.NotUsed; import akka.actor.ActorSystem; import akka.stream.javadsl.Source;

创建ActorSystem后,即可构造区间数据源:

ActorSystem system = ActorSystem.create("Source"); // 形态一:默认步长 1,发射 1, 2, 3, ..., 100 Source<Integer, NotUsed> source = Source.range(1, 100); // 形态二:指定步长 5,发射 1, 6, 11, ..., 96 Source<Integer, NotUsed> sourceStepFive = Source.range(1, 100, 5); // 形态三:负步长实现倒序,发射 100, 99, 98, ..., 1 Source<Integer, NotUsed> sourceStepNegative = Source.range(100, 1, -1);

三种形态对应的签名分别为:

签名说明示例
Source.range(start, end)步长固定为 1 的完整闭区间Source.range(1, 100)→ 1 到 100
Source.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]]))

从源码可以提炼出三个关键事实:

  1. 两参重载是语法糖range(start, end)直接委托给range(start, end, 1),即步长默认值为1
  2. 内部借助 Scala 标准库的Range.inclusive(start, end, step):区间是包含两端(inclusive)的闭区间,这正是Source.range(0, 10)会产出0,1,...,10共 11 个元素、而不是 10 个元素的原因。
  3. Range本身就是惰性迭代器Range.inclusive不会预先分配一个包含全部整数的集合,而是按需逐个产出元素。把它交给Source.apply后,流的发射节奏完全由下游需求驱动,完美契合上文“emits when there is demand”的 Reactive Streams 语义。因此即使Source.range(1, Int.MaxValue)也不会造成内存问题。

此外,从step参数可以推断:步长既可为正(升序),也可为负(降序,如range(100, 1, -1)),这与 ScalaRangeby的语义完全一致。

测试用例验证:闭区间与步长行为

仓库中的单元测试对上述行为给出了直接印证,见 javadsl/SourceTest.java:

@Test public void mustWorkFromRange() throws Exception { CompletionStage<List<Integer>> f = Source.range(0, 10).grouped(20).runWith(Sink.head(), system); final List<Integer> 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 { CompletionStage<List<Integer>> f = Source.range(0, 10, 2).grouped(20).runWith(Sink.head(), system); final List<Integer> 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个元素(010,含两端),顺序严格递增;
  • 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时有几点值得注意:

  1. 闭区间是约定Source.range(1, 100)包含100,若需要“前 100 个整数”应写Source.range(1, 100)Source.range(0, 99)
  2. 惰性安全:由于底层是惰性Range,可放心构造超大区间(如Source.range(1, 1_000_000)),配合groupedtakethrottle等操作符做分页拉取或限速处理,而无需担心一次性占用内存。
  3. 步长与方向step为正数时升序、为负数时降序,但需保证方向与start/end的相对位置一致,否则区间为空(例如Source.range(1, 10, -1)不产生任何元素并立即完成)。
  4. 类型提示:Java 侧返回Source<Integer, NotUsed>NotUsed表示该 Source 物化时不产生有价值的材料化值(materialized value),适用于runForeachrunWith(Sink.xxx)等消费式用法。
  5. 更多整数类数据源:若需要无限递增序列,可改用Source.unfoldSource.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-core

A 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
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询