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

如何在Hazelcast中初始化Java Condition?求具体实现方法

在Hazelcast中实现分布式Condition(ICondition)

完全可以在Hazelcast中实现分布式的Condition逻辑,Hazelcast的CP子系统提供了ICondition接口——这是Java标准Condition的分布式实现,专门用于跨JVM的线程协作场景。

你的代码问题修正

你当前代码里将notEmpty声明为java.util.concurrent.locks.Condition,虽然语法上可行,但建议直接使用Hazelcast的ICondition类型,以便更好地利用其分布式特性:

@Getter
@Service
public class HazelcastService {

    private FencedLock lock;
    private PNCounter counter;
    private ICondition notEmpty; // 改为ICondition类型
    private IList<PushStatementObjectRequestDTO> listData;

    @Autowired
    HazelcastInstance hazelcastInstance;

    @PostConstruct
    public void init(){
        lock = hazelcastInstance.getCPSubsystem().getLock(LIST_LOCK);
        counter = hazelcastInstance.getPNCounter(COUNTER);
        notEmpty = lock.newCondition(); // FencedLock.newCondition()返回ICondition
        listData = hazelcastInstance.getList(LIST_DATA);
    }
}

关键配置前提

使用CP子系统的分布式锁和Condition,必须先启用CP子系统:

方式1:XML配置(hazelcast.xml)

<hazelcast>
    <cp-subsystem>
        <!-- 推荐至少3个CP成员保证一致性,避免单点故障 -->
        <cp-member-count>3</cp-member-count>
    </cp-subsystem>
</hazelcast>

方式2:Java代码配置

Config config = new Config();
config.getCPSubsystemConfig().setCPMemberCount(3);
HazelcastInstance hazelcastInstance = Hazelcast.newHazelcastInstance(config);

分布式Condition的使用示例(生产者-消费者模型)

和本地Condition的使用逻辑一致,但作用于分布式环境:

生产者方法(添加数据并通知等待线程)

public void addData(PushStatementObjectRequestDTO data) {
    lock.lock();
    try {
        listData.add(data);
        counter.increment();
        notEmpty.signalAll(); // 通知所有跨JVM等待的消费者线程
    } finally {
        lock.unlock();
    }
}

消费者方法(等待直到有数据可取)

public PushStatementObjectRequestDTO takeData() throws InterruptedException {
    lock.lock();
    try {
        // 必须用while循环检查条件,避免虚假唤醒(分布式场景同样存在)
        while (listData.isEmpty()) {
            notEmpty.await(); // 释放锁,等待生产者通知
        }
        PushStatementObjectRequestDTO data = listData.remove(0);
        counter.decrement();
        return data;
    } finally {
        lock.unlock();
    }
}

核心注意事项

  • 所有ICondition的方法(await()/signal()/signalAll())必须在FencedLock的锁定范围内调用,否则会抛出异常。
  • 调用await()时会自动释放锁,允许其他节点的线程获取锁;被通知唤醒后,会重新获取锁再继续执行。
  • ICondition继承了Java标准Condition,同时提供了更易用的await(Duration)方法,无需手动转换时间单位。

内容的提问来源于stack exchange,提问作者IdanOst

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 19:35:22