Learn
Java/18-project-multi-threaded-downloader

项目:多线程下载器

本章我们用模拟数据做一个迷你多线程下载器——把第 13 章学到的 ExecutorService / Future / Callable 串起来。真实下载会受网络环境影响;这里我们用字符串标签当「URL」,用 Thread.sleep 模拟下载耗时,能完整跑通并发、重试、进度上报。

1. 需求拆解

我们要做的是一个能并行下载 N 个「文件」的工具:

  • 多线程:固定大小线程池,避免无界创建
  • 进度上报:每个文件下载时回调当前进度
  • 失败重试:单文件失败自动重试最多 3 次
  • 汇总结果:主线程拿到所有 Future 的结果
            ┌──────────┐
  提交 ────►│ thread 1 │────► File("a.txt", bytes)
  任务列表  ├──────────┤
  ────────►│ thread 2 │────► File("b.txt", bytes)
            ├──────────┤
  ────────►│ thread 3 │────► File("c.txt", bytes) (失败重试)
            └──────────┘

2. 项目结构

downloader/
├── pom.xml
└── src/main/java/com/example/dl/
    ├── Main.java              // 启动入口
    ├── Downloader.java        // 线程池 + 任务调度
    ├── DownloadTask.java      // 单文件下载任务(Callable)
    ├── ProgressListener.java  // 进度回调接口
    └── FileResult.java        // 单文件结果

核心类(多文件版本)

public record FileResult(String url, boolean ok, String content, String error) {}
 
public interface ProgressListener {
    void onProgress(String url, int percent);
}
 
public class DownloadTask implements Callable<FileResult> {
    private final String url;
    private final ProgressListener listener;
    public DownloadTask(String url, ProgressListener l) {
        this.url = url; this.listener = l;
    }
    @Override
    public FileResult call() throws Exception {
        int totalSteps = 5;
        for (int i = 1; i <= totalSteps; i++) {
            Thread.sleep(50);             // 模拟网络
            listener.onProgress(url, i * 100 / totalSteps);
        }
        if (url.endsWith("fail")) {
            throw new RuntimeException("simulated 404");
        }
        return new FileResult(url, true, "<<content of " + url + ">>", null);
    }
}
 
public class Downloader {
    private final ExecutorService pool;
    private final int maxRetries;
    public Downloader(int threads, int maxRetries) {
        this.pool = Executors.newFixedThreadPool(threads);
        this.maxRetries = maxRetries;
    }
 
    public List<FileResult> downloadAll(List<String> urls, ProgressListener listener) {
        List<Future<FileResult>> futures = new ArrayList<>();
        for (String url : urls) {
            futures.add(pool.submit(() -> downloadWithRetry(url, listener)));
        }
 
        List<FileResult> results = new ArrayList<>();
        for (Future<FileResult> f : futures) {
            try { results.add(f.get()); }
            catch (Exception e) { results.add(new FileResult("?", false, null, e.getMessage())); }
        }
        return results;
    }
 
    private FileResult downloadWithRetry(String url, ProgressListener listener) {
        Exception last = null;
        for (int attempt = 1; attempt <= maxRetries; attempt++) {
            try {
                return new DownloadTask(url, listener).call();
            } catch (Exception e) {
                last = e;
                System.out.println("retry " + attempt + " for " + url);
            }
        }
        return new FileResult(url, false, null, last.getMessage());
    }
 
    public void shutdown() { pool.shutdown(); }
}
public class Main {
    public static void main(String[] args) {
        List<String> urls = List.of("a.txt", "b.txt", "c-fail.txt", "d.txt");
        Downloader dl = new Downloader(3, 3);
 
        List<FileResult> results = dl.downloadAll(urls, (url, p) ->
            System.out.printf("  [%s] %d%%%n", url, p));
 
        dl.shutdown();
        for (FileResult r : results) {
            System.out.println(r.ok() ? "✓ " + r.url() : "✗ " + r.url() + "  " + r.error());
        }
    }
}

3. Playground:单文件整合版

下面这一份 Playground 把所有类合并到一个 Main,使用内部类,依然保持清晰的层次。

多线程下载器整合
import java.util.*;
import java.util.concurrent.*;
 
public class Main {
    // ====== 数据模型 ======
    record FileResult(String url, boolean ok, String content, String error) {}
 
    @FunctionalInterface
    interface ProgressListener {
        void onProgress(String url, int percent);
    }
 
