
后端并发编程异步编程【免费下载链接】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点击查看免费下载导读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 blockingjava.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.scaladef fromInputStream(in: () InputStream, chunkSize: Int 8192): Source[ByteString, Future[IOResult]]Java DSL 对应的签名定义于 javadsl/StreamConverters.scalapublic static SourceByteString, CompletionStageIOResult fromInputStream(CreatorInputStream in) public static SourceByteString, CompletionStageIOResult fromInputStream(CreatorInputStream in, int chunkSize)关键点参数in是工厂函数而非实例每次Source被物化materialize时都会调用该工厂创建一个新的InputStream。这意味着同一个Source图可以被多次运行例如多次runWith每次都会重新打开输入流。参数chunkSize默认 8192即 8 KiB每个向下游发出的ByteString元素最多为chunkSize字节。物化值Scala 侧为Future[IOResult]Java 侧为CompletionStageIOResult用于在流完成/失败时获知已读取的字节数。分块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让流以错误终止。同时构造函数中有明确的校验源码 L33require(chunkSize 0, schunkSize must be 0 (was $chunkSize))即chunkSize必须大于 0否则在构建Source时立即抛出IllegalArgumentException。物化值 IOResult字节计数与异常语义fromInputStream物化出一个Future[IOResult]/CompletionStageIOResult其中IOResult.count表示从输入流中读取的字节总数InputStreamSource内部用readBytesTotal长整型累加。IOResult定义在 IOResult.scalafinal case class IOResult( count: Long, deprecated(status is always set to Success(Done), 2.6.0) status: Try[Done])需要注意的语义细节成功完成当读到 EOFread返回 -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。因此你无需手动关闭传入的InputStreamAkka 会负责其生命周期但前提是该InputStream的close()方法是可重入且幂等的大多数 JDK 实现满足。调度器配置阻塞读取不能在默认调度器上执行由于InputStream.read是阻塞调用fromInputStream生成的阶段必须运行在专门的阻塞 IO 调度器上避免占用默认调度器导致线程饥饿。全局配置修改akka.stream.materializer.blocking-io-dispatcher默认值为akka.actor.default-blocking-io-dispatcher见 akka-stream 的 reference.confakka.stream.materializer { blocking-io-dispatcher akka.actor.default-blocking-io-dispatcher }默认调度器定义akka.actor.default-blocking-io-dispatcher定义在 akka-actor 的 reference.conf是一个固定线程池大小为 16、吞吐量为 1 的Dispatcherdefault-blocking-io-dispatcher { type Dispatcher executor thread-pool-executor throughput 1 thread-pool-executor { fixed-pool-size 16 } }单 Source 覆盖通过akka.stream.ActorAttributes为某个特定的Source单独指定调度器例如Scalaimport 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。可运行的完整测试代码位于ScalaToFromJavaIOStreams.scalaJavaToFromJavaIOStreams.javaScala 版本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-L42whenReady(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); SourceByteString, CompletionStageIOResult source StreamConverters.fromInputStream(() - inputStream); FlowByteString, ByteString, NotUsed toUpperCase Flow.ByteStringcreate() .map(bs - { String str bs.decodeString(charset).toUpperCase(); return ByteString.fromString(str, charset); }); java.io.OutputStream outputStream new ByteArrayOutputStream(); SinkByteString, CompletionStageIOResult sink StreamConverters.fromOutputStream(() - outputStream); CompletionStageIOResult 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-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 StreamConverters.asOutputStream将阻塞式 java.io.OutputStream 桥接为响应式 SourceAkka 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 PublisherAkka Streams Sink.asPublisher 完全指南将 Akka Stream 桥接到 Reactive Streams Publisher后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考