项目:多线程下载器
本章我们用模拟数据做一个迷你多线程下载器——把第 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 拿到结果
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 发 GET | Thread.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 里的标准实现