Learn
Java/13-concurrency

多线程与并发

现代 CPU 都是多核,业务里「同时做多件事」几乎必学。本章从 Thread 开始,到线程池与 Future,再到同步原语与可见性。

注意:本章 Playground 演示用的线程数都很小(2~4),主线程会 join() 等所有子线程结束再退出,避免「主线程跑完程序就停」的尴尬。

1. 两种创建线程的方式

Thread 与 Runnable
public class Main {
    public static void main(String[] args) throws InterruptedException {
        // 方式 1:继承 Thread
        Thread t1 = new Thread("worker-1") {
            @Override public void run() {
                for (int i = 1; i <= 3; i++) {
                    System.out.println(getName() + " tick " + i);
                }
            }
        };
 
        // 方式 2:实现 Runnable(更推荐,可继承其他类)
        Runnable task = () -> {
            Thread cur = Thread.currentThread();
            for (int i = 1; i <= 3; i++) {
                System.out.println(cur.getName() + " step " + i);
            }
        };
        Thread t2 = new Thread(task, "worker-2");
 
        t1.start();   // 注意:start() 才真正开新线程
        t2.start();
        t1.join();    // 等待 t1 结束
        t2.join();    // 等待 t2 结束
        System.out.println("done");
    }
}
⚠️start() vs run()
  • t.start():启动新线程,由 JVM 调度执行 run()
  • t.run():只是在当前线程调用一个普通方法,根本没并发

2. 线程生命周期

        start()
   NEW ─────────► RUNNABLE ◄────► BLOCKED/WAITING/TIMED_WAITING
                    │                          │
                    └────── run() 结束 ─────────┴───► TERMINATED
状态触发
NEW创建但未 start()
RUNNABLE正在跑或就绪等 CPU
BLOCKED抢不到 monitor 锁(synchronized)
WAITINGwait() / join() / LockSupport.park()
TIMED_WAITINGsleep(ms) / wait(ms)
TERMINATEDrun() 执行完毕或抛未捕获异常

3. synchronized 关键字

多个线程同时改共享变量会竞态。synchronized 给对象/方法加互斥锁。

synchronized 解决竞态
public class Main {
    static class Counter {
        private int n = 0;
        public synchronized void inc() { n++; }   // 等价于 synchronized(this)
        public synchronized int get() { return n; }
    }
 
    public static void main(String[] args) throws InterruptedException {
        Counter c = new Counter();
        int threads = 4, loops = 10000;
 
        Thread[] ts = new Thread[threads];
        for (int i = 0; i < threads; i++) {
            ts[i] = new Thread(() -> {
                for (int j = 0; j < loops; j++) c.inc();
            });
            ts[i].start();
        }
        for (Thread t : ts) t.join();
 
        System.out.printf("期望 %d, 实际 %d%n", threads * loops, c.get());
    }
}
💡锁的是「对象」不是「代码」
  • synchronized 实例方法 → 锁 this
  • synchronized 静态方法 → 锁 Class 对象
  • synchronized(lockObj) 块 → 锁括号里的对象

4. volatile:轻量级可见性

volatile 解决可见性(一个线程的修改能被另一个线程立刻看到),但不解决原子性。

volatile 与可见性
public class Main {
    static volatile boolean running = true;
 
    public static void main(String[] args) throws InterruptedException {
        Thread worker = new Thread(() -> {
            long count = 0;
            while (running) {       // 没 volatile,主线程改 running 可能永远看不见
                count++;
            }
            System.out.println("worker stopped, count = " + count);
        });
        worker.start();
 
        Thread.sleep(50);            // 让 worker 跑一会儿
        running = false;             // 通知 worker 停止
        worker.join();
        System.out.println("main exit");
    }
}
⚠️volatile ≠ 原子

volatile int n = 0; 在并发 n++ 下仍然会丢更新——因为 n++ 是「读-改-写」三步。 需要原子性时用 AtomicInteger / synchronized。

5. wait / notify:生产-消费

wait() 释放锁并阻塞,notify() 唤醒同一个对象上等待的另一个线程。必须在 synchronized 块中调用。

wait / notify
import java.util.*;
 
public class Main {
    static class BoundedQueue {
        private final LinkedList<Integer> data = new LinkedList<>();
        private final int cap;
        BoundedQueue(int cap) { this.cap = cap; }
 
        synchronized void put(int v) throws InterruptedException {
            while (data.size() == cap) wait();     // 满 -> 等
            data.add(v);
            System.out.println("put " + v + " (size=" + data.size() + ")");
            notifyAll();
        }
        synchronized int take() throws InterruptedException {
            while (data.isEmpty()) wait();         // 空 -> 等
            int v = data.removeFirst();
            System.out.println("take " + v + " (size=" + data.size() + ")");
            notifyAll();
            return v;
        }
    }
 
