基于ConcurrentHashMap与Java线程的类生产者/消费者模型问题
针对多线程读写ConcurrentHashMap的等待方案
针对你描述的场景——Thread1持续向ConcurrentHashMap中添加/填充块,Thread2需要等待指定索引的块存在(或填充完成)再读取,我给你两个实用的线程安全方案,完全适配你的需求:
方案一:用CompletableFuture实现异步等待(推荐)
这个方案利用Java 8+的CompletableFuture来天然实现“等待资源就绪”的逻辑,搭配ConcurrentHashMap的线程安全特性,代码简洁且不易出错。
核心思路
把ConcurrentHashMap的value设为CompletableFuture<Block>,而不是直接存Block:
- Thread1准备好块数据后,完成对应的Future,通知等待的Thread2;
- Thread2尝试获取对应索引的Future,若不存在则创建一个,然后调用
get()自动阻塞直到块就绪。
代码实现
首先定义你的Block类(保证数据读写的线程安全):
class Block { private final int index; private byte[] data; public Block(int index) { this.index = index; this.data = new byte[0]; // 初始为空 } // Thread1调用:设置块数据,加同步保证线程安全 public synchronized void setData(byte[] data) { this.data = data; } // Thread2调用:获取块数据,返回副本避免外部修改 public synchronized byte[] getData() { return data.clone(); } public int getIndex() { return index; } }
然后是Blockstore的实现:
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CompletableFuture; class Blockstore { private final ConcurrentHashMap<Integer, CompletableFuture<Block>> blockMap = new ConcurrentHashMap<>(); // Thread1调用:添加并填充块 public void addBlock(int index, byte[] data) { // 原子性获取或创建对应索引的Future,避免重复创建 CompletableFuture<Block> future = blockMap.computeIfAbsent(index, i -> new CompletableFuture<>()); Block block = new Block(index); block.setData(data); future.complete(block); // 标记块已准备就绪,唤醒等待线程 } // Thread2调用:获取指定索引的块,不存在则等待 public Block getBlock(int index) throws InterruptedException { CompletableFuture<Block> future = blockMap.get(index); // 如果Future不存在,创建一个新的(保证后续Thread1添加时能关联上) if (future == null) { future = blockMap.computeIfAbsent(index, i -> new CompletableFuture<>()); } try { return future.get(); // 阻塞直到块完成填充 } catch (Exception e) { Thread.currentThread().interrupt(); throw new InterruptedException("等待块时被中断"); } } }
方案优势
- 无需手动管理锁和通知逻辑,
CompletableFuture自带异步等待机制; computeIfAbsent保证了多线程下不会创建重复的Future,线程安全;- 天然支持超时等待(可以用
get(long timeout, TimeUnit unit)),灵活度高。
方案二:用ReentrantLock + Condition实现传统等待
如果你更倾向于用经典的并发锁机制,这个方案通过ReentrantLock和Condition来实现线程间的等待/通知。
核心思路
- 用
ConcurrentHashMap存储已添加的Block; - 全局锁和Condition用来通知所有等待线程:有新块添加了;
- Thread2在目标块不存在时进入等待,直到Thread1添加块后唤醒。
代码实现
Block类和方案一一致,Blockstore实现如下:
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.ReentrantLock; class Blockstore { private final ConcurrentHashMap<Integer, Block> blockMap = new ConcurrentHashMap<>(); private final ReentrantLock lock = new ReentrantLock(); private final Condition blockAdded = lock.newCondition(); // Thread1调用:添加块 public void addBlock(int index, byte[] data) { Block block = new Block(index); block.setData(data); blockMap.put(index, block); // 唤醒所有等待的线程,告知有新块添加 lock.lock(); try { blockAdded.signalAll(); } finally { lock.unlock(); } } // Thread2调用:获取指定索引的块,不存在则等待 public Block getBlock(int index) throws InterruptedException { Block block = blockMap.get(index); if (block != null) { return block; } // 块不存在,进入等待 lock.lock(); try { // 用while循环避免**虚假唤醒**,这是并发编程的最佳实践 while ((block = blockMap.get(index)) == null) { blockAdded.await(); } return block; } finally { lock.unlock(); } } }
注意事项
- 必须用
while循环检查块是否存在,不能用if,因为可能出现虚假唤醒(线程在没有被显式唤醒的情况下醒来); - 如果你的场景是Thread1先添加空Block再填充数据,记得要在
Block类的读写方法上加同步(如方案一的synchronized),避免Thread2读取到未填充完成的数据。
方案选择建议
- 如果需要等待块完全填充完成(而不仅仅是存在),优先选方案一;
- 如果块是一次性添加完成的,方案二更简单直接;
- 两种方案都完美适配
ConcurrentHashMap的线程安全特性,不会出现读写冲突。
内容的提问来源于stack exchange,提问作者Jordan Mackie
相关产品推荐
相关产品推荐