    // ====== 单任务 ======
    static class DownloadTask implements Callable<FileResult> {
        private final String url;
        private final ProgressListener listener;
        private final int totalSteps;
 
        DownloadTask(String url, ProgressListener l, int totalSteps) {
            this.url = url; this.listener = l; this.totalSteps = totalSteps;
        }
 
        @Override
        public FileResult call() throws Exception {
            for (int i = 1; i <= totalSteps; i++) {
                Thread.sleep(30);   // 模拟网络
                listener.onProgress(url, i * 100 / totalSteps);
            }
            if (url.contains("fail")) {
                throw new RuntimeException("simulated 404");
            }
            return new FileResult(url, true, "<<content of " + url + ">>", null);
        }
    }
 
    // ====== 调度器 ======
    static class Downloader {
        private final ExecutorService pool;
        private final int maxRetries;
        Downloader(int threads, int maxRetries) {
            this.pool = Executors.newFixedThreadPool(threads);
            this.maxRetries = maxRetries;
        }
 
        List<FileResult> downloadAll(List<String> urls, ProgressListener l) {
            List<Future<FileResult>> fs = new ArrayList<>();
            for (String u : urls) fs.add(pool.submit(() -> withRetry(u, l)));
            List<FileResult> out = new ArrayList<>();
            for (Future<FileResult> f : fs) {
                try { out.add(f.get()); }
                catch (Exception e) { out.add(new FileResult("?", false, null, e.getMessage())); }
            }
            return out;
        }
 
        private FileResult withRetry(String url, ProgressListener l) {
            Exception last = null;
            for (int attempt = 1; attempt <= maxRetries; attempt++) {
                try { return new DownloadTask(url, l, 5).call(); }
                catch (Exception e) { last = e; }
            }
            return new FileResult(url, false, null,
                "after " + maxRetries + " retries: " + last.getMessage());
        }
 
        void shutdown() { pool.shutdown(); }
    }
 
    // ====== 入口 ======
    public static void main(String[] args) throws InterruptedException {
        List<String> urls = List.of("a.txt", "b.txt", "c-fail.txt", "d.txt", "e.txt");
        Downloader dl = new Downloader(3, 3);
 
        // 进度回调用 synchronizedList 收集,最后打印汇总
        List<String> progressLog = Collections.synchronizedList(new ArrayList<>());
        ProgressListener listener = (url, p) -> {
            String line = String.format("  [%s] %d%%", url, p);
            progressLog.add(line);
        };
 
        long t0 = System.nanoTime();
        List<FileResult> results = dl.downloadAll(urls, listener);
        double seconds = (System.nanoTime() - t0) / 1e9;
 
        // 等 50ms 让最后一行进度也写完
        Thread.sleep(50);
 
        dl.shutdown();
        System.out.println("=== progress ===");
        progressLog.forEach(System.out::println);
 
        System.out.println("\n=== results ===");
        int ok = 0, fail = 0;
        for (FileResult r : results) {
            if (r.ok()) { ok++; System.out.println("  ✓ " + r.url() + "  " + r.content()); }
            else        { fail++; System.out.println("  ✗ " + r.url() + "  " + r.error()); }
        }
        System.out.printf("%n成功 %d / 失败 %d / 耗时 %.2fs%n", ok, fail, seconds);
    }
}

实际进度顺序因线程调度而异,但 c-fail.txt 一定会有三次(重试三次)进度输出。

4. 分块 Playground:拆解每一步

4.1 基础:Future 拿到结果

Future 收集结果
import java.util.*;
import java.util.concurrent.*;
 
public class Main {
    public static void main(String[] args) throws Exception {
        ExecutorService pool = Executors.newFixedThreadPool(3);
        List<String> urls = List.of("a", "b", "c");
 
        List<Future<String>> futures = new ArrayList<>();
        for (String u : urls) {
            futures.add(pool.submit(() -> {
                Thread.sleep(100);
                return "downloaded " + u;
            }));
        }
 
        for (Future<String> f : futures) {
            System.out.println(f.get());    // 阻塞等到该任务完成
        }
        pool.shutdown();
    }
}

4.2 进度回调

进度回调
import java.util.concurrent.*;
 
