0

[Java Backend Zero to Hello] BÀI 2.7: EXECUTOR FRAMEWORK

Java Backend Zero to Hello

📚 Bài viết thuộc series Java Backend Zero to Hello 📌 Phần: Phase 2: Java Core Nâng Cao | Bài 22/86


BÀI 2.7: EXECUTOR FRAMEWORK

Mục tiêu

  • Hiểu Executor Framework
  • Sử dụng ThreadPool
  • Làm việc với Future và CompletableFuture
  • Scheduled tasks

1. TẠI SAO DÙNG EXECUTOR?

Tạo Thread thủ công tốn tài nguyên. Executor Framework quản lý Thread pool hiệu quả.

So sánh

// ❌ Tạo thread thủ công - tốn tài nguyên
for (int i = 0; i < 1000; i++) {
    new Thread(() -> doWork()).start();
}

// ✅ Dùng Executor - tái sử dụng thread
ExecutorService executor = Executors.newFixedThreadPool(10);
for (int i = 0; i < 1000; i++) {
    executor.submit(() -> doWork());
}
executor.shutdown();

2. CÁC LOẠI EXECUTOR

2.1 newFixedThreadPool

Số thread cố định.

ExecutorService executor = Executors.newFixedThreadPool(5);
// Tối đa 5 thread chạy đồng thời

2.2 newCachedThreadPool

Tạo thread khi cần, tái sử dụng thread rảnh.

ExecutorService executor = Executors.newCachedThreadPool();
// Phù hợp với task ngắn

2.3 newSingleThreadExecutor

Chỉ 1 thread, chạy tuần tự.

ExecutorService executor = Executors.newSingleThreadExecutor();

2.4 newScheduledThreadPool

Thread chạy theo lịch.

ScheduledExecutorService executor = Executors.newScheduledThreadPool(3);

2.5 newWorkStealingPool (Java 8+)

Dùng ForkJoinPool, tối ưu cho CPU-bound.

ExecutorService executor = Executors.newWorkStealingPool();

3. SUBMIT TASK

3.1 execute - Không trả về kết quả

executor.execute(() -> {
    System.out.println("Task running");
});

3.2 submit - Trả về Future

Future<Integer> future = executor.submit(() -> {
    Thread.sleep(1000);
    return 42;
});

// Lấy kết quả (blocking)
Integer result = future.get();

// Lấy với timeout
Integer result2 = future.get(2, TimeUnit.SECONDS);

// Kiểm tra
boolean done = future.isDone();
boolean cancelled = future.isCancelled();
future.cancel(true);

3.3 invokeAll - Nhiều task

List<Callable<String>> tasks = Arrays.asList(
    () -> "Task 1",
    () -> "Task 2",
    () -> "Task 3"
);

List<Future<String>> futures = executor.invokeAll(tasks);
for (Future<String> future : futures) {
    System.out.println(future.get());
}

3.4 invokeAny - Task đầu tiên hoàn thành

String result = executor.invokeAny(tasks);  // Trả về kết quả task đầu tiên

4. SHUTDOWN

executor.shutdown();              // Không nhận task mới, chờ task hiện tại
executor.shutdownNow();           // Cố gắng dừng ngay
executor.awaitTermination(10, TimeUnit.SECONDS);  // Chờ tối đa 10s

// Best practice
executor.shutdown();
try {
    if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
        executor.shutdownNow();
    }
} catch (InterruptedException e) {
    executor.shutdownNow();
    Thread.currentThread().interrupt();
}

5. SCHEDULED TASKS

ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(3);

// Chạy sau 5 giây
scheduler.schedule(() -> {
    System.out.println("Delayed task");
}, 5, TimeUnit.SECONDS);

// Chạy sau 5s, lặp lại mỗi 10s
scheduler.scheduleAtFixedRate(() -> {
    System.out.println("Periodic task");
}, 5, 10, TimeUnit.SECONDS);

// Chạy sau 5s, lặp lại 10s sau khi task trước kết thúc
scheduler.scheduleWithFixedDelay(() -> {
    System.out.println("Periodic task with delay");
}, 5, 10, TimeUnit.SECONDS);

6. COMPLETABLEFUTURE (Java 8+)

Mạnh mẽ hơn Future, hỗ trợ chaining và combining.

6.1 Tạo CompletableFuture

// Từ task
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    return "Hello";
}, executor);

