Zookeeper基本实践_lbclientfactory zookeeper-程序员宅基地

技术标签: zookeeper  分布式  

增删改查节点

package com.org.zookeeper.crud;

import com.org.zookeeper.ClientFactory;
import lombok.extern.slf4j.Slf4j;
import org.apache.curator.framework.CuratorFramework;
import org.apache.zookeeper.AsyncCallback;
import org.apache.zookeeper.CreateMode;
import org.apache.zookeeper.data.Stat;

import java.util.List;

@Slf4j
public class NodeOperation {
    
    public static final String ZK_ADDRESS = "127.0.0.1:2181";

    public static final String ZK_PATH = "/test/CURD/node-1";

    public static final String PARENT_PATH = "/test";


    /**
     * 创建节点
     */
    public void createNode() {
    
        CuratorFramework client = ClientFactory.createSimple(ZK_ADDRESS);

        try {
    
            //启动客户端
            client.start();

            //创建一个ZNode节点,数据为payload
            String data = "hello";
            byte[] payload = data.getBytes("UTF-8");

            /**
             * 1.返回构造者实例
             * 2.如有必要,创建父节点
             * 3.选择节点类型:PERSISTENT、PERSISTENT_SEQUENTIAL、EPHEMERAL、EPHEMERAL_SEQUENTIAL
             * 4.实际创建节点,传入路径与数据
             */
            client.create()
                    .creatingParentsIfNeeded()
                    .withMode(CreateMode.PERSISTENT)
                    .forPath(ZK_PATH, payload);

        } catch (Exception e) {
    
            e.printStackTrace();
        } finally {
    
            client.close();
        }

    }

    /**
     * 读取节点
     */
    public void readNode() {
    
        CuratorFramework client = ClientFactory.createSimple(ZK_ADDRESS);

        try {
    
            client.start();
            //判断节点是否存在
            Stat stat = client.checkExists().forPath(ZK_PATH);
            //节点存在
            if (stat != null) {
    
                //读取节点数据
                byte[] payload = client.getData().forPath(ZK_PATH);
                String data = new String(payload, "UTF-8");
                log.info("read data:" + data);

                //获取父节点的所有子节点,分别打印
                List<String> children = client.getChildren().forPath(PARENT_PATH);
                children.forEach(child -> log.info("child: " + children));
            }

        } catch (Exception e) {
    
            e.printStackTrace();
        } finally {
    
            client.close();
        }
    }

    /**
     * 更新节点
     */
    public void updateNode() {
    
        CuratorFramework client = ClientFactory.createSimple(ZK_ADDRESS);

        try {
    
            client.start();
            String data = "hello world";
            byte[] payload = data.getBytes("UTF-8");

            client.setData().forPath(ZK_PATH, payload);

        } catch (Exception e) {
    
            e.printStackTrace();
        } finally {
    
            client.close();
        }
    }

    /**
     * 异步更新节点
     */
    public void updatNodeAsync() {
    
        CuratorFramework client = ClientFactory.createSimple(ZK_ADDRESS);

        try {
    
            //异步完成更新,回调此实例
            AsyncCallback.StringCallback callback = new AsyncCallback.StringCallback() {
    
                //回调方法
                @Override
                public void processResult(int i, String s, Object o, String s1) {
    
                    log.info(
                            "i = " + i + " | " +
                            "s = " + s + " | " +
                            "o = " + o + " | " +
                            "s1 = " + s1
                    );
                }
            };

            client.start();
            String data = "hello world!";
            byte[] payload = data.getBytes("UTF-8");
            client.setData()
                    .inBackground(callback)
                    .forPath(ZK_PATH, payload);

            Thread.sleep(1000);

        } catch (Exception e) {
    
            e.printStackTrace();
        } finally {
    
            client.start();
        }
    }

    /**
     * 删除节点
     */
    public void deleteNode() {
    
        CuratorFramework client = ClientFactory.createSimple(ZK_ADDRESS);

        try {
    
            client.start();
            //删除节点
            client.delete().forPath(ZK_ADDRESS);

            //查看删除结果
            List<String> children = client.getChildren().forPath(PARENT_PATH);
            for (String child : children) {
    
                log.info("child: " + child);
            }

        } catch (Exception e) {
    
            e.printStackTrace();
        } finally {
    
            client.close();
        }
    }


}

分布式ID生成器

package com.org.zookeeper.distributeIdgenerator;

import com.org.zookeeper.ClientFactory;
import lombok.extern.slf4j.Slf4j;
import org.apache.curator.framework.CuratorFramework;
import org.apache.zookeeper.CreateMode;

@Slf4j
/**
 * id生成器
 */
public class IDMaker {
    

