- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
导读
StreamConverters.fromInputStream是 Akka Streams 提供的阻塞 IO 桥接工具,它把传统的java.io.InputStream(如文件流、网络流、字节数组流)包装成一个类型为Source[ByteString, Future[IOResult]]的响应式数据源,从而让遗留的阻塞式 Java IO 代码可以无缝接入 Akka Streams 的异步、背压(backpressure)管线。读完本文,你将掌握该操作符的完整签名、chunkSize分块语义、物化值IOResult的含义、阻塞调度器的配置方法,以及它与fromOutputStream、asInputStream、asOutputStream等兄弟操作符的配合用法,并深入其底层InputStreamSource阶段的实现原理。
操作符定位:阻塞 IO 世界与响应式世界的桥梁
在 StreamConverters 源码 的对象注释中,官方明确其职责是:
"Converters for interacting with the blocking
java.iostreams APIs and Java 8 Streams"
即专门用于与阻塞式 java.io 流 API及 Java 8 Stream 互操作。整个对象共提供 8 个转换器,fromInputStream是其中最常用、最基础的"读方向"转换器:
| 转换器 | 方向 | 功能 |
|---|---|---|
fromInputStream | 读 | 由InputStream工厂函数创建Source[ByteString, Future[IOResult]] |
fromOutputStream | 写 | 由OutputStream工厂函数创建Sink[ByteString, Future[IOResult]] |
asInputStream | 写 | 物化为InputStream的Sink,反向把流数据暴露给阻塞代码读取 |
asOutputStream | 读 | 物化为OutputStream的Source,让阻塞代码向流内写入 |
fromJavaStream/asJavaStream/javaCollector/javaCollectorParallelUnordered | 双向 | 面向 Java 8 Stream 与 Collector 的转换 |
你可以在 operators 目录 下找到全部 8 个操作符的官方文档页。
函数签名与类型
Scala DSL 的签名定义在 StreamConverters.scala:
def fromInputStream(in: () => InputStream, chunkSize: Int = 8192): Source[ByteString, Future[IOResult]]Java DSL 对应的签名(定义于 javadsl/StreamConverters.scala):
public static Source<ByteString, CompletionStage<IOResult>> fromInputStream(Creator<InputStream> in) public static Source<ByteString, CompletionStage<IOResult>> fromInputStream(Creator<InputStream> in, int chunkSize)关键点:
- 参数
in是工厂函数而非实例:每次Source被物化(materialize)时都会调用该工厂创建一个新的InputStream。这意味着同一个Source图可以被多次运行(例如多次runWith),每次都会重新打开输入流。 - 参数
chunkSize默认 8192(即 8 KiB):每个向下游发出的ByteString元素最多为chunkSize字节。 - 物化值:Scala 侧为
Future[IOResult],Java 侧为CompletionStage<IOResult>,用于在流完成/失败时获知已读取的字节数。
分块(chunk)语义:由底层 read 决定的实际元素大小
fromInputStream发出的每个元素是"最多chunkSize大小的ByteString"。实际大小取决于底层InputStream每次read调用返回的数据量——如果一次read只返回了 3 字节,那么这次发出的元素就是 3 字节,但任何元素都不会超过chunkSize。
这一点在底层实现 InputStreamSource.scala 中一目了然:
override def onPull(): Unit = try { inputStream.read(buffer) match { case -1 => closeStage() case readBytes => readBytesTotal += readBytes push(out, ByteString.fromArray(buffer, 0, readBytes)) } } catch { case NonFatal(t) => failStream(t) failStage(t) }其工作流程为:
- 阶段内部预先分配一个
new ArrayByte的缓冲区(源码 L44); - 当下游发出拉取请求(
onPull)时,调用inputStream.read(buffer)一次性读取最多chunkSize字节; - 返回
-1表示流已到末尾,触发closeStage():关闭输入流、以成功状态完成物化Future并完成整个Source; - 返回正整数
readBytes时,用ByteString.fromArray(buffer, 0, readBytes)精确截取实际读到的字节,推送给下游; - 任何非致命异常都会调用
failStream(t)——先关闭输入流,再以IOOperationIncompleteException失败物化Future,同时failStage让流以错误终止。
同时,构造函数中有明确的校验(源码 L33):
require(chunkSize > 0, s"chunkSize must be > 0 (was $chunkSize)")即chunkSize必须大于 0,否则在构建Source时立即抛出IllegalArgumentException。
物化值 IOResult:字节计数与异常语义
fromInputStream物化出一个Future[IOResult]/CompletionStage<IOResult>,其中IOResult.count表示从输入流中读取的字节总数(InputStreamSource内部用readBytesTotal长整型累加)。IOResult定义在 IOResult.scala:
final case class IOResult( count: Long, @deprecated("status is always set to Success(Done)", "2.6.0") status: Try[Done])需要注意的语义细节:
- 成功完成:当读到 EOF(
read返回 -1)时,以IOResult(readBytesTotal)成功完成,其中count为累计读取的总字节数; - 正常取消:当
Source被下游取消(非失败原因)时,同样以已读取的字节数成功完成物化Future(见源码 L76-L90); - 下游失败:若下游以异常终止(如
SubscriptionWithCancelException的子类之外的原因),物化Future会以IOOperationIncompleteException("Downstream failed before input stream reached end", readBytesTotal, ex)失败,count仍记录已读取字节; - 流被突然终止:
postStop时若 IO 尚未结束,Future以AbruptStageTerminationException失败; - 读取过程中抛异常:以
IOOperationIncompleteException(readBytesTotal, reason)失败。
一个重要警告:字节被Source读取出来,并不保证下游各阶段已经看到这些字节。count只反映"从InputStream读出了多少字节",如果中途流被取消或失败,部分已读字节可能并未被下游消费。设计业务逻辑时(例如处理"已读但未处理"的边界数据)需要对此有明确预期。
生命周期:取消即关闭输入流
文档明确承诺:"The createdInputStreamwill be closed when theSourceis cancelled."(创建的InputStream会在Source被取消时关闭。)
从实现看,关闭逻辑发生在三个时机,统一收敛到closeInputStream()方法(源码 L109-L118):
- 正常读到 EOF 时(
closeStage); - 下游取消 / 失败时(
onDownstreamFinish); - 读取抛出非致命异常时(
failStream)。
此外preStart中若工厂函数本身抛异常,阶段会以IOOperationIncompleteException(0, t)失败物化值并立即failStage(源码 L51-L59)。因此,你无需手动关闭传入的InputStream,Akka 会负责其生命周期;但前提是该InputStream的close()方法是可重入且幂等的(大多数 JDK 实现满足)。
调度器配置:阻塞读取不能在默认调度器上执行
由于InputStream.read是阻塞调用,fromInputStream生成的阶段必须运行在专门的阻塞 IO 调度器上,避免占用默认调度器导致线程饥饿。
- 全局配置:修改
akka.stream.materializer.blocking-io-dispatcher,默认值为"akka.actor.default-blocking-io-dispatcher",见 akka-stream 的 reference.conf:
akka.stream.materializer { blocking-io-dispatcher = "akka.actor.default-blocking-io-dispatcher" }- 默认调度器定义:
akka.actor.default-blocking-io-dispatcher定义在 akka-actor 的 reference.conf,是一个固定线程池大小为 16、吞吐量为 1 的Dispatcher:
default-blocking-io-dispatcher { type = "Dispatcher" executor = "thread-pool-executor" throughput = 1 thread-pool-executor { fixed-pool-size = 16 } }- 单 Source 覆盖:通过
akka.stream.ActorAttributes为某个特定的Source单独指定调度器,例如(Scala):
import akka.stream.ActorAttributes val source = StreamConverters .fromInputStream(() => myInputStream) .withAttributes(ActorAttributes.dispatcher("my-blocking-dispatcher"))Java 侧对应ActorAttributes.dispatcher("my-blocking-dispatcher")与source.withAttributes(...)。这一机制保证了:多个fromInputStream同时运行时,阻塞读取被隔离在专用线程池中,不会拖垮主执行管线。
完整示例:InputStream → 大写转换 → OutputStream
官方文档给出了一个同时使用fromInputStream与fromOutputStream的读写闭环示例:从一个java.io.InputStream读取内容、转大写、再写入java.io.OutputStream。可运行的完整测试代码位于:
- Scala:ToFromJavaIOStreams.scala
- Java:ToFromJavaIOStreams.java
Scala 版本
import java.io.{ ByteArrayInputStream, ByteArrayOutputStream, InputStream, OutputStream } import akka.NotUsed import akka.stream.IOResult import akka.stream.scaladsl.{ Flow, Sink, Source, StreamConverters } import akka.util.ByteString import scala.concurrent.Future val bytes = "Some random input".getBytes val inputStream = new ByteArrayInputStream(bytes) val outputStream = new ByteArrayOutputStream() val source: Source[ByteString, Future[IOResult]] = StreamConverters.fromInputStream(() => inputStream) val toUpperCase: Flow[ByteString, ByteString, NotUsed] = Flow[ByteString].map(_.map(_.toChar.toUpper.toByte)) val sink: Sink[ByteString, Future[IOResult]] = StreamConverters.fromOutputStream(() => outputStream) val eventualResult: Future[IOResult] = source.via(toUpperCase).runWith(sink)当eventualResult完成时,outputStream内部字节数组即为"SOME RANDOM INPUT"——测试断言(ToFromJavaIOStreams.scala#L40-L42):
whenReady(eventualResult) { _ => outputStream.toByteArray.map(_.toChar).mkString should be("SOME RANDOM INPUT") }Java 版本
import akka.NotUsed; import akka.stream.IOResult; import akka.stream.javadsl.*; import akka.util.ByteString; import java.io.*; import java.util.concurrent.CompletionStage; byte[] bytes = "Some random input".getBytes(Charset.defaultCharset()); java.io.InputStream inputStream = new ByteArrayInputStream(bytes); Source<ByteString, CompletionStage<IOResult>> source = StreamConverters.fromInputStream(() -> inputStream); Flow<ByteString, ByteString, NotUsed> toUpperCase = Flow.<ByteString>create() .map(bs -> { String str = bs.decodeString(charset).toUpperCase(); return ByteString.fromString(str, charset); }); java.io.OutputStream outputStream = new ByteArrayOutputStream(); Sink<ByteString, CompletionStage<IOResult>> sink = StreamConverters.fromOutputStream(() -> outputStream); CompletionStage<IOResult> ioResultCompletionStage = source.via(toUpperCase).runWith(sink, system);ioResultCompletionStage完成后,outputStream的字节数组即为输入内容的大写形式。测试断言见 ToFromJavaIOStreams.java#L82-L86。
示例要点拆解
- 工厂函数延迟创建:
() => inputStream在Source物化时才被调用,可借此实现"每次物化重新打开流"; - 纯函数式转换:
Flow[ByteString].map(...)对每个 chunk 独立做大小写转换,无需关心底层分块边界(因为大小写转换不跨 chunk 状态); - 两端都是阻塞桥接:读端
fromInputStream阻塞读取,写端fromOutputStream阻塞写入,中间的处理逻辑完全异步; runWith(sink)同时物化 Source 与 Sink,返回 Sink 侧(Keep.right)的物化值Future[IOResult]。
底层实现:GraphStage 视角
fromInputStream的实现极为简洁——它直接委托给内部阶段:
def fromInputStream(in: () => InputStream, chunkSize: Int = 8192): Source[ByteString, Future[IOResult]] = { Source.fromGraph(new InputStreamSource(in, chunkSize)) }InputStreamSource是一个GraphStageWithMaterializedValue[SourceShape[ByteString], Future[IOResult]](源码 InputStreamSource.scala#L30-L31),其关键设计:
- 背压驱动读取:只有在下游
onPull时才调用一次read,天然实现了"下游要多少、上游读多少"的背压语义,不会提前把整个流读进内存; - 惰性创建:
InputStream在preStart(阶段启动、物化完成)时才通过工厂创建,失败则立刻以 0 字节失败物化值; - 单元素缓冲:内部只有一个
chunkSize的字节数组,内存占用恒定,不会随流大小增长; - 精确的失败传播:所有异常都包装为
IOOperationIncompleteException(携带已读字节数),便于调用方进行部分数据的补偿处理,其定义见 IOResult.scala#L82-L88。
从源码结构可以推断,InputStreamSource属于akka.stream.impl.io包(与OutputStreamGraphStage、InputStreamSinkStage、OutputStreamSourceStage并列),是 Akka Streams 内置 IO 桥接图阶段体系的一员。
兄弟操作符与选型建议
fromInputStream常常与以下操作符搭配使用:
fromOutputStream:读写的镜像操作符,创建Sink[ByteString, Future[IOResult]],由工厂函数提供OutputStream,支持autoFlush参数,详见 fromOutputStream 文档;asInputStream:反向桥接——把Source[ByteString]物化为一个阻塞的java.io.InputStream,供遗留代码读取流中的数据,默认读超时 5 秒;asOutputStream:把Source物化为java.io.OutputStream,遗留代码往里写数据即注入流中,默认写超时 5 秒。
选型建议:
- 你的异步管线需要从 InputStream 拉数据→
fromInputStream; - 需要把流结果写回一个 OutputStream→
fromOutputStream; - 需要让遗留同步代码消费流→
asInputStream(遗留代码持有InputStream主动读取); - 需要让遗留同步代码生产数据→
asOutputStream。
四个操作符共同构成了 Akka Streams 与 Java IO 世界双向、可组合的互操作面。
实践注意事项
chunkSize并非越大越好:增大它意味着每次read更大的分配与更大的元素粒度,会减少元素数量但增加单元素延迟;对高吞吐文件流可适当调大(如 64 KiB),对低延迟场景保持默认 8192 即可;- 不要在流处理链中做跨 chunk 有状态解析:因为 chunk 边界由底层
InputStream.read行为决定(可能不是固定大小),跨 chunk 的状态(如拆行、半包)需要用Framing或自定义状态化Flow处理; count不等于"下游处理量":IOResult.count是"从输入流读出的字节数",中间有缓冲/丢弃时两者会不一致;- 阻塞调度器容量有限:
default-blocking-io-dispatcher固定 16 线程,若并发打开大量fromInputStream实例,考虑自定义更大的阻塞调度器并通过ActorAttributes.dispatcher指定; - 每次物化都会重新打开流:如果输入流只能打开一次(如一次性网络连接),需要保证工厂函数在多次物化场景下的行为符合预期,或复用同一个
Source仅物化一次。
- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
相关推荐
Akka Streams StreamConverters.asOutputStream:将阻塞式 java.io.OutputStream 桥接为响应式 Source
Akka Streams StreamConverters.asOutputStream:将阻塞式 java.io.OutputStream 桥接为响应式 So
后端并发编程异步编程Akka Streams StreamConverters.asJavaStream 详解:将 Akka Sink 物化为 Java 8 Stream 的桥接之道
Akka Streams StreamConverters.asJavaStream 详解:将 Akka Sink 物化为 Java 8 Stream 的桥接之
后端并发编程异步编程Akka Streams Sink.asPublisher 完全指南:将 Akka Stream 桥接到 Reactive Streams Publisher
Akka Streams Sink.asPublisher 完全指南:将 Akka Stream 桥接到 Reactive Streams Publisher
后端并发编程异步编程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考