You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 08:11:53