Java线程通信案例详解:从管道流到CompletableFuture的完整实践
目录导读
- 为什么线程通信是并发编程的基石?
- 经典通信方式:wait/notify机制深度剖析
- 高级通信工具:BlockingQueue与Condition
- 现代异步通信:CompletableFuture实战
- 常见问题问答(FAQ)
- 总结与性能对比
为什么线程通信是并发编程的基石?
在多线程环境中,线程之间往往需要协作完成任务,如生产者-消费者、主从协作等场景,如果线程间无法有效通信,就会出现数据不一致、死锁或忙等待(CPU空转)等问题,Java提供了多种线程通信机制,从最底层的wait/notify到高层次的并发工具类,选择合适的方式能显著提升系统吞吐量和代码可读性,本文将结合真实生产级案例,逐一演示不同通信方式的适用场景与代码实现。

经典通信方式:wait/notify机制深度剖析
案例需求:实现一个容量为1的缓冲区,生产者线程放入数据,消费者线程取出数据,两者必须交替执行。
public class WaitNotifyExample {
private final Object lock = new Object();
private int data;
private boolean empty = true;
public void produce(int value) throws InterruptedException {
synchronized (lock) {
while (!empty) { lock.wait(); } // 缓冲非空则等待
data = value;
empty = false;
lock.notifyAll(); // 唤醒消费者
}
}
public int consume() throws InterruptedException {
synchronized (lock) {
while (empty) { lock.wait(); } // 缓冲为空则等待
empty = true;
lock.notifyAll(); // 唤醒生产者
return data;
}
}
}
关键点:wait()会释放锁并进入等待集,notifyAll()唤醒所有等待线程,但必须用while而非if检查条件,以避免虚假唤醒,此方式虽基础,但代码繁琐,且容易因忘记唤醒而死锁。
高级通信工具:BlockingQueue与Condition
1 BlockingQueue——企业级首选
案例:使用LinkedBlockingQueue实现多生产者-多消费者模型,无需手动同步。
ExecutorService producers = Executors.newFixedThreadPool(2);
ExecutorService consumers = Executors.newFixedThreadPool(2);
BlockingQueue<Integer> queue = new LinkedBlockingQueue<>(10);
// 生产者任务
for (int i = 0; i < 2; i++) {
final int id = i;
producers.execute(() -> {
try {
for (int j = 0; j < 100; j++) {
queue.put(id * 100 + j); // 阻塞式放入
TimeUnit.MILLISECONDS.sleep(10);
}
} catch (InterruptedException e) { Thread.currentThread().interrupt(); }
});
}
// 消费者任务
for (int i = 0; i < 2; i++) {
consumers.execute(() -> {
try {
while (true) {
Integer val = queue.take(); // 阻塞式取出
System.out.println("消费: " + val);
}
} catch (InterruptedException e) { Thread.currentThread().interrupt(); }
});
}
BlockingQueue内部使用ReentrantLock和Condition实现线程安全,自动处理等待与唤醒,大幅降低出错率。
2 Condition——更细粒度的可控通信
当需要“多个等待条件”时,Condition比synchronized更灵活,例如读写锁场景,可分离“读条件”和“写条件”:
ReentrantLock lock = new ReentrantLock();
Condition notFull = lock.newCondition();
Condition notEmpty = lock.newCondition();
// 生产者
lock.lock();
try {
while (count == size) notFull.await();
// 插入数据
notEmpty.signal();
} finally { lock.unlock(); }
现代异步通信:CompletableFuture实战
案例:一个订单服务需要同时调用“库存服务”和“折扣服务”,两者独立,最后合并结果。
CompletableFuture<Inventory> invFuture = CompletableFuture
.supplyAsync(() -> queryInventory(orderId)); // 异步查库存
CompletableFuture<Discount> discFuture = CompletableFuture
.supplyAsync(() -> getDiscount(orderId)); // 异步查折扣
CompletableFuture<OrderResult> result = invFuture
.thenCombine(discFuture, (inv, disc) ->
new OrderResult(inv.stock, disc.rate)); // 合并结果
result.thenAccept(System.out::println); // 非阻塞打印
优势:无需显式创建线程或使用锁,基于回调(Callback)和事件驱动,天然支持多任务编排(如串行、并行、超时控制),是异步微服务通信的利器。
常见问题问答(FAQ)
Q1:wait()和sleep()有什么区别?
A:wait()释放锁并依赖notify唤醒,必须在synchronized块内调用;sleep()不释放锁,到期自动唤醒,可随时调用。
Q2:什么时候用BlockingQueue,什么时候用CompletableFuture?
A:若需要数据缓冲(如任务队列),选BlockingQueue;若需要异步结果传递(如RPC调用),选CompletableFuture。
Q3:使用notify还是notifyAll?
A:单生产者单消费者用notify;多生产者多消费者必须用notifyAll,否则可能唤醒同类型线程导致“信号丢失”死锁。
Q4:Condition和内置锁的wait/notify性能差异?
A:Condition支持多条件队列和多路等待,且支持公平锁和非公平锁,在高竞争场景吞吐量更高。
总结与性能对比
| 通信方式 | 适用场景 | 代码复杂度 | 线程阻塞 |
|---|---|---|---|
| wait/notify | 简单单缓冲 | 高(易错) | 是(释放锁) |
| BlockingQueue | 生产-消费队列 | 低 | 是(自动) |
| Condition | 多条件复杂等待 | 中 | 是(精细) |
| CompletableFuture | 异步结果组合 | 低(声明式) | 否(回调) |
最佳实践建议:
- 日常开发优先使用
BlockingQueue或CompletableFuture - 若需要维护复杂状态(如多路闸门),可选用
Condition - 永远避免“忙等待”(死循环检查变量),那会浪费CPU资源
延伸思考:在Java 9+中,FlowAPI提供了响应式流(Reactive Streams)支持,适用于背压场景,可视为线程通信的下一代演进方向,掌握以上案例,您已能应对90%以上的线程协作需求,但务必结合压测工具(如JMH)验证吞吐量,避免过早优化。