Reactor响应式编程中如何实现带优先级和可控缓冲的生产者-消费者模式

Reactor响应式编程中如何实现带优先级和可控缓冲的生产者-消费者模式

java reactor的生产者-消费者模式中,当内置sinks无法满足任务优先级、队列监控及清空需求时,可利用`sinks.many().unicast().onbackpressurebuffer()`结合外部`priorityqueue`实现高效、可控的异步任务处理,避免阻塞式操作,从而构建一个功能更强大的响应式任务处理系统。

1. Reactor生产者-消费者模式中的挑战

在基于Reactor的应用程序中,生产者-消费者模式常用于异步任务处理。通常,我们会使用Sinks.Many来在生产者和消费者之间传递数据,例如:

Sinks.Many taskSink = Sinks.many().multicast().onBackpressureBuffer(1000, false);// 生产者Flux dates = loadDates();dates.filterWhen(...)     .concatMap(date -> taskManager.getTaskByDate(date))     .doOnNext(taskSink::tryEmitNext)     .subscribe();// 消费者taskProcessor.process(taskSink.asFlux())             .subscribeOn(Schedulers.boundedElastic())             .subscribe();

这种实现方式在大多数情况下运行良好。然而,当系统在高负载下运行时,我们可能会遇到以下痛点:

队列可见性: 无法直接获取Sink中当前待处理任务的数量。队列清空: 无法方便地清空Sink中所有待处理任务。任务优先级: 无法对Sink中的任务进行优先级排序。

由于标准Sinks.Many不提供对内部缓冲队列的直接访问,上述需求难以满足。

2. 避免阻塞式操作:为什么poll()在响应式编程中是问题

为了解决上述问题,一种常见的尝试是引入自定义的包装类,其中包含一个PriorityBlockingQueue,并通过Flux.create结合poll()方法从队列中获取元素:

// 自定义任务队列MergingQueue taskQueue = new PriorityMergingQueue();// 生产者Flux dates = loadDates();dates.filterWhen(...)     .concatMap(date -> taskManager.getTaskByDate(date))     .doOnNext(taskQueue::enqueue) // 将任务加入自定义队列     .subscribe();// 消费者taskProcessor.process(Flux.create((sink) -> {     sink.onRequest(n -> {          Task task;          try {                while(!sink.isCancel() && n > 0) {                    // 潜在的阻塞操作                    if((task = taskQueue.poll(1, TimeUnit.SECOND))  != null) {                        sink.next(task);                        n--;                    }                }          } catch(InterruptedException e) {                Thread.currentThread().interrupt();                sink.error(e);          }     });})).subscribeOn(Schedulers.boundedElastic()).subscribe();

尽管这种方法似乎解决了优先级和队列访问的问题,但其中使用PriorityBlockingQueue.poll(1, TimeUnit.SECOND)是一个阻塞式操作。响应式编程的核心目标之一就是避免阻塞,因为阻塞操作会占用线程并等待,这与Reactor的非阻塞、异步特性相悖。在长时间运行测试中,阻塞式poll()可能导致线程挂起,严重影响系统的响应性和吞吐量。

3. Reactor原生解决方案:结合Sinks.many().unicast()与外部PriorityQueue

Reactor提供了一个更优雅、更符合响应式编程原则的解决方案:利用Sinks.many().unicast().onBackpressureBuffer(Queue)。这个方法允许我们提供一个外部的Queue作为Sink的背压缓冲区。通过使用PriorityQueue作为这个外部队列,我们可以完美地解决任务优先级、队列可见性和清空的问题,同时保持非阻塞特性。

3.1 核心原理

Sinks.many().unicast().onBackpressureBuffer(Queue)的unicast特性意味着只有一个订阅者可以消费Sink发出的元素。当这个订阅者无法及时处理元素时,新发出的元素会被存储到我们提供的外部Queue中。

通过这种方式,我们可以:

实现优先级: 将一个PriorityQueue作为外部队列,并定义好任务的比较器,Sink将自动按照优先级从队列中取出任务。监控队列: 直接访问外部PriorityQueue实例,调用其size()方法即可获取当前待处理任务数量。清空队列: 直接调用外部PriorityQueue实例的clear()方法即可清空所有待处理任务。

3.2 示例代码

