zookeeper学习七-分布式计数器&分布式屏障

1 分布式计数器

@Test
    public void testDistributedAtomicInteger() throws Exception {
        DistributedAtomicInteger atomicInteger = new DistributedAtomicInteger(curatorFramework, Constant.ROOT_PATH,
                        new RetryNTimes(3, 1000));
        AtomicValue<Integer> value = atomicInteger.add(1);
        log.info("current value:{}", value.postValue());
        log.info("before value:{}", value.preValue());
    }

2 分布式屏障

2.1 单屏障

@Test
    public void testDistributedBarrier() throws Exception {
        final DistributedBarrier barrier = new DistributedBarrier(curatorFramework, Constant.ROOT_PATH);
        for (int i = 0; i < 5; i++) {
            executor.execute(() -> {
                try {
                    System.out.println(Thread.currentThread().getName() + "设置barrier!");
                    // 设置
                    barrier.setBarrier();
                    // 等待
                    barrier.waitOnBarrier();
                    System.out.println("---------开始执行程序----------");
                } catch (Exception e) {
                    e.printStackTrace();
                }
            });
        }
        Thread.sleep(1000);
        //释放
        barrier.removeBarrier();
        Thread.sleep(10000);
    }

输出

pool-1-thread-1设置barrier!
pool-1-thread-2设置barrier!
pool-1-thread-3设置barrier!
pool-1-thread-4设置barrier!
pool-1-thread-5设置barrier!
---------开始执行程序----------
---------开始执行程序----------
---------开始执行程序----------
---------开始执行程序----------
---------开始执行程序----------

2.2 双屏障

@Test
    public void testDistributedDoubleBarrier() throws InterruptedException {
        for (int i = 0; i < 5; i++) {
            executor.execute(() -> {
                DistributedDoubleBarrier barrier = new DistributedDoubleBarrier(curatorFramework, Constant.ROOT_PATH, 5);
                try {
                    Thread.sleep(1000 * (new Random()).nextInt(3));
                    System.out.println(Thread.currentThread().getName() + "已经准备");
                    barrier.enter();
                    System.out.println("同时开始运行...");
                    Thread.sleep(1000 * (new Random()).nextInt(3));
                    System.out.println(Thread.currentThread().getName() + "运行完毕");
                    barrier.leave();
                    System.out.println("同时退出运行...");
                } catch (Exception e) {
                    e.printStackTrace();
                }
            });
        }
        Thread.sleep(5000L);
    }

输出

pool-1-thread-2已经准备
pool-1-thread-4已经准备
pool-1-thread-3已经准备
pool-1-thread-5已经准备
pool-1-thread-1已经准备
同时开始运行...
同时开始运行...
同时开始运行...
同时开始运行...
同时开始运行...
pool-1-thread-2运行完毕
pool-1-thread-3运行完毕
pool-1-thread-5运行完毕
pool-1-thread-1运行完毕
pool-1-thread-4运行完毕
同时退出运行...
同时退出运行...
同时退出运行...
同时退出运行...
同时退出运行...
经验分享 程序员 微信小程序 职场和发展