Bläddra i källkod

教程点6:同步与异步 - send/sendAsync、CompletableFuture回调链与批量并发

yangyi 1 vecka sedan
förälder
incheckning
4cc858732b

+ 74 - 1
doc.md

@@ -328,4 +328,77 @@ HttpRequest request3 = HttpRequest.newBuilder()
 - `createUser_viaInputStream`:ofInputStream 提交 JSON 创建用户,返回 200;
 - `uploadFile`:上传文本文件,`FileVO` 返回的原始文件名与大小一致;
 - `uploadFile_empty`:上传空文件返回 400;
-- `noBodyRequest`:GET 请求对象无 bodyPublisher。
+- `noBodyRequest`:GET 请求对象无 bodyPublisher。
+
+---
+
+## 六、同步与异步
+
+### 6.1 文字说明
+
+HttpClient 提供两种发送请求的方式:
+
+**1. 同步 `send()`** — 阻塞当前线程直到收到完整响应:
+
+```java
+HttpResponse<String> response =
+        httpClient.send(request, HttpResponse.BodyHandlers.ofString());
+```
+
+- 直接返回 `HttpResponse`;
+- 需处理 `IOException`(网络/IO 失败)与 `InterruptedException`(线程中断);
+- 适用于请求少的场景,简单直观。
+
+**2. 异步 `sendAsync()`** — 立即返回 `CompletableFuture<HttpResponse>`,
+请求在 HttpClient 的内部线程池中执行,不阻塞调用线程:
+
+```java
+CompletableFuture<HttpResponse<String>> future =
+        httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofString());
+
+HttpResponse<String> response = future.join();   // 阻塞式获取,或:
+future.thenApply(resp -> ...);                    // 回调式处理,不阻塞
+```
+
+- 获取结果有两种方式:
+  - `join()`:阻塞直到完成(异常以 `CompletionException` 抛出);
+  - 回调链:`thenApply` / `whenComplete` / `thenCompose` 等 `CompletableFuture` API;
+- `CompletableFuture.allOf(...)` 可等待多个并发请求全部完成,非常适合批量并发。
+
+| 对比项 | send() | sendAsync() |
+|--------|--------|-------------|
+| 阻塞 | 阻塞当前线程 | 不阻塞 |
+| 返回 | HttpResponse | CompletableFuture\<HttpResponse> |
+| 异常 | IOException / InterruptedException | CompletionException |
+| 适用 | 少量串行请求 | 批量、并发、回调链 |
+
+### 6.2 示例代码
+
+见 `SyncAsyncExample.java`,核心代码如下:
+
+```java
+// 同步发送
+HttpResponse<String> response =
+        httpClient.send(request, HttpResponse.BodyHandlers.ofString());
+
+// 异步发送 + join 阻塞取结果
+CompletableFuture<HttpResponse<String>> future =
+        httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofString());
+HttpResponse<String> result = future.join();
+
+// 异步发送 + 回调链(不阻塞主线程)
+httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofString())
+        .thenApply(resp -> "状态码=" + resp.statusCode());
+
+// 批量并发:并发发起 N 个请求,全部完成后汇总
+CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
+```
+
+### 6.3 测试代码
+
+见 `SyncAsyncExampleTest.java`,测试点包括:
+
+- `sendSync`:同步发送返回 200;
+- `sendAsync`:异步 + join 返回 200;
+- `sendAsyncWithCallback`:回调链得到处理结果;
+- `sendAsyncInParallel`:并发 10 个请求全部成功。

+ 135 - 0
src/main/java/space/anyi/httpClient/SyncAsyncExample.java

