后端并发编程异步编程【免费下载链接】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.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对象上包括asInputStreamSink 物化为InputStream本文主题asOutputStreamSource 物化为OutputStreamfromInputStream包装InputStream的 SourcefromOutputStream包装OutputStream的 Sink以及asJavaStream、fromJavaStream、javaCollector、javaCollectorParallelUnordered等 Java 8 Stream/Collector 转换器由于这些操作符本质上是阻塞 API官方文档明确指出其实现运行在独立的调度器dispatcher上通过akka.stream.blocking-io-dispatcher配置详见下文阻塞语义与调度配置一节。二、核心概念创建一个物化为 InputStream 的 SinkasInputStream的目标很明确创建一个 Sink物化后得到一个InputStream通过读取该InputStream来触发流经 Sink 的需求。也就是说数据流向为Source[ByteString] → (转换/过滤) → Sink[ByteString, InputStream] ← 外部代码读取该 InputStream流经该 Sink 的字节会缓存在内部队列中等待消费者通过InputStream.read()拉取每次read()都会向流的上游传递需求从而形成读取驱动消费的模式。签名SignatureScala DSLscaladsl/StreamConverters.scaladef asInputStream(readTimeout: FiniteDuration 5.seconds): Sink[ByteString, InputStream]Java DSLjavadsl/StreamConverters.scala提供两个重载public static SinkByteString, InputStream asInputStream() // 默认超时 public static SinkByteString, InputStream asInputStream(java.time.Duration readTimeout) // 显式指定超时关键点输入元素类型固定为ByteString即流经该 Sink 的元素必须是akka.util.ByteString这与 Akka Streams 面向字节流 IO 的统一约定一致物化值为java.io.InputStream可在runWith后直接交给任何传统 Java IO 代码使用readTimeout默认 5 秒单次read()操作在等待新数据时最多阻塞的时间超时抛出IOExceptionTimeout on waiting for new data。Scala 版本参数类型为scala.concurrent.duration.FiniteDurationJava 版本接受java.time.Duration。生命周期语义文档明确规定了两个方向的终止规则源码实现完全对应流完成 → InputStream 结束当流入该 Sink 的上游流完成complete时InputStream也会结束此时再调用read()返回-1EOF 语义。关闭 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.scalaScala与 ToFromJavaIOStreams.javaJava演示了读取 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(); FlowByteString, ByteString, NotUsed toUpperCase Flow.ByteStringcreate() .map( bs - { String str bs.decodeString(charset).toUpperCase(); return ByteString.fromString(str, charset); }); final SinkByteString, 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 与 InputStreamAdapterasInputStream的底层实现是内部 APIInputStreamSinkStage一个GraphStageWithMaterializedValue[SinkShape[ByteString], InputStream]位于 akka-stream/src/main/scala/akka/stream/impl/io/InputStreamSinkStage.scala。整个协作机制分为三部分stage 与适配器之间的消息协议见 InputStreamSinkStage.scala#L22-L38stage → 适配器Initialized初始化完成、Data(data)非空数据块、Finished上游完成、Failed(cause)上游失败适配器 → stageReadElementAcknowledgementchunk 已被读完可继续拉取、CloseInputStream 被关闭取消流。GraphStageLogicpreStart()时向队列放入Initialized并立即pull(in)onPush()时把非空的ByteString放入队列且仅在队列剩余容量大于 1 时继续向上游拉取实现背压onUpstreamFinish()/onUpstreamFailure()分别追加Finished/Failed并完成/失败 stagepostStop()在异常终止如 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.confakka.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()返回 -1ignore an empty ByteString空的ByteString元素被忽略不产生数据throw error when reactive stream is closed关闭InputStream后上游收到取消再次read()抛IOExceptionthrow exception when call read with wrong parameters非法参数负偏移、越界长度等抛IllegalArgumentExceptionreturn IOException when stream is failed上游流失败时read()抛IOException原始异常作为causeread next byte as an int from InputStream单字节read()语义返回 0-255 的无符号值EOF 返回 -1fail to materialize with zero sized input buffer缓冲大小为 0 时物化失败throw from inputstream read if terminated abruptlymaterializer 关闭导致异常终止时read()抛IOExceptiona 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 中同时演示了fromInputStreamfromOutputStream的完整 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-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 Sink.asPublisher 完全指南将 Akka Stream 桥接到 Reactive Streams PublisherAkka 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 桥接为响应式 SourceAkka Streams StreamConverters.asOutputStream将阻塞式 java.io.OutputStream 桥接为响应式 So后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
