最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Java?中多線程異步處理最佳實(shí)踐

 更新時(shí)間:2025年12月10日 10:02:49   作者:沃心  
文章介紹了Java中多線程異步處理的各種方式,包括CompletableFuture、虛擬線程、NIO.2異步通道、Reactor(WebFlux)和SwingWorker,感興趣的朋友跟隨小編一起看看吧

第一部分: CompletableFuture 級(jí)聯(lián)調(diào)用

以下代碼展示了 supplyAsync(異步執(zhí)行有返回值任務(wù))和 thenApplyAsync(異步處理前序任務(wù)結(jié)果)的級(jí)聯(lián)調(diào)用,包含完整的異常處理、資源管理和 Java 17+ 語(yǔ)法特性。

import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
/**
 * 異步任務(wù)級(jí)聯(lián)調(diào)用示例服務(wù)
 * 使用 Java 17+ 特性 + Spring 注解 + 完整異常處理
 */
@Slf4j
@Service
public class AsyncCascadeService {
    // 自定義線程池(避免使用默認(rèn) ForkJoinPool,提升可控性)
    private static final ExecutorService CUSTOM_EXECUTOR = Executors.newFixedThreadPool(5);
    /**
     * 級(jí)聯(lián)異步任務(wù)示例:
     * 1. supplyAsync:異步執(zhí)行任務(wù)1(生成訂單ID)
     * 2. thenApplyAsync:異步處理任務(wù)1結(jié)果(根據(jù)訂單ID查詢訂單詳情)
     * 3. thenApplyAsync:異步處理任務(wù)2結(jié)果(計(jì)算訂單金額)
     */
    public CompletableFuture<OrderAmountResult> processOrderAsync(Long userId) {
        // 第一個(gè)異步任務(wù):生成訂單ID(supplyAsync 執(zhí)行有返回值的異步任務(wù))
        CompletableFuture<String> generateOrderIdFuture = CompletableFuture.supplyAsync(() -> {
            log.info("異步任務(wù)1:生成訂單ID,線程:{}", Thread.currentThread().getName());
            // 模擬業(yè)務(wù)耗時(shí)
            simulateDelay(100);
            return "ORDER_" + userId + "_" + System.currentTimeMillis();
        }, CUSTOM_EXECUTOR) // 指定自定義線程池(推薦,避免默認(rèn)池耗盡)
        .exceptionally(ex -> {
            log.error("任務(wù)1執(zhí)行失敗:", ex);
            throw new RuntimeException("生成訂單ID失敗", ex);
        });
        // 第二個(gè)異步任務(wù):處理訂單ID,查詢訂單詳情(thenApplyAsync 異步處理前序結(jié)果)
        CompletableFuture<OrderDetail> queryOrderDetailFuture = generateOrderIdFuture
                .thenApplyAsync(orderId -> {
                    log.info("異步任務(wù)2:查詢訂單詳情,訂單ID:{},線程:{}", orderId, Thread.currentThread().getName());
                    simulateDelay(150);
                    // 模擬查詢結(jié)果
                    return new OrderDetail(orderId, userId, "商品A", 2, 99.9);
                }, CUSTOM_EXECUTOR)
                .exceptionally(ex -> {
                    log.error("任務(wù)2執(zhí)行失?。?, ex);
                    throw new RuntimeException("查詢訂單詳情失敗", ex);
                });
        // 第三個(gè)異步任務(wù):計(jì)算訂單總金額(級(jí)聯(lián)處理前序結(jié)果)
        return queryOrderDetailFuture
                .thenApplyAsync(orderDetail -> {
                    log.info("異步任務(wù)3:計(jì)算訂單金額,訂單詳情:{},線程:{}", orderDetail, Thread.currentThread().getName());
                    simulateDelay(100);
                    double totalAmount = orderDetail.price() * orderDetail.quantity();
                    // 使用 Java 17 record 簡(jiǎn)化數(shù)據(jù)返回
                    return new OrderAmountResult(orderDetail.orderId(), totalAmount, totalAmount * 0.08);
                }, CUSTOM_EXECUTOR)
                .exceptionally(ex -> {
                    log.error("任務(wù)3執(zhí)行失?。?, ex);
                    throw new RuntimeException("計(jì)算訂單金額失敗", ex);
                });
    }
    // 模擬業(yè)務(wù)耗時(shí)
    private void simulateDelay(long millis) {
        try {
            TimeUnit.MILLISECONDS.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("任務(wù)被中斷", e);
        }
    }
    // Java 17 Record:訂單詳情
    public record OrderDetail(String orderId, Long userId, String productName, int quantity, double price) {}
    // Java 17 Record:訂單金額計(jì)算結(jié)果
    public record OrderAmountResult(String orderId, double totalAmount, double tax) {}
    // 資源關(guān)閉:Spring 銷毀時(shí)關(guān)閉線程池(try-with-resources 思想)
    @Override
    protected void finalize() throws Throwable {
        try (ExecutorService executor = CUSTOM_EXECUTOR) {
            executor.shutdown();
            if (!executor.awaitTermination(5, TimeUnit.SECONDS)) {
                executor.shutdownNow();
            }
            log.info("異步線程池已關(guān)閉");
        } finally {
            super.finalize();
        }
    }
    // 測(cè)試入口
    public static void main(String[] args) {
        AsyncCascadeService service = new AsyncCascadeService();
        // 執(zhí)行級(jí)聯(lián)異步任務(wù)
        CompletableFuture<OrderAmountResult> resultFuture = service.processOrderAsync(1001L);
        // 阻塞獲取結(jié)果(實(shí)際業(yè)務(wù)中建議使用 thenAccept/thenRun 非阻塞處理)
        resultFuture.whenComplete((result, ex) -> {
            if (ex != null) {
                log.error("級(jí)聯(lián)任務(wù)執(zhí)行失敗", ex);
            } else {
                log.info("級(jí)聯(lián)任務(wù)執(zhí)行完成:{}", result);
            }
            // 關(guān)閉線程池
            CUSTOM_EXECUTOR.shutdown();
        }).join();
    }
}

核心知識(shí)點(diǎn)解釋

  • supplyAsync 作用
    • 異步執(zhí)行有返回值的任務(wù),參數(shù)是 Supplier<T>(無(wú)入?yún)?、有返回值?/li>
    • 第二個(gè)參數(shù)指定自定義線程池(推薦,避免使用 JVM 默認(rèn)的 ForkJoinPool.commonPool(),防止核心線程被耗盡)
  • thenApplyAsync 作用
    • 異步處理前序 CompletableFuture 的結(jié)果,參數(shù)是 Function<T, R>(接收前序結(jié)果、返回新結(jié)果)
    • thenApply 的區(qū)別:thenApply 使用前序任務(wù)的線程執(zhí)行,thenApplyAsync 提交到線程池異步執(zhí)行,真正實(shí)現(xiàn)“級(jí)聯(lián)異步”
  • Java 17+ 特性使用
    • record:簡(jiǎn)化數(shù)據(jù)載體(OrderDetail/OrderAmountResult),自動(dòng)生成構(gòu)造器、getter、equals/hashCode 等
    • 異常處理:exceptionally 捕獲前序任務(wù)異常并兜底,whenComplete 統(tǒng)一處理最終結(jié)果/異常
  • 資源管理
    • 線程池使用 try-with-resources 思想關(guān)閉(finalize 方法中通過(guò) try-with-resources 語(yǔ)法確保線程池關(guān)閉)
    • 中斷處理:捕獲 InterruptedException 后恢復(fù)線程中斷狀態(tài),避免線程狀態(tài)丟失
  • 級(jí)聯(lián)調(diào)用邏輯
  • 生成訂單ID(supplyAsync)→ 查詢訂單詳情(thenApplyAsync)→ 計(jì)算金額(thenApplyAsync)
    • 每個(gè)步驟都異步執(zhí)行,且后序任務(wù)依賴前序任務(wù)的結(jié)果,實(shí)現(xiàn)“流水線式”異步處理。

最佳實(shí)踐

  • 避免在 thenApplyAsync 中執(zhí)行耗時(shí)過(guò)長(zhǎng)的任務(wù),建議拆分多個(gè)小任務(wù)
  • 優(yōu)先使用自定義線程池,通過(guò) @Configuration 配置線程池參數(shù)(核心線程數(shù)、最大線程數(shù)、隊(duì)列等)
  • 異常處理:每個(gè)異步步驟都建議添加 exceptionallyhandle,避免異常穿透導(dǎo)致整個(gè)鏈路失敗
  • 非阻塞處理結(jié)果:使用 thenAccept(消費(fèi)結(jié)果)/thenRun(無(wú)結(jié)果處理)替代 join() 阻塞調(diào)用

第二部分: 非阻塞處理結(jié)果

以下示例除阻塞的 join(),完全基于 thenAccept(消費(fèi)結(jié)果)、thenRun(無(wú)結(jié)果收尾)實(shí)現(xiàn)非阻塞處理,并補(bǔ)充 whenComplete 統(tǒng)一異常處理,同時(shí)優(yōu)化 Spring 規(guī)范的資源關(guān)閉(替換不推薦的 finalize()@PreDestroy)。

import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.beans.factory.annotation.PreDestroy;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
/**
 * 非阻塞處理 CompletableFuture 結(jié)果示例
 * 核心:thenAccept(消費(fèi)結(jié)果)/thenRun(收尾動(dòng)作)替代 join() 阻塞調(diào)用
 */
@Slf4j
@Service
public class NonBlockingAsyncService {
    // 自定義線程池(Spring 中建議通過(guò) @Bean 配置,此處簡(jiǎn)化)
    private final ExecutorService customExecutor = Executors.newFixedThreadPool(5);
    // 訂單狀態(tài)枚舉(Java 17+ 增強(qiáng)枚舉,可結(jié)合 switch 表達(dá)式)
    public enum OrderStatus { SUCCESS, FAILURE }
    /**
     * 級(jí)聯(lián)異步任務(wù)(非阻塞核心邏輯)
     */
    public void processOrderNonBlocking(Long userId) {
        // 步驟1:異步生成訂單ID(supplyAsync)
        CompletableFuture<String> orderIdFuture = CompletableFuture.supplyAsync(() -> {
            log.info("【異步任務(wù)1】生成訂單ID | 線程:{}", Thread.currentThread().getName());
            simulateDelay(100);
            return "ORDER_" + userId + "_" + System.currentTimeMillis();
        }, customExecutor)
        .exceptionally(ex -> {
            log.error("【任務(wù)1失敗】生成訂單ID異常", ex);
            throw new RuntimeException("生成訂單ID失敗", ex);
        });
        // 步驟2:異步查詢訂單詳情(thenApplyAsync)
        CompletableFuture<OrderDetail> orderDetailFuture = orderIdFuture
                .thenApplyAsync(orderId -> {
                    log.info("【異步任務(wù)2】查詢訂單詳情 | 訂單ID:{} | 線程:{}", orderId, Thread.currentThread().getName());
                    simulateDelay(150);
                    return new OrderDetail(orderId, userId, "Java編程實(shí)戰(zhàn)", 2, 89.9);
                }, customExecutor)
                .exceptionally(ex -> {
                    log.error("【任務(wù)2失敗】查詢訂單詳情異常", ex);
                    throw new RuntimeException("查詢訂單詳情失敗", ex);
                });
        // 步驟3:非阻塞處理最終結(jié)果(核心:thenAccept + thenRun)
        orderDetailFuture
                // 1. whenComplete:統(tǒng)一處理結(jié)果/異常(非阻塞)
                .whenComplete((detail, ex) -> {
                    if (ex != null) {
                        log.error("【級(jí)聯(lián)任務(wù)失敗】訂單處理異常", ex);
                        // 非阻塞處理失敗邏輯(如通知、落庫(kù))
                        handleOrderFailure(userId, ex.getMessage());
                    } else {
                        log.info("【級(jí)聯(lián)任務(wù)成功】訂單詳情:{}", detail);
                    }
                })
                // 2. thenAccept:消費(fèi)結(jié)果(有返回值時(shí)用,非阻塞)
                .thenAccept(detail -> {
                    log.info("【非阻塞消費(fèi)】計(jì)算訂單金額 | 訂單ID:{}", detail.orderId());
                    double total = detail.quantity() * detail.price();
                    // 模擬:非阻塞更新訂單金額到數(shù)據(jù)庫(kù)
                    updateOrderAmount(detail.orderId(), total);
                })
                // 3. thenRun:無(wú)結(jié)果收尾動(dòng)作(非阻塞)
                .thenRun(() -> {
                    log.info("【非阻塞收尾】訂單處理流程結(jié)束 | 清理臨時(shí)資源/發(fā)送通知");
                    // 模擬:非阻塞發(fā)送處理完成通知
                    sendOrderNotify(userId, OrderStatus.SUCCESS);
                });
        // 關(guān)鍵:當(dāng)前線程(如主線程/Controller線程)不會(huì)阻塞,直接繼續(xù)執(zhí)行后續(xù)邏輯
        log.info("【主線程不阻塞】訂單處理請(qǐng)求已提交,當(dāng)前線程繼續(xù)執(zhí)行其他任務(wù)...");
    }
    // ------------------- 模擬業(yè)務(wù)方法 -------------------
    private void simulateDelay(long millis) {
        try {
            TimeUnit.MILLISECONDS.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("任務(wù)被中斷", e);
        }
    }
    private void handleOrderFailure(Long userId, String errorMsg) {
        log.warn("【失敗處理】用戶{}訂單處理失敗:{}", userId, errorMsg);
        // 實(shí)際場(chǎng)景:非阻塞發(fā)送失敗通知、記錄異常日志等
    }
    private void updateOrderAmount(String orderId, double total) {
        log.info("【更新金額】訂單{}總金額:{}元", orderId, total);
        // 實(shí)際場(chǎng)景:非阻塞調(diào)用DAO更新數(shù)據(jù)庫(kù)
    }
    private void sendOrderNotify(Long userId, OrderStatus status) {
        // Java 17+ switch表達(dá)式(簡(jiǎn)化分支邏輯)
        String msg = switch (status) {
            case SUCCESS -> "用戶" + userId + "訂單處理完成";
            case FAILURE -> "用戶" + userId + "訂單處理失敗";
        };
        log.info("【發(fā)送通知】{}", msg);
    }
    // ------------------- Java 17+ 特性 -------------------
    // Record:簡(jiǎn)化數(shù)據(jù)載體(自動(dòng)生成構(gòu)造器、getter、equals等)
    public record OrderDetail(String orderId, Long userId, String productName, int quantity, double price) {}
    // ------------------- 資源管理(Spring規(guī)范) -------------------
    // Spring 銷毀Bean時(shí)關(guān)閉線程池(try-with-resources 確保資源釋放)
    @PreDestroy
    public void destroy() {
        log.info("【關(guān)閉線程池】開始釋放異步線程池資源...");
        // try-with-resources 自動(dòng)關(guān)閉 ExecutorService(Java 9+ 支持)
        try (ExecutorService executor = customExecutor) {
            executor.shutdown();
            if (!executor.awaitTermination(3, TimeUnit.SECONDS)) {
                executor.shutdownNow();
                log.warn("【強(qiáng)制關(guān)閉】線程池未正常終止,強(qiáng)制關(guān)閉");
            }
        } catch (InterruptedException e) {
            customExecutor.shutdownNow();
            Thread.currentThread().interrupt();
            log.error("【關(guān)閉異?!烤€程池關(guān)閉被中斷", e);
        }
        log.info("【關(guān)閉完成】異步線程池已釋放");
    }
    // ------------------- 測(cè)試入口(非阻塞效果演示) -------------------
    public static void main(String[] args) {
        NonBlockingAsyncService service = new NonBlockingAsyncService();
        // 提交異步任務(wù)(非阻塞)
        service.processOrderNonBlocking(1001L);
        // 主線程繼續(xù)執(zhí)行(不會(huì)等待異步任務(wù)完成)
        log.info("【主線程】已提交訂單處理任務(wù),開始執(zhí)行其他業(yè)務(wù)邏輯...");
        simulateMainThreadWork();
        // 等待異步任務(wù)完成(僅測(cè)試用,實(shí)際業(yè)務(wù)無(wú)需此步驟)
        try {
            TimeUnit.SECONDS.sleep(2);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
    private static void simulateMainThreadWork() {
        log.info("【主線程】執(zhí)行其他任務(wù):檢查庫(kù)存、驗(yàn)證用戶權(quán)限...");
        try {
            TimeUnit.MILLISECONDS.sleep(500);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

核心知識(shí)點(diǎn):非阻塞處理的關(guān)鍵區(qū)別

方法作用是否阻塞適用場(chǎng)景
join()阻塞當(dāng)前線程,等待異步任務(wù)返回結(jié)果測(cè)試/必須同步獲取結(jié)果的場(chǎng)景
thenAccept()異步消費(fèi)結(jié)果(入?yún)榻Y(jié)果,無(wú)返回值)有結(jié)果需要處理(如更新/通知)
thenRun()異步執(zhí)行收尾動(dòng)作(無(wú)入?yún)ⅰo(wú)返回值)結(jié)果處理完成后的清理/日志
whenComplete異步處理結(jié)果+異常(入?yún)榻Y(jié)果+異常)統(tǒng)一捕獲異常+分支處理

非阻塞處理的核心優(yōu)勢(shì)

  1. 提升系統(tǒng)吞吐量:提交異步任務(wù)的線程(如 Controller 線程)無(wú)需等待,可立即處理下一個(gè)請(qǐng)求,避免線程阻塞導(dǎo)致的資源浪費(fèi);
  2. 避免線程死鎖/饑餓:阻塞式 join() 可能導(dǎo)致線程池耗盡,非阻塞回調(diào)更符合異步編程的設(shè)計(jì)初衷;
  3. 邏輯解耦:結(jié)果消費(fèi)、異常處理、收尾動(dòng)作拆分到不同回調(diào)中,代碼更清晰(符合單一職責(zé))。

關(guān)鍵注意事項(xiàng)

  1. 異常處理:非阻塞場(chǎng)景下,必須通過(guò) whenComplete/exceptionally 捕獲異常,否則異常會(huì)被“吞掉”(無(wú)感知失?。?;
  2. 線程池選擇thenAccept/thenRun 默認(rèn)使用前序任務(wù)的線程池,建議統(tǒng)一指定自定義線程池(避免默認(rèn)池耗盡);
  3. 資源生命周期:Spring 中通過(guò) @PreDestroy 關(guān)閉線程池(替代 finalize(),finalize() 已被標(biāo)記為過(guò)時(shí)),結(jié)合 try-with-resources 確保資源釋放;
  4. Java 17+ 增強(qiáng):示例中使用 switch 表達(dá)式(替代傳統(tǒng) switch)、record(簡(jiǎn)化數(shù)據(jù)類),符合最新語(yǔ)法規(guī)范。

輸出效果(非阻塞特征)

運(yùn)行代碼后,控制臺(tái)會(huì)先打印主線程不阻塞的日志,再異步打印各任務(wù)的日志,體現(xiàn)“提交任務(wù)后主線程立即繼續(xù)執(zhí)行”的非阻塞特性:

【主線程不阻塞】訂單處理請(qǐng)求已提交,當(dāng)前線程繼續(xù)執(zhí)行其他任務(wù)...
【主線程】執(zhí)行其他任務(wù):檢查庫(kù)存、驗(yàn)證用戶權(quán)限...
【異步任務(wù)1】生成訂單ID | 線程:pool-1-thread-1
【異步任務(wù)2】查詢訂單詳情 | 訂單ID:ORDER_1001_17338xxxxxx | 線程:pool-1-thread-2
【級(jí)聯(lián)任務(wù)成功】訂單詳情:OrderDetail[orderId=ORDER_1001_17338xxxxxx, userId=1001, productName=Java編程實(shí)戰(zhàn), quantity=2, price=89.9]
【非阻塞消費(fèi)】計(jì)算訂單金額 | 訂單ID:ORDER_1001_17338xxxxxx
【更新金額】訂單ORDER_1001_17338xxxxxx總金額:179.8元
【非阻塞收尾】訂單處理流程結(jié)束 | 清理臨時(shí)資源/發(fā)送通知
【發(fā)送通知】用戶1001訂單處理完成

第三部分 其它異步編程方式全解析

Java 異步編程體系從基礎(chǔ)線程模型到現(xiàn)代響應(yīng)式編程、虛擬線程逐步演進(jìn),以下是主流方式的代碼示例+特性解析,均基于 Java 17+ 語(yǔ)法(部分如虛擬線程需 Java 21+ 正式支持,17+ 可通過(guò)預(yù)覽特性啟用),并包含異常處理、資源管理規(guī)范。

一、基礎(chǔ)核心:線程池 + Future(Java 5+)

ExecutorService 結(jié)合 Future 是 CompletableFuture 之前的核心異步方式,通過(guò)提交 Callable/Runnable 實(shí)現(xiàn)異步任務(wù),缺點(diǎn)是結(jié)果獲取阻塞、無(wú)鏈?zhǔn)秸{(diào)用,但勝在簡(jiǎn)單可控。

代碼示例(Java 17+)

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
/**
 * 線程池 + Future 異步示例
 * 核心:提交Callable獲取Future,通過(guò)get()阻塞獲取結(jié)果(可設(shè)置超時(shí))
 */
public class FutureAsyncDemo {
    // 自定義線程池(Java 17+ 推薦使用Executors靜態(tài)工廠+try-with-resources關(guān)閉)
    private static final ExecutorService EXECUTOR = Executors.newVirtualThreadPerTaskExecutor(); // Java 21+ 虛擬線程池(17+需預(yù)覽)
    // Java 17 Record:異步任務(wù)結(jié)果
    public record TaskResult(Long taskId, String content, long costTimeMs) {}
    // 異步任務(wù):生成業(yè)務(wù)數(shù)據(jù)
    private static Callable<TaskResult> generateDataTask(Long taskId) {
        return () -> {
            long start = System.currentTimeMillis();
            log("任務(wù){(diào)}開始執(zhí)行,線程:{}", taskId, Thread.currentThread().getName());
            // 模擬耗時(shí)
            TimeUnit.MILLISECONDS.sleep(200);
            long cost = System.currentTimeMillis() - start;
            return new TaskResult(taskId, "異步生成數(shù)據(jù)_" + taskId, cost);
        };
    }
    public static void main(String[] args) {
        // 提交異步任務(wù)
        Future<TaskResult> future1 = EXECUTOR.submit(generateDataTask(1L));
        Future<TaskResult> future2 = EXECUTOR.submit(generateDataTask(2L));
        // 非阻塞提交后,主線程可執(zhí)行其他邏輯
        log("主線程執(zhí)行其他任務(wù)...");
        // 阻塞獲取結(jié)果(可設(shè)置超時(shí),避免永久阻塞)
        try {
            TaskResult result1 = future1.get(1, TimeUnit.SECONDS); // 超時(shí)時(shí)間1秒
            TaskResult result2 = future2.get(1, TimeUnit.SECONDS);
            log("任務(wù)1結(jié)果:{}", result1);
            log("任務(wù)2結(jié)果:{}", result2);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            log("任務(wù)被中斷:{}", e.getMessage());
        } catch (ExecutionException e) {
            log("任務(wù)執(zhí)行失?。簕}", e.getCause().getMessage());
        } catch (TimeoutException e) {
            log("任務(wù)超時(shí):{}", e.getMessage());
            future1.cancel(true); // 超時(shí)取消任務(wù)
            future2.cancel(true);
        } finally {
            // try-with-resources 關(guān)閉線程池(Java 9+ 支持ExecutorService自動(dòng)關(guān)閉)
            try (ExecutorService executor = EXECUTOR) {
                executor.shutdown();
                if (!executor.awaitTermination(3, TimeUnit.SECONDS)) {
                    executor.shutdownNow();
                    log("線程池強(qiáng)制關(guān)閉");
                }
            } catch (InterruptedException e) {
                EXECUTOR.shutdownNow();
                Thread.currentThread().interrupt();
            }
        }
    }
    private static void log(String msg, Object... args) {
        System.out.printf("[%s] %s%n", Thread.currentThread().getName(), String.format(msg, args));
    }
}

核心特性

  • 優(yōu)點(diǎn):基礎(chǔ)通用,兼容所有Java版本,線程池可控;
  • 缺點(diǎn)Future.get() 阻塞線程,無(wú)鏈?zhǔn)秸{(diào)用、無(wú)異步異常處理,無(wú)法組合多個(gè)任務(wù);
  • 適用場(chǎng)景:簡(jiǎn)單異步任務(wù),無(wú)需結(jié)果聯(lián)動(dòng)的場(chǎng)景。

二、現(xiàn)代輕量:虛擬線程(VirtualThread,Java 21+ 正式/17+ 預(yù)覽)

Project Loom 推出的虛擬線程(用戶態(tài)線程)是 Java 異步編程的重大升級(jí),相比傳統(tǒng)平臺(tái)線程,虛擬線程輕量(百萬(wàn)級(jí)創(chuàng)建)、無(wú)上下文切換開銷,天然適配異步場(chǎng)景。

代碼示例(Java 21+)

import java.time.Duration;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
/**
 * 虛擬線程(VirtualThread)異步示例
 * Java 21+ 正式支持,17+ 需添加 --enable-preview 啟動(dòng)參數(shù)
 */
public class VirtualThreadAsyncDemo {
    // Java 17 Record:用戶數(shù)據(jù)
    public record UserData(Long userId, String username, String email) {}
    // 異步任務(wù):查詢用戶信息
    private static UserData queryUser(Long userId) {
        try {
            // 模擬IO阻塞(虛擬線程阻塞時(shí)會(huì)釋放載體線程,無(wú)性能損耗)
            TimeUnit.MILLISECONDS.sleep(300);
            log("查詢用戶{}完成,線程:{}", userId, Thread.currentThread().getName());
            return new UserData(userId, "user_" + userId, userId + "@example.com");
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("查詢用戶中斷", e);
        }
    }
    public static void main(String[] args) throws InterruptedException {
        // 虛擬線程池(Java 21+ 推薦)
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            // 提交1000個(gè)虛擬線程任務(wù)(輕量,無(wú)性能壓力)
            for (long i = 1; i <= 1000; i++) {
                long userId = i;
                executor.submit(() -> queryUser(userId));
            }
            // 主線程等待任務(wù)完成(僅測(cè)試用)
            executor.awaitTermination(Duration.ofSeconds(5));
            log("所有虛擬線程任務(wù)執(zhí)行完成");
        } // try-with-resources 自動(dòng)關(guān)閉線程池
    }
    private static void log(String msg, Object... args) {
        System.out.printf("[%s] %s%n", Thread.currentThread().getName(), String.format(msg, args));
    }
}

核心特性

  • 優(yōu)點(diǎn)
    • 輕量:?jiǎn)蝹€(gè)JVM可創(chuàng)建百萬(wàn)級(jí)虛擬線程,遠(yuǎn)超平臺(tái)線程(千級(jí));
    • 無(wú)阻塞損耗:虛擬線程遇到IO阻塞時(shí),會(huì)釋放底層載體線程(OS線程),避免資源浪費(fèi);
    • 編程模型簡(jiǎn)單:無(wú)需回調(diào)地獄,同步寫法實(shí)現(xiàn)異步效果;
  • 缺點(diǎn):Java 21+ 正式支持,低版本需升級(jí);
  • 適用場(chǎng)景:IO密集型異步任務(wù)(如數(shù)據(jù)庫(kù)、網(wǎng)絡(luò)調(diào)用),替代傳統(tǒng)線程池+回調(diào)。

三、異步IO:NIO.2 異步通道(Java 7+)

Java NIO.2 提供 AsynchronousFileChannel/AsynchronousSocketChannel 等異步IO通道,基于事件驅(qū)動(dòng)實(shí)現(xiàn)文件/網(wǎng)絡(luò)IO的異步操作,無(wú)需線程阻塞等待IO完成。

代碼示例(異步文件讀寫,Java 17+)

import java.nio.ByteBuffer;
import java.nio.channels.AsynchronousFileChannel;
import java.nio.channels.CompletionHandler;
import java.nio.file.Path;
import java.nio.file.StandardOpenOption;
import java.nio.charset.StandardCharsets;
/**
 * NIO.2 異步文件IO示例
 * 核心:CompletionHandler 回調(diào)處理異步結(jié)果,無(wú)線程阻塞
 */
public class AsyncFileIODemo {
    private static final Path FILE_PATH = Path.of("async_demo.txt");
    public static void main(String[] args) throws InterruptedException {
        // try-with-resources 管理異步文件通道(自動(dòng)關(guān)閉資源)
        try (AsynchronousFileChannel fileChannel = AsynchronousFileChannel.open(
                FILE_PATH,
                StandardOpenOption.READ,
                StandardOpenOption.WRITE,
                StandardOpenOption.CREATE
        )) {
            // 步驟1:異步寫入文件
            String content = "Java 17 異步IO示例:" + System.currentTimeMillis();
            ByteBuffer writeBuffer = ByteBuffer.wrap(content.getBytes(StandardCharsets.UTF_8));
            fileChannel.write(writeBuffer, 0, null, new CompletionHandler<Integer, Void>() {
                @Override
                public void completed(Integer bytesWritten, Void attachment) {
                    log("異步寫入完成,寫入字節(jié)數(shù):{}", bytesWritten);
                    // 寫入完成后異步讀取文件
                    readFileAsync(fileChannel);
                }
                @Override
                public void failed(Throwable exc, Void attachment) {
                    log("異步寫入失敗:{}", exc.getMessage());
                }
            });
            // 主線程等待異步操作完成(僅測(cè)試用)
            Thread.sleep(2000);
        } catch (Exception e) {
            log("文件通道操作異常:{}", e.getMessage());
        }
    }
    // 異步讀取文件
    private static void readFileAsync(AsynchronousFileChannel fileChannel) {
        ByteBuffer readBuffer = ByteBuffer.allocate(1024);
        fileChannel.read(readBuffer, 0, null, new CompletionHandler<Integer, Void>() {
            @Override
            public void completed(Integer bytesRead, Void attachment) {
                if (bytesRead == -1) {
                    log("文件讀取完成(無(wú)數(shù)據(jù))");
                    return;
                }
                readBuffer.flip();
                String content = StandardCharsets.UTF_8.decode(readBuffer).toString();
                log("異步讀取結(jié)果:{}", content);
            }
            @Override
            public void failed(Throwable exc, Void attachment) {
                log("異步讀取失?。簕}", exc.getMessage());
            }
        });
    }
    private static void log(String msg, Object... args) {
        System.out.printf("[%s] %s%n", Thread.currentThread().getName(), String.format(msg, args));
    }
}

核心特性

  • 優(yōu)點(diǎn):底層基于操作系統(tǒng)異步IO(AIO),無(wú)用戶態(tài)線程阻塞,IO效率極高;
  • 缺點(diǎn):編程模型繁瑣(依賴 CompletionHandler 回調(diào)),僅適用于IO場(chǎng)景;
  • 適用場(chǎng)景:高并發(fā)文件/網(wǎng)絡(luò)IO(如文件服務(wù)器、網(wǎng)絡(luò)通信)。

四、響應(yīng)式編程:Reactor(Spring WebFlux,Java 8+)

Reactor 是 Java 響應(yīng)式編程的主流框架(實(shí)現(xiàn) Reactive Streams 規(guī)范),基于 Mono(單值異步)/Flux(多值異步)實(shí)現(xiàn)非阻塞、背壓可控的異步編程,是 Spring WebFlux 的核心。

代碼示例(Spring Boot 3.x + Java 17+)

import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
import java.time.Duration;
/**
 * Reactor 響應(yīng)式異步示例(Spring Boot 3.x + Java 17+)
 */
@Slf4j
@Service
public class ReactorAsyncDemo {
    // Java 17 Record:商品數(shù)據(jù)
    public record Product(Long id, String name, double price) {}
    // 異步查詢商品(Mono 表示單值異步結(jié)果)
    public Mono<Product> queryProductAsync(Long productId) {
        return Mono.fromSupplier(() -> {
                    log.info("查詢商品{},線程:{}", productId, Thread.currentThread().getName());
                    // 模擬IO耗時(shí)
                    try {
                        Thread.sleep(200);
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                        throw new RuntimeException("查詢中斷", e);
                    }
                    return new Product(productId, "商品_" + productId, 99.9 + productId);
                })
                .subscribeOn(Schedulers.boundedElastic()) // 指定異步線程池
                .timeout(Duration.ofSeconds(1)) // 超時(shí)控制
                .onErrorResume(ex -> { // 異常兜底
                    log.error("查詢商品失?。簕}", ex.getMessage());
                    return Mono.just(new Product(productId, "默認(rèn)商品", 0.0));
                });
    }
    // 測(cè)試入口(Spring Boot 環(huán)境)
    public static void main(String[] args) {
        ReactorAsyncDemo service = new ReactorAsyncDemo();
        // 非阻塞消費(fèi)結(jié)果
        service.queryProductAsync(1001L)
                .doOnNext(product -> log.info("查詢結(jié)果:{}", product)) // 消費(fèi)結(jié)果
                .doOnTerminate(() -> log.info("異步任務(wù)結(jié)束")) // 收尾動(dòng)作
                .subscribe(); // 訂閱(觸發(fā)異步執(zhí)行)
        // 主線程等待(僅測(cè)試用)
        try {
            Thread.sleep(1000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

核心特性

  • 優(yōu)點(diǎn)
    • 非阻塞+背壓:可控制數(shù)據(jù)生產(chǎn)速率,避免消費(fèi)者過(guò)載;
    • 豐富的操作符:map/flatMap/filter 等實(shí)現(xiàn)復(fù)雜異步邏輯;
    • 與 Spring 生態(tài)深度集成:Spring WebFlux 支持全響應(yīng)式Web開發(fā);
  • 缺點(diǎn):學(xué)習(xí)曲線陡峭,調(diào)試難度高;
  • 適用場(chǎng)景:高并發(fā)微服務(wù)、流式數(shù)據(jù)處理、實(shí)時(shí)響應(yīng)系統(tǒng)。

五、補(bǔ)充:SwingWorker(GUI 異步場(chǎng)景)

SwingWorker 是 Java 專門為 Swing GUI 設(shè)計(jì)的異步工具,用于在后臺(tái)線程執(zhí)行耗時(shí)任務(wù),避免阻塞UI線程(EDT),屬于小眾但專用的異步方式。

核心特性

  • 優(yōu)點(diǎn):專為GUI異步設(shè)計(jì),內(nèi)置進(jìn)度回調(diào)、結(jié)果發(fā)布;
  • 缺點(diǎn):僅適用于Swing/AWT場(chǎng)景,通用性差;
  • 適用場(chǎng)景:桌面GUI應(yīng)用的后臺(tái)任務(wù)(如數(shù)據(jù)加載、文件導(dǎo)出)。

六、各異步方式對(duì)比與選型建議

方式核心優(yōu)勢(shì)核心劣勢(shì)適用場(chǎng)景
ExecutorService+Future基礎(chǔ)通用,兼容性好阻塞獲取結(jié)果,無(wú)鏈?zhǔn)秸{(diào)用簡(jiǎn)單異步任務(wù),無(wú)結(jié)果聯(lián)動(dòng)
VirtualThread輕量高效,同步寫法異步Java 21+ 依賴IO密集型異步任務(wù)(數(shù)據(jù)庫(kù)/網(wǎng)絡(luò))
NIO.2 異步通道底層AIO,IO效率極高編程繁瑣,僅適用于IO高并發(fā)文件/網(wǎng)絡(luò)IO
Reactor(WebFlux)非阻塞+背壓,生態(tài)完善學(xué)習(xí)成本高,調(diào)試難高并發(fā)微服務(wù)、流式數(shù)據(jù)處理
SwingWorker專為GUI設(shè)計(jì),適配EDT通用性差Swing/AWT桌面應(yīng)用
CompletableFuture鏈?zhǔn)秸{(diào)用,任務(wù)組合線程池需手動(dòng)管理通用異步任務(wù),結(jié)果聯(lián)動(dòng)/組合

選型核心原則

  1. IO密集型任務(wù):優(yōu)先選 虛擬線程(Java 21+),次選 CompletableFuture;
  2. 高并發(fā)流式處理:選 Reactor(WebFlux)
  3. 文件/網(wǎng)絡(luò)IO:選 NIO.2 異步通道;
  4. 簡(jiǎn)單異步無(wú)聯(lián)動(dòng):選 ExecutorService+Future;
  5. GUI應(yīng)用:選 SwingWorker。

所有方式均需遵循:

  • 資源管理:使用 try-with-resources 關(guān)閉線程池/IO通道;
  • 異常處理:捕獲中斷異常并恢復(fù)線程中斷狀態(tài),避免狀態(tài)丟失;
  • 線程池規(guī)范:自定義線程池,避免使用默認(rèn)池(如 ForkJoinPool.commonPool())。

到此這篇關(guān)于Java 中多線程異步處理的文章就介紹到這了,更多相關(guān)java多線程異步內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

最新評(píng)論

彝良县| 滁州市| 察雅县| 阳谷县| 福建省| 寻乌县| 上思县| 营山县| 东兰县| 江北区| 马公市| 温州市| 武夷山市| 双桥区| 屏南县| 宁武县| 墨玉县| 玛沁县| 綦江县| 海宁市| 襄垣县| 麟游县| 色达县| 宜黄县| 马公市| 什邡市| 衢州市| 天水市| 江阴市| 喀喇| 永兴县| 石台县| 永寿县| 佳木斯市| 徐州市| 平原县| 德江县| 潞城市| 丰县| 安龙县| 奎屯市|