// Từ giá trị
CompletableFuture<String> completed = CompletableFuture.completedFuture("Done");

6.2 Chaining

CompletableFuture.supplyAsync(() -> "Hello")
    .thenApply(s -> s + " World")           // Biến đổi
    .thenApply(String::toUpperCase)         // "HELLO WORLD"
    .thenAccept(System.out::println);       // In ra

6.3 Combining

CompletableFuture<String> future1 = CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture<String> future2 = CompletableFuture.supplyAsync(() -> "World");

// Combine 2 future
CompletableFuture<String> combined = future1.thenCombine(future2, (a, b) -> a + " " + b);
// "Hello World"

// Chờ cả 2
CompletableFuture<Void> all = CompletableFuture.allOf(future1, future2);

// Chờ 1 trong 2
CompletableFuture<Object> any = CompletableFuture.anyOf(future1, future2);

6.4 Xử lý lỗi

CompletableFuture.supplyAsync(() -> {
    if (true) throw new RuntimeException("Error");
    return "Result";
})
.exceptionally(ex -> "Default value")  // Fallback khi lỗi
.thenAccept(System.out::println);

6.5 Ví dụ thực tế

public CompletableFuture<UserProfile> getUserProfile(Long userId) {
    CompletableFuture<User> userFuture = CompletableFuture
        .supplyAsync(() -> userRepository.findById(userId));

    CompletableFuture<List<Order>> ordersFuture = CompletableFuture
        .supplyAsync(() -> orderRepository.findByUserId(userId));

    return userFuture.thenCombine(ordersFuture, (user, orders) -> {
        return new UserProfile(user, orders);
    });
}

7. CUSTOM THREADPOOLEXECUTOR

ThreadPoolExecutor executor = new ThreadPoolExecutor(
    5,                      // corePoolSize
    10,                     // maximumPoolSize
    60,                     // keepAliveTime
    TimeUnit.SECONDS,       // unit
    new LinkedBlockingQueue<>(100),  // workQueue
    Executors.defaultThreadFactory(),
    new ThreadPoolExecutor.AbortPolicy()  // rejection policy
);

Rejection Policies

  • AbortPolicy - Ném exception (mặc định)
  • CallerRunsPolicy - Chạy trên thread gọi
  • DiscardPolicy - Bỏ qua
  • DiscardOldestPolicy - Bỏ task cũ nhất

8. BÀI TẬP THỰC HÀNH

Bài 1: Download song song

public List<String> downloadAll(List<String> urls) throws InterruptedException {
    ExecutorService executor = Executors.newFixedThreadPool(10);
    List<Future<String>> futures = urls.stream()
        .map(url -> executor.submit(() -> download(url)))
        .collect(Collectors.toList());

    List<String> results = new ArrayList<>();
    for (Future<String> future : futures) {
        try {
            results.add(future.get(30, TimeUnit.SECONDS));
        } catch (TimeoutException e) {
            future.cancel(true);
        }
    }
    executor.shutdown();
    return results;
}

Bài 2: Scheduled Task

ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);

scheduler.scheduleAtFixedRate(() -> {
    System.out.println("Cleanup at " + LocalDateTime.now());
    // Cleanup logic
}, 0, 1, TimeUnit.HOURS);

Bài 3: CompletableFuture chain

CompletableFuture.supplyAsync(() -> fetchUser(userId))
    .thenApply(this::enrichUser)
    .thenAccept(this::sendNotification)
    .exceptionally(ex -> {
        log.error("Error", ex);
        return null;
    });

9. TÓM TẮT

Khái niệm Mô tả
ExecutorService Quản lý thread pool
FixedThreadPool Số thread cố định
CachedThreadPool Tạo thread khi cần
ScheduledExecutor Chạy theo lịch
Future Kết quả async
CompletableFuture Future mạnh mẽ, chaining
shutdown Đóng executor

Bài tiếp theo: 2.8 I/O & File


🧭 Điều Hướng Series

⬅️ Bài trước: BÀI 2.6: ĐA LUỒNG (MULTITHREADING)

📋 Lộ trình tổng quan: Xem Toàn Bộ Series

➡️ Bài tiếp theo: BÀI 2.8: I/O & FILE


All rights reserved

Viblo
Hãy đăng ký một tài khoản Viblo để nhận được nhiều bài viết thú vị hơn.
Đăng kí