Appearance
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 中,你可以使用 Observable、Single、Completable 和 Flowable 类型来表示异步数据流,并通过操作符来处理这些数据流。
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通过Mono和Flux类型提供了对响应式流的支持。
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操作
}
}