ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Akka Streams Source.fromPublisher:接入 java.util.concurrent.Flow.Publisher 的响应式流集成指南

Akka Streams Source.fromPublisher:接入 java.util.concurrent.Flow.Publisher 的响应式流集成指南 后端并发编程异步编程【免费下载链接】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点击查看免费下载导读Source.fromPublisher是 Akka Streams 中用于与 Reactive Streams 生态包括 JDK 9 的java.util.concurrent.Flow无缝互操作的入口操作符。本文以官方文档 akka-docs/src/main/paradox/stream/operators/Source/fromPublisher.md 为主体结合仓库源码剖析其签名、底层实现、背压协调机制与多次物化语义并通过完整的 Scala/Java 数据库客户端示例演示如何把一个支持 Reactive Streams 的第三方库如数据库驱动接入 Akka Streams 处理管道。读完本文你将掌握JavaFlowSupport.Source.fromPublisher与org.reactivestreams.Publisher两种接入方式的选择原则并能直接将其用于实际项目集成。一、操作符定位Source 的 Reactive Streams 集成入口在 Akka Streams 内置的 Source 操作符索引akka-docs/src/main/paradox/stream/operators/index.md中fromPublisher被定位为Integration with Reactive Streams, subscribes to ajava.util.concurrent.Flow.Publisher.它的核心价值在于当你希望从一个支持 Reactive Streams 规范的第三方库例如响应式数据库驱动、HTTP 客户端、消息中间件客户端中拉取元素时不必自行编写适配层只需把对方提供的Publisher交给fromPublisher就能得到标准的 Akka StreamsSource进而接入map、filter、buffer等所有下游操作符。这一能力在 JavaFlowSupport.scala 中体现得十分直白JavaFlowSupport对象专门为 JDK 9 的java.util.concurrent.Flow.*接口提供Source、Flow、Sink三组工厂方法fromPublisher只是其中与Source相关的一例。二、方法签名Scala 与 Java 双 APIScala APIScala 侧签名定义于 akka-stream/src/main/scala/akka/stream/scaladsl/JavaFlowSupport.scaladef fromPublisherT: Source[T, NotUsed]注意它是JavaFlowSupport.Source.fromPublisher与位于akka.stream.scaladsl.Source中处理org.reactivestreams.Publisher的Source.fromPublisher同名但参数类型不同// akka-stream/src/main/scala/akka/stream/scaladsl/Source.scala def fromPublisherT: Source[T, NotUsed] // org.reactivestreams.Publisher两者共享相同的语义区别仅在于接入的Publisher类型来自哪个规范族。Java APIJava 侧签名摘自 FromPublisher.javastatic T akka.stream.javadsl.SourceT, NotUsed fromPublisher(PublisherT publisher)其中Publisher为java.util.concurrent.Flow.Publisher返回值为akka.stream.javadsl.Source物化值为akka.NotUsed——这意味着该 Source 本身不产生有意义的物化结果全部控制权都在上游 Publisher 一侧。三、底层实现从 j.u.c.Flow 到内部 PublisherSourcefromPublisher之所以能零成本接入两种规范关键在于仓库内部维护了一套双向转换器。查看 JavaFlowSupport.scala 的实现def fromPublisherT: Source[T, NotUsed] scaladsl.Source.fromPublisher(publisher.asRs)调用链可以分为两层类型适配层publisher.asRs通过隐式转换把java.util.concurrent.Flow.Publisher[T]包装为org.reactivestreams.Publisher[T]。转换逻辑位于 akka-stream/src/main/scala/akka/stream/impl/JavaFlowAndRsConverters.scala其中asRs方法会判断入参如果是本仓库产出的RsPublisherToJavaFlowAdapter实例则直接解包避免重复包装否则新建JavaFlowPublisherToRsAdapter包装器。所有适配器都只是对subscribe、request、cancel、onNext等信号的透传转发不引入任何额外缓冲或语义变化。图构建层转换后的org.reactivestreams.Publisher进入akka.stream.scaladsl.Source.fromPublisher最终构建为new PublisherSource(publisher, DefaultAttributes.publisherSource, shape(PublisherSource))见 Source.scalaPublisherSource定义于akka-stream/src/main/scala/akka/stream/impl/Modules.scala。PublisherSource负责在物化时订阅上游Publisher并将上游发出的元素与背压信号桥接进 Akka Streams 的图执行引擎。值得注意的是JavaFlowAndRsConverters.scala 明确标注为InternalApi这两组接口本是设计给共享库如数据库驱动做互操作用的应用层不应直接触碰转换器而是统一走JavaFlowSupport门面——这正是文档推荐JavaFlowSupport.Source.fromPublisher的原因。四、实战示例接入响应式数据库客户端官方文档以使用支持 Reactive Streams 的数据库客户端查询行数据为例场景非常典型数据库驱动是Publisher的生产者Akka Streams 是消费方两者都遵守 Reactive Streams 规范因此背压可以贯穿整条链路。Scala 示例完整代码见 akka-docs/src/test/scala-jdk9-only/docs/stream/operators/source/FromPublisher.scalaimport java.util.concurrent.Flow.Subscriber; import java.util.concurrent.Flow.Publisher; import akka.NotUsed; import akka.stream.scaladsl.Source; import akka.stream.scaladsl.JavaFlowSupport; case class Row(name: String) class DatabaseClient { def fetchRows(): Publisher[Row] ??? } val databaseClient: DatabaseClient ??? val names: Source[String, NotUsed] // A new subscriber will subscribe to the supplied publisher for each // materialization, so depending on whether the database client supports // this the Source can be materialized more than once. JavaFlowSupport.Source.fromPublisher(databaseClient.fetchRows()) .map(row row.name);Java 示例完整代码见 akka-docs/src/test/java-jdk9-only/jdocs/stream/operators/source/FromPublisher.javaimport java.util.concurrent.Flow.Publisher; import akka.NotUsed; import akka.stream.javadsl.Source; import akka.stream.javadsl.JavaFlowSupport; class Example { public SourceString, NotUsed names() { // A new subscriber will subscribe to the supplied publisher for each // materialization, so depending on whether the database client supports // this the Source can be materialized more than once. return JavaFlowSupport.Source.RowfromPublisher(databaseClient.fetchRows()) .map(row - row.getField(name)); } }示例解读来源类型databaseClient.fetchRows()返回java.util.concurrent.Flow.Publisher[Row]fromPublisher直接消费它无需任何包装代码。下游加工得到的Source[Row, NotUsed]与普通 Source 无异可直接链式调用.map(row row.name)提取字段再交给runForeach、Sink或其它操作符。物化语义源码注释明确指出——每次物化都会有一个新的订阅者去订阅传入的 Publisher。因此如果数据库客户端支持多次订阅该 Source 可以被物化多次反之若 Publisher 只允许单次订阅重复物化会失败。这决定了该 Source 能否安全复用比如在多个流中共享蓝图。背压贯通数据库驱动与 Akka Streams 都实现了 Reactive Streams 规范PublisherSource会把下游的需求信号demand逐级传递给数据库的Publisher。当消费速度慢于生产速度时上游会暂停产出从而避免数据库行数据在内存中无限堆积导致 OOM。五、关键行为与注意事项1. 背压协调文档强调coordinate backpressure as needed。背压的传递依赖 Reactive Streams 的Subscription.request(n)协议PublisherSource内部依据下游缓冲区和需求状态向上游请求元素上游按请求数量分批产出。这正是数据库行被消费得比产出慢时不会内存溢出的机制保证。2. 上游失败与完成作为规范实现Publisher通过onError或onComplete信号结束流onError会使流以该异常失败可被下游的recover等操作符捕获onComplete则正常完成。这些信号由PublisherSource桥接进图执行引擎与 Akka Streams 的流生命周期保持一致。3. JDK 8 兼容性org.reactivestreams.Publisher 路线由于java.util.concurrent.Flow在 JDK 9 才引入文档对 JDK 8 用户给出了明确的替代方案使用 org.reactivestreams 库的org.reactivestreams.Publisher配合akka.stream.scaladsl.Source.fromPublisherScala或akka.stream.javadsl.Source.fromPublisherJava。该 API 早在 JDK 8 时代就已存在语义与JavaFlowSupport版本完全一致且依赖同一套PublisherSource实现。两者的关系在 JavaFlowSupport.scala 的 Scaladoc 中亦有说明org.reactivestreams版本先于 Java 9 存在两者承载相同语义。4. 与 asSubscriber 的互补关系若你面对的不是提供 Publisher的 API而是接收 Subscriber的 API例如某些库要求传入回调订阅者则应使用JavaFlowSupport.Source.asSubscriber——它在每次物化时产出一个java.util.concurrent.Flow.Subscriber作为物化值可挂接到外部 Publisher 上为其填充元素。详见 asSubscriber.md。两者一拉一推共同覆盖了响应式库互操作的两种常见接口形态。5. JavaFlowSupport 的完整能力从 JavaFlowSupport.scala 可以看到与fromPublisher配套的还有Source.asSubscriber[T]: Source[T, java.util.concurrent.Flow.Subscriber[T]]Flow.fromProcessor/Flow.fromProcessorMat/Flow.toProcessorSink.asPublisher(fanout: Boolean)/Sink.fromSubscriber这意味着 j.u.c.Flow 生态的 Publisher、Subscriber、Processor 三种角色都能与 Akka Streams 的 Source、Flow、Sink 一一对应构成完整的互操作矩阵。六、小结Source.fromPublisher及 JDK 9 的JavaFlowSupport.Source.fromPublisher是 Akka Streams 与响应式数据库驱动、消息客户端等第三方库对接的标准化入口。其背后是仓库内JavaFlowAndRsConverters适配器与PublisherSource图节点的协同前者负责两种规范类型的零开销互转后者负责订阅、背压与生命周期信号的桥接。实际使用时只需把握三点按 JDK 版本选择JavaFlowSupportJDK 9或org.reactivestreams变体JDK 8明确每次物化都会重新订阅 Publisher据此判断 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 Source.asSubscriber 实战将 java.util.concurrent.Flow.Subscriber 无缝接入响应式流Akka Streams Source.asSubscriber 实战将 java.util.concurrent.Flow.Subscriber 无缝接入响后端并发编程异步编程Akka Streams Flow.completionStageFlow基于 CompletionStage 的延迟 Flow 创建与流式接入指南Akka Streams Flow.completionStageFlow基于 CompletionStage 的延迟 Flow 创建与流式接入指南 本篇技术后端并发编程异步编程Akka Streams Sink.asPublisher 完全指南将 Akka Stream 桥接到 Reactive Streams PublisherAkka Streams Sink.asPublisher 完全指南将 Akka Stream 桥接到 Reactive Streams Publisher后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表