下面是一个演示如何使用外部PriorityQueue与Sinks.many().unicast()实现带优先级、可监控的生产者-消费者模式的例子。

import reactor.core.publisher.Sinks;import reactor.core.scheduler.Schedulers;import java.time.Duration;import java.time.LocalTime;import java.time.ZoneOffset;import java.time.temporal.ChronoUnit;import java.util.Comparator;import java.util.Queue;import java.util.PriorityQueue;import java.util.concurrent.TimeUnit;// 任务记录类,包含优先级record Task(int prio, String name) {}public class PriorityTaskProcessor {    private static void log(Object message) {        System.out.println(LocalTime.now(ZoneOffset.UTC).truncatedTo(ChronoUnit.MILLIS) + ": " + message);    }    public void externalBufferDemo() throws InterruptedException {        // 1. 创建一个PriorityQueue作为外部缓冲区        // 优先级高的(prio值大)的任务先处理,所以使用reversed()        Queue taskQueue = new PriorityQueue(Comparator.comparingInt(Task::prio).reversed());        // 2. 创建unicast Sink,并指定使用外部的PriorityQueue作为背压缓冲区        Sinks.Many taskSink = Sinks.many().unicast().onBackpressureBuffer(taskQueue);        // 3. 消费者:订阅Sink发出的Flux        // 为了演示效果,这里模拟一个处理延迟        taskSink.asFlux()                .delayElements(Duration.ofMillis(100)) // 模拟每个任务处理需要100ms                .doOnNext(task -> log("处理任务: " + task))                .subscribeOn(Schedulers.boundedElastic()) // 在弹性调度器上执行处理逻辑                .subscribe(                        task -> {}, // onNext                        error -> log("消费者发生错误: " + error.getMessage()), // onError                        () -> log("消费者完成") // onComplete                );        // 4. 生产者:向Sink发射任务        log("开始发射任务...");        for (int i = 0; i < 10; i++) {            // 发射不同优先级的任务            taskSink.tryEmitNext(new Task(i, "Task-" + i));            // 模拟生产者快速生产            Thread.sleep(10);        }        log("任务发射完毕.");        // 5. 检查Sink中任务数量(直接访问外部队列)        log("当前Sink中待处理任务数量: " + taskQueue.size());        // 6. 模拟一段时间后清空队列        Thread.sleep(350); // 等待一些任务被处理        log("准备清空Sink中所有待处理任务...");        taskQueue.clear(); // 直接清空外部PriorityQueue        log("清空后Sink中待处理任务数量: " + taskQueue.size());        // 7. 继续等待,观察清空后的处理情况        Thread.sleep(1500);        log("演示结束.");    }    public static void main(String[] args) throws InterruptedException {        new PriorityTaskProcessor().externalBufferDemo();    }}

3.3 运行输出分析

运行上述代码,你可能会看到类似以下的输出(时间戳会有所不同):

09:41:11.347: 开始发射任务...09:41:11.437: 任务发射完毕.09:41:11.437: 当前Sink中待处理任务数量: 909:41:11.539: 处理任务: Task[prio=0, name=Task-0] // Task-0先被处理,因为delayElements的内部队列09:41:11.642: 处理任务: Task[prio=9, name=Task-9] // 队列中的最高优先级任务09:41:11.745: 处理任务: Task[prio=8, name=Task-8] // 接下来是优先级8的任务09:41:11.787: 准备清空Sink中所有待处理任务...09:41:11.787: 清空后Sink中待处理任务数量: 009:41:11.848: 处理任务: Task[prio=7, name=Task-7] // 注意:此任务在清空后仍被处理09:41:12.051: 演示结束.

输出解释:

