商城首页欢迎来到中国正版软件门户

您的位置: 首页 > 文章列表 > 编程开发 > 响应式编程(反应式编程)的来龙去脉(同步编程、多线程编程、异步编程再到响应式编程)

响应式编程(反应式编程)的来龙去脉(同步编程、多线程编程、异步编程再到响应式编程)

  发布于2026-07-18 阅读(0)

扫一扫,手机访问

响应式编程的来龙去脉(同步编程、多线程编程、异步编程再到响应式编程)

如果要讲清楚响应式编程的来龙去脉,最好从最基础的同步编程开始,一步步看到多线程、异步编程,最后再进入响应式编程的世界。这个过程本身就是编程思维不断进化的缩影。

简介

这篇文章会从头梳理这一脉络。从同步编程讲起,到多线程编程的出现,再到为了解决线程阻塞、资源浪费而诞生的异步编程,以及为了让异步编程更好写、更好维护而出现的响应式编程。

1. 示例

我们用同一个实际编程的例子来感受这几种编程方式的区别。

假设要执行一个任务A:先调用功能M算出一个结果,同时调用功能N算出另一个结果,然后把两个结果做某种计算,得到最终结果。M和N都需要一定的时间(比如有阻塞的IO操作)。

任务A示意图

接下来我们会用不同的编程模式来实现同样的功能,重点感受每种模式的特点。

2. 同步编程

最简单也最直接的思路:单线程、顺序执行。

代码写出来就是一步一步的:先通过功能M算出m,再通过功能N算出n,然后计算m+n。M和N各自需要1000ms的等待时间。

public static int functionM(){
    try {
        // calculate m...
        Thread.sleep(1000);
    } catch (InterruptedException e) {
        e.printStackTrace();
    }
    return 1;
}

public static int functionN(){
    try {
        // calculate n...
        Thread.sleep(1000);
    } catch (InterruptedException e) {
        e.printStackTrace();
    }
    return 2;
}

主线程执行如下:

public static void main(String[] args) throws InterruptedException {
    long start = System.currentTimeMillis();
    int m = functionM();
    int n = functionN();
    //Complex calculations
    Thread.sleep(1000);
    System.out.println("计算结果:"+(m+n));
    System.out.println("耗时:"+(System.currentTimeMillis()-start));
}

结果是:耗时3036ms,接近3秒。这就是单线程顺序执行的结果,M和N依次等待,浪费了大量时间。

3. 多线程编程

仔细观察任务A:M和N的执行其实是相互独立的,完全可以同时进行。于是我们想到了多线程并行。

任务A可以拆解成B、C、D三个子任务:B执行M,C执行N,D等待B和C的结果再进行计算。B和C之间没有依赖关系。

任务分解示意图

实现方式就是给B和C各自分配一个子线程,主线程通过Future等待结果,然后执行D。

public class MultithreadingDemo {
    public static void main(String[] args) throws InterruptedException, ExecutionException {
        long start = System.currentTimeMillis();
        ExecutorService threadpool = Executors.newCachedThreadPool();
        Future FutureTaskB = threadpool.submit(MultithreadingDemo::TaskB);
        Future FutureTaskC = threadpool.submit(MultithreadingDemo::TaskC);
        int m = FutureTaskB.get();
        int n = FutureTaskC.get();
        int result = TaskD(m, n);
        System.out.println("计算结果:" + result);
        System.out.println("耗时:" + (System.currentTimeMillis() - start));
        threadpool.shutdown();
    }
    public static int TaskB() { return functionM(); }
    public static int TaskC() { return functionN(); }
    public static int TaskD(int m, int n) throws InterruptedException {
        Thread.sleep(1000);
        return m + n;
    }
}

运行结果:耗时2078ms,比同步方式快了将近1秒。这就是多线程并行的效果。

4. 异步编程

