- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
导读
StreamConverters.asInputStream是 Akka Streams 中用于与阻塞式java.ioAPI 互操作的核心转换器:它创建一个Sink,物化(materialize)后返回一个java.io.InputStream,通过读取该InputStream触发下游需求(demand),从而把响应式流"反向"暴露给传统、同步的 Java IO 消费方。读完本文,你将掌握asInputStream的签名与参数、生命周期与背压语义、Scala/Java 双端用法、内部实现原理,以及阻塞式 IO 在 Akka Streams 中的调度配置与常见陷阱。本文以 asInputStream.md 为骨架,并结合 StreamConverters.scala、InputStreamSinkStage.scala 与 InputStreamSinkSpec.scala 等仓库源码展开。
一、操作符定位:Additional Sink and Source converters
在 Akka Streams 操作符索引中,asInputStream归属于 Additional Sink and Source converters 分类。该类转换器用于与java.io.InputStream/java.io.OutputStream集成,全部集中在StreamConverters对象上,包括:
asInputStream:Sink 物化为InputStream(本文主题)asOutputStream:Source 物化为OutputStreamfromInputStream:包装InputStream的 SourcefromOutputStream:包装OutputStream的 Sink- 以及
asJavaStream、fromJavaStream、javaCollector、javaCollectorParallelUnordered等 Java 8 Stream/Collector 转换器
由于这些操作符本质上是阻塞 API,官方文档明确指出其实现运行在独立的调度器(dispatcher)上,通过akka.stream.blocking-io-dispatcher配置,详见下文"阻塞语义与调度配置"一节。
二、核心概念:创建一个物化为 InputStream 的 Sink
asInputStream的目标很明确:创建一个 Sink,物化后得到一个InputStream,通过读取该InputStream来触发流经 Sink 的需求。
也就是说,数据流向为:
Source[ByteString] → (转换/过滤) → Sink[ByteString, InputStream] ← 外部代码读取该 InputStream流经该 Sink 的字节会缓存在内部队列中,等待消费者通过InputStream.read()拉取;每次read()都会向流的上游传递需求,从而形成"读取驱动消费"的模式。
签名(Signature)
Scala DSL(scaladsl/StreamConverters.scala):
def asInputStream(readTimeout: FiniteDuration = 5.seconds): Sink[ByteString, InputStream]Java DSL(javadsl/StreamConverters.scala)提供两个重载:
public static Sink<ByteString, InputStream> asInputStream() // 默认超时 public static Sink<ByteString, InputStream> asInputStream(java.time.Duration readTimeout) // 显式指定超时关键点:
- 输入元素类型固定为
ByteString,即流经该 Sink 的元素必须是akka.util.ByteString,这与 Akka Streams 面向字节流 IO 的统一约定一致; - 物化值为
java.io.InputStream,可在runWith后直接交给任何传统 Java IO 代码使用; readTimeout默认 5 秒:单次read()操作在等待新数据时最多阻塞的时间,超时抛出IOException("Timeout on waiting for new data")。Scala 版本参数类型为scala.concurrent.duration.FiniteDuration,Java 版本接受java.time.Duration。
生命周期语义
文档明确规定了两个方向的终止规则,源码实现完全对应:
- 流完成 → InputStream 结束:当流入该 Sink 的上游流完成(complete)时,
InputStream也会结束;此时再调用read()返回-1(EOF 语义)。 - 关闭 InputStream → 取消流:外部调用
InputStream.close()会取消(cancel)流入该 Sink 的流,上游随之终止。
在 InputStreamSinkStage.scala 中可以看到:上游完成时 stage 向共享队列追加Finished消息并completeStage();上游失败时追加Failed(ex)并failStage(ex)。而在 InputStreamAdapter 的read实现中,取到Finished即返回-1、取到Failed则包装为IOException抛出(原始异常作为cause保留),close()则通过sendToStage(Close)通知 stage 取消整个流。
三、Reactive Streams 语义
asInputStream的背压/取消行为可以用一句话概括:"读取驱动需求,关闭即取消"。文档给出的 Reactive Streams 语义为:
| 场景 | 行为 |
|---|---|
InputStream被关闭 | cancels:取消流入该 Sink 的流 |
| 没有挂起的读取 | backpressures:上游因无人读取而背压 |
从实现层面看,InputStreamSinkStage 使用一个容量为maxBuffer + 2的LinkedBlockingDeque作为 stage 与InputStreamAdapter之间的共享缓冲:stage 在sendPullIfAllowed()中仅当队列剩余容量大于 1 时才向上游pull(in),从而保证缓冲不溢出、实现精确背压;InputStreamAdapter每次读完一个完整 chunk 后会发送ReadElementAcknowledgement消息,通知 stage 释放容量并继续拉取。对应地,akka.stream.Attributes.InputBuffer属性可调节内部缓冲大小(默认 16),详见下文。
四、完整示例:Scala 与 Java 双版本
文档附带的示例来自仓库中的 ToFromJavaIOStreams.scala(Scala)与 ToFromJavaIOStreams.java(Java),演示了"读取 Source 内容 → 转为大写 → 物化为 InputStream"的完整链路。
Scala 示例
import akka.NotUsed import akka.stream.scaladsl.{ Flow, Sink, Source, StreamConverters } import akka.util.ByteString import java.io.InputStream // 将每个 ByteString 中的字节转为大写 val toUpperCase: Flow[ByteString, ByteString, NotUsed] = Flow[ByteString].map(_.map(_.toChar.toUpper.toByte)) val source: Source[ByteString, NotUsed] = Source.single(ByteString("some random input")) // 物化后得到 InputStream val sink: Sink[ByteString, InputStream] = StreamConverters.asInputStream() val inputStream: InputStream = source.via(toUpperCase).runWith(sink)随后即可用标准 Java IO 方式读取:
inputStream.read() should be('S') // 读到大写 'S' inputStream.close() // 关闭将取消上游流测试中对结果的断言为读到的首字节是'S',即"some random input"经toUpperCase转换后变成"SOME RANDOM INPUT"。
Java 示例
import akka.NotUsed; import akka.actor.ActorSystem; import akka.stream.javadsl.*; import akka.util.ByteString; import java.io.InputStream; import java.nio.charset.Charset; ActorSystem system = ActorSystem.create("ToFromJavaIOStreams"); Charset charset = Charset.defaultCharset(); Flow<ByteString, ByteString, NotUsed> toUpperCase = Flow.<ByteString>create() .map( bs -> { String str = bs.decodeString(charset).toUpperCase(); return ByteString.fromString(str, charset); }); final Sink<ByteString, InputStream> sink = StreamConverters.asInputStream(); final InputStream stream = Source.single(ByteString.fromString("Some random input")) .via(toUpperCase) .runWith(sink, system); // 读取 17 个字节并断言内容为 "SOME RANDOM INPUT" byte[] a = new byte[17]; stream.read(a);注意 Java 版本的runWith(sink, system)需要显式传入ActorSystem。Java 端断言使用assertArrayEquals("SOME RANDOM INPUT".getBytes(), a)。
五、参数详解与内部缓冲配置
readTimeout:单次读取的最大阻塞时间
readTimeout是asInputStream唯一的直接参数,控制InputStream上单次读取等待数据的最大时长。在 InputStreamSinkStage.scala 中体现为:
sharedBuffer.poll(readTimeout.toMillis, TimeUnit.MILLISECONDS) match { case Data(data) => detachedChunk = Some(data); readBytes(a, begin, length) case Finished => isStageAlive = false; -1 case Failed(ex) => isStageAlive = false; throw new IOException(ex) case null => throw new IOException("Timeout on waiting for new data") case Initialized => throw new IllegalStateException("message 'Initialized' must come first") }即:在超时时间内没有新数据到达时,read()会抛出IOException("Timeout on waiting for new data")。此外,stage 初始化时还会等待Initialized消息,同样受readTimeout约束。默认值为 5 秒;如果你的消费方可能长时间不读取(比如高频轮询场景),应适当调大该值。
InputBuffer:内部缓冲大小
文档与源码注释指出,内部缓冲大小可通过akka.stream.ActorAttributes(即Attributes.inputBuffer)配置。在 InputStreamSinkStage 中:
val maxBuffer = inheritedAttributes.getInputBuffer).max require(maxBuffer > 0, "Buffer size must be greater than 0") val dataQueue = new LinkedBlockingDequeStreamToAdapterMessage- 缓冲队列容量为
maxBuffer + 2(额外的 1 个位置预留给Finished/Failed消息); - 缓冲大小必须大于 0,否则物化时抛出
IllegalArgumentException(见测试用例 "fail to materialize with zero sized input buffer"); - 可通过
withAttributes(Attributes.inputBuffer(initial, max))调整,例如StreamConverters.asInputStream().withAttributes(inputBuffer(8, 8))。
六、实现原理:InputStreamSinkStage 与 InputStreamAdapter
asInputStream的底层实现是内部 APIInputStreamSinkStage(一个GraphStageWithMaterializedValue[SinkShape[ByteString], InputStream]),位于 akka-stream/src/main/scala/akka/stream/impl/io/InputStreamSinkStage.scala。整个协作机制分为三部分:
stage 与适配器之间的消息协议(见 InputStreamSinkStage.scala#L22-L38):
- stage → 适配器:
Initialized(初始化完成)、Data(data)(非空数据块)、Finished(上游完成)、Failed(cause)(上游失败); - 适配器 → stage:
ReadElementAcknowledgement(chunk 已被读完,可继续拉取)、Close(InputStream 被关闭,取消流)。
- stage → 适配器:
GraphStageLogic:
preStart()时向队列放入Initialized并立即pull(in);onPush()时把非空的ByteString放入队列,且仅在队列剩余容量大于 1 时继续向上游拉取(实现背压);onUpstreamFinish()/onUpstreamFailure()分别追加Finished/Failed并完成/失败 stage;postStop()在异常终止(如 materializer 关闭)时补充Failed(AbruptStageTerminationException)。InputStreamAdapter:真正的
InputStream实现。它在read时从共享队列poll数据块,支持三种读取形态:read()(单字节)、read(byte[])、read(byte[], off, len);单次读取可以跨多个上游 chunk 拼接(见getData的尾递归实现),也可以只消费 chunk 的一部分(剩余部分保留在detachedChunk中);读取长度为 0 的数组直接返回 0 且不向上游请求数据(测试用例 "a read of length 0 should not request bytes from upstream" 专门验证了这一点)。
七、阻塞语义与调度配置(重要陷阱)
asInputStream物化的InputStream是阻塞式实现:read()会阻塞当前线程直到上游有数据可用。官方文档(operators/index.md)给出了明确警告:
asInputStream和asOutputStream物化的InputStream/OutputStream是阻塞 API 实现,会阻塞线程直到上游数据可用。由于其阻塞本质,这些对象不能用于mapMaterializedValue回调中,否则会导致流物化过程的死锁。
例如下面的代码会因超时而失败:
// 反例:在物化回调里直接 read() 可能永远阻塞,导致死锁 .toMat(StreamConverters.asInputStream().mapMaterializedValue { inputStream => inputStream.read() // 这可能会永远阻塞 // ... }).run()正确做法是把读取逻辑放到独立线程(或消费者 actor/线程池)中,不要在物化阶段同步读取。
调度器配置
与所有阻塞式 IO 操作符一样,其实现运行在独立的阻塞 IO 调度器上。仓库默认配置见 akka-stream/src/main/resources/reference.conf:
akka.stream.materializer { # 执行阻塞操作的流操作符使用的调度器 blocking-io-dispatcher = "akka.actor.default-blocking-io-dispatcher" }即默认使用akka.actor.default-blocking-io-dispatcher。可以通过修改该配置项,或为单个流使用ActorAttributes.dispatcher(...)覆盖调度器,从而避免阻塞线程影响流处理主线程。
八、测试验证:行为边界一览
仓库中的 InputStreamSinkSpec.scala 对asInputStream的行为做了全面验证,可作为理解语义边界的权威参考:
| 测试用例 | 验证的行为 |
|---|---|
| "read bytes from InputStream" | 基础读取:Source.single→runWith(asInputStream())可读回全部字节 |
| "read bytes correctly if requested by InputStream not in chunk size" | 读取请求大小与上游 chunk 大小不一致时正确拼接/切分 |
| "block read until get requested number of bytes from upstream" | 上游无数据时read阻塞,数据到达后返回;上游完成后read()返回 -1 |
| "ignore an empty ByteString" | 空的ByteString元素被忽略,不产生数据 |
| "throw error when reactive stream is closed" | 关闭InputStream后,上游收到取消;再次read()抛IOException |
| "throw exception when call read with wrong parameters" | 非法参数(负偏移、越界长度等)抛IllegalArgumentException |
| "return IOException when stream is failed" | 上游流失败时,read()抛IOException,原始异常作为cause |
| "read next byte as an int from InputStream" | 单字节read()语义:返回 0-255 的无符号值,EOF 返回 -1 |
| "fail to materialize with zero sized input buffer" | 缓冲大小为 0 时物化失败 |
| "throw from inputstream read if terminated abruptly" | materializer 关闭导致异常终止时,read()抛IOException |
| "a read of length 0 should not request bytes from upstream" | 长度为 0 的读取不向上游请求数据 |
这些测试用例与文档描述完全一致,同时补充了 EOF(-1)、超时(IOException)、上游失败传播、参数校验等文档未展开的边界细节,是深入理解该操作符行为的最佳资料。
九、与同族转换器的配合使用
asInputStream常与同一族的其他转换器配合,构建完整的 java.io 互操作链路:
fromInputStream:方向相反的转换器,创建包装InputStream的Source(默认 chunk 大小 8192 字节,物化为Future[IOResult]);asOutputStream:创建物化为OutputStream的Source(默认写超时 5 秒);fromOutputStream:创建写入OutputStream的Sink(默认不自动 flush,可通过autoFlush开启)。
例如 ToFromJavaIOStreams.scala 中同时演示了fromInputStream+fromOutputStream的完整 IO 转换链路。当你需要把 Akka Streams 中的数据导出给只接受java.io.InputStream的遗留库(如 HTTP 客户端、压缩工具、数据库驱动等)时,asInputStream就是那座桥。
十、使用建议小结
- 确认元素类型:
asInputStream只接受ByteString,上游若非字节流需先用map转换为ByteString; - 管理读取超时:消费方读取间隔可能大于默认 5 秒时,务必通过
asInputStream(Duration)/asInputStream(FiniteDuration)显式调大readTimeout; - 避开物化回调死锁:不要在
mapMaterializedValue中同步调用read(),应把阻塞读取放到独立线程; - 善用生命周期:
close()会取消上游流,用完即关;上游流正常完成时read()返回-1,据此判断流式数据结束; - 按需调整缓冲:高频小批量读取可调大
InputBuffer减少上游拉取次数;注意缓冲必须大于 0。
- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
相关推荐
Akka Streams Sink.asPublisher 完全指南:将 Akka Stream 桥接到 Reactive Streams Publisher
Akka Streams Sink.asPublisher 完全指南:将 Akka Stream 桥接到 Reactive Streams Publisher
后端并发编程异步编程Akka Streams StreamConverters.asJavaStream 详解:将 Akka Sink 物化为 Java 8 Stream 的桥接之道
Akka Streams StreamConverters.asJavaStream 详解:将 Akka Sink 物化为 Java 8 Stream 的桥接之
后端并发编程异步编程Akka Streams StreamConverters.asOutputStream:将阻塞式 java.io.OutputStream 桥接为响应式 Source
Akka Streams StreamConverters.asOutputStream:将阻塞式 java.io.OutputStream 桥接为响应式 So
后端并发编程异步编程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考