Reactor详解之:异常处理
lipiwang 2024-11-26 06:06 8 浏览 0 评论
简介
不管是在响应式编程还是普通的程序设计中,异常处理都是一个非常重要的方面。今天将会给大家介绍Reactor中异常的处理流程。
Reactor的异常一般处理方法
先举一个例子,我们创建一个Flux,在这个Flux中,我们产生一个异常,看看是什么情况:
Flux flux2= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i));
flux2.subscribe(System.out::println);
我们会得到一个异常ErrorCallbackNotImplemented:
100 / 1 = 100
100 / 2 = 50
reactor.core.Exceptions$ErrorCallbackNotImplemented: java.lang.ArithmeticException: / by zero
那怎么处理这个异常呢?
有两种方式,第一种方式就是我们之前文章讲过的,在subscribe的时候指定onError方法:
Flux flux2= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i));
flux2.subscribe(System.out::println,
error -> System.err.println("Error: " + error));
还是刚才的代码,但是这次我们在subscribe的时候,添加了onError处理器,看下运行结果:
Divided by zero :(
100 / 1 = 100
100 / 2 = 50
Error: java.lang.ArithmeticException: / by zero
可以看到异常已经被我们捕获了,并且进行了合适的处理。
除了在subscribe中进行处理,我们还可以在publish的时候,就指定异常的处理模式,这就是我们要介绍的第二种方法:
Flux flux= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i))
.onErrorReturn("Divided by zero :(");
flux.subscribe(System.out::println);
上面的例子中,在创建Flux的时候,手动指定了其onErrorReturn方法,我们看下输出结果:
100 / 1 = 100
100 / 2 = 50
Divided by zero :(
注意,对于Flux或者Mono来说,所有的异常都是一个终止的操作,即使你使用了异常处理,原生成序列也不会继续。
但是如果你对异常进行了处理,那么它会将oneError信号转换成为新的序列的开始,并将替换掉之前上游产生的序列。
各种异常处理方式详解
在一般的程序中,我们的异常应该怎么处理呢?大家很容易想到的是try catch。而Reactor中subscribe的onError方法,就是try catch的一个具体应用:
Flux flux2= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i));
flux2.subscribe(System.out::println,
error -> System.err.println("Error: " + error));
还是上的例子,我们在onError方法中,对异常进行了处理。
如果转换成为常规代码,应该是下面的样子:
public void normalErrorHandle(){
try{
Arrays.asList(1,2,0).stream().map(i -> "100 / " + i + " = " + (100 / i)).forEach(System.out::println);
}catch (Exception e){
System.err.println("Error: " + e);
}
}
除了这种最基本的异常处理方法之外,Reactor还提供了很多种不同的异常处理方法,下面我们来一一介绍一下。
Static Fallback Value
Static Fallback Value的意思是,在遇到异常的时候会fallback到一个静态的默认值。比如我们之前讲到的onErrorReturn。
Flux flux= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i))
.onErrorReturn("Divided by zero :(");
当然onErrorReturn还支持一个Predicate参数,用来判断要falback的异常是否满足条件。
public final Flux<T> onErrorReturn(Predicate<? super Throwable> predicate, T fallbackValue)
Fallback Method
除了fallback Value之外,还支持Fallback Method。也就是说如果你想在捕获异常之后调用其他的方法,就可以使用Fallback Method。
这里Fallback Method是用onErrorResume来表示的。
public void useFallbackMethod(){
Flux flux= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i))
.onErrorResume(e -> System.out::println);
flux.subscribe(System.out::println);
}
Dynamic Fallback Value
所谓的动态Fallback Value就是根据你抛出的异常进行判断,通过定位不同的Error从而fallback到不同的值:
public void useDynamicFallback(){
Flux flux= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i))
.onErrorResume(error -> Mono.just(
MyWrapper.fromError(error)));
}
public static class MyWrapper{
public static String fromError(Throwable error){
return "That is a new Error";
}
}
Catch and Rethrow
同样的,我们可以在捕获异常之后进行rethrow:
Flux flux= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i))
.onErrorResume(error -> Flux.error(
new RuntimeException("oops, ArithmeticException!", error)));
Flux flux2= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i))
.onErrorMap(error -> new RuntimeException("oops, ArithmeticException!", error));
有两种方式,第一种就是在onErrorResume中使用Flux.error构建一个新的Flux,另外一种就是直接在onErrorMap中进行处理。
Log or React on the Side
有时候你只是想记录一下异常信息,并不想破坏原来的React结构,那么可以试着使用doOnError。
public void useDoOnError(){
Flux flux= Flux.just(1, 2, 0)
.map(i -> "100 / " + i + " = " + (100 / i))
.doOnError(error -> System.out.println("we got the error: "+ error));
}
Finally Block
如果我们在代码中使用了某些资源,一般情况下我们需要在finally中对其进行关闭,或者使用JDK7中引入的 try-with-resource 。
举个例子,下面的是使用finally的方式:
Stats stats = new Stats();
stats.startTimer();
try {
doSomethingDangerous();
}
finally {
stats.stopTimerAndRecordTiming();
}
下面是使用try-with-resource的方式:
try (SomeAutoCloseable disposableInstance = new SomeAutoCloseable()) {
return disposableInstance.toString();
}
那么在Reactor中,我们也有两种方式和其对应。
第一种就是doFinally方法:
Stats stats = new Stats();
LongAdder statsCancel = new LongAdder();
Flux<String> flux =
Flux.just("foo", "bar")
.doOnSubscribe(s -> stats.startTimer())
.doFinally(type -> {
stats.stopTimerAndRecordTiming();
if (type == SignalType.CANCEL)
statsCancel.increment();
})
.take(1);
上面的例子中,doFinally实际上做的就是finally block做的事情。
第二种是使用using,我们先看一个using的定义:
public static <T, D> Flux<T> using(Callable<? extends D> resourceSupplier, Function<? super D, ? extends
Publisher<? extends T>> sourceSupplier, Consumer<? super D> resourceCleanup)
可以看到using支持三个参数,resourceSupplier是一个生成器,用来在subscribe的时候生成要发送的resource对象。
sourceSupplier是一个生成Publisher的工厂,接收resourceSupplier传过来的resource,然后生成Publisher对象。
resourceCleanup用来对resource进行收尾操作。
那么我们怎么用呢?
举个例子:
public void useUsing(){
AtomicBoolean isDisposed = new AtomicBoolean();
Disposable disposableInstance = new Disposable() {
@Override
public void dispose() {
isDisposed.set(true);
}
@Override
public String toString() {
return "DISPOSABLE";
}
};
Flux<String> flux =
Flux.using(
() -> disposableInstance,
disposable -> Flux.just(disposable.toString()),
Disposable::dispose);
}
上面的例子中,我们创建了一个Disposable对象,作为resource,然后对这个resource进行加工,返回一个Flux对象,最后通过调用Disposable::dispose方法,对resource进行销毁。
Retrying
有时候我们遇到了异常,可能需要重试几次,Reactor为我们提供了retry方法,先看一个例子:
public void testRetry(){
Flux.interval(Duration.ofMillis(250))
.map(input -> {
if (input < 3){
return "tick " + input;
}
throw new RuntimeException("boom");
})
.retry(1)
.elapsed()
.subscribe(System.out::println, System.err::println);
try {
Thread.sleep(2100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
看下输出结果:
[264,tick 0]
[255,tick 1]
[241,tick 2]
[506,tick 0]
[252,tick 1]
[253,tick 2]
java.lang.RuntimeException: boom
retry的作用就是当遇到异常的时候,重启一个新的序列。
elapsed是用来展示产生的value时间之间的duration。
从结果我们可以看到,retry之前是不会产生异常信息的。
本文的例子learn-reactive
本文作者:flydean程序那些事
本文链接:http://www.flydean.com/reactor-handle-errors/
本文来源:flydean的博客
欢迎关注我的公众号:「程序那些事」最通俗的解读,最深刻的干货,最简洁的教程,众多你不知道的小技巧等你来发现!
相关推荐
- linux实例之设置时区的方式有哪些
-
linux系统下的时间管理是一个复杂但精细的功能,而时区又是时间管理非常重要的一个辅助功能。时区解决了本地时间和UTC时间的差异,从而确保了linux系统下时间戳和时间的准确性和一致性。比如文件的时间...
- Linux set命令用法(linux cp命令的用法)
-
Linux中的set命令用于设置或显示系统环境变量。1.设置环境变量:-setVAR=value:设置环境变量VAR的值为value。-exportVAR:将已设置的环境变量VAR导出,使其...
- python环境怎么搭建?小白看完就会!简简单单
-
很多小伙伴安装了python不会搭建环境,看完这个你就会了Python可应用于多平台包括Linux和MacOSX。你可以通过终端窗口输入"python"命令来查看本地是否...
- Linux环境下如何设置多个交叉编译工具链?
-
常见的Linux操作系统都可以通过包管理器安装交叉编译工具链,比如Ubuntu环境下使用如下命令安装gcc交叉编译器:sudoapt-getinstallgcc-arm-linux-gnueab...
- JMeter环境变量配置技巧与注意事项
-
通过给JMeter配置环境变量,可以快捷的打开JMeter:打开终端。执行jmeter。配置环境变量的方法如下。Mac和Linux系统在~/.bashrc中加如下内容:export...
- C/C++|头文件、源文件分开写的源起及作用
-
1C/C++编译模式通常,在一个C++程序中,只包含两类文件——.cpp文件和.h文件。其中,.cpp文件被称作C++源文件,里面放的都是C++的源代码;而.h文件则被称...
- linux中内部变量,环境变量,用户变量的区别
-
unixshell的变量分类在Shell中有三种变量:内部变量,环境变量,用户变量。内部变量:系统提供,不用定义,不能修改环境变量:系统提供,不用定义,可以修改,可以利用export将用户变量转为环...
- 在Linux中输入一行命令后究竟发生了什么?
-
Linux,这个开源的操作系统巨人,以其强大的命令行界面而闻名。无论你是初学者还是经验丰富的系统管理员,理解在Linux终端输入一条命令并按下回车后发生的事情,都是掌握Linux核心的关键。从表面上看...
- Nodejs安装、配置与快速入门(node. js安装)
-
Nodejs是现代JavaScript语言产生革命性变化的一个主要框架,它使得JavaScript从一门浏览器语言成为可以在服务器端运行、开发各种各样应用的通用语言。在不同的平台下,Nodejs的安装...
- Ollama使用指南【超全版】(olaplex使用方法图解)
-
一、Ollama快速入门Ollama是一个用于在本地运行大型语言模型的工具,下面将介绍如何在不同操作系统上安装和使用Ollama。官网:https://ollama.comGithub:http...
- linux移植(linux移植lvgl)
-
1uboot移植l移植linux之前需要先移植一个bootlader代码,主要用于启动linux内核,lLinux系统包括u-boot、内核、根文件系统(rootfs)l引导程序的主要作用将...
- Linux日常小技巧参数优化(linux参数调优)
-
Linux系统参数优化可以让系统更加稳定、高效、安全,提高系统的性能和使用体验。下面列出一些常见的Linux系统参数优化示例,包括修改默认配置、网络等多方面。1.修改默认配置1.1修改默认编辑器默...
- Linux系统编程—条件变量(linux 条件变量开销)
-
条件变量是用来等待线程而不是上锁的,条件变量通常和互斥锁一起使用。条件变量之所以要和互斥锁一起使用,主要是因为互斥锁的一个明显的特点就是它只有两种状态:锁定和非锁定,而条件变量可以通过允许线程阻塞和等...
- 面试题-Linux系统优化进阶学习(linux系统的优化)
-
一.基础必备优化:1.关闭SElinux2.FirewalldCenetOS7Iptables(C6)安全组(阿里云)3.网络管理服务||NetworkManager|network...
- 嵌入式Linux开发教程:Linux Shell
-
本章重点介绍Linux的常用操作和命令。在介绍命令之前,先对Linux的Shell进行了简单介绍,然后按照大多数用户的使用习惯,对各种操作和相关命令进行了分类介绍。对相关命令的介绍都力求通俗易懂,都给...
你 发表评论:
欢迎- 一周热门
- 最近发表
- 标签列表
-
- maven镜像 (69)
- undefined reference to (60)
- zip格式 (63)
- oracle over (62)
- date_format函数用法 (67)
- 在线代理服务器 (60)
- shell 字符串比较 (74)
- x509证书 (61)
- localhost (65)
- java.awt.headless (66)
- syn_sent (64)
- settings.xml (59)
- 弹出窗口 (56)
- applicationcontextaware (72)
- my.cnf (73)
- httpsession (62)
- pkcs7 (62)
- session cookie (63)
- java 生成uuid (58)
- could not initialize class (58)
- beanpropertyrowmapper (58)
- word空格下划线不显示 (73)
- jar文件 (60)
- jsp内置对象 (58)
- makefile编写规则 (58)