    public static final String ZK_ADDRESS = "";
    CuratorFramework client = ClientFactory.createSimple(ZK_ADDRESS);

    /**
     * 创建临时顺序节点
     *
     * @param pathPerfix 节点路径
     * @return 创建后的完整路径名称
     */
    private String createSeqNode(String pathPerfix) {
    
        try {
    
            client.start();
            String destPath = client.create()
                    .creatingParentsIfNeeded()
                    .withMode(CreateMode.EPHEMERAL_SEQUENTIAL)
                    .forPath(pathPerfix);

            return destPath;
        } catch (Exception e) {
    
            log.error(e.getMessage(), e);
        } finally {
    
            client.close();
        }

        return null;
    }

    public String makeId(String nodeName) {
    
        String str = createSeqNode(nodeName);
        if (null == str) {
    
            return null;
        }

        //取得zk节点的末尾序号
        int index = str.lastIndexOf(nodeName);
        if (index > 0) {
    
            index += str.length();
            return index <= str.length() ? str.substring(index) : "";
        }

        return str;
    }

}


分布式锁

package com.org.zookeeper.distributelock;

import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
import org.apache.zookeeper.AsyncCallback;
import org.apache.zookeeper.CreateMode;
import org.apache.zookeeper.KeeperException;
import org.apache.zookeeper.WatchedEvent;
import org.apache.zookeeper.Watcher;
import org.apache.zookeeper.ZooDefs;
import org.apache.zookeeper.ZooKeeper;
import org.apache.zookeeper.data.Stat;

import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;
@Getter
@Setter
@Slf4j
public class ZkLock implements AsyncCallback.StringCallback, AsyncCallback.Children2Callback, Watcher {
    
    ZooKeeper zk;
    CountDownLatch cd = new CountDownLatch(1);
    String ZK_PATH = "/lock";
    String lockName;
    String threadName;

    public void tryLock() {
    

        try {
    
            zk.create(ZK_PATH, threadName.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL,
                    this, threadName);
            cd.await();
        } catch (Exception e) {
    
            e.printStackTrace();
        }
    }

    public void unlock() {
    
        try {
    
            zk.delete("/" + lockName, -1);
        } catch (InterruptedException e) {
    
            e.printStackTrace();
        } catch (KeeperException e) {
    
            e.printStackTrace();
        }
    }

    //create
    @Override
    public void processResult(int i, String path, Object ctx, String name) {
    
        //单个线程启动后创建锁,然后get锁目录的所有子节点,不注册watch所在目录
        log.info(ctx.toString() + " create path: " + name);
        lockName = name.substring(1);
        zk.getChildren("/", false, this, ctx);
    }

    //getChildren
    @Override
    public void processResult(int rc, String path, Object ctx, List<String> children, Stat stat) {
    
        //获得锁目录的所有有序节点,然后排序,取自己在有序list中的index
        if (children == null) {
    
            log.error(ctx.toString() + "children null");
        } else {
    
            try {
    
                Collections.sort(children);
                int i = children.indexOf(lockName);
                if (i < 1) {
    
                    log.info(threadName + " is first...");
                    zk.setData("/", threadName.getBytes(), -1);
                    cd.countDown();
                } else {
    
                    //监听前一个节点的变化
                    log.info(threadName + " watch " + children.get(i -1));
                    zk.exists("/" + children.get(i - 1), this);
                }
            } catch (KeeperException e) {
    
                e.printStackTrace();
            } catch (InterruptedException e) {
    
                e.printStackTrace();
            }
        }


    }

    //exists
    @Override
    public void process(WatchedEvent watchedEvent) {
    
        Event.EventType type = watchedEvent.getType();
        switch (type) {
    
            case NodeDeleted:
                zk.getChildren("/", false, this, "");
            case NodeChildrenChanged:
                break;
        }

    }
}
版权声明:本文为博主原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接和本声明。
本文链接:https://blog.csdn.net/blade1122/article/details/107737320

智能推荐

while循环&CPU占用率高问题深入分析与解决方案_main函数使用while(1)循环cpu占用99-程序员宅基地

