爆款云主机2核4G限时秒杀,88元/年起!
查看详情

活动

天翼云最新优惠活动,涵盖免费试用,产品折扣等,助您降本增效!
热门活动
  • 618智算钜惠季 爆款云主机2核4G限时秒杀,88元/年起!
  • 免费体验DeepSeek,上天翼云息壤 NEW 新老用户均可免费体验2500万Tokens,限时两周
  • 云上钜惠 HOT 爆款云主机全场特惠,更有万元锦鲤券等你来领!
  • 算力套餐 HOT 让算力触手可及
  • 天翼云脑AOne NEW 连接、保护、办公,All-in-One!
  • 中小企业应用上云专场 产品组合下单即享折上9折起,助力企业快速上云
  • 息壤高校钜惠活动 NEW 天翼云息壤杯高校AI大赛,数款产品享受线上订购超值特惠
  • 天翼云电脑专场 HOT 移动办公新选择,爆款4核8G畅享1年3.5折起,快来抢购!
  • 天翼云奖励推广计划 加入成为云推官,推荐新用户注册下单得现金奖励
免费活动
  • 免费试用中心 HOT 多款云产品免费试用,快来开启云上之旅
  • 天翼云用户体验官 NEW 您的洞察,重塑科技边界

智算服务

打造统一的产品能力,实现算网调度、训练推理、技术架构、资源管理一体化智算服务
智算云(DeepSeek专区)
科研助手
  • 算力商城
  • 应用商城
  • 开发机
  • 并行计算
算力互联调度平台
  • 应用市场
  • 算力市场
  • 算力调度推荐
一站式智算服务平台
  • 模型广场
  • 体验中心
  • 服务接入
智算一体机
  • 智算一体机
大模型
  • DeepSeek-R1-昇腾版(671B)
  • DeepSeek-R1-英伟达版(671B)
  • DeepSeek-V3-昇腾版(671B)
  • DeepSeek-R1-Distill-Llama-70B
  • DeepSeek-R1-Distill-Qwen-32B
  • Qwen2-72B-Instruct
  • StableDiffusion-V2.1
  • TeleChat-12B

应用商城

天翼云精选行业优秀合作伙伴及千余款商品,提供一站式云上应用服务
进入甄选商城进入云市场创新解决方案
办公协同
  • WPS云文档
  • 安全邮箱
  • EMM手机管家
  • 智能商业平台
财务管理
  • 工资条
  • 税务风控云
企业应用
  • 翼信息化运维服务
  • 翼视频云归档解决方案
工业能源
  • 智慧工厂_生产流程管理解决方案
  • 智慧工地
建站工具
  • SSL证书
  • 新域名服务
网络工具
  • 翼云加速
灾备迁移
  • 云管家2.0
  • 翼备份
资源管理
  • 全栈混合云敏捷版(软件)
  • 全栈混合云敏捷版(一体机)
行业应用
  • 翼电子教室
  • 翼智慧显示一体化解决方案

合作伙伴

天翼云携手合作伙伴,共创云上生态,合作共赢
天翼云生态合作中心
  • 天翼云生态合作中心
天翼云渠道合作伙伴
  • 天翼云代理渠道合作伙伴
天翼云服务合作伙伴
  • 天翼云集成商交付能力认证
天翼云应用合作伙伴
  • 天翼云云市场合作伙伴
  • 天翼云甄选商城合作伙伴
天翼云技术合作伙伴
  • 天翼云OpenAPI中心
  • 天翼云EasyCoding平台
天翼云培训认证
  • 天翼云学堂
  • 天翼云市场商学院
天翼云合作计划
  • 云汇计划
天翼云东升计划
  • 适配中心
  • 东升计划
  • 适配互认证

开发者

开发者相关功能入口汇聚
技术社区
  • 专栏文章
  • 互动问答
  • 技术视频
资源与工具
  • OpenAPI中心
开放能力
  • EasyCoding敏捷开发平台
培训与认证
  • 天翼云学堂
  • 天翼云认证
魔乐社区
  • 魔乐社区

支持与服务