    public static void main(String[] args) throws InterruptedException {
        BoundedQueue q = new BoundedQueue(2);
        Thread producer = new Thread(() -> {
            try { for (int i = 1; i <= 4; i++) { q.put(i); Thread.sleep(20); } }
            catch (InterruptedException ignored) {}
        });
        Thread consumer = new Thread(() -> {
            try { for (int i = 0; i < 4; i++) { q.take(); Thread.sleep(50); } }
            catch (InterruptedException ignored) {}
        });
        producer.start(); consumer.start();
        producer.join();  consumer.join();
    }
}

6. 线程池 ExecutorService

频繁 new Thread() 浪费资源。线程池复用线程,控制并发上限。

线程池
import java.util.concurrent.*;
 
public class Main {
    public static void main(String[] args) throws InterruptedException {
        // 固定大小线程池
        ExecutorService pool = Executors.newFixedThreadPool(3);
 
        for (int i = 1; i <= 6; i++) {
            final int id = i;
            pool.submit(() -> {
                String name = Thread.currentThread().getName();
                System.out.printf("[%s] task %d start%n", name, id);
                try { Thread.sleep(100); } catch (InterruptedException ignored) {}
                System.out.printf("[%s] task %d end%n", name, id);
            });
        }
 
        pool.shutdown();                          // 不再接受新任务
        pool.awaitTermination(5, TimeUnit.SECONDS); // 等所有任务结束
        System.out.println("all done");
    }
}
⚠️生产里别用 Executors.newFixedThreadPool

它的内部队列是无界 LinkedBlockingQueue,任务堆积会 OOM。生产用 new ThreadPoolExecutor(...) 显式指定队列上限和拒绝策略。

7. Future 与 Callable

Runnable 没返回值,Callable<T> 可以返回结果并抛异常。

Future / Callable
import java.util.concurrent.*;
 
public class Main {
    public static void main(String[] args) throws Exception {
        ExecutorService pool = Executors.newFixedThreadPool(2);
 
        Callable<Integer> square = (Callable<Integer>) () -> {
            Thread.sleep(100);
            return 7 * 7;
        };
        Future<Integer> f = pool.submit(square);
        System.out.println("任务已提交,主线程做点别的...");
        Thread.sleep(50);
        System.out.println("isDone? " + f.isDone());
        Integer r = f.get();            // 阻塞直到拿到结果
        System.out.println("isDone? " + f.isDone());
        System.out.println("7² = " + r);
 
        pool.shutdown();
    }
}

8. AtomicInteger:无锁原子

AtomicInteger
import java.util.concurrent.atomic.*;
 
public class Main {
    public static void main(String[] args) throws InterruptedException {
        AtomicInteger n = new AtomicInteger(0);
        int threads = 4, loops = 50000;
 
        Thread[] ts = new Thread[threads];
        for (int i = 0; i < threads; i++) {
            ts[i] = new Thread(() -> {
                for (int j = 0; j < loops; j++) n.incrementAndGet();
            });
            ts[i].start();
        }
        for (Thread t : ts) t.join();
        System.out.printf("期望 %d, AtomicInteger 实际 %d%n", threads * loops, n.get());
    }
}

🎯 练习

并发求和
// 任务:用 ExecutorService + Future 并行计算多个数组之和,再合并
// 期望输出:sum = 1000
// (把 1..1000 分给 4 个线程并行累加)
 
import java.util.*;
import java.util.concurrent.*;
 
public class Main {
    public static void main(String[] args) throws Exception {
        // 补全:把 1..1000 切成 4 段,每段起一个 Callable 计算和
        // 最后用 Future.get() 累加
        ExecutorService pool = Executors.newFixedThreadPool(4);
        // ...
        pool.shutdown();
        // System.out.println("sum = " + total);
    }
}

小结

  • ✅ Thread.start() 开新线程;run() 只是普通调用
  • ✅ synchronized 互斥(锁对象);volatile 仅保证可见性
  • ✅ wait() 释放锁;notifyAll() 唤醒;必须配 synchronized
  • ✅ 线程池复用线程;生产用 ThreadPoolExecutor 而非 Executors
  • ✅ Future<T> 拿到 Callable 的返回值;get() 阻塞
  • ✅ AtomicInteger 用 CAS 实现无锁原子操作

下一章 反射:运行时检视与操作类与对象。