public class Main {
    public static void main(String[] args) throws Exception {
        ExecutorService pool = Executors.newSingleThreadExecutor();
        Future<Integer> f = pool.submit(() -> {
            for (int p = 25; p <= 100; p += 25) {
                Thread.sleep(40);
                System.out.println("  progress: " + p + "%");
            }
            return 1024;
        });
        Integer bytes = f.get();
        pool.shutdown();
        System.out.println("done, bytes = " + bytes);
    }
}

4.3 失败重试

失败重试
import java.util.concurrent.atomic.AtomicInteger;
 
public class Main {
    public static void main(String[] args) throws Exception {
        AtomicInteger attempts = new AtomicInteger(0);
        int maxRetries = 3;
 
        for (int i = 1; i <= maxRetries; i++) {
            attempts.incrementAndGet();
            System.out.println("尝试第 " + attempts.get() + " 次");
            try {
                if (attempts.get() < 2) {
                    throw new RuntimeException("还没好");
                }
                System.out.println("成功!");
                break;
            } catch (Exception e) {
                System.out.println("失败: " + e.getMessage());
                if (i == maxRetries) {
                    System.out.println("超过最大重试次数,放弃");
                }
            }
        }
    }
}
💡真实场景的重试要点
  • 退避策略:指数退避 + 抖动(sleep(base * 2^attempt + random))
  • 区分可重试错误(5xx、超时)和不可重试错误(404、400)
  • 设上限 + 设总超时(避免无限重试把资源耗尽)
  • 配合断路器(Resilience4j、Hystrix)

5. 真实下载器对比

真实我们的迷你版
HttpClient 发 GETThread.sleep
InputStream.read 流式读固定 5 步
CompletableFuture 链式回调Future.get() 阻塞
RetryTemplate / Resilience4j手写 for 循环
进度用 Content-Length 计算按 step 推算百分比
写盘暂存为字符串

🎯 练习

限流下载
// 任务:在前面的下载器基础上,给 Downloader 增加
// 「最多同时进行 N 个任务」的能力。
// 提示:用 Semaphore(N),每个任务开始前 acquire,结束 release
// 用 3 个线程但同时只能跑 2 个任务(Semaphore = 2),
// 5 个任务中一个会失败
// 期望输出大致:
//   start A
//   start B
//   progress 100% for A
//   start C        (A 完成后)
//   progress 100% for B
//   start D
//   ...
 
import java.util.*;
import java.util.concurrent.*;
 
public class Main {
    record FileResult(String url, boolean ok) {}
 
    public static void main(String[] args) throws Exception {
        Semaphore gate = new Semaphore(2);   // 限流:最多并发 2 个
        ExecutorService pool = Executors.newFixedThreadPool(5);
 
        List<String> urls = List.of("A", "B", "C", "D", "E");
        List<Future<FileResult>> fs = new ArrayList<>();
 
        for (String u : urls) {
            fs.add(pool.submit(() -> {
                // 补全:
                // 1. gate.acquire()
                // 2. 打印 "start u"
                // 3. Thread.sleep(80)
                // 4. 打印 "done u"
                // 5. gate.release()
                // 6. 返回 FileResult(u, true)
                return null;
            }));
        }
        for (Future<FileResult> f : fs) f.get();
        pool.shutdown();
    }
}

小结

  • ✅ ExecutorService + Future 是并发下载的最小骨架
  • ✅ 进度回调用接口注入;线程安全用 Collections.synchronizedList 或 CopyOnWriteArrayList
  • ✅ 重试要带退避和上限;区分可重试 vs 不可重试错误
  • ✅ 真实项目优先用 CompletableFuture(链式、可组合)替代裸 Future
  • ✅ 限流用 Semaphore 或专门的 RateLimiter
  • ✅ 框架层面:Resilience4j(重试/熔断/限流)、OkHttp(HTTP 客户端)

🎉 Java 进阶 10 章完结!到这里你已经掌握:

  • 继承多态 / 泛型 / 注解 / 枚举
  • 多线程 / 反射 / JVM / GC / Maven 测试
  • 两个真实项目:HTTP 服务 + 并发下载器

接下来推荐学习方向:

  • Spring 全家桶:Spring Data JPA、Spring Security、Spring Cloud
  • JVM 调优:GC 日志分析、arthas 在线诊断
  • 并发工具:java.util.concurrent 高级组件(CompletableFuture / StampedLock / Phaser)
  • 设计模式:策略、状态、装饰器在 Java 里的标准实现