18并发编程综合应用
经过前面系列的学习,我们已经掌握了线程、锁、线程池、JUC 工具、并发集合等知识。作为应用篇,也是本系列的收尾篇,本文不再纠结底层原理,而是回答一个工程问题:什么场景下该用什么并发手段。全篇围绕七类典型应用展开——效率(利用多核)、限制(限流与资源控制)、互斥、同步与异步、缓存一致性、分治(Fork/Join)、统筹协调与定时任务。
使用多线程充分利用 CPU
环境搭建
基准测试工具选择比较靠谱的 JMH(Java Microbenchmark Harness),它会执行程序预热、多次测试并取平均,避免手写计时带来的误差。
CPU 核数限制有两种思路:
- 使用虚拟机,分配合适的核;
- 使用
msconfig分配合适的核,需要重启,比较麻烦。
并行计算方式的选择:
- 最初想直接使用 parallel stream,后来发现它有自己的问题(公共线程池共享、不可控);
- 改为自己手动控制 thread,实现简单的并行计算。
使用 JMH 官方 archetype 生成基准测试工程:
mvn archetype:generate -DinteractiveMode=false \
-DarchetypeGroupId=org.openjdk.jmh \
-DarchetypeArtifactId=jmh-java-benchmark-archetype \
-DgroupId=org.sample -DartifactId=test -Dversion=1.0测试代码:对一个长度为 1 亿的数组求和。方法 c() 把数组拆成 4 段,交给 4 个线程并行计算;方法 d() 只开 1 个线程串行求和,作为对照组。
package org.sample;
import java.util.Arrays;
import java.util.concurrent.FutureTask;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Mode;
import org.openjdk.jmh.annotations.Warmup;
@Fork(1)
@BenchmarkMode(Mode.AverageTime)
@Warmup(iterations = 3)
@Measurement(iterations = 5)
public class MyBenchmark {
static int[] ARRAY = new int[1000_000_00];
static {
Arrays.fill(ARRAY, 1);
}
@Benchmark
public int c() throws Exception {
int[] array = ARRAY;
FutureTask<Integer> t1 = new FutureTask<>(() -> {
int sum = 0;
for (int i = 0; i < 250_000_00; i++) {
sum += array[0 + i];
}
return sum;
});
FutureTask<Integer> t2 = new FutureTask<>(() -> {
int sum = 0;
for (int i = 0; i < 250_000_00; i++) {
sum += array[250_000_00 + i];
}
return sum;
});
FutureTask<Integer> t3 = new FutureTask<>(() -> {
int sum = 0;
for (int i = 0; i < 250_000_00; i++) {
sum += array[500_000_00 + i];
}
return sum;
});
FutureTask<Integer> t4 = new FutureTask<>(() -> {
int sum = 0;
for (int i = 0; i < 250_000_00; i++) {
sum += array[750_000_00 + i];
}
return sum;
});
new Thread(t1).start();
new Thread(t2).start();
new Thread(t3).start();
new Thread(t4).start();
return t1.get() + t2.get() + t3.get() + t4.get();
}
@Benchmark
public int d() throws Exception {
int[] array = ARRAY;
FutureTask<Integer> t1 = new FutureTask<>(() -> {
int sum = 0;
for (int i = 0; i < 1000_000_00; i++) {
sum += array[0 + i];
}
return sum;
});
new Thread(t1).start();
return t1.get();
}
}双核 CPU(4 个逻辑 CPU)
打包后执行:
java -jar target/benchmarks.jar并行方法 c() 的结果(预热 3 轮、度量 5 轮):
Result: 0.020 ±(99.9%) 0.001 s/op [Average]
Statistics: (min, avg, max) = (0.020, 0.020, 0.020), stdev = 0.000
Confidence interval (99.9%): [0.019, 0.021]串行方法 d() 的结果:
Result: 0.043 ±(99.9%) 0.003 s/op [Average]
Statistics: (min, avg, max) = (0.042, 0.043, 0.044), stdev = 0.001
Confidence interval (99.9%): [0.040, 0.045]汇总:
| Benchmark | Mode | Samples | Score (s/op) | 误差 |
|---|---|---|---|---|
| MyBenchmark.c(4 线程并行) | avgt | 5 | 0.020 | 0.001 |
| MyBenchmark.d(单线程串行) | avgt | 5 | 0.043 | 0.003 |
可以看到,多核下效率提升非常明显,并行版本快了一倍左右。这类纯 CPU 密集型、且数据可均匀切分的任务,正是多线程最能发挥价值的场景。
单核 CPU
同样的测试放到单核环境下:
| Benchmark | Mode | Samples | Score (s/op) | 误差 |
|---|---|---|---|---|
| MyBenchmark.c(4 线程并行) | avgt | 5 | 0.061 | 0.060 |
| MyBenchmark.d(单线程串行) | avgt | 5 | 0.064 | 0.071 |
两者性能几乎一样,并行版本甚至因为线程切换略有波动。没有多余的 CPU 核心可供利用时,拆线程并不会带来加速——线程不是银弹,先确认瓶颈是不是「CPU 没吃饱」。
限制对资源的使用
并发不是一味地「快」,很多时候还要「稳」:高峰期限流、保护有限资源,防止系统被压垮。
限制对 CPU 的使用
在没有利用 CPU 来计算时,不要让 while(true) 空转浪费 CPU,这时可以使用 yield 或 sleep 把 CPU 的使用权让给其他程序。
sleep 实现,适用于无需锁同步的场景:
while (true) {
try {
Thread.sleep(50);
} catch (InterruptedException e) {
e.printStackTrace();
}
}也可以用 wait 或条件变量达到类似的效果。不同的是,后两种都需要加锁,并且需要相应的唤醒操作,一般适用于需要进行同步的场景;sleep 适用于无需锁同步的场景。
wait 实现:
synchronized (锁对象) {
while (条件不满足) {
try {
锁对象.wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
// do sth...
}条件变量实现:
lock.lock();
try {
while (条件不满足) {
try {
条件变量.await();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
// do sth...
} finally {
lock.unlock();
}限制对共享资源的使用:Semaphore
使用 Semaphore 限流:在访问高峰期让请求线程阻塞,高峰期过去后再释放许可。需要注意:
- 它只适合限制单机线程数量;
- 它限制的仅是线程数,而不是资源数(例如连接数,可对比 Tomcat LimitLatch 的实现)。
下面用 Semaphore 实现一个简单连接池。对比「享元模式」下用 wait/notify 的实现,它的性能和可读性显然更好。注意实现中许可数与连接(资源)数保持相等。
@Slf4j(topic = "c.Pool")
class Pool {
// 1. 连接池大小
private final int poolSize;
// 2. 连接对象数组
private Connection[] connections;
// 3. 连接状态数组:0 表示空闲,1 表示繁忙
private AtomicIntegerArray states;
private Semaphore semaphore;
// 4. 构造方法初始化
public Pool(int poolSize) {
this.poolSize = poolSize;
// 让许可数与资源数一致
this.semaphore = new Semaphore(poolSize);
this.connections = new Connection[poolSize];
this.states = new AtomicIntegerArray(new int[poolSize]);
for (int i = 0; i < poolSize; i++) {
connections[i] = new MockConnection("连接" + (i + 1));
}
}
// 5. 借连接
public Connection borrow() {
// 获取许可
try {
semaphore.acquire(); // 没有许可的线程,在此等待
} catch (InterruptedException e) {
e.printStackTrace();
}
for (int i = 0; i < poolSize; i++) {
// 获取空闲连接
if (states.get(i) == 0) {
if (states.compareAndSet(i, 0, 1)) {
log.debug("borrow {}", connections[i]);
return connections[i];
}
}
}
// 不会执行到这里
return null;
}
// 6. 归还连接
public void free(Connection conn) {
for (int i = 0; i < poolSize; i++) {
if (connections[i] == conn) {
states.set(i, 0);
log.debug("free {}", conn);
semaphore.release();
break;
}
}
}
}单位时间内限流:Guava RateLimiter
Semaphore 控制的是「同时有多少个线程在跑」,而有些场景需要控制「单位时间内允许多少个请求」,即限流(QPS)。Guava 的 RateLimiter 基于令牌桶实现,一行代码即可完成。
@RestController
public class TestController {
private RateLimiter limiter = RateLimiter.create(50);
@GetMapping("/test")
public String test() {
limiter.acquire();
return "ok";
}
}用 ApacheBench 压测 10 秒、并发 10:
ab -c 10 -t 10 http://localhost:8080/test没有限流之前(注释掉 limiter.acquire()),10 秒处理了 24706 个请求,QPS 约 2469:
Concurrency Level: 10
Time taken for tests: 10.005 seconds
Complete requests: 24706
Failed requests: 0
Requests per second: 2469.42 [#/sec] (mean)
Time per request: 4.050 [ms] (mean)限流之后(每秒 50 个许可),10 秒只处理了 545 个请求,QPS 被稳定压在 54 左右,超出的请求排队等待:
Concurrency Level: 10
Time taken for tests: 10.007 seconds
Complete requests: 545
Failed requests: 0
Requests per second: 54.46 [#/sec] (mean)
Time per request: 183.621 [ms] (mean)代价是平均响应时间从 4ms 上升到 183ms,但系统始终处于可控范围,不会被突发流量击垮。
互斥
当多个线程访问同一份共享可变数据时,必须保证操作的正确性。思路有两种:悲观互斥与乐观重试。
悲观互斥
互斥实际是悲观锁的思想:总假设并发冲突会发生,先加锁把共享资源保护起来,同一时刻只允许一个线程进入临界区。以取款需求为例,先定义账户接口与测试方法:
interface Account {
// 获取余额
Integer getBalance();
// 取款
void withdraw(Integer amount);
/**
* 方法内会启动 1000 个线程,每个线程做 -10 元 的操作
* 如果初始余额为 10000 那么正确的结果应当是 0
*/
static void demo(Account account) {
List<Thread> ts = new ArrayList<>();
for (int i = 0; i < 1000; i++) {
ts.add(new Thread(() -> {
account.withdraw(10);
}));
}
long start = System.nanoTime();
ts.forEach(Thread::start);
ts.forEach(t -> {
try {
t.join();
} catch (InterruptedException e) {
e.printStackTrace();
}
});
long end = System.nanoTime();
System.out.println(account.getBalance()
+ " cost: " + (end - start) / 1000_000 + " ms");
}
}用 synchronized 保护余额的读和写:
class AccountSync implements Account {
private Integer balance;
public AccountSync(Integer balance) {
this.balance = balance;
}
@Override
public Integer getBalance() {
synchronized (this) {
return this.balance;
}
}
@Override
public void withdraw(Integer amount) {
synchronized (this) {
this.balance -= amount;
}
}
}乐观重试
另一种是乐观锁思想,它其实不是互斥:不加锁,先尝试修改,如果修改期间没人动过数据就成功,否则重试。底层依赖 CAS(Compare And Swap)。
class AccountCas implements Account {
private AtomicInteger balance;
public AccountCas(int balance) {
this.balance = new AtomicInteger(balance);
}
@Override
public Integer getBalance() {
return balance.get();
}
@Override
public void withdraw(Integer amount) {
while (true) {
// 获取余额的最新值
int prev = balance.get();
// 要修改的余额
int next = prev - amount;
// 真正修改
if (balance.compareAndSet(prev, next)) {
break;
}
}
}
}CAS 的典型套路是「读旧值 → 算新值 → compareAndSet」三步,失败就自旋重试。冲突激烈时重试会变多;数据库中的版本号机制(update ... where version = ?)也是同样的思路。悲观锁适合冲突多、临界区长的场景,乐观重试适合冲突少、操作短的场景。
同步与异步
调用一个耗时操作(如查库、读文件、远程调用)后,调用方要不要等它的结果?这决定了该用同步还是异步。
需要等待结果
这时既可以用同步处理,也可以用异步处理,关键看结果由谁来接收。
join 实现(同步)
static int result = 0;
private static void test1() throws InterruptedException {
log.debug("开始");
Thread t1 = new Thread(() -> {
log.debug("开始");
sleep(1);
log.debug("结束");
result = 10;
}, "t1");
t1.start();
t1.join();
log.debug("结果为:{}", result);
}输出:
20:30:40.453 [main] c.TestJoin - 开始
20:30:40.541 [Thread-0] c.TestJoin - 开始
20:30:41.543 [Thread-0] c.TestJoin - 结束
20:30:41.551 [main] c.TestJoin - 结果为:10评价:
- 需要外部共享变量,不符合面向对象封装的思想;
- 必须等待线程结束,不能配合线程池使用。
Future 实现(同步)
FutureTask 本身既是任务又是结果凭证,规避了共享变量:
private static void test2() throws InterruptedException, ExecutionException {
log.debug("开始");
FutureTask<Integer> result = new FutureTask<>(() -> {
log.debug("开始");
sleep(1);
log.debug("结束");
return 10;
});
new Thread(result, "t1").start();
log.debug("结果为:{}", result.get());
}输出:
10:11:57.880 c.TestSync [main] - 开始
10:11:57.942 c.TestSync [t1] - 开始
10:11:58.943 c.TestSync [t1] - 结束
10:11:58.943 c.TestSync [main] - 结果为:10配合线程池使用
FutureTask 可以方便地配合线程池使用。线程池 submit 返回的 Future,其实际类型也是 FutureTask:
private static void test3() throws InterruptedException, ExecutionException {
ExecutorService service = Executors.newFixedThreadPool(1);
log.debug("开始");
Future<Integer> result = service.submit(() -> {
log.debug("开始");
sleep(1);
log.debug("结束");
return 10;
});
log.debug("结果为:{}, result 的类型:{}", result.get(), result.getClass());
service.shutdown();
}输出:
10:17:40.090 c.TestSync [main] - 开始
10:17:40.150 c.TestSync [pool-1-thread-1] - 开始
10:17:41.151 c.TestSync [pool-1-thread-1] - 结束
10:17:41.151 c.TestSync [main] - 结果为:10, result 的类型:class java.util.concurrent.FutureTask评价:结果仍然由 main 线程接收,get() 方法会让调用线程同步等待。
自定义实现(同步)
手写「一个线程等待另一个线程的结果」,即模式篇中的保护性暂停模式(Guarded Suspension):用一个不可变的结果对象配合 wait/notify,在结果未就绪时等待,就绪后被唤醒并返回。此处不再展开。
CompletableFuture 实现(异步)
前面几种方式,调用线程都要同步 get() 等待。CompletableFuture 可以让调用线程异步处理结果——计算与结果回调分别派发给不同线程池:
private static void test4() {
// 进行计算的线程池
ExecutorService computeService = Executors.newFixedThreadPool(1);
// 接收结果的线程池
ExecutorService resultService = Executors.newFixedThreadPool(1);
log.debug("开始");
CompletableFuture.supplyAsync(() -> {
log.debug("开始");
sleep(1);
log.debug("结束");
return 10;
}, computeService).thenAcceptAsync((result) -> {
log.debug("结果为:{}", result);
}, resultService);
}输出:
10:36:28.114 c.TestSync [main] - 开始
10:36:28.164 c.TestSync [pool-1-thread-1] - 开始
10:36:29.165 c.TestSync [pool-1-thread-1] - 结束
10:36:29.165 c.TestSync [pool-2-thread-1] - 结果为:10评价:
- 可以让调用线程异步处理结果,实际是其他线程去同步等待;
- 可以方便地分离不同职责的线程池;
- 以任务为中心,而不是以线程为中心。
BlockingQueue 实现(异步)
生产者把结果放入队列,消费者从队列取结果,SynchronousQueue 不存储元素,放与取直接握手:
private static void test6() {
ExecutorService consumer = Executors.newFixedThreadPool(1);
ExecutorService producer = Executors.newFixedThreadPool(1);
BlockingQueue<Integer> queue = new SynchronousQueue<>();
log.debug("开始");
producer.submit(() -> {
log.debug("开始");
sleep(1);
log.debug("结束");
try {
queue.put(10);
} catch (InterruptedException e) {
e.printStackTrace();
}
});
consumer.submit(() -> {
try {
Integer result = queue.take();
log.debug("结果为:{}", result);
} catch (InterruptedException e) {
e.printStackTrace();
}
});
}不需等待结果
这时最好使用异步处理:调用方提交任务后立刻返回,去做自己的事。以读取一个大文件为例:
@Slf4j(topic = "c.FileReader")
public class FileReader {
public static void read(String filename) {
int idx = filename.lastIndexOf(File.separator);
String shortName = filename.substring(idx + 1);
try (FileInputStream in = new FileInputStream(filename)) {
long start = System.currentTimeMillis();
log.debug("read [{}] start ...", shortName);
byte[] buf = new byte[1024];
int n = -1;
do {
n = in.read(buf);
} while (n != -1);
long end = System.currentTimeMillis();
log.debug("read [{}] end ... cost: {} ms", shortName, end - start);
} catch (IOException e) {
e.printStackTrace();
}
}
}没有用线程时,方法调用是同步的,main 读完文件(耗时约 4 秒)后才能打印「do other things」:
@Slf4j(topic = "c.Sync")
public class Sync {
public static void main(String[] args) {
String fullPath = "E:\\1.mp4";
FileReader.read(fullPath);
log.debug("do other things ...");
}
}18:39:15 [main] c.FileReader - read [1.mp4] start ...
18:39:19 [main] c.FileReader - read [1.mp4] end ... cost: 4090 ms
18:39:19 [main] c.Sync - do other things ...普通线程实现,提交任务后立即返回:
private static void test1() {
new Thread(() -> FileReader.read(Constants.MP4_FULL_PATH)).start();
log.debug("do other things ...");
}18:41:53 [main] c.Async - do other things ...
18:41:53 [Thread-0] c.FileReader - read [1.mp4] start ...
18:41:57 [Thread-0] c.FileReader - read [1.mp4] end ... cost: 4197 ms线程池实现,复用线程、避免频繁创建销毁:
private static void test2() {
ExecutorService service = Executors.newFixedThreadPool(1);
service.execute(() -> FileReader.read(Constants.MP4_FULL_PATH));
log.debug("do other things ...");
service.shutdown();
}CompletableFuture 实现,默认使用 ForkJoinPool 公共池:
private static void test3() throws IOException {
CompletableFuture.runAsync(() -> FileReader.read(Constants.MP4_FULL_PATH));
log.debug("do other things ...");
System.in.read();
}11:09:38.145 c.TestAsyc [main] - do other things ...
11:09:38.145 c.FileReader [ForkJoinPool.commonPool-worker-1] - read [1.mp4] start ...
11:09:40.514 c.FileReader [ForkJoinPool.commonPool-worker-1] - read [1.mp4] end ... cost: 2369 ms小结:要结果用 join/Future(同步)或 CompletableFuture 回调(异步);不要结果直接丢给线程或线程池异步执行即可。
缓存
缓存更新策略
引入缓存后,更新数据时到底是「先清缓存」还是「先更新数据库」?两种顺序在并发下都可能产生数据不一致。
先清缓存的风险时序:
- 线程 A 清空缓存;
- 线程 B 查询缓存未命中,查询数据库得到旧值(x=1);
- A 将查询结果放入缓存(x=1);
- A 将新数据存入库(x=2);
- 后续查询将一直读到旧值(x=1)。
先更新数据库的风险时序:
- 线程 A 将新数据存入库(x=2);
- 线程 B 查询缓存命中旧值(x=1);
- A 清空缓存;
- 后续查询数据库可以得到新值(x=2)。
对比可见,「先更新数据库、再删缓存」的最坏情况只是其他线程短时间内读到旧值,缓存一旦被删就能恢复,因此实践中一般采用这种方式。
再补充一种极端情况:查询线程 A 查询数据时,恰好缓存由于时间到期失效(或第一次查询)——
- A 缓存没有命中,查询数据库得到旧值(x=1);
- B 将新数据存入库(x=2);
- B 清空缓存;
- A 将查询结果放入缓存(x=1);
- 后续查询将一直是旧值(x=1)。
这种情况出现的几率非常小(要求读比写慢、且写恰好夹在读库与回填之间),相关分析可参见 Facebook 关于缓存一致性的论文。
读写锁实现一致性缓存
下面使用读写锁实现一个简单的按需加载缓存:读操作加读锁,允许多线程并发读;更新操作加写锁,独占并清空缓存;缓存未命中时再升级为写锁查库回填(双重检查,避免重复查库)。
class GenericCachedDao<T> {
// HashMap 作为缓存非线程安全, 需要保护
HashMap<SqlPair, T> map = new HashMap<>();
ReentrantReadWriteLock lock = new ReentrantReadWriteLock();
GenericDao genericDao = new GenericDao();
public int update(String sql, Object... params) {
SqlPair key = new SqlPair(sql, params);
// 加写锁, 防止其它线程对缓存读取和更改
lock.writeLock().lock();
try {
int rows = genericDao.update(sql, params);
map.clear();
return rows;
} finally {
lock.writeLock().unlock();
}
}
public T queryOne(Class<T> beanClass, String sql, Object... params) {
SqlPair key = new SqlPair(sql, params);
// 加读锁, 防止其它线程对缓存更改
lock.readLock().lock();
try {
T value = map.get(key);
if (value != null) {
return value;
}
} finally {
lock.readLock().unlock();
}
// 加写锁, 防止其它线程对缓存读取和更改
lock.writeLock().lock();
try {
// get 方法上面部分是可能多个线程进来的, 可能已经向缓存填充了数据
// 为防止重复查询数据库, 再次验证
T value = map.get(key);
if (value == null) {
// 如果没有, 查询数据库
value = genericDao.queryOne(beanClass, sql, params);
map.put(key, value);
}
return value;
} finally {
lock.writeLock().unlock();
}
}
// 作为 key 保证其是不可变的
class SqlPair {
private String sql;
private Object[] params;
public SqlPair(String sql, Object[] params) {
this.sql = sql;
this.params = params;
}
@Override
public boolean equals(Object o) {
if (this == o) {
return true;
}
if (o == null || getClass() != o.getClass()) {
return false;
}
SqlPair sqlPair = (SqlPair) o;
return sql.equals(sqlPair.sql) &&
Arrays.equals(params, sqlPair.params);
}
@Override
public int hashCode() {
int result = Objects.hash(sql);
result = 31 * result + Arrays.hashCode(params);
return result;
}
}
}注意:以上实现体现的是读写锁的应用,能保证缓存和数据库的一致性,但还有以下问题没有考虑:
- 适合读多写少,如果写操作比较频繁,以上实现性能低;
- 没有考虑缓存容量;
- 没有考虑缓存过期;
- 只适合单机;
- 并发性还是低,目前只会用一把锁;
- 更新方法太过简单粗暴,清空了所有 key(可考虑按类型分区或重新设计 key);
- 乐观锁实现:用 CAS 去更新。
分治:Fork/Join
分治思想:把一个大任务拆成若干小任务并行执行,再把各小任务的结果合并。JDK 中的 Fork/Join 框架就是它的载体——RecursiveTask(有返回值)负责 fork(拆分)与 join(合并),ForkJoinPool 负责调度。
案例:单词计数
26 个线程分别读取 26 个文件并统计单词次数。先给出通用的统计框架:
private static <V> void demo(Supplier<Map<String, V>> supplier,
BiConsumer<Map<String, V>, List<String>> consumer) {
Map<String, V> counterMap = supplier.get();
List<Thread> ts = new ArrayList<>();
for (int i = 1; i <= 26; i++) {
int idx = i;
Thread thread = new Thread(() -> {
List<String> words = readFromFile(idx);
consumer.accept(counterMap, words);
});
ts.add(thread);
}
ts.forEach(t -> t.start());
ts.forEach(t -> {
try {
t.join();
} catch (InterruptedException e) {
e.printStackTrace();
}
});
System.out.println(counterMap);
}
public static List<String> readFromFile(int i) {
ArrayList<String> words = new ArrayList<>();
try (BufferedReader in = new BufferedReader(new InputStreamReader(
new FileInputStream("tmp/" + i + ".txt")))) {
while (true) {
String word = in.readLine();
if (word == null) {
break;
}
words.add(word);
}
return words;
} catch (IOException e) {
throw new RuntimeException(e);
}
}解法 1:ConcurrentHashMap + LongAdder,用 computeIfAbsent 保证计数器只创建一次,LongAdder 保证高并发自增的性能:
demo(
() -> new ConcurrentHashMap<String, LongAdder>(),
(map, words) -> {
for (String word : words) {
map.computeIfAbsent(word, (key) -> new LongAdder()).increment();
}
}
);解法 2:parallel stream 一把梭,groupingBy 分组、summingInt 求和,函数式风格更简洁:
Map<String, Integer> collect = IntStream.range(1, 27).parallel()
.mapToObj(idx -> readFromFile(idx))
.flatMap(list -> list.stream())
.collect(Collectors.groupingBy(Function.identity(),
Collectors.summingInt(w -> 1)));案例:求和
对 1~10 求和。定义 RecursiveTask:当区间只剩 1 个或 2 个数时直接计算(终止条件),否则从中间一分为二,两个子任务分别 fork(),再 join() 合并结果。
class AddTask3 extends RecursiveTask<Integer> {
int begin;
int end;
public AddTask3(int begin, int end) {
this.begin = begin;
this.end = end;
}
@Override
public String toString() {
return "{" + begin + "," + end + '}';
}
@Override
protected Integer compute() {
// 5, 5
if (begin == end) {
log.debug("join() {}", begin);
return begin;
}
// 4, 5
if (end - begin == 1) {
log.debug("join() {} + {} = {}", begin, end, end + begin);
return end + begin;
}
// 1 5
int mid = (end + begin) / 2; // 3
AddTask3 t1 = new AddTask3(begin, mid); // 1,3
t1.fork();
AddTask3 t2 = new AddTask3(mid + 1, end); // 4,5
t2.fork();
log.debug("fork() {} + {} = ?", t1, t2);
int result = t1.join() + t2.join();
log.debug("join() {} + {} = {}", t1, t2, result);
return result;
}
}然后提交给 ForkJoinPool 执行:
public static void main(String[] args) {
ForkJoinPool pool = new ForkJoinPool(4);
System.out.println(pool.invoke(new AddTask3(1, 10)));
}任务拆分树为 {1,10} → {1,5} + {6,10} → …… 一直拆到叶子节点并行计算、逐层向上合并。运行输出:
[ForkJoinPool-1-worker-0] - join() 1 + 2 = 3
[ForkJoinPool-1-worker-3] - join() 4 + 5 = 9
[ForkJoinPool-1-worker-0] - join() 3
[ForkJoinPool-1-worker-1] - fork() {1,3} + {4,5} = ?
[ForkJoinPool-1-worker-2] - fork() {1,2} + {3,3} = ?
[ForkJoinPool-1-worker-2] - join() {1,2} + {3,3} = 6
[ForkJoinPool-1-worker-1] - join() {1,3} + {4,5} = 15
15统筹:烧水泡茶
并发编程不只有「快慢」,还有「先后」。烧水泡茶是经典的工序统筹题:老王洗水壶(1 分钟)、烧开水(15 分钟);小王洗茶壶(1 分钟)、洗茶杯(2 分钟)、拿茶叶(1 分钟);水开且茶叶到位后才能泡茶。两条工序线并行,最后汇合。
解法 1:join
小王完成自己的工序后,join 等待老王把水烧开,再泡茶:
Thread t1 = new Thread(() -> {
log.debug("洗水壶");
sleep(1);
log.debug("烧开水");
sleep(15);
}, "老王");
Thread t2 = new Thread(() -> {
log.debug("洗茶壶");
sleep(1);
log.debug("洗茶杯");
sleep(2);
log.debug("拿茶叶");
sleep(1);
try {
t1.join();
} catch (InterruptedException e) {
e.printStackTrace();
}
log.debug("泡茶");
}, "小王");
t1.start();
t2.start();输出:
19:19:37.547 [小王] c.TestMakeTea - 洗茶壶
19:19:37.547 [老王] c.TestMakeTea - 洗水壶
19:19:38.552 [小王] c.TestMakeTea - 洗茶杯
19:19:38.552 [老王] c.TestMakeTea - 烧开水
19:19:40.553 [小王] c.TestMakeTea - 拿茶叶
19:19:53.553 [小王] c.TestMakeTea - 泡茶缺陷:
- 上面模拟的是小王等老王的水烧开了由小王泡茶,如果反过来要实现老王等小王的茶叶拿来了由老王泡茶,代码最好能适应两种情况;
- 两个线程其实是各执行各的,如果要模拟老王把水壶交给小王泡茶,或小王把茶叶交给老王泡茶,join 难以表达。
解法 2:wait/notify
两个线程互相等待对方的中间产物(开水与花茶),用同一个锁对象配合 wait/notifyAll 协调,谁最后发现两个条件都满足谁泡茶:
class S2 {
static String kettle = "冷水";
static String tea = null;
static final Object lock = new Object();
static boolean maked = false;
public static void makeTea() {
new Thread(() -> {
log.debug("洗水壶");
sleep(1);
log.debug("烧开水");
sleep(5);
synchronized (lock) {
kettle = "开水";
lock.notifyAll();
while (tea == null) {
try {
lock.wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
if (!maked) {
log.debug("拿({})泡({})", kettle, tea);
maked = true;
}
}
}, "老王").start();
new Thread(() -> {
log.debug("洗茶壶");
sleep(1);
log.debug("洗茶杯");
sleep(2);
log.debug("拿茶叶");
sleep(1);
synchronized (lock) {
tea = "花茶";
lock.notifyAll();
while (kettle.equals("冷水")) {
try {
lock.wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
if (!maked) {
log.debug("拿({})泡({})", kettle, tea);
maked = true;
}
}
}, "小王").start();
}
}输出:
20:04:48.179 c.S2 [小王] - 洗茶壶
20:04:48.179 c.S2 [老王] - 洗水壶
20:04:49.185 c.S2 [老王] - 烧开水
20:04:49.185 c.S2 [小王] - 洗茶杯
20:04:51.185 c.S2 [小王] - 拿茶叶
20:04:54.185 c.S2 [老王] - 拿(开水)泡(花茶)解法 2 解决了解法 1 的问题,不过老王和小王需要相互等待。不如让他们只负责各自的任务,泡茶交给第三人来做。
解法 3:第三者协调
老王只管烧好水、小王只管拿茶叶,完成后通知;「王夫人」等待两个条件都齐备,统一负责泡茶,职责最清晰:
class S3 {
static String kettle = "冷水";
static String tea = null;
static final Object lock = new Object();
public static void makeTea() {
new Thread(() -> {
log.debug("洗水壶");
sleep(1);
log.debug("烧开水");
sleep(5);
synchronized (lock) {
kettle = "开水";
lock.notifyAll();
}
}, "老王").start();
new Thread(() -> {
log.debug("洗茶壶");
sleep(1);
log.debug("洗茶杯");
sleep(2);
log.debug("拿茶叶");
sleep(1);
synchronized (lock) {
tea = "花茶";
lock.notifyAll();
}
}, "小王").start();
new Thread(() -> {
synchronized (lock) {
while (kettle.equals("冷水") || tea == null) {
try {
lock.wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
log.debug("拿({})泡({})", kettle, tea);
}
}, "王夫人").start();
}
}输出:
20:13:18.202 c.S3 [小王] - 洗茶壶
20:13:18.202 c.S3 [老王] - 洗水壶
20:13:19.206 c.S3 [小王] - 洗茶杯
20:13:19.206 c.S3 [老王] - 烧开水
20:13:21.206 c.S3 [小王] - 拿茶叶
20:13:24.207 c.S3 [王夫人] - 拿(开水)泡(花茶)三个版本演进的背后,正是把「线程间的等待与协作」从硬编码的 join,抽象成条件等待,再抽象成独立协调者的过程。
定时任务
定期执行
如何让任务在每周四 18:00:00 定时执行?思路:先算出最近一个周四 18:00 的时间点作为首次延时(如果今天周四且已过 18 点,则顺延到下周),再以一周为固定周期调度。
// 获得当前时间
LocalDateTime now = LocalDateTime.now();
// 获取本周四 18:00:00.000
LocalDateTime thursday =
now.with(DayOfWeek.THURSDAY).withHour(18).withMinute(0).withSecond(0).withNano(0);
// 如果当前时间已经超过本周四 18:00:00.000,那么找下周四 18:00:00.000
if (now.compareTo(thursday) >= 0) {
thursday = thursday.plusWeeks(1);
}
// 计算时间差,即延时执行时间
long initialDelay = Duration.between(now, thursday).toMillis();
// 计算间隔时间,即 1 周的毫秒值
long oneWeek = 7 * 24 * 3600 * 1000;
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
System.out.println("开始时间:" + new Date());
executor.scheduleAtFixedRate(() -> {
System.out.println("执行时间:" + new Date());
}, initialDelay, oneWeek, TimeUnit.MILLISECONDS);scheduleAtFixedRate 按固定速率触发,适合这类与日历对齐的周期任务;更复杂的 cron 表达式调度则可以交给 Quartz、Spring @Scheduled 或分布式调度框架。
写在最后
至此,JUC 并发编程系列就全部结束了。回头看这条学习路线:我们从线程与进程的基本概念起步,理解了 synchronized 与 Monitor、wait/notify、JMM 内存模型(可见性、有序性、原子性);然后掌握了 CAS 与原子类、各种锁(ReentrantLock、读写锁、StampedLock)、AQS 与 JUC 工具类、线程池原理与参数设计、不可变设计与线程安全集合、CompletableFuture 与 Fork/Join;最后在本篇把这些知识放回真实场景——要提速就利用多核算 CPU 密集任务,要保护系统就用 Semaphore、RateLimiter 限流,要保证正确就用悲观锁或乐观 CAS,要解耦耗时操作就按「要不要结果」选择 Future、回调或线程池,要保证缓存一致就用读写锁加双重检查,要处理大任务就分治拆细,要协调多道工序就用条件等待,要周期执行就交给调度线程池。
并发工具会不断更新,但背后的主题始终不变:安全、性能与可维护性之间的权衡。愿你在日后的开发中,既能写出跑得更快的程序,也能守住多线程下那份最难的正确性。