返回文章列表
JUC并发编程
JUC并发应用CompletableFutureForkJoin限流

18并发编程综合应用

经过前面系列的学习,我们已经掌握了线程、锁、线程池、JUC 工具、并发集合等知识。作为应用篇,也是本系列的收尾篇,本文不再纠结底层原理,而是回答一个工程问题:什么场景下该用什么并发手段。全篇围绕七类典型应用展开——效率(利用多核)、限制(限流与资源控制)、互斥、同步与异步、缓存一致性、分治(Fork/Join)、统筹协调与定时任务。

使用多线程充分利用 CPU

环境搭建

基准测试工具选择比较靠谱的 JMH(Java Microbenchmark Harness),它会执行程序预热、多次测试并取平均,避免手写计时带来的误差。

CPU 核数限制有两种思路:

  1. 使用虚拟机,分配合适的核;
  2. 使用 msconfig 分配合适的核,需要重启,比较麻烦。

并行计算方式的选择:

  1. 最初想直接使用 parallel stream,后来发现它有自己的问题(公共线程池共享、不可控);
  2. 改为自己手动控制 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 回调(异步);不要结果直接丢给线程或线程池异步执行即可。

缓存

缓存更新策略

引入缓存后,更新数据时到底是「先清缓存」还是「先更新数据库」?两种顺序在并发下都可能产生数据不一致。

先清缓存的风险时序:

  1. 线程 A 清空缓存;
  2. 线程 B 查询缓存未命中,查询数据库得到旧值(x=1);
  3. A 将查询结果放入缓存(x=1);
  4. A 将新数据存入库(x=2);
  5. 后续查询将一直读到旧值(x=1)。

先更新数据库的风险时序:

  1. 线程 A 将新数据存入库(x=2);
  2. 线程 B 查询缓存命中旧值(x=1);
  3. A 清空缓存;
  4. 后续查询数据库可以得到新值(x=2)。

对比可见,「先更新数据库、再删缓存」的最坏情况只是其他线程短时间内读到旧值,缓存一旦被删就能恢复,因此实践中一般采用这种方式。

再补充一种极端情况:查询线程 A 查询数据时,恰好缓存由于时间到期失效(或第一次查询)——

  1. A 缓存没有命中,查询数据库得到旧值(x=1);
  2. B 将新数据存入库(x=2);
  3. B 清空缓存;
  4. A 将查询结果放入缓存(x=1);
  5. 后续查询将一直是旧值(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、回调或线程池,要保证缓存一致就用读写锁加双重检查,要处理大任务就分治拆细,要协调多道工序就用条件等待,要周期执行就交给调度线程池。

并发工具会不断更新,但背后的主题始终不变:安全、性能与可维护性之间的权衡。愿你在日后的开发中,既能写出跑得更快的程序,也能守住多线程下那份最难的正确性。