多线程虽然解决了并行,但仔细看上面的代码:主线程在调用FutureTaskB.get()时是阻塞的——它一直在等子线程的结果,这段时间内什么都干不了。如果业务流量很高,线程是一种宝贵资源,大量线程这样空等着,很容易耗尽系统内存。

于是异步编程的思路出现了:主线程发出任务指令后直接返回,不等结果。等B和C执行完了,再安排其他线程(不一定是原来的线程)去执行后续操作。我们的关注点从“谁执行、什么时候执行”转移到“任务之间的依赖关系”上。

Ja va 8的CompletableFuture是一个很好的工具。我们用它来对这个过程建模:主线程只负责编排任务,然后就可以去做别的事情,不必等着。

任务依赖关系图:

任务依赖关系图

代码示例如下:

public static void main(String[] args) throws InterruptedException, ExecutionException {
    triggerA();
    //处理其他任务,完全不用管理任务A。
    Thread.sleep(3000);
}
public static void triggerA() {
    long start = System.currentTimeMillis();
    System.out.println("编排任务A线程ID为:"+Thread.currentThread().getId());
    CompletableFuture taskBFuture = CompletableFuture.supplyAsync(()->{
        System.out.println("执行任务B的当前线程ID:"+Thread.currentThread().getId());
        return taskB();
    });
    CompletableFuture taskCFuture = CompletableFuture.supplyAsync(()->{
        System.out.println("执行任务C的当前线程ID:"+Thread.currentThread().getId());
        return taskC();
    });
    //当任务C、D完成之后,回调执行任务D
    taskBFuture.thenCombineAsync(taskCFuture,(m,n)->{
        try {
            System.out.println("执行任务C的当前线程ID:"+Thread.currentThread().getId());
            Thread.sleep(1000);
            System.out.println("耗时:"+(System.currentTimeMillis()-start));
            System.out.println("taskA计算结果:"+(m+n));
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        return null;
    });
}
public static int taskB() { return functionM(); }
public static int taskC() { return functionN(); }

运行结果:耗时2042ms,线程1(主线程)编排完任务后就去忙别的了,B和C分别在11、12线程执行,然后由11线程执行D。主线程没有被阻塞。

这里还引出一个更直观的例子:C#的异步编程。三个任务,每个任务分为前、中、后三部分,中间部分是等待操作(如IO)。传统思路下线程会阻塞在等待操作上;异步编程则让线程直接返回,去执行其他任务,等等待完成了再派一个线程去执行后续操作。

C#异步执行结果示例

Ja va中的非阻塞IO也是类似思想:线程不等待IO操作,而是通过事件机制在IO完成后唤醒执行后续操作。这种模式在高并发场景下效果显著,极大解放了线程资源。

5. 响应式编程

5.1 响应式的由来

从上面的演变可以看到,编码的视角在逐渐从“控制任务执行”转向“描述任务本身和任务之间的依赖关系”。我们越来越关注“事情本身”。

但之前举的例子都很简单。如果业务逻辑复杂,任务的编排也会变得复杂、难以实现。这就是响应式编程要解决的问题。

回过头看我们拆解的任务B、C、D:B和C完成之后回调D执行。与其叫“回调”,不如叫“触发”——B、C的执行结束事件触发了D的执行。换句话说,我们可以看作D订阅了B、C的结束事件。当事件发生时,D开始响应并执行。

事件订阅发布示意图

发布者发布各种事件(包括BC执行结束的事件),订阅者订阅这些事件。BC执行结束后发布事件通知订阅者,订阅者响应执行D。这就是用订阅-发布模式来表现异步编程。

这么做的好处是:在响应式编程中,一系列事件(产生事件的地方)被看作流或管道。利用函数式编程中对流的处理,可以把异步编程中的许多操作转化为对数据流的操作,尤其是在对流进行转换时,效率和优雅度都会提升一个台阶。

Project Reactor举个例子。有一个管道每隔100ms发布一个数字(相当于事件),一个订阅者负责打印这些数字。传统异步操作肯定会阻塞等待,浪费线程资源。改用响应式编程:

Flux.range(1,10)
    .delayElements(Duration.of(100,ChronoUnit.MILLIS))
    .subscribe(System.out::println);

这段代码的意思是:声明一个可以生成1到10数字的管道(Flux.range(1,10)),每隔100ms生成一个数字(delayElements),打印服务订阅这个管道(subscribe),有元素生成时就通知打印服务进行打印(System.out::println)。

为了更直观地对比不同风格的异步编程,再看两个Project Reactor官方的例子。

5.2 回调 vs 响应式

需求:展示五个商品——优先展示用户最喜爱的五个商品,如果没有最喜爱商品,就用推荐的商品展示。

使用回调实现(可用Gua va实现):

userService.getFa vorites(userId, new Callback>() {
    public void onSuccess(List list) {
        if (list.isEmpty()) {
            suggestionService.getSuggestions(new Callback>() {
                public void onSuccess(List list) {
                    UiUtils.submitOnUiThread(() -> {
                        list.stream().limit(5).forEach(uiList::show);
                    });
                }
                public void onError(Throwable error) {
                    UiUtils.errorPopup(error);
                }
            });
        } else {
            list.stream()
                .limit(5)
                .forEach(fa vId -> fa voriteService.getDetails(fa vId, 
                    new Callback() {
                        public void onSuccess(Fa vorite details) {
                            UiUtils.submitOnUiThread(() -> uiList.show(details));
                        }
                        public void onError(Throwable error) {
                            UiUtils.errorPopup(error);
                        }
                    }));
        }
    }
    public void onError(Throwable error) {
        UiUtils.errorPopup(error);
    }
});

一层层的回调嵌套,对于开发和维护来说简直是地狱——这被称为“回调地狱”(Callback Hell)。命令式编程可读性极差,层层缩进,难以维护,容易出错。

使用响应式编程实现:

userService.getFa vorites(userId)
    .flatMap(fa voriteService::getDetails)
    .switchIfEmpty(suggestionService.getSuggestions())
    .take(5)
    .publishOn(UiUtils.uiThreadScheduler())
    .subscribe(uiList::show, UiUtils::errorPopup);

流程清晰,转换操作简单,底层异步实现被完全屏蔽。只需要把要做什么表达出来,描述好任务之间的关系,功能就实现了。

5.3 CompletableFuture 响应式

再看一个例子:用一个ID队列,分别获取每个ID的name和statistic,然后组合起来。这是典型的异步编排场景。

CompletableFuture实现:

CompletableFuture> ids = ifhIds();
CompletableFuture> result = ids.thenComposeAsync(l -> {
    Stream> zip = l.stream().map(i -> {
        CompletableFuture nameTask = ifhName(i);
        CompletableFuture statTask = ifhStat(i);
        return nameTask.thenCombineAsync(statTask, 
            (name, stat) -> "Name " + name + " has stats " + stat);
    });
    List> combinationList = zip.collect(Collectors.toList());
    CompletableFuture[] combinationArray = 
        combinationList.toArray(new CompletableFuture[combinationList.size()]);
    CompletableFuture allDone = CompletableFuture.allOf(combinationArray);
    return allDone.thenApply(v -> 
        combinationList.stream().map(CompletableFuture::join)
            .collect(Collectors.toList()));
});
List results = result.join();
assertThat(results).contains(...);

虽然比原生回调简化了一些,但读起来依然比较吃力。任务依赖关系图可以帮助理解:

CompletableFuture任务依赖图

响应式编程实现:

Flux ids = ifhrIds();
Flux combinations = ids.flatMap(id -> {
    Mono nameTask = ifhrName(id);
    Mono statTask = ifhrStat(id);
    return nameTask.zipWith(statTask, 
        (name, stat) -> "Name " + name + " has stats " + stat);
});
Mono> result = combinations.collectList();
List results = result.block();
assertThat(results).containsExactly(...);

FLUX是可以生成0到N个元素的管道,Mono是可以生成0到1个元素的管道。代码的声明性更强,对任务编排的API也更友好,开发和理解的难度明显降低。

从回调 → CompletableFuture → 响应式编程,代码的可维护性、声明性和面对复杂业务的灵活性都在逐步提高。

最后用响应式编程改造一下一开始那个异步编程的例子:

public class ReactiveProgramDemo {
    public static void main(String[] args) throws InterruptedException {
        long start = System.currentTimeMillis();
        Mono monoB = Mono.fromCallable(ReactiveProgramDemo::taskB)
            .publishOn(Schedulers.boundedElastic());
        Mono monoC = Mono.fromCallable(ReactiveProgramDemo::taskC)
            .publishOn(Schedulers.boundedElastic());
        Mono.zip(monoB, monoC)
            .subscribeOn(Schedulers.boundedElastic())
            .subscribe(nums -> {
                System.out.println("计算结果为:" + taskD(nums.getT1(), nums.getT2()));
                System.out.println("tasA执行用时:" + (System.currentTimeMillis() - start));
            });
        //主线程做其他的事情
        Thread.sleep(5000);
    }
    // ... taskB, taskC, taskD 同上
}

结果:耗时2249ms,一切异步执行,主线程不被阻塞。

5.4 响应式编程的特性

随着响应式编程的普及,reactive-streams规范也逐渐成形,现在的响应式框架都是按照这个规范进行接口实现的。概括一下它的核心特性:

  • 异步编程的实现:响应式编程本质上就是一种异步编程,与基于回调、CompletableFuture实现的异步编程属于同一范畴。
  • 更好的表达性,不易出错:用声明式的方式描述业务逻辑,而不是控制流程,自然减少了出错概率。
  • 更灵活的操作方式:得益于数据流的函数式编程,提供了丰富的转换和操作符。
  • 更好的异常处理:其他方式的异步编程很难做到优雅的异常处理,响应式编程可以轻松实现。
Flux.range(1,100)
    .map(x->1/(x%10))
    .doOnNext(System.out::println)
    .doOnError(System.err::println)
    .subscribe()
  • 支持背压:可以控制上游管道的发送速度,避免下游消费不过来。

这些特性仅仅是入门,后续的文章会针对响应式编程进行更详细的讲解。

6. 总结

从最初的问题出发,我们经历了几个阶段:

先是同步编程——顺序执行,简单直接,但不能避免等待。

然后发现有些地方可以并行执行,于是出现了多线程编程,但依然没有摆脱“同步等待”的低效率问题。

随着非阻塞IO和异步接口(如Ja va的NIO和CompletableFuture)的普及,异步编程流行起来。它大大提高了线程的利用率,让服务在更少的时间内完成更多计算,性能大幅提升。

但新的问题随之而来:异步编程在面对复杂业务时的设计难度大,任务之间的关系复杂。用传统方式实现异步编程,容易出现冗长的嵌套代码,维护成本高,而且对复杂业务的任务编排十分困难。

这时候,封装更完善的响应式编程出现了。它从事件的角度出发,将异步编程转化为对事件流的操作,利用函数式编程对数据流进行处理,让异步编程变得更加优雅。同时提供了灵活的转换操作、异常处理和背压机制。对于异步编程来说,这确实是令人兴奋的进步。

7. 思考

语言或框架的接口提供了异步的能力,它们到底是怎么实现的?

8. 参考

一文带你彻底了解ja va异步编程
ja va-8-completablefuture-tutorial
Reactor 3 Reference Guide
Difference Between Asynchronous Programming and Multithreading in C#
What is the difference between asynchronous programming and multithreading?

本文转载于:https://blog.csdn.net/tianpengfei123/article/details/127752227 如有侵犯,请联系zhengruancom@outlook.com删除。
免责声明:正软商城发布此文仅为传递信息,不代表正软商城认同其观点或证实其描述。

热门关注