Skip to content
1.CompletableFuture

可以使用CompletableFuture.allOf方法。这个方法接收一个CompletableFuture数组,当数组中的所有CompletableFuture都正常完成后,它返回一个新的CompletableFuture

java
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

public class AsyncRequestsExample {

	public static void main(String[] args) {
		// 创建一个线程池
		ExecutorService executor = Executors.newFixedThreadPool(3);

		// 模拟三个异步请求
		CompletableFuture<String> future1 = CompletableFuture.supplyAsync(() -> {
			// 模拟耗时的异步操作
			sleep(1);
			return "Result of Future 1";
		}, executor);

		CompletableFuture<String> future2 = CompletableFuture.supplyAsync(() -> {
			// 模拟耗时的异步操作
			sleep(2);
			return "Result of Future 2";
		}, executor);

		CompletableFuture<String> future3 = CompletableFuture.supplyAsync(() -> {
			// 模拟耗时的异步操作
			sleep(3);
			return "Result of Future 3";
		}, executor);

		// 使用 allOf 等待所有异步操作结果
		CompletableFuture<Void> combinedFuture = CompletableFuture.allOf(future1, future2, future3);

		// 当所有的异步操作完成后,执行一些操作
		combinedFuture.thenRun(() -> {
			try {
				System.out.println("All futures completed.");
				System.out.println(future3.get()); // 获取第三个异步操作的结果
				System.out.println(future1.get()); // 获取第一个异步操作的结果
				System.out.println(future2.get()); // 获取第二个异步操作的结果

			} catch (InterruptedException | ExecutionException e) {
				e.printStackTrace();
			}
		});

		System.out.println("do something.");

		// 关闭线程池
		executor.shutdown();
	}

	private static void sleep(int seconds) {
		try {
			TimeUnit.SECONDS.sleep(seconds);
		} catch (InterruptedException e) {
			Thread.currentThread().interrupt(); // 设置中断标志
			throw new IllegalStateException(e);
		}
	}
}
2.CountDownLatch

CountDownLatch是一个同步辅助类,在完成一组正在其他线程中执行的操作之前,它允许一个或多个线程等待。

java
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

public class CountDownLatchExample {
	public static void main(String[] args) throws InterruptedException {
		int numberOfTasks = 3;
		CountDownLatch latch = new CountDownLatch(numberOfTasks);
		ExecutorService executor = Executors.newFixedThreadPool(numberOfTasks);

		for (int i = 1; i <= numberOfTasks; i++) {
			int finalI = i;
			executor.submit(() -> {
				try {
					// 模拟耗时操作
					TimeUnit.SECONDS.sleep(finalI);
					System.out.println("Result of Task " + finalI);
				} catch (InterruptedException e) {
					Thread.currentThread().interrupt();
				} finally {
					latch.countDown();
				}
			});
		}

		System.out.println("do something.");

		// 等待所有任务完成
		latch.await();

		System.out.println("All tasks completed.");

		executor.shutdown();
	}
}
3.Rxjava

RxJava 是一个实现了响应式编程的库,它允许你使用可观察序列来编写异步和基于事件的程序。在 RxJava 中,你可以使用 ObservableSingleCompletableFlowable 类型来表示异步数据流,并通过操作符来处理这些数据流。

java
import io.reactivex.rxjava3.core.Observable;
import io.reactivex.rxjava3.schedulers.Schedulers;

import java.util.concurrent.TimeUnit;

public class RxJavaAsyncRequestsExample {

	public static void main(String[] args) {
		// 创建三个异步请求的Observable
		Observable<String> observable1 = Observable.fromCallable(() -> {
			// 模拟耗时的异步操作
			sleep(1);
			return "Result of Observable 1";
		}).subscribeOn(Schedulers.io()); // 指定在IO调度器上执行

		Observable<String> observable2 = Observable.fromCallable(() -> {
			// 模拟耗时的异步操作
			sleep(2);
			return "Result of Observable 2";
		}).subscribeOn(Schedulers.io()); // 指定在IO调度器上执行

		Observable<String> observable3 = Observable.fromCallable(() -> {
			// 模拟耗时的异步操作
			sleep(3);
			return "Result of Observable 3";
		}).subscribeOn(Schedulers.io()); // 指定在IO调度器上执行

		// 使用zip操作符等待所有Observable完成,并合并结果
		Observable.zip(observable1, observable2, observable3, (result1, result2, result3) -> {
				// 当所有的Observable都发出了数据项时,这个函数会被调用,参数就是每个Observable发出的数据项
				return result1 + ", " + result2 + ", " + result3;
			})
			.subscribe(
				result -> System.out.println("All observables completed with results: " + result),
				Throwable::printStackTrace, // 错误处理
				() -> System.out.println("This will be printed upon completion of all Observables") // 完成处理
			);

		System.out.println("do something.");

		// 等待足够长的时间以确保异步操作完成
		sleep(5);
	}

	private static void sleep(int seconds) {
		try {
			TimeUnit.SECONDS.sleep(seconds);
		} catch (InterruptedException e) {
			Thread.currentThread().interrupt(); // 设置中断标志
			throw new IllegalStateException(e);
		}
	}
}
4.Reactor

Reactor是另一个响应式编程库,与RxJava类似,它也是基于Reactive Streams规范。在Spring WebFlux中,Reactor通过MonoFlux类型提供了对响应式流的支持。

java
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import java.util.concurrent.TimeUnit;

public class WebFluxAsyncRequestsExample {

	public static void main(String[] args) {
		Mono<String> task1 = Mono.fromCallable(() -> {
			// 模拟耗时操作
			TimeUnit.SECONDS.sleep(1);
			return "Result of Task 1";
		});

		Mono<String> task2 = Mono.fromCallable(() -> {
			// 模拟耗时操作
			TimeUnit.SECONDS.sleep(2);
			return "Result of Task 2";
		});

		Mono<String> task3 = Mono.fromCallable(() -> {
			// 模拟耗时操作
			TimeUnit.SECONDS.sleep(3);
			return "Result of Task 3";
		});

		System.out.println("do something.");

		Flux.merge(task1, task2, task3)
			.doOnNext(System.out::println)
			.blockLast(); // 在实际应用中避免使用block操作
	}
}