Java 多线程
多线程
线程是程序执行的最小单位,是进程中的一个执行路径。Java 语言内置了对多线程编程的支持,使得开发者可以轻松地创建和管理线程,从而实现并发执行任务,提高程序的性能和响应能力。
并发和并行
并发:在单个处理器上,通过时间片轮转等方式实现多个线程交替执行。 并行:在多处理器系统上,多个线程可以同时执行。
实现方式
继承 Thread 类
Thread 类,是 Java 中用于创建和管理线程的核心类。通过继承 Thread 类,可以创建自定义的线程类,并重写其 run() 方法来定义线程执行的任务。但是在启动线程时,不要直接调用 run() 方法,而是调用 start() 方法来启动线程。
▼java复制代码public class MyThread extends Thread { @Override public void run() { for (int i = 0; i < 100; i++) { System.out.println(getName() + " " +i); } } }
▼java复制代码MyThread myThread1 = new MyThread(); myThread1.setName("Thread 1"); // 不要调用 run 方式,而是执行 start 方法来创建线程 myThread1.start(); MyThread myThread2 = new MyThread(); myThread2.setName("Thread 2"); myThread2.start();
实现 Runnable 接口
Runnable 接口,是 Java 中另一种创建线程的方式。通过实现 Runnable 接口,可以将线程任务与线程本身分离,使得同一个 Runnable 对象可以被多个线程共享。实现 Runnable 接口需要重写 run() 方法,然后将 Runnable 对象传递给 Thread 类的构造函数来创建线程。
▼java复制代码public class MyRunnable implements Runnable { @Override public void run() { Thread thread = Thread.currentThread(); for (int i = 0; i < 100; i++) { System.out.println(thread.getName() + " " + i); } } }
▼java复制代码// 创建 Runnable 对象 MyRunnable myRunnable = new MyRunnable(); // 将 Runnable 对象传给 Thread Thread thread1 = new Thread(myRunnable); thread1.setName("Thread 1"); // 启动线程 thread1.start(); Thread thread2 = new Thread(myRunnable); thread2.setName("Thread 2"); thread2.start();
实现 Callable 接口
Callable 接口,是 Java 5 引入的一种创建线程的方式。与 Runnable 接口不同,Callable 接口的 call() 方法可以返回结果,并且可以抛出异常。实现 Callable 接口需要重写 call() 方法,然后将 Callable 对象传递给 FutureTask 类,再将 FutureTask 对象传递给 Thread 类来创建线程。
▼java复制代码public class MyCallable implements Callable<Integer> { @Override public Integer call() { Thread thread = Thread.currentThread(); // 求和并返回 int sum = 0; for (int i = 0; i < 100; i++) { sum += i; System.out.println(thread.getName() + " " + i); } return sum; } }
▼java复制代码// 创建 Callable 对象 MyCallable callable = new MyCallable(); // 将 Callable 对象交给 FutureTask FutureTask<Integer> futureTask = new FutureTask<Integer>(callable); // 将 FutureTask 交给 Thread 并开启线程 new Thread(futureTask).start(); // 通过 FutureTask 获取线程的返回值 Integer sum = futureTask.get(); System.out.println(sum);
Thread 构造方法
| 构造方法 | 说明 |
|---|---|
| Thread() | 创建一个新的线程对象,默认线程组为当前线程的线程组,线程名称为 "Thread-x" |
| Thread(Runnable task) | 创建一个新的线程对象,指定线程执行的任务 |
| Thread(ThreadGroup group, Runnable task) | 创建一个新的线程对象,指定线程组和线程执行的任务 |
| Thread(String name) | 创建一个新的线程对象,指定线程名称 |
| Thread(ThreadGroup group, String name) | 创建一个新的线程对象,指定线程组和线程名称 |
| Thread(Runnable task, String name) | 创建一个新的线程对象,指定线程执行的任务和线程名称 |
| Thread(ThreadGroup group, Runnable task, String name) | 创建一个新的线程对象,指定线程组、线程执行的任务和线程名称 |
| Thread(ThreadGroup group, Runnable task, String name, long stackSize) | 创建一个新的线程对象,指定线程组、线程执行的任务、线程名称和线程栈大小 |
| Thread(ThreadGroup group, Runnable task, String name, long stackSize, boolean inheritInheritableThreadLocals) | 创建一个新的线程对象,指定线程组、线程执行的任务、线程名称、线程栈大小和是否继承可继承的线程局部变量 |
Thread 常用方法
| 方法名 | 说明 |
|---|---|
currentThread() | 获取当前线程对象 |
yield() | 让出当前 CPU 执行权,进入就绪状态 |
sleep(xxx) | 使当前线程休眠指定时间,进入阻塞状态 |
| ofPlatform() | 创建一个与平台相关的线程工厂 |
| ofVirtual() | 创建一个虚拟线程工厂 |
| startVirtualThread(Runnable task) | 启动一个虚拟线程来执行指定任务 |
| isVirtual() | 判断当前线程是否为虚拟线程 |
| start() | 启动线程,调用线程的 run() 方法 |
| run() | 线程执行的任务代码 |
| interrupt() | 中断线程 |
| interrupted() | 静态方法,判断当前线程是否被中断,并清除中断状态 |
| isInterrupted() | 判断线程是否被中断,不清除中断状态 |
| isAlive() | 判断线程是否存活 |
set / getPriority(int newPriority) | 设置 / 获取线程优先级 |
set / getName(String name) | 设置 / 获取线程名称,线程默认名为:Thread-x |
| getThreadGroup() | 获取线程所属的线程组 |
| activeCount() | 获取线程组中活动线程的数量 |
| enumerate(Thread[] tarray) | 将线程组中的活动线程复制到指定数组中 |
join(xxx) | 等待线程终止,可以指定等待时间 |
| dumpStack() | 打印当前线程的堆栈信息 |
setDaemon(boolean on) | 设置线程为守护线程或用户线程 |
| isDaemon() | 判断线程是否为守护线程 |
| getContextClassLoader(xxx) | 获取线程的上下文类加载器 |
| holdsLock(Object obj) | 判断当前线程是否持有指定对象的锁 |
| getStackTrace() | 获取线程的堆栈跟踪信息 |
| getAllStackTraces() | 获取所有线程的堆栈跟踪信息 |
| threadId() | 获取线程的唯一标识符 |
| getState() | 获取线程的状态 |
| set / getDefaultUncaughtExceptionHandler(UncaughtExceptionHandler ueh) | 设置 / 获取线程的默认未捕获异常处理器 |
| set / getUncaughtExceptionHandler() | 设置 / 获取线程的未捕获异常处理器 |
Thread 代码示例
▼java复制代码// 获取当前线程 Thread thread = Thread.currentThread(); // 获取线程名 String name = thread.getName(); System.out.println(name); // main // 修改线程名 thread.setName("thread test"); System.out.println(thread.getName()); // thread test // 睡眠 System.out.println(System.currentTimeMillis()); // 1761309398420 Thread.sleep(5000); System.out.println(System.currentTimeMillis()); // 1761309403421
线程优先级
线程优先级是指线程在 CPU 调度中的优先级别。Java 中的线程优先级范围从 1 到 10,数值越大,优先级越高,默认优先级为 5。可以通过 setPriority() 方法设置线程的优先级,通过 getPriority() 方法获取线程的优先级。需要注意的是,线程优先级只是对线程调度的一种建议,具体的调度行为取决于操作系统的实现。
▼java复制代码MyThread thread1 = new MyThread(); MyThread thread2 = new MyThread(); System.out.println(thread1.getPriority()); // 5 System.out.println(thread2.getPriority()); // 5 // 设置线程优先级 thread1.setPriority(1); thread2.setPriority(10); // 启动线程,但是按照优先级,按理说 thread2 应该先执行,但实际结果不一定,毕竟优先级高只代表有更高的概率被 CPU 调度,而不是 100% 会被调度 thread1.start(); thread2.start();
守护线程
守护线程是一种特殊的线程,它的生命周期依赖于其他非守护线程。当所有非守护线程结束时,JVM 会自动终止所有守护线程。守护线程通常用于执行后台任务,如垃圾回收、日志记录等。可以通过 setDaemon(true) 方法将线程设置为守护线程。
▼java复制代码public class MyThread1 extends Thread { @Override public void run() { for (int i = 0; i < 10; i++) { System.out.println(getName() + " " + i); } } }
▼java复制代码public class MyThread2 extends Thread { @Override public void run() { for (int i = 0; i < 1000; i++) { System.out.println(getName() + " " + i); } } }
▼java复制代码MyThread1 myThread1 = new MyThread1(); MyThread2 myThread2 = new MyThread2(); // 设置 myThread2 为守护线程 myThread2.setDaemon(true); // 启动线程,当 myThread1 结束后,myThread2 也会自动结束 myThread1.start(); myThread2.start();
礼让线程
礼让线程是指在多线程环境中,某个线程主动放弃 CPU 使用权,让其他线程有机会执行。可以通过调用 Thread.yield() 方法来实现线程的礼让。需要注意的是,礼让线程并不一定会导致当前线程立即被挂起,具体行为依赖于操作系统的线程调度策略。
▼java复制代码public class MyThread extends Thread { @Override public void run() { for (int i = 0; i < 100; i++) { System.out.println(getName() + " " + i); // 当 i 为 50 时,礼让线程 if (i == 50) { System.out.println(getName() + " 礼让线程"); Thread.yield(); } } } }
▼java复制代码MyThread myThread1 = new MyThread(); MyThread myThread2 = new MyThread(); myThread1.setName("Thread 1"); myThread2.setName("Thread 2"); myThread1.start(); myThread2.start();
插队线程
插队线程是指在多线程环境中,某个线程通过调用 Thread.join() 方法,等待另一个线程执行完毕后再继续执行。这样可以确保某个线程在另一个线程完成之前不会继续执行,从而实现线程间的协调和同步。
▼java复制代码MyThread thread = new MyThread(); thread.start(); // 让主线程等待 thread 线程执行完毕 thread.join(); String name = Thread.currentThread().getName(); for (int i = 0; i < 100; i++) { System.out.println(name + i); }
线程的生命周期
线程的生命周期包括以下几个状态:
- 新建状态(New):线程对象被创建,但尚未启动。
- 通过
start()方法启动线程后,线程进入就绪状态。
- 通过
- 就绪状态(Runnable):线程已经启动,等待 CPU 调度执行。
- 通过
CPU 调度后,线程进入运行状态。
- 通过
- 运行状态(Running):线程正在执行其任务。
- 通过
sleep()、wait()等方法,线程可以进入阻塞状态。 - 当 CPU 被其他线程抢走后,回到就绪状态。
- 当线程执行完毕后进入终止状态。
- 通过
- 阻塞状态(Blocked):线程因等待某个条件而暂停执行。
- 通过
notify()或notifyAll()方法,线程可以从阻塞状态恢复到就绪状态。 sleep()方法结束后,线程也会从阻塞状态恢复到就绪状态。
- 通过
- 等待状态(Waiting):线程无限期等待另一个线程的通知。
- 通过
notify()或notifyAll()方法,线程可以从等待状态恢复到就绪状态。 - 通过
join()方法,线程可以等待另一个线程执行完毕后恢复到就绪状态。
- 通过
- 超时等待状态(Timed Waiting):线程等待指定时间后自动恢复执行。
- 终止状态(Terminated):线程执行完毕或被强制终止。
安全问题
在多线程环境中,多个线程可能会同时访问共享资源,导致数据不一致或其他安全问题。为了确保线程安全,可以使用同步机制,如 synchronized 关键字、Lock 接口等。
▼java复制代码public class MyThread extends Thread { static int ticked = 0; @Override public void run() { while (true) { if (ticked < 100) { // 模拟延迟 try { Thread.sleep(100); } catch (InterruptedException e) { throw new RuntimeException(e); } ticked++; System.out.println(getName() + "正在售卖第" + ticked + "张票"); } else { break; } } } }
▼java复制代码// 模拟三个窗口同时售票 MyThread myThread1 = new MyThread(); MyThread myThread2 = new MyThread(); MyThread myThread3 = new MyThread(); myThread1.start(); myThread2.start(); myThread3.start(); // 控制台打印输出结果,会发现,Thread-1 和 Thread-2 卖的票数会重复,导致数据不一致 // Thread-0正在售卖第1张票 // ...... // Thread-2正在售卖第96张票 // Thread-1正在售卖第97张票 // Thread-0正在售卖第95张票 // Thread-1正在售卖第99张票 // Thread-2正在售卖第99张票 // Thread-0正在售卖第100张票
解决线程安全问题
线程安全问题的解决方法有很多种,如使用 synchronized 关键字、Lock 接口、volatile 关键字、原子类等。
synchronized 关键字
synchronized 关键字可以用于方法或代码块,用于保证同一时间只有一个线程可以访问该方法或代码块。
▼java复制代码// synchronized 用于代码块 public class MyThread extends Thread { static int ticked = 0; // 创建一个锁对象,用于锁定资源,必须保证该对象是唯一的 static final Object lock = new Object(); @Override public void run() { while (true) { // 锁定资源 synchronized (lock) { if (ticked < 100) { try { Thread.sleep(100); } catch (InterruptedException e) { throw new RuntimeException(e); } ticked++; System.out.println(getName() + "正在售卖第" + ticked + "张票"); } else { break; } } } } }
▼java复制代码// synchronized 用于方法 public class MyThread extends Thread { static int ticked = 0; @Override public void run() { while (true) { if (method()) break; } } private synchronized boolean method() { if (ticked < 100) { try { Thread.sleep(100); } catch (InterruptedException e) { throw new RuntimeException(e); } ticked++; System.out.println(getName() + "正在售卖第" + ticked + "张票"); return false; } return true; } }
Lock 接口
Lock 接口提供了比 synchronized 更加细粒度的控制,可以更灵活地控制线程访问资源。提供了 tryLock() 方法,该方法尝试获取锁, lock() 方法,该方法获取锁, unlock() 方法,该方法释放锁。
Lock 是一个接口,所以我们需要使用它的实现类:ReentrantLock
▼java复制代码public class MyThread extends Thread { static int ticked = 0; static final Lock lock = new ReentrantLock(); @Override public void run() { while (true) { // 锁定资源 lock.lock(); if (ticked < 100) { try { Thread.sleep(1); } catch (InterruptedException e) { throw new RuntimeException(e); } ticked++; System.out.println(getName() + "正在售卖第" + ticked + "张票"); } else { // 当线程执行完毕,确保锁释放掉 lock.unlock(); break; } // 执行完毕,确保锁释放掉 lock.unlock(); } } }
死锁
当一个线程执行是需要获取多个锁,但是只获取了部分锁,其他锁被其他线程占用,自己没有获取全部锁,自己无法执行,而其他线程也无法执行,就会发生死锁。死锁问题出现的原因:多个线程同时访问同一资源,多个线程都等待其他线程释放资源,从而形成僵局,无法继续执行。
▼java复制代码public class MyThread extends Thread { static int ticked = 0; // 创建两个锁对象 static final Lock lockA = new ReentrantLock(); static final Lock lockB = new ReentrantLock(); @Override public void run() { while (true) { if (getName().equals("Thread-1")) { // 线程 1 需要获取 A 和 B lockA.lock(); System.out.println("Thread-1 获取 A"); // 当线程 1 获取 A 之后,可能Thread-2已持有B,那么需要等待Thread-2释放B lockB.lock(); System.out.println("Thread-1 获取 B"); } else { // 线程 2 需要获取 B 和 A lockB.lock(); System.out.println("Thread-2 获取 B"); // 当线程 2 获取 B 之后,可能Thread-1已持有A,那么需要等待Thread-1释放A lockA.lock(); System.out.println("Thread-2 获取 A"); } lockA.unlock(); lockB.unlock(); } } } // 控制台打印结果:线程1和线程2都进入死锁状态 // Thread-2 获取 B // Thread-1 获取 A
等待唤醒机制
等待唤醒机制,也叫同步机制,是多线程并发访问共享资源的一种解决方案。典型案例就是生产者-消费者问题。消费者负责获取数据,生产者负责生产数据。当队列已满时,生产者等待,当队列已空时,消费者等待。消费者获取数据时,会调用 wait() 方法,该方法会释放锁,并进入等待状态,当生产者生产数据时,会调用 notify() 方法,该方法会唤醒等待的消费者线程,消费者线程会重新获取锁,并继续执行。
消费者
- 判断是否存在需要消费的数据
- 没有数据,调用
wait()方法,释放锁,并进入等待状态,等待生产者生产数据 - 存在数据,消费数据,结束之后,调用
notify()方法,唤醒等待的生产者线程
生产者
- 判断数据是否被消费
- 没有被消费,调用
wait()方法,释放锁,并进入等待状态,等待消费者消费数据 - 存在被消费,生产数据,结束之后,调用
notify()方法,唤醒消费者消费数据
等待唤醒实现方式
| 机制 | synchronized | ReentrantLock |
|---|---|---|
| 等待方法 | obj.wait() | condition.await() |
| 唤醒单个线程 | obj.notify() | condition.signal() |
| 唤醒所有线程 | obj.notifyAll() | condition.signalAll() |
| 使用位置 | 必须在 synchronized 块中 | 必须在 lock.lock() 和 lock.unlock() 之间 |
Synchronized
▼java复制代码public class ConsumerQueue { // 最多生产/消费 10 个数据 public static int dataCount = 10; // 记录数据状态 public static boolean hasData = false; // 锁 public static final Object lock = new Object(); }
▼java复制代码// 生产者 public class Producer extends Thread { @Override public void run() { while (true) { synchronized (ConsumerQueue.lock) { // 只生产 10 个数据,生产完就结束线程 if (ConsumerQueue.dataCount == 0) { break; } // 1. 判断队列中是否有数据 if (ConsumerQueue.hasData) { try { // 2. 有数据,就进入等待状态 ConsumerQueue.lock.wait(); } catch (InterruptedException e) { throw new RuntimeException(e); } } else { // 3. 没有数据,就生产数据 System.out.println(getName() + "生产了一条数据,还可以生产" + ConsumerQueue.dataCount + "条数据"); // 4. 唤醒消费者(只唤醒和这把锁相关的线程) ConsumerQueue.lock.notifyAll(); // 5. 修改数据状态 ConsumerQueue.hasData = true; } } } } }
▼java复制代码// 消费者 public class Consumer extends Thread { @Override public void run() { while (true) { // 上锁 synchronized (ConsumerQueue.lock) { // 只消费 10 个数据,消费完就结束线程 if (ConsumerQueue.dataCount == 0) { break; } // 1. 判断队列中是否有数据 if (ConsumerQueue.hasData) { // 2. 有数据,就消费数据 ConsumerQueue.dataCount--; System.out.println(getName() + "消费了一条数据,还可以消费" + ConsumerQueue.dataCount + "条数据"); // 3. 唤醒生产者(只唤醒和这把锁相关的线程) ConsumerQueue.lock.notifyAll(); // 4. 修改数据状态 ConsumerQueue.hasData = false; } else { try { // 4. 没有数据,进入等待状态 ConsumerQueue.lock.wait(); } catch (InterruptedException e) { throw new RuntimeException(e); } } } } } }
▼java复制代码Producer producer = new Producer(); Consumer consumer = new Consumer(); producer.setName("生产者"); consumer.setName("消费者"); producer.start(); consumer.start(); // 生产者生产了一条数据,还可以生产10条数据 // 消费者消费了一条数据,还可以消费9条数据 // 生产者生产了一条数据,还可以生产9条数据 // 消费者消费了一条数据,还可以消费8条数据 // 生产者生产了一条数据,还可以生产8条数据 // 消费者消费了一条数据,还可以消费7条数据 // 生产者生产了一条数据,还可以生产7条数据 // 消费者消费了一条数据,还可以消费6条数据 // 生产者生产了一条数据,还可以生产6条数据 // 消费者消费了一条数据,还可以消费5条数据 // 生产者生产了一条数据,还可以生产5条数据 // 消费者消费了一条数据,还可以消费4条数据 // 生产者生产了一条数据,还可以生产4条数据 // 消费者消费了一条数据,还可以消费3条数据 // 生产者生产了一条数据,还可以生产3条数据 // 消费者消费了一条数据,还可以消费2条数据 // 生产者生产了一条数据,还可以生产2条数据 // 消费者消费了一条数据,还可以消费1条数据 // 生产者生产了一条数据,还可以生产1条数据 // 消费者消费了一条数据,还可以消费0条数据
Lock
▼java复制代码// 消费队列 public class ConsumerQueue { // 最多生产/消费 10 个数据 public static int dataCount = 10; // 记录数据状态 public static boolean hasData = false; // 锁 public static final Lock lock = new ReentrantLock(); // Condition 对象 public static final Condition condition = lock.newCondition(); }
▼java复制代码// 生产者 public class Producer extends Thread { @Override public void run() { while (true) { // 上锁 ConsumerQueue.lock.lock(); // 只生产 10 个数据,生产完就结束线程 if (ConsumerQueue.dataCount == 0) { // 解锁 ConsumerQueue.lock.unlock(); break; } // 1. 判断队列中是否有数据 if (ConsumerQueue.hasData) { try { // 2. 有数据,就进入等待状态 ConsumerQueue.condition.await(); } catch (InterruptedException e) { throw new RuntimeException(e); } } else { // 3. 没有数据,就生产数据 System.out.println(getName() + "生产了一条数据,还可以生产" + ConsumerQueue.dataCount + "条数据"); // 4. 唤醒消费者(只唤醒和这把锁相关的线程) ConsumerQueue.condition.signal(); // 5. 修改数据状态 ConsumerQueue.hasData = true; } // 解锁 ConsumerQueue.lock.unlock(); } } }
▼java复制代码// 消费者 public class Consumer extends Thread { @Override public void run() { while (true) { // 上锁 ConsumerQueue.lock.lock(); // 只消费 10 个数据,消费完就结束线程 if (ConsumerQueue.dataCount == 0) { // 解锁 ConsumerQueue.lock.unlock(); break; } // 1. 判断队列中是否有数据 if (ConsumerQueue.hasData) { // 2. 有数据,就消费数据 ConsumerQueue.dataCount--; System.out.println(getName() + "消费了一条数据,还可以消费" + ConsumerQueue.dataCount + "条数据"); // 3. 唤醒生产者(只唤醒和这把锁相关的线程) ConsumerQueue.condition.signal(); // 4. 修改数据状态 ConsumerQueue.hasData = false; } else { try { // 4. 没有数据,进入等待状态 ConsumerQueue.condition.await(); } catch (InterruptedException e) { throw new RuntimeException(e); } } // 解锁 ConsumerQueue.lock.unlock(); } } }
▼java复制代码Producer producer = new Producer(); Consumer consumer = new Consumer(); producer.setName("生产者"); consumer.setName("消费者"); producer.start(); consumer.start(); // 运行结果: // 生产者生产了一条数据,还可以生产10条数据 // 消费者消费了一条数据,还可以消费9条数据 // 生产者生产了一条数据,还可以生产9条数据 // 消费者消费了一条数据,还可以消费8条数据 // 生产者生产了一条数据,还可以生产8条数据 // 消费者消费了一条数据,还可以消费7条数据 // 生产者生产了一条数据,还可以生产7条数据 // 消费者消费了一条数据,还可以消费6条数据 // 生产者生产了一条数据,还可以生产6条数据 // 消费者消费了一条数据,还可以消费5条数据 // 生产者生产了一条数据,还可以生产5条数据 // 消费者消费了一条数据,还可以消费4条数据 // 生产者生产了一条数据,还可以生产4条数据 // 消费者消费了一条数据,还可以消费3条数据 // 生产者生产了一条数据,还可以生产3条数据 // 消费者消费了一条数据,还可以消费2条数据 // 生产者生产了一条数据,还可以生产2条数据 // 消费者消费了一条数据,还可以消费1条数据 // 生产者生产了一条数据,还可以生产1条数据 // 消费者消费了一条数据,还可以消费0条数据
阻塞队列
阻塞队列是 Java 集合框架提供的一个接口,它继承了 Queue 接口,并且提供了一些额外的方法,使得队列的使用更加方便。
阻塞队列分类
| 阻塞队列 | 描述 |
|---|---|
| ArrayBlockingQueue | 基于数组的有界阻塞队列,容量固定且 FIFO(先进先出)。适用于需要严格控制资源占用的场景,如固定大小的线程池任务队列,避免内存溢出。 |
| LinkedBlockingQueue | 基于链表的可选有界阻塞队列(默认容量为 Integer.MAX_VALUE,可设为有界),FIFO 顺序。常用于高吞吐量的生产者-消费者模型,如 Web 服务器请求队列,兼顾性能与灵活性。 |
| PriorityBlockingQueue | 无界优先级阻塞队列,元素按自然顺序或 Comparator 排序。适用于需要优先级调度的任务,如实时系统中的紧急任务处理(例如告警系统)。 |
| DelayQueue | 无界延迟阻塞队列,元素必须实现 Delayed 接口,仅当延迟过期后才能被取出。专用于定时任务调度,如缓存过期清理、定时任务触发(类似 Timer 的替代方案)。 |
| SynchronousQueue | 不存储元素的同步阻塞队列,每个插入操作必须等待移除操作(反之亦然)。适用于直接传递(hand-off)场景,如 Executors.newCachedThreadPool()中的任务传递,减少中间缓冲开销。 |
| LinkedTransferQueue | 基于链表的无界阻塞队列,支持 transfer()方法实现生产者直接移交元素给消费者。用于高性能数据传递场景,如低延迟交易系统,避免额外的队列存储步骤。 |
▼java复制代码public class Producer extends Thread { private ArrayBlockingQueue<String> queue; public Producer(ArrayBlockingQueue<String> queue) { this.queue = queue; } @Override public void run() { while (true) { try { queue.put("data"); } catch (InterruptedException e) { throw new RuntimeException(e); } } } }
▼java复制代码public class Consumer extends Thread { private ArrayBlockingQueue<String> queue; public Consumer(ArrayBlockingQueue<String> queue) { this.queue = queue; } @Override public void run() { while (true) { try { queue.take(); } catch (InterruptedException e) { throw new RuntimeException(e); } } } }
▼java复制代码// 创建阻塞队列 ArrayBlockingQueue<String> queue = new ArrayBlockingQueue<>(1); Producer producer = new Producer(queue); Consumer consumer = new Consumer(queue); producer.setName("生产者"); consumer.setName("消费者"); producer.start(); consumer.start();
综合练习
模拟卖票
▼java复制代码public class MyThread extends Thread { public static int ticket = 1000; @Override public void run() { while (true) { synchronized (MyThread.class) { if (ticket > 0) { try { Thread.sleep(3000); } catch (InterruptedException e) { throw new RuntimeException(e); } ticket--; System.out.println(getName() + "卖出一张票,还剩下" + ticket + "张票"); } else { break; } } } } }
▼java复制代码MyThread myThread1 = new MyThread(); MyThread myThread2 = new MyThread(); myThread1.setName("窗口 A"); myThread2.setName("窗口 B"); myThread1.start(); myThread2.start();
分发礼品
▼java复制代码public class MyThread extends Thread { public static int gift = 100; @Override public void run() { while (true) { synchronized (MyThread.class) { if (gift > 10) { gift--; System.out.println(getName() + "送出一个礼物,还剩下" + gift + "个礼物"); } else { break; } } } } }
▼java复制代码MyThread myThread1 = new MyThread(); MyThread myThread2 = new MyThread(); myThread1.setName("张三"); myThread2.setName("李四"); myThread1.start(); myThread2.start();
打印奇数
▼java复制代码public class MyThread extends Thread { public static int num = 0; @Override public void run() { while (true) { synchronized (MyThread.class) { if (num > 100) { break; } if (num % 2 == 1) { System.out.println(num); } num++; } } } }
▼java复制代码MyThread myThread1 = new MyThread(); MyThread myThread2 = new MyThread(); myThread1.start(); myThread2.start();
抢红包
▼java复制代码public class MyThread extends Thread { public static double allMoney = 100.0; public static int count = 3; public static double minMoney = 0.01; @Override public void run() { synchronized (MyThread.class) { if (count == 0) { System.out.println(getName() + "没有抢到红包"); return; } double v; if (count == 1) v = allMoney; else v = new Random().nextDouble(allMoney - (count - 1) * minMoney); if (v < minMoney) v = minMoney; allMoney -= v; count--; System.out.println(getName() + "抢到了" + v + "元,还剩下" + count + "个红包, 还剩下" + allMoney + "元"); } } }
▼java复制代码MyThread myThread1 = new MyThread(); MyThread myThread2 = new MyThread(); MyThread myThread3 = new MyThread(); MyThread myThread4 = new MyThread(); MyThread myThread5 = new MyThread(); myThread1.setName("张三"); myThread2.setName("李四"); myThread3.setName("王五"); myThread4.setName("赵六"); myThread5.setName("孙七"); myThread1.start(); myThread2.start(); myThread3.start(); myThread4.start(); myThread5.start();
抽奖
▼java复制代码public class MyThread extends Thread { private List<Integer> list; public MyThread(List<Integer> list) { this.list = list; } @Override public void run() { while (true) { synchronized (MyThread.class) { if (list.isEmpty()) return; Collections.shuffle(list); Integer remove = list.removeFirst(); System.out.println(getName() + remove); } try { Thread.sleep(1000); } catch (InterruptedException e) { throw new RuntimeException(e); } } } }
▼java复制代码List<Integer> list = new ArrayList<Integer>(); Collections.addAll(list, 100, 200, 300, 400, 500, 600, 700); MyThread myThread1 = new MyThread(list); MyThread myThread2 = new MyThread(list); myThread1.setName("张三"); myThread2.setName("李四"); myThread1.start(); myThread2.start(); System.out.println(list);
抽奖结束后打印结果
▼java复制代码public class MyThread extends Thread { private List<Integer> list; public MyThread(List<Integer> list) { this.list = list; } public List<Integer> result = new ArrayList<>(); @Override public void run() { while (true) { synchronized (MyThread.class) { if (list.isEmpty()) { System.out.println(this.result); return; } Collections.shuffle(list); Integer remove = list.removeFirst(); this.result.add(remove); } try { Thread.sleep(1000); } catch (InterruptedException e) { throw new RuntimeException(e); } } } }
▼java复制代码public class ThreadTest01 { public static void main(String[] args) { List<Integer> list = new ArrayList<>(); Collections.addAll(list, 100, 200, 300, 400, 500, 600, 700); MyThread myThread1 = new MyThread(list); MyThread myThread2 = new MyThread(list); myThread1.setName("张三"); myThread2.setName("李四"); myThread1.start(); myThread2.start(); System.out.println(list); } }
最大抽奖结果
▼java复制代码public class MyCallable implements Callable<Integer> { private List<Integer> list; public MyCallable(List<Integer> list) { this.list = list; } @Override public Integer call() throws Exception { List<Integer> result = new ArrayList<>(); while (true) { synchronized (MyCallable.class) { if (list.isEmpty()) { System.out.println(result); Optional<Integer> reduce = result.stream().reduce(Integer::max); if (reduce.isPresent()) return reduce.get(); } Collections.shuffle(list); Integer remove = list.removeFirst(); result.add(remove); } Thread.sleep(1000); } } }
▼java复制代码List<Integer> list = new ArrayList<>(); Collections.addAll(list, 100, 200, 300, 400, 500, 600, 700); MyCallable callable1 = new MyCallable(list); FutureTask<Integer> futureTask1 = new FutureTask<>(callable1); FutureTask<Integer> futureTask2 = new FutureTask<>(callable1); Thread myThread1 = new Thread(futureTask1); Thread myThread2 = new Thread(futureTask2); myThread1.setName("张三"); myThread2.setName("李四"); myThread1.start(); myThread2.start(); Integer integer1 = futureTask1.get(); System.out.println(integer1); Integer integer2 = futureTask2.get(); System.out.println(integer2); System.out.println(Integer.max(integer1, integer2));
线程池
线程池:线程复用,提高效率,避免频繁创建线程。
在 Java 中,线程池的实现类有:Executors(不推荐),ThreadPoolExecutor(推荐)。
线程池分类
Executors
| 线程池 | 描述 |
|---|---|
newFixedThreadPool | 创建一个定长线程池,可控制线程最大并发数,超出的线程会在队列中等待。 |
| newWorkStealingPool | 创建一个拥有并行级别为 CPU 核数的线程池,该线程池的线程数会根据 CPU 的核数动态变化。 |
| newSingleThreadExecutor | 创建一个单线程化的线程池,它只会用唯一的工作线程来执行任务,保证所有任务按照指定顺序执行。 |
newCachedThreadPool | 创建一个可缓存线程池,如果线程池长度超过处理需要,可灵活回收空闲线程,若无可回收,则新建线程。 |
| newThreadPerTaskExecutor | 创建一个线程池,它使用一个单独的线程来处理任务,该线程会一直运行,直到被中断。 |
| newVirtualThreadPerTaskExecutor | 创建一个虚拟线程池,它使用一个单独的虚拟线程来处理任务,该虚拟线程会一直运行,直到被中断。 |
| newSingleThreadScheduledExecutor() | 创建一个单线程的定时执行器,该线程会按照指定的时间间隔重复执行任务。 |
| newScheduledThreadPool(int corePoolSize) | 创建一个定长线程池,支持定时及周期性任务执行。 |
ThreadPoolExecutor
| 线程池 | 描述 |
|---|---|
| ThreadPoolExecutor(int corePoolSize,int maximumPoolSize,long keepAliveTime,TimeUnit unit,BlockingQueue workQueue,ThreadFactory threadFactory,RejectedExecutionHandler handler) | 创建一个线程池,指定核心线程数、最大线程数、空闲线程的存活时间、任务队列、线程工厂和拒绝策略。 |
Executors 类
Executors 底层实现类是 ThreadPoolExecutor,只不过 Executors 对象进行了封装,使其更易用。
▼java复制代码public class MyRunnable implements Runnable { @Override public void run() { String name = Thread.currentThread().getName(); System.out.println(name + "线程执行了"); } }
▼java复制代码// 使用 try-with-resources 来确保线程池被正确关闭 // CachedThreadPool 没有线程数限制,当没有线程空闲,且有任务时,会创建新的线程执行任务。 try (ExecutorService executorService = Executors.newCachedThreadPool()) { executorService.submit(new MyRunnable()); Thread.sleep(100); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); } catch (InterruptedException e) { throw new RuntimeException(e); }
▼java复制代码// FixedThreadPool 创建一个定长线程池,当有任务提交时,会根据线程池的容量来分配线程执行任务。若线程池已满,则任务会进入队列等待。 try (ExecutorService executorService = Executors.newFixedThreadPool(3)) { executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); executorService.submit(new MyRunnable()); }
ThreadPoolExecutor 类
创建一个线程池,指定核心线程数、最大线程数、空闲线程的存活时间、存活时间单位、任务队列、线程工厂和拒绝策略。
- 当有任务提交时,核心线程会执行任务。
- 当核心线程数已满,任务会进入队列等待。
- 当队列已满,任务会创建非核心线程执行任务,不是去执行队列中的任务,而是执行新提交的任务。
- 当非核心线程数已满,任务会进入拒绝策略处理。
拒绝策略:
- AbortPolicy:默认策略,抛出 RejectedExecutionException 异常。
- DiscardPolicy:忽略新任务,不抛出异常。
- DiscardOldestPolicy:丢弃队列中靠前的任务,并执行新任务。
- CallerRunsPolicy:将新任务执行在调用者线程中。
线程池多大比较合适:
- CPU 密集型:线程数 = CPU 线程数 + 1
- IO 密集型:线程数 = CPU 线程数 * CPU 利用率 * (CPU 计算时间 + 等待时间) / CPU 计算时间
▼java复制代码ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor( 3, // 核心线程 6, // 最大线程数,临时线程数 = 最大线程数 - 核心线程数 = 6 - 3 = 3 60, // 临时线程最大存活时间为 60 s TimeUnit.SECONDS, // 秒 new ArrayBlockingQueue<>(3), // 任务队列,最多 3 个任务排队 Executors.defaultThreadFactory(), // 默认线程工厂 new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略,将新任务执行在调用者线程中 ); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.submit(new MyRunnable()); threadPoolExecutor.close(); // pool-1-thread-3线程执行了 // pool-1-thread-1线程执行了 // pool-1-thread-1线程执行了 // pool-1-thread-2线程执行了 // pool-1-thread-3线程执行了 // pool-1-thread-1线程执行了 // pool-1-thread-6线程执行了 // main线程执行了 // 该任务被拒绝,执行在调用者线程中 // pool-1-thread-5线程执行了 // pool-1-thread-4线程执行了