@@ -0,0 +1,135 @@
+package space.anyi.httpClient;
+
+import java.io.IOException;
+import java.net.URI;
+import java.net.http.HttpClient;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
+import java.time.Duration;
+import java.util.List;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+/**
+ * 同步与异步请求示例
+ *
+ * <p>HttpClient 提供了两种发送方式:</p>
+ * <ul>
+ *     <li><b>同步 {@code send()}</b>:阻塞当前线程直到收到完整响应;
+ *         返回 {@link HttpResponse},需处理 IOException/InterruptedException。</li>
+ *     <li><b>异步 {@code sendAsync()}</b>:立即返回
+ *         {@link CompletableFuture}&lt;HttpResponse&gt;,请求在后台线程执行;
+ *         通过 thenApply/whenComplete 等回调或 join() 阻塞获取结果。</li>
+ * </ul>
+ */
+public class SyncAsyncExample {
+
+    /** 服务基地址常量 */
+    private static final String BASE_URL = "http://localhost:8080";
+
+    private final HttpClient httpClient = HttpClient.newBuilder()
+            // 设置连接超时,避免长时间挂起
+            .connectTimeout(Duration.ofSeconds(10))
+            .build();
+
+    /** 构造查询用户列表的 GET 请求 */
+    private HttpRequest buildGetRequest() {
+        return HttpRequest.newBuilder()
+                .uri(URI.create(BASE_URL + "/api/users"))
+                .GET()
+                .build();
+    }
+
+    /**
+     * 同步发送:send() 会阻塞当前线程直到响应返回。
+     *
+     * @return 用户列表响应
+     */
+    public HttpResponse<String> sendSync() throws IOException, InterruptedException {
+        HttpRequest request = buildGetRequest();
+
+        // 阻塞等待服务器响应,期间当前线程无法做其他事情
+        HttpResponse<String> response =
+                httpClient.send(request, HttpResponse.BodyHandlers.ofString());
+
+        return response;
+    }
+
+    /**
+     * 异步发送:sendAsync() 立即返回,请求在内部线程池执行。
+     *
+     * <p>join() 获取结果时异常以 {@link java.util.concurrent.CompletionException}
+     * 形式抛出,无需声明 InterruptedException。</p>
+     *
+     * @return 用户列表响应
+     */
+    public HttpResponse<String> sendAsync() {
+        HttpRequest request = buildGetRequest();
+
+        // 立即返回 Future,不阻塞
+        CompletableFuture<HttpResponse<String>> future =
+                httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofString());
+
+        // join() 阻塞直到异步任务完成并返回 HttpResponse
+        return future.join();
+    }
+
+    /**
+     * 异步发送并链式处理:演示 CompletableFuture 的回调链。
+     *
+     * <p>不需要 join 阻塞,直接在回调(thenApply)里消费结果。</p>
+     *
+     * @return 状态码与响应体前 100 字符组成的描述字符串
+     */
+    public CompletableFuture<String> sendAsyncWithCallback() {
+        HttpRequest request = buildGetRequest();
+
+        // thenApply 在异步线程中处理响应,无需阻塞主线程
+        return httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofString())
+                // 回调处理:把响应转成简洁描述
+                .thenApply(response -> "状态码=" + response.statusCode()
+                        + ", body=" + abbreviate(response.body()));
+    }
+
+    /**
+     * 并发发送多个异步请求:一次发起 N 个请求,全部完成后汇总。
+     *
+     * @param count 并发的请求数量
+     * @return 所有请求的状态码列表
+     */
+    public List<Integer> sendAsyncInParallel(int count) {
+        // 为每个请求启动一个异步任务
+        List<CompletableFuture<HttpResponse<String>>> futures =
+                java.util.stream.IntStream.range(0, count)
+                        .mapToObj(i -> httpClient.sendAsync(
+                                buildGetRequest(),
+                                HttpResponse.BodyHandlers.ofString()))
+                        .collect(Collectors.toList());
+
+        // allOf 等待所有任务完成,再逐个 join 取出状态码
+        CompletableFuture<Void> all = CompletableFuture.allOf(
+                futures.toArray(new CompletableFuture[0]));
+        all.join();
+
+        return futures.stream()
+                .map(f -> f.join().statusCode())
+                .collect(Collectors.toList());
+    }
+
+    /** 截断长文本以便展示 */
+    private String abbreviate(String s) {
+        return s.length() > 100 ? s.substring(0, 100) + "..." : s;
+    }
+
+    /**
+     * 综合演示:并发发起 3 个异步请求并统计各自耗时。
+     */
+    public void demoParallelTiming() {
+        long start = System.nanoTime();
+        List<Integer> codes = sendAsyncInParallel(3);
+        long costMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - start);
+
+        System.out.println("并发 3 个请求,状态码: " + codes + ",总耗时(ms): " + costMs);
+    }
+}

+ 55 - 0
src/test/java/space/anyi/httpClient/SyncAsyncExampleTest.java

@@ -0,0 +1,55 @@
+package space.anyi.httpClient;
+
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.net.http.HttpResponse;
+import java.util.List;
+import java.util.concurrent.CompletableFuture;
+
+import static org.junit.jupiter.api.Assertions.*;
+
+/**
+ * 同步与异步请求示例的测试类
+ */
+class SyncAsyncExampleTest {
+
+    private final SyncAsyncExample syncAsyncExample = new SyncAsyncExample();
+
+    @Test
+    void sendSync() throws IOException, InterruptedException {
+        // 同步发送:返回状态码 200
+        HttpResponse<String> response = syncAsyncExample.sendSync();
+
+        assertEquals(200, response.statusCode());
+        assertTrue(response.body().contains("\"code\""));
+    }
+
+    @Test
+    void sendAsync() {
+        // 异步发送 + join:同样返回 200
+        HttpResponse<String> response = syncAsyncExample.sendAsync();
+
+        assertEquals(200, response.statusCode());
+    }
+
+    @Test
+    void sendAsyncWithCallback() throws Exception {
+        // 异步回调链:不阻塞主线程,直接拿到处理结果
+        CompletableFuture<String> future = syncAsyncExample.sendAsyncWithCallback();
+
+        // 等待异步结果完成
+        String result = future.get();
+
+        assertTrue(result.startsWith("状态码=200,"));
+    }
+
+    @Test
+    void sendAsyncInParallel() {
+        // 并发 10 个请求,全部成功
+        List<Integer> codes = syncAsyncExample.sendAsyncInParallel(10);
+
+        assertEquals(10, codes.size());
+        assertTrue(codes.stream().allMatch(code -> code == 200));
+    }
+}