为您提供全方位支持与服务,全流程技术保障,助您轻松上云,安全无忧
文档与工具
  • 文档中心
  • 新手上云
  • 自助服务
  • OpenAPI中心
定价
  • 价格计算器
  • 定价策略
基础服务
  • 售前咨询
  • 在线支持
  • 在线支持
  • 工单服务
  • 建议与反馈
  • 用户体验官
  • 服务保障
  • 客户公告
  • 会员中心
增值服务
  • 红心服务
  • 首保服务
  • 客户支持计划
  • 专家技术服务
  • 备案管家

了解天翼云

天翼云秉承央企使命,致力于成为数字经济主力军,投身科技强国伟大事业,为用户提供安全、普惠云服务
品牌介绍
  • 关于天翼云
  • 智算云
  • 天翼云4.0
  • 新闻资讯
  • 天翼云APP
基础设施
  • 全球基础设施
  • 信任中心
最佳实践
  • 精选案例
  • 超级探访
  • 云杂志
  • 分析师和白皮书
  • 天翼云·创新直播间
市场活动
  • 2025智能云生态大会
  • 2024智算云生态大会
  • 2023云生态大会
  • 2022云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心

      Reactor:深入理解reactor core

      首页 知识中心 其他 文章详情页

      Reactor:深入理解reactor core

      2023-03-24 10:31:28 阅读次数:150

      core

      目录

      • ​​简介​​
      • ​​自定义Subscriber​​
      • ​​Backpressure处理​​
      • ​​创建Flux​​
      • ​​使用generate​​
      • ​​使用create​​
      • ​​使用push​​
      • ​​使用Handle​​

      简介

      上篇文章我们简单的介绍了Reactor的发展史和基本的Flux和Mono的使用,本文将会进一步挖掘Reactor的高级用法,一起来看看吧。

      自定义Subscriber

      之前的文章我们提到了4个Flux的subscribe的方法:

      Disposable subscribe(); 

      Disposable subscribe(Consumer<? super T> consumer);

      Disposable subscribe(Consumer<? super T> consumer,
      Consumer<? super Throwable> errorConsumer);

      Disposable subscribe(Consumer<? super T> consumer,
      Consumer<? super Throwable> errorConsumer,
      Runnable completeConsumer);

      Disposable subscribe(Consumer<? super T> consumer,
      Consumer<? super Throwable> errorConsumer,
      Runnable completeConsumer,
      Consumer<? super Subscription> subscriptionConsumer);

      这四个方法,需要我们使用lambda表达式来自定义consumer,errorConsumer,completeSonsumer和subscriptionConsumer这四个Consumer。

      写起来比较复杂,看起来也不太方便,我们考虑一下,这四个Consumer是不是和Subscriber接口中定义的4个方法是一一对应的呢?

      public static interface Subscriber<T> {

      public void onSubscribe(Subscription subscription);

      public void onNext(T item);

      public void onError(Throwable throwable);

      public void onComplete();
      }

      对的,所以我们有一个更加简单点的subscribe方法:

      public final void subscribe(Subscriber<? super T> actual)

      这个subscribe方法直接接收一个Subscriber类。从而实现了所有的功能。

      自己写Subscriber太麻烦了,Reactor为我们提供了一个BaseSubscriber的类,它实现了Subscriber中的所有功能,还附带了一些其他的方法。

      我们看下BaseSubscriber的定义:

      public abstract class BaseSubscriber<T> implements CoreSubscriber<T>, Subscription,
      Disposable

      注意,BaseSubscriber是单次使用的,这就意味着,如果它首先subscription到Publisher1,然后subscription到Publisher2,那么将会取消对第一个Publisher的订阅。

      因为BaseSubscriber是一个抽象类,所以我们需要继承它,并且重写我们需要自己实现的方法。

      下面看一个自定义的Subscriber:

      public class CustSubscriber<T> extends BaseSubscriber<T> {

      public void hookOnSubscribe(Subscription subscription) {
      System.out.println("Subscribed");
      request(1);
      }

      public void hookOnNext(T value) {
      System.out.println(value);
      request(1);
      }
      }

      BaseSubscriber中有很多以hook开头的方法,这些方法都是我们可以重写的,而Subscriber原生定义的on开头的方法,在BaseSubscriber中都是final的,都是不能重写的。

      我们看一个定义:

      @Override
      public final void onSubscribe(Subscription s) {
      if (Operators.setOnce(S, this, s)) {
      try {
      hookOnSubscribe(s);
      }
      catch (Throwable throwable) {
      onError(Operators.onOperatorError(s, throwable, currentContext()));
      }
      }
      }

      可以看到,它内部实际上调用了hook的方法。

      上面的CustSubscriber中,我们重写了两个方法,一个是hookOnSubscribe,在建立订阅的时候调用,一个是hookOnNext,在收到onNext信号的时候调用。

      在这些方法中,给了我们足够的自定义空间,上面的例子中我们调用了request(1),表示再请求一个元素。

      其他的hook方法还有: hookOnComplete, hookOnError, hookOnCancel 和 hookFinally。

      Backpressure处理

      我们之前讲过了,reactive stream的最大特征就是可以处理Backpressure。

      什么是Backpressure呢?就是当consumer处理过不来的时候,可以通知producer来减少生产速度。

      我们看下BaseSubscriber中默认的hookOnSubscribe实现:

      protected void hookOnSubscribe(Subscription subscription){
      subscription.request(Long.MAX_VALUE);
      }

      可以看到默认是request无限数目的值。 也就是说默认情况下没有Backpressure。

      通过重写hookOnSubscribe方法,我们可以自定义处理速度。

      除了request之外,我们还可以在publisher中限制subscriber的速度。

      public final Flux<T> limitRate(int prefetchRate) {
      return onAssembly(this.publishOn(Schedulers.immediate(), prefetchRate));
      }

      在Flux中,我们有一个limitRate方法,可以设定publisher的速度。

      比如subscriber request(100),然后我们设置limitRate(10),那么最多producer一次只会产生10个元素。

      创建Flux

      接下来,我们要讲解一下怎么创建Flux,通常来讲有4种方法来创建Flux。

      使用generate

      第一种方法就是最简单的同步创建的generate.

      先看一个例子:

      public void useGenerate(){
      Flux<String> flux = Flux.generate(
      () -> 0,
      (state, sink) -> {
      sink.next("3 x " + state + " = " + 3*state);
      if (state == 10) sink.complete();
      return state + 1;
      });

      flux.subscribe(System.out::println);
      }

      输出结果:

      3 x 0 = 0
      3 x 1 = 3
      3 x 2 = 6
      3 x 3 = 9
      3 x 4 = 12
      3 x 5 = 15
      3 x 6 = 18
      3 x 7 = 21
      3 x 8 = 24
      3 x 9 = 27
      3 x 10 = 30

      上面的例子中,我们使用generate方法来同步的生成元素。

      generate接收两个参数:

      public static <T, S> Flux<T> generate(Callable<S> stateSupplier, BiFunction<S, SynchronousSink<T>, S> generator)

      第一个参数是stateSupplier,用来指定初始化的状态。

      第二个参数是一个generator,用来消费SynchronousSink,并生成新的状态。

      上面的例子中,我们每次将state+1,一直加到10。

      然后使用subscribe来将所有的生成元素输出。

      使用create

      Flux也提供了一个create方法来创建Flux,create可以是同步也可以是异步的,并且支持多线程操作。

      因为create没有初始的state状态,所以可以用在多线程中。

      create的一个非常有用的地方就是可以将第三方的异步API和Flux关联起来,举个例子,我们有一个自定义的EventProcessor,当处理相应的事件的时候,会去调用注册到Processor中的listener的一些方法。

      interface MyEventListener<T> {
      void onDataChunk(List<T> chunk);
      void processComplete();
      }

      我们怎么把这个Listener的响应行为和Flux关联起来呢?

      public void useCreate(){
      EventProcessor myEventProcessor = new EventProcessor();
      Flux<String> bridge = Flux.create(sink -> {
      myEventProcessor.register(
      new MyEventListener<String>() {
      public void onDataChunk(List<String> chunk) {
      for(String s : chunk) {
      sink.next(s);
      }
      }
      public void processComplete() {
      sink.complete();
      }
      });
      });
      }

      使用create就够了,create接收一个consumer参数:

      public static <T> Flux<T> create(Consumer<? super FluxSink<T>> emitter)

      这个consumer的本质是去消费FluxSink对象。

      上面的例子在MyEventListener的事件中对FluxSink对象进行消费。

      使用push

      push和create一样,也支持异步操作,但是同时只能有一个线程来调用next, complete 或者 error方法,所以它是单线程的。

      使用Handle

      Handle和上面的三个方法不同,它是一个实例方法。

      它和generate很类似,也是消费SynchronousSink对象。

      Flux<R> handle(BiConsumer<T, SynchronousSink<R>>);

      不同的是它的参数是一个BiConsumer,是没有返回值的。

      看一个使用的例子:

      public void useHandle(){
      Flux<String> alphabet = Flux.just(-1, 30, 13, 9, 20)
      .handle((i, sink) -> {
      String letter = alphabet(i);
      if (letter != null)
      sink.next(letter);
      });

      alphabet.subscribe(System.out::println);
      }

      public String alphabet(int letterNumber) {
      if (letterNumber < 1 || letterNumber > 26) {
      return null;
      }
      int letterIndexAscii = 'A' + letterNumber - 1;
      return "" + (char) letterIndexAscii;
      }

      本文的例子​​learn-reactive​​

      版权声明:本文内容来自第三方投稿或授权转载,原文地址:https://blog.51cto.com/flydean/5689738,作者:程序那些事,版权归原作者所有。本网站转在其作品的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如因作品内容、版权等问题需要同本网站联系,请发邮件至ctyunbbs@chinatelecom.cn沟通。

      上一篇:Docker registry 你有深入的玩过吗?

      下一篇:Pandas高级教程之:window操作

      相关文章

      2025-05-13 09:50:48

      error: "net.bridge.bridge-nf-call-ip6tables" is an unknown key

      error: "net.bridge.bridge-nf-call-ip6tables" is an unknown key

      2025-05-13 09:50:48
      core , kernel
      2025-04-01 10:28:07

      git的学习笔记之设置(忽略换行 权限)

      git的学习笔记之设置(忽略换行 权限)

      2025-04-01 10:28:07
      config , core , false , git , merge , true
      2025-02-26 07:20:49

      【性能】Linux服务器低延迟技术

      【性能】Linux服务器低延迟技术

      2025-02-26 07:20:49
      core , 中断 , 延迟 , 线程
      2024-12-24 10:19:43

      使用gcore生成当前崩溃进程生成dump文件并定位错误

      使用gcore生成当前崩溃进程生成dump文件并定位错误

      2024-12-24 10:19:43
      core , dump , gdb , 文件 , 生成
      2024-12-11 06:42:09

      Linux下从一个线程获取另一个线程的函数调用堆栈并生成转储文件

      Linux下从一个线程获取另一个线程的函数调用堆栈并生成转储文件

      2024-12-11 06:42:09
      core , cpp , gdb , pthread
      2024-12-11 06:20:04

      详细分析Linux中的core dump异常(附 Demo排查)

      Core dump 是指在程序异常终止时,操作系统将程序的内存映像保存到磁盘上的一种机制。

      2024-12-11 06:20:04
      core , dump , 内存 , 文件 , 生成 , 程序 , 错误
      2024-12-10 06:59:17

      出现 Error querying database. Cause: com.baomidou.mybatisplus.core.exceptions 等解决方法

      出现 Error querying database. Cause: com.baomidou.mybatisplus.core.exceptions 等解决方法

      2024-12-10 06:59:17
      com , core , 问题
      2024-06-17 10:03:58

      bash检测CPU的物理个数、逻辑个数、core个数等信息

      bash检测CPU的物理个数、逻辑个数、core个数等信息

      2024-06-17 10:03:58
      core , CPU
      2024-06-06 08:24:01

      02jumpserver之部署core

      02jumpserver之部署core

      2024-06-06 08:24:01
      core , jumpserver , 镜像
      查看更多
      推荐标签

      作者介绍

      天翼云小翼
      天翼云用户

      文章

      33561

      阅读量

      5248934

      查看更多

      最新文章

      error: "net.bridge.bridge-nf-call-ip6tables" is an unknown key

      2025-05-13 09:50:48

      git的学习笔记之设置(忽略换行 权限)

      2025-04-01 10:28:07

      使用gcore生成当前崩溃进程生成dump文件并定位错误

      2024-12-24 10:19:43

      出现 Error querying database. Cause: com.baomidou.mybatisplus.core.exceptions 等解决方法

      2024-12-10 06:59:17

      bash检测CPU的物理个数、逻辑个数、core个数等信息

      2024-06-17 10:03:58

      查看更多

      热门文章

      bash检测CPU的物理个数、逻辑个数、core个数等信息

      2024-06-17 10:03:58

      出现 Error querying database. Cause: com.baomidou.mybatisplus.core.exceptions 等解决方法

      2024-12-10 06:59:17

      使用gcore生成当前崩溃进程生成dump文件并定位错误

      2024-12-24 10:19:43

      git的学习笔记之设置(忽略换行 权限)

      2025-04-01 10:28:07

      error: "net.bridge.bridge-nf-call-ip6tables" is an unknown key

      2025-05-13 09:50:48

      查看更多

      热门标签

      linux java python javascript 数组 前端 docker Linux vue 函数 shell git 节点 容器 示例
      查看更多

      相关产品

      弹性云主机

      随时自助获取、弹性伸缩的云服务器资源

      天翼云电脑(公众版)

      便捷、安全、高效的云电脑服务

      对象存储

      高品质、低成本的云上存储服务

      云硬盘

      为云上计算资源提供持久性块存储

      查看更多

      随机文章

      出现 Error querying database. Cause: com.baomidou.mybatisplus.core.exceptions 等解决方法

      error: "net.bridge.bridge-nf-call-ip6tables" is an unknown key

      使用gcore生成当前崩溃进程生成dump文件并定位错误

      bash检测CPU的物理个数、逻辑个数、core个数等信息

      git的学习笔记之设置(忽略换行 权限)

      • 7*24小时售后
      • 无忧退款
      • 免费备案
      • 专家服务
      售前咨询热线
      400-810-9889转1
      关注天翼云
      • 旗舰店
      • 天翼云APP
      • 天翼云微信公众号
      服务与支持
      • 备案中心
      • 售前咨询
      • 智能客服
      • 自助服务
      • 工单管理
      • 客户公告
      • 涉诈举报
      账户管理
      • 管理中心
      • 订单管理
      • 余额管理
      • 发票管理
      • 充值汇款
      • 续费管理
      快速入口
      • 天翼云旗舰店
      • 文档中心
      • 最新活动
      • 免费试用
      • 信任中心
      • 天翼云学堂
      云网生态
      • 甄选商城
      • 渠道合作
      • 云市场合作
      了解天翼云
      • 关于天翼云
      • 天翼云APP
      • 服务案例
      • 新闻资讯
      • 联系我们
      热门产品
      • 云电脑
      • 弹性云主机
      • 云电脑政企版
      • 天翼云手机
      • 云数据库
      • 对象存储
      • 云硬盘
      • Web应用防火墙
      • 服务器安全卫士
      • CDN加速
      热门推荐
      • 云服务备份
      • 边缘安全加速平台
      • 全站加速
      • 安全加速
      • 云服务器
      • 云主机
      • 智能边缘云
      • 应用编排服务
      • 微服务引擎
      • 共享流量包
      更多推荐
      • web应用防火墙
      • 密钥管理
      • 等保咨询
      • 安全专区
      • 应用运维管理
      • 云日志服务
      • 文档数据库服务
      • 云搜索服务
      • 数据湖探索
      • 数据仓库服务
      友情链接
      • 中国电信集团
      • 189邮箱
      • 天翼企业云盘
      • 天翼云盘
      ©2025 天翼云科技有限公司版权所有 增值电信业务经营许可证A2.B1.B2-20090001
      公司地址:北京市东城区青龙胡同甲1号、3号2幢2层205-32室
      • 用户协议
      • 隐私政策
      • 个人信息保护
      • 法律声明
      备案 京公网安备11010802043424号 京ICP备 2021034386号