当前Sink中待处理任务数量: 9: 在所有任务发射完毕后,由于消费者有100ms的延迟,大部分任务都进入了taskQueue。第一个任务Task-0可能在生产者循环结束前就被delayElements的内部队列捕获并开始处理,所以外部队列中剩下9个。处理任务: Task[prio=0, name=Task-0]: 尽管PriorityQueue通常会先处理优先级最高的任务,但delayElements操作符本身有一个内部队列(通常大小为1)。这意味着当Task-0被Sink发出后,它可能立即进入delayElements的内部队列并开始计时,在Task-1被发出之前就已经被消费者处理。处理任务: Task[prio=9, name=Task-9]: 紧接着,PriorityQueue的优先级特性开始生效。Task-9(优先级最高)被取出并处理。处理任务: Task[prio=8, name=Task-8]: 随后是Task-8。清空后Sink中待处理任务数量: 0: 调用taskQueue.clear()后,外部队列被清空。处理任务: Task[prio=7, name=Task-7]: 尽管外部队列已清空,但Task-7仍然被处理了。这是因为在taskQueue.clear()被调用之前,Task-7可能已经从taskQueue中取出,并进入了delayElements操作符的内部队列中等待处理。

4. 关于多播(Multicast)的需求

上述解决方案使用了unicast Sink,这意味着只有一个订阅者可以消费其发出的元素。如果您的业务场景确实需要多个消费者订阅同一个Flux,并让他们都能接收到外部PriorityQueue中按优先级取出的任务,您可以在taskSink.asFlux()之后,利用Reactor提供的多播操作符来实现:

// 如果需要多播,可以在unicast Sink的Flux上应用多播操作符taskSink.asFlux()        .publish() // 或 .share(), .replay() 等        .autoConnect(2) // 示例:等待2个订阅者连接后开始发射        .delayElements(Duration.ofMillis(100))        .subscribe(consumer1); // 第一个消费者taskSink.asFlux()        .publish() // 再次强调,多播操作符应作用在原始Flux上,而不是创建多个Flux        .autoConnect(2)        .delayElements(Duration.ofMillis(100))        .subscribe(consumer2); // 第二个消费者

重要提示: 在这种多播场景下,虽然外部PriorityQueue确保了任务在进入Sink时的优先级,但一旦任务被Sink发出并进入多播管道,每个订阅者会独立接收到这些任务。如果多个消费者需要独立地、按照自己的节奏处理任务,且每个消费者都需要完整的优先级队列功能,那么可能需要为每个消费者维护一个独立的unicast Sink和PriorityQueue,或者重新评估多播的必要性。

5. 总结与最佳实践

通过将Sinks.many().unicast().onBackpressureBuffer()与外部PriorityQueue结合使用,我们能够:

实现任务优先级: 确保高优先级任务优先处理。增强可观测性: 轻松获取待处理任务的数量。提供控制能力: 允许动态清空待处理任务。保持响应式特性: 避免了阻塞式poll()操作,符合Reactor的非阻塞编程范式。

这种模式为构建高效、可控且符合响应式原则的生产者-消费者系统提供了一个强大的工具,尤其适用于需要精细化任务调度和监控的场景。在设计响应式系统时,应始终优先考虑Reactor提供的原生操作符和机制,以充分利用其非阻塞和异步的优势。

以上就是Reactor响应式编程中如何实现带优先级和可控缓冲的生产者-消费者模式的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/11972.html

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
上一篇 2025年11月18日 07:30:26
下一篇 2025年11月18日 07:52:51

