Reactor Mono应用

2024-06-22 08:20
文章标签 应用 reactor mono

本文主要是介绍Reactor Mono应用,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

使用案例

创建Mono

使用静态工厂方法创建Mono
import reactor.core.publisher.Mono;public class MonoExample {public static void main(String[] args) {// 创建一个包含值的MonoMono<String> monoWithValue = Mono.just("Hello, Reactor!");// 创建一个空的MonoMono<String> emptyMono = Mono.empty();// 创建一个包含错误的MonoMono<String> errorMono = Mono.error(new RuntimeException("Something went wrong"));}
}
从Callable、Supplier、CompletableFuture等创建Mono
import reactor.core.publisher.Mono;import java.util.concurrent.CompletableFuture;public class MonoExample {public static void main(String[] args) {// 从Callable创建MonoMono<String> callableMono = Mono.fromCallable(() -> "Hello from Callable");// 从Supplier创建MonoMono<String> supplierMono = Mono.fromSupplier(() -> "Hello from Supplier");// 从CompletableFuture创建MonoCompletableFuture<String> future = CompletableFuture.supplyAsync(() -> "Hello from CompletableFuture");Mono<String> futureMono = Mono.fromFuture(future);}
}

订阅Mono

订阅是开始数据流处理的关键步骤

import reactor.core.publisher.Mono;public class MonoExample {public static void main(String[] args) {Mono<String> mono = Mono.just("Hello, Reactor!");// 订阅并消费数据mono.subscribe(value -> System.out.println("Received: " + value), // onNexterror -> System.err.println("Error: " + error),   // onError() -> System.out.println("Completed")             // onComplete);}
}

操作Mono

转换和操作数据
import reactor.core.publisher.Mono;public class MonoExample {public static void main(String[] args) {Mono<String> mono = Mono.just("hello");// 转换数据Mono<String> transformedMono = mono.map(String::toUpperCase);// 链式操作transformedMono.flatMap(value -> Mono.just(value + " World")).subscribe(System.out::println);  // 输出: HELLO World}
}

异常处理

import reactor.core.publisher.Mono;public class MonoExample {public static void main(String[] args) {Mono<String> monoWithError = Mono.error(new RuntimeException("Original Error"));// 捕获并处理错误monoWithError.onErrorReturn("Fallback value").subscribe(System.out::println);  // 输出: Fallback value// 使用 onErrorResume 提供备用的 MonomonoWithError.onErrorResume(error -> {System.err.println("Error: " + error);return Mono.just("Recovered value");}).subscribe(System.out::println);  // 输出: Recovered value}
}

组合Mono

import reactor.core.publisher.Mono;public class MonoExample {public static void main(String[] args) {Mono<String> mono1 = Mono.just("Hello");Mono<String> mono2 = Mono.just("World");// 组合两个MonoMono<String> combinedMono = Mono.zip(mono1, mono2, (s1, s2) -> s1 + " " + s2);combinedMono.subscribe(System.out::println);  // 输出: Hello World}
}

调试和日志

使用日志功能可以帮助调试和监控数据流。

import reactor.core.publisher.Mono;
import reactor.util.Logger;
import reactor.util.Loggers;public class MonoExample {public static void main(String[] args) {Logger logger = Loggers.getLogger(MonoExample.class);Mono<String> mono = Mono.just("Hello, Reactor!").log();  // 默认日志mono.subscribe(System.out::println);}
}

调度器(Schedulers)

使用调度器来控制Mono的执行线程。

import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;public class MonoExample {public static void main(String[] args) {Mono<String> mono = Mono.fromCallable(() -> {// 在独立的线程池中执行Thread.sleep(1000);return "Hello from another thread";});mono.subscribeOn(Schedulers.boundedElastic())  // 指定订阅时的调度器.publishOn(Schedulers.parallel())         // 指定发布时的调度器.subscribe(System.out::println);}
}

应用场景

异步计算

Mono可以用来表示和处理异步计算的结果。例如,当你需要从一个异步操作中获取一个值时,可以使用Mono。

Mono<String> asyncResult = Mono.fromCallable(() -> {// 模拟异步计算Thread.sleep(1000);return "Result";
});asyncResult.subscribe(result -> System.out.println("Received: " + result));

调用远程服务

在微服务架构中,调用远程服务(如REST API或gRPC)时,通常会返回一个单一的结果。这是Mono的一个典型应用场景。

WebClient webClient = WebClient.create("http://example.com");Mono<String> response = webClient.get().uri("/resource").retrieve().bodyToMono(String.class);response.subscribe(body -> System.out.println("Response: " + body));

数据库查询

Mono非常适合表示数据库查询返回的单个结果。例如,查询一个用户的信息

Mono<User> userMono = reactiveUserRepository.findById(userId);userMono.subscribe(user -> System.out.println("User: " + user));

事件驱动的处理

在事件驱动架构中,某些事件处理结果可能是单一的值。例如,处理某个事件并返回一个处理结果

Mono<EventResult> eventResultMono = processEvent(event);eventResultMono.subscribe(result -> System.out.println("Event processed: " + result));

错误处理

使用Mono可以优雅地处理异步操作中的错误。例如,如果某个操作可能会失败,可以返回一个错误的Mono并在订阅时处理错误

Mono<String> result = performOperation().onErrorReturn("Fallback value");result.subscribe(value -> System.out.println("Received: " + value),error -> System.err.println("Error: " + error)
);

延迟操作

Mono可以用于表示一个延迟操作,执行某些延迟逻辑

Mono<Long> delayMono = Mono.delay(Duration.ofSeconds(3));delayMono.subscribe(time -> System.out.println("Delayed for 3 seconds"));

条件逻辑

使用Mono可以在异步流中进行条件判断和逻辑处理。例如,根据某个条件返回不同的结果

Mono<String> conditionalMono = Mono.just("data").flatMap(data -> {if (data.equals("condition")) {return Mono.just("Condition met");} else {return Mono.just("Condition not met");}});conditionalMono.subscribe(System.out::println);

转换与映射

Mono可以用于对单个值进行转换或映射。例如,将一个值转换为另一个类型的值

Mono<String> originalMono = Mono.just("original");Mono<Integer> transformedMono = originalMono.map(String::length);transformedMono.subscribe(length -> System.out.println("Length: " + length));

资源管理

Mono可以用于在异步操作中管理资源,如文件或连接的打开和关闭

Mono.using(() -> new BufferedReader(new FileReader("data.txt")), reader -> Mono.fromCallable(() -> reader.readLine()), BufferedReader::close
).subscribe(line -> System.out.println("Read line: " + line));

组合多个异步操作

Mono可以用于组合多个异步操作,构建复杂的异步数据流

Mono<String> combinedMono = Mono.zip(Mono.just("Hello"),Mono.just("World"),(s1, s2) -> s1 + " " + s2
);combinedMono.subscribe(System.out::println);  // 输出: Hello World

这篇关于Reactor Mono应用的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1083753

相关文章

Python中随机休眠技术原理与应用详解

《Python中随机休眠技术原理与应用详解》在编程中,让程序暂停执行特定时间是常见需求,当需要引入不确定性时,随机休眠就成为关键技巧,下面我们就来看看Python中随机休眠技术的具体实现与应用吧... 目录引言一、实现原理与基础方法1.1 核心函数解析1.2 基础实现模板1.3 整数版实现二、典型应用场景2

Python Dash框架在数据可视化仪表板中的应用与实践记录

《PythonDash框架在数据可视化仪表板中的应用与实践记录》Python的PlotlyDash库提供了一种简便且强大的方式来构建和展示互动式数据仪表板,本篇文章将深入探讨如何使用Dash设计一... 目录python Dash框架在数据可视化仪表板中的应用与实践1. 什么是Plotly Dash?1.1

Android Kotlin 高阶函数详解及其在协程中的应用小结

《AndroidKotlin高阶函数详解及其在协程中的应用小结》高阶函数是Kotlin中的一个重要特性,它能够将函数作为一等公民(First-ClassCitizen),使得代码更加简洁、灵活和可... 目录1. 引言2. 什么是高阶函数?3. 高阶函数的基础用法3.1 传递函数作为参数3.2 Lambda

Java中&和&&以及|和||的区别、应用场景和代码示例

《Java中&和&&以及|和||的区别、应用场景和代码示例》:本文主要介绍Java中的逻辑运算符&、&&、|和||的区别,包括它们在布尔和整数类型上的应用,文中通过代码介绍的非常详细,需要的朋友可... 目录前言1. & 和 &&代码示例2. | 和 ||代码示例3. 为什么要使用 & 和 | 而不是总是使

Python循环缓冲区的应用详解

《Python循环缓冲区的应用详解》循环缓冲区是一个线性缓冲区,逻辑上被视为一个循环的结构,本文主要为大家介绍了Python中循环缓冲区的相关应用,有兴趣的小伙伴可以了解一下... 目录什么是循环缓冲区循环缓冲区的结构python中的循环缓冲区实现运行循环缓冲区循环缓冲区的优势应用案例Python中的实现库

SpringBoot整合MybatisPlus的基本应用指南

《SpringBoot整合MybatisPlus的基本应用指南》MyBatis-Plus,简称MP,是一个MyBatis的增强工具,在MyBatis的基础上只做增强不做改变,下面小编就来和大家介绍一下... 目录一、MyBATisPlus简介二、SpringBoot整合MybatisPlus1、创建数据库和

python中time模块的常用方法及应用详解

《python中time模块的常用方法及应用详解》在Python开发中,时间处理是绕不开的刚需场景,从性能计时到定时任务,从日志记录到数据同步,时间模块始终是开发者最得力的工具之一,本文将通过真实案例... 目录一、时间基石:time.time()典型场景:程序性能分析进阶技巧:结合上下文管理器实现自动计时

Java逻辑运算符之&&、|| 与&、 |的区别及应用

《Java逻辑运算符之&&、||与&、|的区别及应用》:本文主要介绍Java逻辑运算符之&&、||与&、|的区别及应用的相关资料,分别是&&、||与&、|,并探讨了它们在不同应用场景中... 目录前言一、基本概念与运算符介绍二、短路与与非短路与:&& 与 & 的区别1. &&:短路与(AND)2. &:非短

Spring AI集成DeepSeek三步搞定Java智能应用的详细过程

《SpringAI集成DeepSeek三步搞定Java智能应用的详细过程》本文介绍了如何使用SpringAI集成DeepSeek,一个国内顶尖的多模态大模型,SpringAI提供了一套统一的接口,简... 目录DeepSeek 介绍Spring AI 是什么?Spring AI 的主要功能包括1、环境准备2

Spring AI与DeepSeek实战一之快速打造智能对话应用

《SpringAI与DeepSeek实战一之快速打造智能对话应用》本文详细介绍了如何通过SpringAI框架集成DeepSeek大模型,实现普通对话和流式对话功能,步骤包括申请API-KEY、项目搭... 目录一、概述二、申请DeepSeek的API-KEY三、项目搭建3.1. 开发环境要求3.2. mav