文章浏览阅读3.8k次,点赞9次,收藏28次。直接上一个工作中碰到的问题,另外一个系统开启多线程调用我这边的接口,然后我这边会开启多线程批量查询第三方接口并且返回给调用方。使用的是两三年前别人遗留下来的方法,放到线上后发现确实是可以正常取到结果,但是一旦调用,CPU占用就直接100%(部署环境是win server服务器)。因此查看了下相关的老代码并使用JProfiler查看发现是在某个while循环的时候有问题。具体项目代码就不贴了,类似于下面这段代码。​​​​​​while(flag) {//your code;}这里的flag._main函数使用while(1)循环cpu占用99

【无标题】jetbrains idea shift f6不生效_idea shift +f6快捷键不生效-程序员宅基地

文章浏览阅读347次。idea shift f6 快捷键无效_idea shift +f6快捷键不生效

node.js学习笔记之Node中的核心模块_node模块中有很多核心模块,以下不属于核心模块,使用时需下载的是-程序员宅基地

文章浏览阅读135次。Ecmacript 中没有DOM 和 BOM核心模块Node为JavaScript提供了很多服务器级别,这些API绝大多数都被包装到了一个具名和核心模块中了,例如文件操作的 fs 核心模块 ,http服务构建的http 模块 path 路径操作模块 os 操作系统信息模块// 用来获取机器信息的var os = require('os')// 用来操作路径的var path = require('path')// 获取当前机器的 CPU 信息console.log(os.cpus._node模块中有很多核心模块,以下不属于核心模块,使用时需下载的是

数学建模【SPSS 下载-安装、方差分析与回归分析的SPSS实现(软件概述、方差分析、回归分析)】_化工数学模型数据回归软件-程序员宅基地

文章浏览阅读10w+次,点赞435次,收藏3.4k次。SPSS 22 下载安装过程7.6 方差分析与回归分析的SPSS实现7.6.1 SPSS软件概述1 SPSS版本与安装2 SPSS界面3 SPSS特点4 SPSS数据7.6.2 SPSS与方差分析1 单因素方差分析2 双因素方差分析7.6.3 SPSS与回归分析SPSS回归分析过程牙膏价格问题的回归分析_化工数学模型数据回归软件

利用hutool实现邮件发送功能_hutool发送邮件-程序员宅基地

文章浏览阅读7.5k次。如何利用hutool工具包实现邮件发送功能呢?1、首先引入hutool依赖<dependency> <groupId>cn.hutool</groupId> <artifactId>hutool-all</artifactId> <version>5.7.19</version></dependency>2、编写邮件发送工具类package com.pc.c..._hutool发送邮件

docker安装elasticsearch,elasticsearch-head,kibana,ik分词器_docker安装kibana连接elasticsearch并且elasticsearch有密码-程序员宅基地

文章浏览阅读867次,点赞2次,收藏2次。docker安装elasticsearch,elasticsearch-head,kibana,ik分词器安装方式基本有两种,一种是pull的方式,一种是Dockerfile的方式,由于pull的方式pull下来后还需配置许多东西且不便于复用,个人比较喜欢使用Dockerfile的方式所有docker支持的镜像基本都在https://hub.docker.com/docker的官网上能找到合..._docker安装kibana连接elasticsearch并且elasticsearch有密码

随便推点

Python 攻克移动开发失败!_beeware-程序员宅基地

文章浏览阅读1.3w次,点赞57次,收藏92次。整理 | 郑丽媛出品 | CSDN(ID:CSDNnews)近年来,随着机器学习的兴起,有一门编程语言逐渐变得火热——Python。得益于其针对机器学习提供了大量开源框架和第三方模块,内置..._beeware

Swift4.0_Timer 的基本使用_swift timer 暂停-程序员宅基地

文章浏览阅读7.9k次。//// ViewController.swift// Day_10_Timer//// Created by dongqiangfei on 2018/10/15.// Copyright 2018年 飞飞. All rights reserved.//import UIKitclass ViewController: UIViewController { ..._swift timer 暂停

元素三大等待-程序员宅基地

文章浏览阅读986次,点赞2次,收藏2次。1.硬性等待让当前线程暂停执行,应用场景:代码执行速度太快了,但是UI元素没有立马加载出来,造成两者不同步,这时候就可以让代码等待一下,再去执行找元素的动作线程休眠,强制等待 Thread.sleep(long mills)package com.example.demo;import org.junit.jupiter.api.Test;import org.openqa.selenium.By;import org.openqa.selenium.firefox.Firefox.._元素三大等待

Java软件工程师职位分析_java岗位分析-程序员宅基地

文章浏览阅读3k次,点赞4次,收藏14次。Java软件工程师职位分析_java岗位分析

Java:Unreachable code的解决方法_java unreachable code-程序员宅基地

文章浏览阅读2k次。Java:Unreachable code的解决方法_java unreachable code

标签data-*自定义属性值和根据data属性值查找对应标签_如何根据data-*属性获取对应的标签对象-程序员宅基地

文章浏览阅读1w次。1、html中设置标签data-*的值 标题 11111 222222、点击获取当前标签的data-url的值$('dd').on('click', function() { var urlVal = $(this).data('ur_如何根据data-*属性获取对应的标签对象

推荐文章

热门文章

相关标签