- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
StreamConverters.fromOutputStream是 Akka Streams 提供的用于与阻塞式java.io.OutputStream互操作的 Sink 转换器。本文将基于 官方操作符文档 并结合仓库源码,讲解它的签名、生命周期、autoFlush参数、调度器配置与错误处理,并给出 Scala / Java 双语言可运行示例,帮助你将任何传统 OutputStream(文件、网络、压缩流等)无缝接入响应式流管线。
操作符概览
fromOutputStream创建一个 Sink,将流入的ByteString写入由给定工厂函数创建的java.io.OutputStream。它属于 Additional Sink and Source converters 一组,与 fromInputStream 互为读写两侧的镜像操作符。
Scala DSL 签名
定义于 akka.stream.scaladsl.StreamConverters:
def fromOutputStream(out: () => OutputStream, autoFlush: Boolean = false): Sink[ByteString, Future[IOResult]]Java DSL 签名
定义于 akka.stream.javadsl.StreamConverters,提供两个重载,一个使用默认autoFlush = false,另一个显式传入:
Sink<ByteString, CompletionStage<IOResult>> fromOutputStream(Creator<OutputStream> f) Sink<ByteString, CompletionStage<IOResult>> fromOutputStream(Creator<OutputStream> f, boolean autoFlush)Java 版本底层委托给 Scala 实现(scaladsl.StreamConverters.fromOutputStream(() => f.create(), autoFlush).toCompletionStage()),因此两个 API 的行为完全一致,仅物化值类型不同:Scala 返回Future[IOResult],Java 返回CompletionStage[IOResult]。
物化值与 IOResult
该 Sink 物化(materialize)为一个IOResult的异步结果:
- 在流成功完成时,以已写入的字节总数完成该 Future/CompletionStage;
- 若 IO 操作失败,则以携带已写入字节数与底层异常的
IOOperationIncompleteException完成。
IOResult定义于 akka.stream.IOResult,包含count: Long与status: Try[Unit]两个字段,并提供了便捷方法wasSuccessful: Boolean与getError: Throwable,可用于检查写入是否完全成功。
底层实现:OutputStreamGraphStage
从源码结构看,该 Sink 由内部GraphStage实现,即 akka.stream.impl.io.OutputStreamGraphStage(标注为@InternalApi,仅供内部使用)。它通过GraphStageWithMaterializedValue[SinkShape[ByteString], Future[IOResult]]管理一个Promise[IOResult]作为物化值。
整个写入生命周期分为以下几个阶段:
- 创建(preStart):在 stage 启动时调用工厂函数
factory()创建 OutputStream,成功后立即pull(in)请求上游元素;若创建失败,则用IOOperationIncompleteException(bytesWritten, t)失败物化值并failStage(t)。 - 写入(onPush):每次收到一个
ByteString,调用outputStream.write(next.toArrayUnsafe())写入字节数组;若autoFlush为 true 则随后outputStream.flush();累加bytesWritten += next.size后继续pull(in)请求下一个元素。 - 上游完成(onUpstreamFinish):流正常结束时先执行一次
flush(),将缓冲数据刷出。 - 停止清理(postStop):若 OutputStream 非空,再次
flush()并close(),然后以IOResult(bytesWritten)成功完成物化值。
值得注意的几点行为:
- 失败传播:写入过程中任何
NonFatal异常都会导致物化值以IOOperationIncompleteException失败,同时 stage 以该异常失败,从而触发下游取消,保证阻塞 IO 错误能进入 Akka Streams 的失败通道。 - 自动关闭:OutputStream 由该 Sink 负责关闭,用户无需(也不应)自行关闭;文档明确"当流入该 Sink 的流完成时 OutputStream 会被关闭"。
- 取消语义:当 OutputStream 不再可写(例如底层通道已关闭)时,Sink 会取消上游流。
参数详解:autoFlush
autoFlush是唯一的行为开关,默认值为false:
| 取值 | 行为 |
|---|---|
false(默认) | 仅在流完成、stop 清理时统一flush(),批量写入吞吐更高 |
true | 每写入一个ByteString字节数组后立即flush(),数据及时落盘/发送,但频繁 flush 会降低吞吐 |
选择建议:需要实时性(如边写边给外部读取的管道、日志追加场景)时开启autoFlush;追求批量写入性能(如一次性导出大文件)时保持默认关闭。需要说明的是,源码中的 flush 调用发生在每次onPush之后、pull(in)之前,因此autoFlush = true时每个上游元素都会触发一次 flush。
完整可运行示例
文档提供了同时使用fromInputStream与fromOutputStream的示例:从java.io.InputStream读取内容,转大写后写回java.io.OutputStream。仓库中对应的可执行测试位于 Scala 测试 与 Java 测试。
Scala 版本
import java.io.{ ByteArrayInputStream, ByteArrayOutputStream, InputStream, OutputStream } import akka.NotUsed import akka.stream.IOResult import akka.stream.scaladsl.{ Flow, Keep, 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"Java 版本
import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.util.concurrent.CompletionStage; import akka.NotUsed; import akka.actor.ActorSystem; import akka.stream.IOResult; import akka.stream.javadsl.Flow; import akka.stream.javadsl.Sink; import akka.stream.javadsl.Source; import akka.stream.javadsl.StreamConverters; import akka.util.ByteString; ActorSystem system = ActorSystem.create("ToFromJavaIOStreams"); 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 -> ByteString.fromString(bs.decodeString(charset).toUpperCase(), 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 底层字节数组 // 将包含从 inputStream 读出并转大写后的内容两个测试类还分别断言了outputStream.toByteArray结果为"SOME RANDOM INPUT",可直接作为运行验证的基准。注意fromInputStream与fromOutputStream的工厂函数都是() => OutputStream,这意味着每次物化都会重新调用一次工厂——若同一个 Sink 被多次runWith,会为每次运行创建独立的 OutputStream 实例。
调度器配置:阻塞 IO 专用 dispatcher
由于OutputStream.write是阻塞操作,Akka 为这类 IO 转换器预设了专用 dispatcher。该 Sink 的默认属性定义于 akka.stream.impl.Stages:
val outputStreamSink = name("outputStreamSink") and IODispatcher其中IODispatcher指向 reference.conf 中的配置项:
akka.stream.materializer { blocking-io-dispatcher = "akka.actor.default-blocking-io-dispatcher" }即默认使用akka.actor.default-blocking-io-dispatcher(一个带线程池上限、允许阻塞的 dispatcher),避免阻塞调用占据普通 actor 派发线程。有两种方式自定义:
- 全局配置:修改
akka.stream.materializer.blocking-io-dispatcher指向自定义 dispatcher 名称; - 局部覆盖:通过
akka.stream.ActorAttributes为单个 Sink 覆盖调度器,例如:
sink.withAttributes(ActorAttributes.dispatcher("my-blocking-dispatcher"))错误处理实践
由于物化值为Future[IOResult]/CompletionStage[IOResult],读取写入结果时需同时处理成功与失败两种路径:
- 成功时通过
IOResult.count获取写入字节数; - 失败时异常类型为
IOOperationIncompleteException,它同时携带了失败前已写入的字节数与根因Throwable,可用于部分写入诊断。
一个典型的防御式写法(Scala):
eventualResult.foreach { ioResult => if (ioResult.wasSuccessful) println(s"written ${ioResult.count} bytes") else ioResult.getError.printStackTrace() }需要留意的是,Future本身失败(如磁盘满、通道关闭)时,上述foreach不会执行,需配合recover/onComplete捕获IOOperationIncompleteException。
与 asOutputStream 的辨析
StreamConverters中还提供了方向相反的 asOutputStream,两者容易混淆,对比如下:
| 维度 | fromOutputStream | asOutputStream |
|---|---|---|
| 类型 | Sink(写入方) | Source(读取方) |
| 数据流向 | 流中的ByteString→ OutputStream | 外部代码写入 OutputStream → 流入流 |
| 物化值 | Future[IOResult](字节数) | OutputStream(可直接write) |
| 典型场景 | 把响应式流结果落到传统写入 API | 把传统写入 API 的数据引入流内处理 |
简单记忆:fromXxx把"外部 Java IO 对象"作为流的终点;asXxx把"外部 Java IO 对象"作为流的起点。
适用场景总结
fromOutputStream的核心价值是桥接阻塞 IO 与响应式流,常见用法包括:
- 将流计算结果写入
FileOutputStream、ZipOutputStream等压缩流、BufferedOutputStream; - 对接只接受 OutputStream 的第三方库(如某些加密、编码、上报 SDK);
- 与
fromInputStream配对,实现传统InputStream → OutputStream处理链的流式化,避免一次性加载全部字节到内存。
由于写入全程在专用阻塞 dispatcher 上执行、每次只处理一个ByteString元素,即使面对超大数据量也能保持有界内存占用;配合背压机制,上游生产速率会自动适配磁盘/网络的实际写入能力。
- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
相关推荐
Akka Streams 中的 StreamConverters.fromJavaStream:将 Java 8 Stream 接入响应式流
Akka Streams 中的 StreamConverters.fromJavaStream:将 Java 8 Stream 接入响应式流 导读 Stream
后端并发编程异步编程Akka Streams `Source.asSubscriber` 实战:将 `java.util.concurrent.Flow.Subscriber` 无缝接入响应式流
Akka Streams Source.asSubscriber 实战:将 java.util.concurrent.Flow.Subscriber 无缝接入响
后端并发编程异步编程Akka Streams 的 Source.fromJavaStream:将 Java 8 Stream 按需接入响应式流
Akka Streams 的 Source.fromJavaStream:将 Java 8 Stream 按需接入响应式流 本指南围绕 Akka Streams
后端并发编程异步编程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考