zookeeper实现分布式应用系统服务器上下线动态感知程序、监听机制与守护线程

简介: zookeeper实现分布式应用系统服务器上下线动态感知程序、监听机制与守护线程

需求


在分布式系统中存在多个服务器,这些服务器可以动态上下线,而客户端可以连接任意服务器,但是如果连接的服务器突然下线那么客户端需要重新连接其他服务器,这就需要在服务器上下线的时候客户端能感知,获取哪些可以连接的服务器。


解决思路


每次服务器启动的时候去zookeeper上进行注册(注册规则自由指定,比如简单使用/servers/server001 hostname),而客户端上线就获取服务器列表,并对节点进行监听,一旦有服务器下线那么就能监听到事件从而重新获取服务器列表。

image.png


程序简单实现

服务器端:

/**
 * 服务端程序
 * @author
 *
 */
public class DistributedServer {
    private static final String connectionString = "192.168.47.141:2181";
    public static final Integer sessionTimeout = 2000;
    public static ZooKeeper zkClient = null;
    /**
     * 获取zookeeper连接
     * @throws Exception
     */
    public void getConnection() throws Exception{
        zkClient = new ZooKeeper(connectionString, sessionTimeout, new Watcher(){
            //收到事件通知后的回调函数(应该是我们自己的事件处理逻辑)
            public void process(WatchedEvent event) {
                System.out.println(event.getType()+","+event.getPath());
                try {
                    //为了能一直监听,调用一次注册一次
                    zkClient.getChildren("/", true);
                }catch(Exception e){
                }
            }});
    }
    /**
     * 注册服务器信息
     * @param hostname 注册的服务器名
     * @throws Exception
     */
    public void registerServer(String hostname) throws Exception{
        //创建的是带序号的临时节点 生成的节点像/servers/server000001,/servers/server000002等
        //节点数据即为注册的主机名
        String path = zkClient.create("/server", hostname.getBytes(),
                ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL);
        System.out.println(hostname+" --上线了-- "+path);
    }
    /**
     * 服务器注册完后,执行业务逻辑
     * @param hostname
     * @throws IOException
     */
    public void executeBusiness(String hostname) throws IOException {
        System.out.println(hostname+"开始工作了!");
        System.in.read();
    }
    public static void main(String[] args) throws Exception {
        //获取zookeeper连接
        DistributedServer server = new DistributedServer();
        server.getConnection();
        //服务器上线,完成注册
        Scanner scanner = new Scanner(System.in);
        System.out.println("输入hostname");
        String hostname = scanner.nextLine();
        server.registerServer(hostname);
        //执行业务逻辑
        server.executeBusiness(hostname);
    }
}

客户端:

/*
 * 客户端程序
 */
public class DistributeClient {
    private static final String connectionString = "192.168.47.141:2181";
    public static final Integer sessionTimeout = 2000;
    public static ZooKeeper zkClient = null;
    public static final String parentNode = "/";
    //注意:加volatile的意义何在?使得多线程看到的服务器列表一致而不会拷贝到自己的工作空间
    public volatile List<String> serverList = new ArrayList<String>();
    /**
     * 获取zookeeper连接
     * @throws Exception
     */
    public void getConnection() throws Exception{
        zkClient = new ZooKeeper(connectionString, sessionTimeout, new Watcher(){
            //收到事件通知后的回调函数(应该是我们自己的事件处理逻辑)
            public void process(WatchedEvent event) {
                System.out.println(event.getType()+","+event.getPath());
                try {
                    //重新获取(更新)服务器列表,并进行监听
                    getServerList();
                }catch(Exception e){
                }
            }});
    }
    /**
     * 获取服务器列表信息,并对父节点进行监听
     * @throws Exception
     */
    public void getServerList() throws Exception{
        //获取服务器列表,并对父节点进行监听
        //getChildren()相对于命令行 ls /znode,对子节点进行监听
        List<String> children = zkClient.getChildren(parentNode, true);
        //创建临时集合,将子节点存入
        List<String> childrenList = new ArrayList<String>();
        for (String child : children) {
            byte[] data = zkClient.getData(parentNode+child, false, null);
            childrenList.add(new String(data));
        }
        //将临时集合中的节点赋给服务器列表serverList,以便业务线程使用
        serverList = childrenList;
        System.out.println(serverList);
    }
    /**
     * 业务功能
     * @throws Exception
     */
    public void executeBusiness() throws Exception{
        System.out.println("获取的服务器列表:"+serverList);
        System.out.println("客户端开始工作了...");
        System.in.read();
    }
    public static void main(String[] args) throws Exception {
        //获取zookeeper连接
        DistributeClient client = new DistributeClient();
        client.getConnection();
        //获取服务器列表
        client.getServerList();
        //业务功能
        client.executeBusiness();
    }
}

测试

运行三次服务器端程序,输入的hostname分别为mini1,mini2,mini3当成注册了三个服务器


image.png

image.png



运行客户端程序(可以启动多次,简单起见这里就一次)


image.png


关闭其中2个(mini1,mini2)连接zookeeper的客户端(关闭后注册的服务器也就消失了),查看客户端输出


image.png


一旦服务器下线了,客户端能监听到并且重新获取服务器列表。


