后端并发编程异步编程【免费下载链接】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),仅供参考
企业数字化 ERP 产品动态
相关推荐
车道线检测数据集详解:YOLO与VOC双格式训练全流程 简介:面向车道线检测与自动驾驶视觉场景,这份YOLO格式目标检测数据集为开发者提供可直接用于主流YOLO系列模型训练的高质量素材。压缩包共2000个文件,包含1161个XML标签与839个TXT标签,关联1659张带标注图像,整体74.72… · 2026/9/24 18:49:36
C++ Win32塔防游戏教学框架:纯标准库实现可调试游戏骨架 简介:本资源是一套基于C开发的塔防类游戏源码,完整复刻《王国保卫战》核心玩法,专为计算机、自动化等专业本科生课程设计与毕业设计实践打造。代码结构清晰,涵盖游戏主循环、关卡管理、塔与怪物基类、UI界面及音效系统等模块&… · 2026/9/24 18:49:15
EasyWeChat 4.x 快速上手指南:PHP 微信 SDK 的环境要求、安装配置与模块全景 后端即时通讯 【免费下载链接】easywechat 📦 一个 PHP 微信 SDK 项目地址: https://gitcode.com/gh_mirrors/ea/easywechat 点击查看 免费下载 本文以 EasyWeChat 4.x 版本文档为核心,系统讲解该 PHP 微信 SDK 的定位、运行环境、Composer … · 2026/9/24 19:30:47
RunAnywhere 本地 RAG 端到端实践:从 RAG 测试语料看检索增强生成的完整链路 AI模型推理服务推理引擎本地部署多模态 【免费下载链接】runanywhere-sdks Production ready toolkit to run AI locally 项目地址: https://gitcode.com/gh_mirrors/ru/runanywhere-sdks 点击查看 免费下载 导读
本文以 core/tests/data/rag_sample.md 这份位于 … · 2026/9/24 19:30:47
微PE与Ventoy协同实战:系统急救与多ISO启动的底层逻辑 1. 为什么现在还在用微PE?——一个被低估的“系统急救员”真实价值 微PE不是过时的古董,而是我过去三年在27家中小IT服务商、14所高校机房、8个社区维修点反复验证过的“最小可靠解”。它不炫技,不联网,不依赖硬件抽象层ÿ… · 2026/9/24 19:30:33
Python模板注入检测工具源码解析:SSTI检测与利用实战 简介:这是一套面向Web安全研究人员与渗透测试学习者的Server-Side模板注入与代码注入检测利用工具源码,采用Python开发,可帮助读者理解模板引擎漏洞的检测逻辑与利用方式,适合具备一定安全基础的中高级人员研究参考。资源包共103个… · 2026/9/24 19:30:33
Java银行排号系统源码与数据库设计:并发取号、队列调度及论文框架 简介:这份资源是面向高校计算机专业学生与Java初学者的一套银行排号系统完整毕业设计资料,围绕服务器端与客户端双模块架构展开,可用于课程设计、毕设选题或Java桌面应用练手。系统功能划分清晰:服务器端涵盖取号、统计、删除、查… · 2026/9/24 19:30:33
MySQL主从数据一致性校验:pt-table-checksum实战与踩坑总结 说到底,MySQL主从数据一致性这事儿,绝大多数人一开始都会觉得“只要复制没断,数据就肯定一样”。我以前也是这么想的,直到线上因为一个漏加主键的同步脚本,主从数据悄悄分裂了一个多月才被发现,当时配合业务… · 2026/9/24 19:30:27
基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程 简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为… · 2026/9/24 0:00:13
1D-CNN时间序列建模实战:从Conv1d原理到工业落地 简介:面向时间序列数据建模的一维卷积神经网络完整实现,适合深度学习入门者及需要快速验证时序模型的研究者,能够从音频、文本、传感器或股价等序列中挖掘局部特征与时间依赖。压缩包体积很小,只有3KB,内含3个Python脚… · 2026/9/24 0:00:26
柔软的L:汉语语流中被忽视的舌肌张力控制 1. 这个“L”不是字母表里的L,而是舌尖上的L最近在几个方言群和语音教学社群里,反复看到有人发一句:“也说字母L:柔软的长舌”。初看以为是英语发音课笔记,点开才发现全是方言爱好者、播音系学生、语言康复师甚至戏曲演… · 2026/9/24 0:00:44