相关推荐

  • soul怎么发长视频瞬间_Soul长视频瞬间发布方法

    可通过分段发布、格式转换或剪辑压缩三种方法在Soul上传长视频。一、将长视频用相册编辑功能拆分为多个30秒内片段,依次发布并标注“Part 1”“Part 2”保持连贯;二、使用“格式工厂”等工具将视频转为MP4(H.264)、分辨率≤1080p、帧率≤30fps、大小≤50MB,适配平台要求;三、…

    2025年12月6日 软件教程
    400
  • 天猫app淘金币抵扣怎么使用

    在天猫app购物时,淘金币是一项能够帮助你节省开支的实用功能。掌握淘金币的抵扣使用方法,能让你以更实惠的价格买到心仪商品。 当你选好商品并准备下单时,记得查看商品页面是否支持淘金币抵扣。如果该商品支持此项功能,在提交订单的页面会明确显示相关提示。你会看到淘金币的具体抵扣比例——通常情况下,淘金币可按…

    2025年12月6日 软件教程
    500
  • Pboot插件缓存机制的详细解析_Pboot插件缓存清理的命令操作

    插件功能异常或页面显示陈旧内容可能是缓存未更新所致。PbootCMS通过/runtime/cache/与/runtime/temp/目录缓存插件配置、模板解析结果和数据库查询数据,提升性能但影响调试。解决方法包括:1. 手动删除上述目录下所有文件;2. 后台进入“系统工具”-“缓存管理”,勾选插件、…

    2025年12月6日 软件教程
    100
  • Word2013如何插入SmartArt图形_Word2013SmartArt插入的视觉表达

    答案:可通过四种方法在Word 2013中插入SmartArt图形。一、使用“插入”选项卡中的“SmartArt”按钮,选择所需类型并插入;二、从快速样式库中选择常用模板如组织结构图直接应用;三、复制已有SmartArt图形到目标文档后调整内容与格式;四、将带项目符号的文本选中后右键转换为Smart…

    2025年12月6日 软件教程
    000
  • 《kk键盘》一键发图开启方法

    如何在kk键盘中开启一键发图功能? 1、打开手机键盘,找到并点击“kk”图标。 2、进入工具菜单后,选择“一键发图”功能入口。 3、点击“去开启”按钮,跳转至无障碍服务设置页面。 4、在系统通用设置中,进入“已下载的应用”列表。 j2me3D游戏开发简单教程 中文WORD版 本文档主要讲述的是j2m…

    2025年12月6日 软件教程
    100
  • 怎样用免费工具美化PPT_免费美化PPT的实用方法分享

    利用KIMI智能助手可免费将PPT美化为科技感风格,但需核对文字准确性;2. 天工AI擅长优化内容结构,提升逻辑性,适合高质量内容需求;3. SlidesAI支持语音输入与自动排版,操作便捷,利于紧急场景;4. Prezo提供多种模板,自动生成图文并茂幻灯片,适合学生与初创团队。 如果您有一份内容完…

    2025年12月6日 软件教程
    000
  • Pages怎么协作编辑同一文档 Pages多人实时协作的流程

    首先启用Pages共享功能,点击右上角共享按钮并选择“添加协作者”,设置为可编辑并生成链接;接着复制链接通过邮件或社交软件发送给成员,确保其使用Apple ID登录iCloud后即可加入编辑;也可直接在共享菜单中输入邮箱地址定向邀请,设定编辑权限后发送;最后在共享面板中管理协作者权限,查看实时在线状…

    2025年12月6日 软件教程
    100
  • 哔哩哔哩的视频卡在加载中怎么办_哔哩哔哩视频加载卡顿解决方法

    视频加载停滞可先切换网络或重启路由器,再清除B站缓存并重装应用,接着调低播放清晰度并关闭自动选分辨率,随后更改播放策略为AVC编码,最后关闭硬件加速功能以恢复播放。 如果您尝试播放哔哩哔哩的视频,但进度条停滞在加载状态,无法继续播放,这通常是由于网络、应用缓存或播放设置等因素导致。以下是解决此问题的…

    2025年12月6日 软件教程
    000
  • REDMI K90系列正式发布,售价2599元起!

    10月23日,redmi k90系列正式亮相,推出redmi k90与redmi k90 pro max两款新机。其中,redmi k90搭载骁龙8至尊版处理器、7100mah大电池及100w有线快充等多项旗舰配置,起售价为2599元,官方称其为k系列迄今为止最完整的标准版本。 图源:REDMI红米…

    2025年12月6日 行业动态
    200
  • 买家网购苹果手机仅退款不退货遭商家维权,法官调解后支付货款

    10 月 24 日消息,据央视网报道,近年来,“仅退款”服务逐渐成为众多网购平台的常规配置,但部分消费者却将其当作“免费试用”的手段,滥用规则谋取私利。 江苏扬州市民李某在某电商平台购买了一部苹果手机,第二天便以“不想要”为由在线申请“仅退款”,当时手机尚在物流运输途中。第三天货物送达后,李某签收了…

    2025年12月6日 行业动态
    000
  • Linux中如何安装Nginx服务_Linux安装Nginx服务的完整指南

    首先更新系统软件包,然后通过对应包管理器安装Nginx,启动并启用服务,开放防火墙端口,最后验证欢迎页显示以确认安装成功。 在Linux系统中安装Nginx服务是搭建Web服务器的第一步。Nginx以高性能、低资源消耗和良好的并发处理能力著称,广泛用于静态内容服务、反向代理和负载均衡。以下是在主流L…

    2025年12月6日 运维
    000
  • 当贝X5S怎样看3D

    当贝X5S观看3D影片无立体效果时,需开启3D模式并匹配格式:1. 播放3D影片时按遥控器侧边键,进入快捷设置选择3D模式;2. 根据片源类型选左右或上下3D格式;3. 可通过首页下拉进入电影专区选择3D内容播放;4. 确认片源为Side by Side或Top and Bottom格式,并使用兼容…

    2025年12月6日 软件教程
    100
  • Linux journalctl与systemctl status结合分析

    先看 systemctl status 确认服务状态,再用 journalctl 查看详细日志。例如 nginx 启动失败时,systemctl status 显示 Active: failed,journalctl -u nginx 发现端口 80 被占用,结合两者可快速定位问题根源。 在 Lin…

    2025年12月6日 运维
    100
  • 华为新机发布计划曝光:Pura 90系列或明年4月登场

    近日,有数码博主透露了华为2025年至2026年的新品规划,其中pura 90系列预计在2026年4月发布,有望成为华为新一代影像旗舰。根据路线图,华为将在2025年底至2026年陆续推出mate 80系列、折叠屏新机mate x7系列以及nova 15系列,而pura 90系列则将成为2026年上…

    2025年12月6日 行业动态
    100
  • TikTok视频无法下载怎么办 TikTok视频下载异常修复方法

    先检查链接格式、网络设置及工具版本。复制以https://www.tiktok.com/@或vm.tiktok.com开头的链接,删除?后参数,尝试短链接;确保网络畅通,可切换地区节点或关闭防火墙;更新工具至最新版,优先选用yt-dlp等持续维护的工具。 遇到TikTok视频下载不了的情况,别急着换…

    2025年12月6日 软件教程
    100
  • Linux如何防止缓冲区溢出_Linux防止缓冲区溢出的安全措施

    缓冲区溢出可通过栈保护、ASLR、NX bit、安全编译选项和良好编码实践来防范。1. 使用-fstack-protector-strong插入canary检测栈破坏;2. 启用ASLR(kernel.randomize_va_space=2)随机化内存布局;3. 利用NX bit标记不可执行内存页…

    2025年12月6日 运维
    000
  • 2025年双十一买手机选直板机还是选折叠屏?建议看完这篇再做决定

    随着2025年双十一购物节的临近,许多消费者在选购智能手机时都会面临一个共同的问题:是选择传统的直板手机,还是尝试更具科技感的折叠屏设备?其实,这个问题的答案早已在智能手机行业的演进中悄然浮现——如今的手机市场已不再局限于“拼参数、堆配置”的初级竞争,而是迈入了以形态革新驱动用户体验升级的新时代。而…

    2025年12月6日 行业动态
    000
  • Linux如何优化系统性能_Linux系统性能优化的实用方法

    优化Linux性能需先监控资源使用,通过top、vmstat等命令分析负载,再调整内核参数如TCP优化与内存交换,结合关闭无用服务、选用合适文件系统与I/O调度器,持续按需调优以提升系统效率。 Linux系统性能优化的核心在于合理配置资源、监控系统状态并及时调整瓶颈环节。通过一系列实用手段,可以显著…

    2025年12月6日 运维
    000
  • Pboot插件数据库连接的配置教程_Pboot插件数据库备份的自动化脚本

    首先配置PbootCMS数据库连接参数,确保插件正常访问;接着创建auto_backup.php脚本实现备份功能;然后通过Windows任务计划程序或Linux Cron定时执行该脚本,完成自动化备份流程。 如果您正在开发或维护一个基于PbootCMS的网站,并希望实现插件对数据库的连接配置以及自动…

    2025年12月6日 软件教程
    000
  • 今日头条官方主页入口 今日头条平台直达网址官方链接

    今日头条官方主页入口是www.toutiao.com,该平台通过个性化信息流推送图文、短视频等内容,具备分类导航、便捷搜索及跨设备同步功能。 今日头条官方主页入口在哪里?这是不少网友都关注的,接下来由PHP小编为大家带来今日头条平台直达网址官方链接,感兴趣的网友一起随小编来瞧瞧吧! www.tout…

    2025年12月6日 软件教程
    000

发表回复

登录后才能评论
关注微信