目录
相关文章
|
Kubernetes 大数据 调度
Airflow vs Argo Workflows:分布式任务调度系统的“华山论剑”
本文对比了Apache Airflow与Argo Workflows两大分布式任务调度系统。两者均支持复杂的DAG任务编排、社区支持及任务调度功能,且具备优秀的用户界面。Airflow以Python为核心语言,适合数据科学家使用,拥有丰富的Operator库和云服务集成能力;而Argo Workflows基于Kubernetes设计,支持YAML和Python双语定义工作流,具备轻量化、高性能并发调度的优势,并通过Kubernetes的RBAC机制实现多用户隔离。在大数据和AI场景中,Airflow擅长结合云厂商服务,Argo则更适配Kubernetes生态下的深度集成。
1394 34
|
9月前
|
消息中间件 分布式计算 资源调度
《聊聊分布式》ZooKeeper与ZAB协议:分布式协调的核心引擎
ZooKeeper是一个开源的分布式协调服务,基于ZAB协议实现数据一致性,提供分布式锁、配置管理、领导者选举等核心功能,具有高可用、强一致和简单易用的特点,广泛应用于Kafka、Hadoop等大型分布式系统中。
|
10月前
|
存储 算法 安全
“卧槽,系统又崩了!”——别慌,这也许是你看过最通俗易懂的分布式入门
本文深入解析分布式系统核心机制:数据分片与冗余副本实现扩展与高可用,租约、多数派及Gossip协议保障一致性与容错。探讨节点故障、网络延迟等挑战,揭示CFT/BFT容错原理,剖析规模与性能关系,为构建可靠分布式系统提供理论支撑。
450 2
|
10月前
|
机器学习/深度学习 算法 安全
新型电力系统下多分布式电源接入配电网承载力评估方法研究(Matlab代码实现)
新型电力系统下多分布式电源接入配电网承载力评估方法研究(Matlab代码实现)
313 3
|
数据采集 缓存 NoSQL
分布式新闻数据采集系统的同步效率优化实战
本文介绍了一个针对高频新闻站点的分布式爬虫系统优化方案。通过引入异步任务机制、本地缓存池、Redis pipeline 批量写入及身份池策略,系统采集效率提升近两倍,数据同步延迟显著降低,实现了分钟级热点追踪能力,为实时舆情监控与分析提供了高效、稳定的数据支持。
533 1
分布式新闻数据采集系统的同步效率优化实战
|
存储 SpringCloudAlibaba Java
【SpringCloud Alibaba系列】一文全面解析Zookeeper安装、常用命令、JavaAPI操作、Watch事件监听、分布式锁、集群搭建、核心理论
一文全面解析Zookeeper安装、常用命令、JavaAPI操作、Watch事件监听、分布式锁、集群搭建、核心理论。
【SpringCloud Alibaba系列】一文全面解析Zookeeper安装、常用命令、JavaAPI操作、Watch事件监听、分布式锁、集群搭建、核心理论
|
存储 运维 安全
盘古分布式存储系统的稳定性实践
本文介绍了阿里云飞天盘古分布式存储系统的稳定性实践。盘古作为阿里云的核心组件,支撑了阿里巴巴集团的众多业务,确保数据高可靠性、系统高可用性和安全生产运维是其关键目标。文章详细探讨了数据不丢不错、系统高可用性的实现方法,以及通过故障演练、自动化发布和健康检查等手段保障生产安全。总结指出,稳定性是一项系统工程,需要持续迭代演进,盘古经过十年以上的线上锤炼,积累了丰富的实践经验。
1482 7
|
存储 分布式计算 Hadoop
基于Java的Hadoop文件处理系统:高效分布式数据解析与存储
本文介绍了如何借鉴Hadoop的设计思想,使用Java实现其核心功能MapReduce,解决海量数据处理问题。通过类比图书馆管理系统,详细解释了Hadoop的两大组件:HDFS(分布式文件系统)和MapReduce(分布式计算模型)。具体实现了单词统计任务,并扩展支持CSV和JSON格式的数据解析。为了提升性能,引入了Combiner减少中间数据传输,以及自定义Partitioner解决数据倾斜问题。最后总结了Hadoop在大数据处理中的重要性,鼓励Java开发者学习Hadoop以拓展技术边界。
617 7
|
存储 运维 负载均衡
构建高可用性GraphRAG系统:分布式部署与容错机制
【10月更文挑战第28天】作为一名数据科学家和系统架构师,我在构建和维护大规模分布式系统方面有着丰富的经验。最近,我负责了一个基于GraphRAG(Graph Retrieval-Augmented Generation)模型的项目,该模型用于构建一个高可用性的问答系统。在这个过程中,我深刻体会到分布式部署和容错机制的重要性。本文将详细介绍如何在生产环境中构建一个高可用性的GraphRAG系统,包括分布式部署方案、负载均衡、故障检测与恢复机制等方面的内容。
897 4
构建高可用性GraphRAG系统:分布式部署与容错机制
|
存储 运维 NoSQL
分布式读写锁的奥义:上古世代 ZooKeeper 的进击
本文作者将介绍女娲对社区 ZooKeeper 在分布式读写锁实践细节上的思考,希望帮助大家理解分布式读写锁背后的原理。
476 11

热门文章

最新文章