利用ZooKeeper开发分布式应用系统案例--服务端与客户端实现

本文主要是介绍利用ZooKeeper开发分布式应用系统案例--服务端与客户端实现,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

利用ZooKeeper开发分布式应用系统案例--服务端与客户端实现

服务端代码:

package cn.edu360.zk.distributesystem;import java.io.IOException;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.Ids;
import org.apache.zookeeper.ZooKeeper;
import org.apache.zookeeper.data.Stat;public class TimeQueryServer {ZooKeeper zk = null;//启动zk客户端连接public void connectZK() throws Exception {zk = new ZooKeeper("hadoop1:2181,hadoop2:2181,hadoop3:2181", 2000, null);}//注册服务器信息public void registerServerInfo(String hostname,String port) throws Exception, InterruptedException {/** 先判断注册节点是否存在,如果不存在,则创建*/Stat stat = zk.exists("/servers", false);if(stat == null) {zk.create("/servers", null, Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);}//注册服务器数据到zk的约定注册节点下String create = zk.create("/servers/server", (hostname + ":" + port).getBytes(), Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL);System.out.println(hostname + "服务器向zk注册信息成功,注册的节点为:" + create);}//启动业务线程开始处理业务public static void main(String[] args) throws Exception, Exception {TimeQueryServer timeQueryServer = new TimeQueryServer();//启动zk客户端连接timeQueryServer.connectZK();//注册服务器信息timeQueryServer.registerServerInfo(args[0], args[1]);//启动业务线程开始处理业务new TimeQueryService(Integer.parseInt(args[1])).start();}}

服务端线程代码:

package cn.edu360.zk.distributesystem;import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.net.ServerSocket;
import java.net.Socket;
import java.util.Date;public class TimeQueryService extends Thread{int port = 0;public TimeQueryService(int port) {this.port = port;}@Overridepublic void run() {try {ServerSocket ss = new ServerSocket(port);System.out.println("业务线程已绑定端口"+ port + "准备接受消费端请求了...");while(true) {Socket sc = ss.accept();InputStream inputStream = sc.getInputStream();OutputStream outputStream = sc.getOutputStream();outputStream.write(new Date().toString().getBytes());}} catch (IOException e) {// TODO Auto-generated catch blocke.printStackTrace();}}}

客户端代码:

package cn.edu360.zk.distributesystem;import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.net.ServerSocket;
import java.net.Socket;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;import org.apache.zookeeper.KeeperException;
import org.apache.zookeeper.WatchedEvent;
import org.apache.zookeeper.Watcher;
import org.apache.zookeeper.Watcher.Event.EventType;
import org.apache.zookeeper.Watcher.Event.KeeperState;
import org.apache.zookeeper.ZooKeeper;public class Consumer {//定义一个list用于存放最新的在线服务器列表private volatile ArrayList<String> onlineServers = new ArrayList<String>();//构造zk连接对象ZooKeeper zk = null;public void connectZK() throws Exception {zk = new ZooKeeper("hadoop1:2181,hadoop2:2181,hadoop3:2181", 2000, new Watcher() {@Overridepublic void process(WatchedEvent event) {if(event.getState() == KeeperState.SyncConnected && event.getType() == EventType.NodeChildrenChanged) {try {//事件回调逻辑中,再次查询zk上的在线服务器节点即可,查询逻辑中又再次注册子节点事件监听。getOnlineServers();} catch (Exception e) {// TODO Auto-generated catch blocke.printStackTrace();}}}});}//查询在线服务器列表public void getOnlineServers() throws Exception, InterruptedException {List<String> children = zk.getChildren("/servers", true);ArrayList<String> list = new ArrayList<String>();for (String child : children) {byte[] data = zk.getData("/servers/"+child, false, null);String serverInfo = new String(data);list.add(serverInfo);}onlineServers = list;System.out.println("查询了一次zk,当前在线的服务器有:"+list);}public void sendRequest() throws Exception {Random random = new Random();while(true) {try {//挑选一台当前在线的服务器	int nextInt = random.nextInt(onlineServers.size());String server = onlineServers.get(nextInt);String hostname = server.split(":")[0];int port = Integer.parseInt(server.split(":")[1]);System.out.println("本次请求挑选的服务器为:" + server);Socket socket = new Socket(hostname, port);OutputStream outputStream = socket.getOutputStream();outputStream.write("haha".getBytes());outputStream.flush();InputStream inputStream = socket.getInputStream();byte[] buf = new byte[256];int read = inputStream.read(buf);System.out.println("服务器相应的时间为" + new String(buf,0,read));outputStream.close();inputStream.close();socket.close();Thread.sleep(2000);}catch(Exception e) {e.printStackTrace();}}}public static void main(String[] args) throws Exception {Consumer consumer = new Consumer();//构造zk连接对象consumer.connectZK();//查询在线服务器列表consumer.getOnlineServers();//处理业务(向一台服务器发送时间查询请求)consumer.sendRequest();}}

这篇关于利用ZooKeeper开发分布式应用系统案例--服务端与客户端实现的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/176901

相关文章

Python位移操作和位运算的实现示例

《Python位移操作和位运算的实现示例》本文主要介绍了Python位移操作和位运算的实现示例,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一... 目录1. 位移操作1.1 左移操作 (<<)1.2 右移操作 (>>)注意事项:2. 位运算2.1

如何在 Spring Boot 中实现 FreeMarker 模板

《如何在SpringBoot中实现FreeMarker模板》FreeMarker是一种功能强大、轻量级的模板引擎,用于在Java应用中生成动态文本输出(如HTML、XML、邮件内容等),本文... 目录什么是 FreeMarker 模板?在 Spring Boot 中实现 FreeMarker 模板1. 环

Qt实现网络数据解析的方法总结

《Qt实现网络数据解析的方法总结》在Qt中解析网络数据通常涉及接收原始字节流,并将其转换为有意义的应用层数据,这篇文章为大家介绍了详细步骤和示例,感兴趣的小伙伴可以了解下... 目录1. 网络数据接收2. 缓冲区管理(处理粘包/拆包)3. 常见数据格式解析3.1 jsON解析3.2 XML解析3.3 自定义

SpringMVC 通过ajax 前后端数据交互的实现方法

《SpringMVC通过ajax前后端数据交互的实现方法》:本文主要介绍SpringMVC通过ajax前后端数据交互的实现方法,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价... 在前端的开发过程中,经常在html页面通过AJAX进行前后端数据的交互,SpringMVC的controll

Java Stream流使用案例深入详解

《JavaStream流使用案例深入详解》:本文主要介绍JavaStream流使用案例详解,本文通过实例代码给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下吧... 目录前言1. Lambda1.1 语法1.2 没参数只有一条语句或者多条语句1.3 一个参数只有一条语句或者多

Spring Security自定义身份认证的实现方法

《SpringSecurity自定义身份认证的实现方法》:本文主要介绍SpringSecurity自定义身份认证的实现方法,下面对SpringSecurity的这三种自定义身份认证进行详细讲解,... 目录1.内存身份认证(1)创建配置类(2)验证内存身份认证2.JDBC身份认证(1)数据准备 (2)配置依

利用python实现对excel文件进行加密

《利用python实现对excel文件进行加密》由于文件内容的私密性,需要对Excel文件进行加密,保护文件以免给第三方看到,本文将以Python语言为例,和大家讲讲如何对Excel文件进行加密,感兴... 目录前言方法一:使用pywin32库(仅限Windows)方法二:使用msoffcrypto-too

C#使用StackExchange.Redis实现分布式锁的两种方式介绍

《C#使用StackExchange.Redis实现分布式锁的两种方式介绍》分布式锁在集群的架构中发挥着重要的作用,:本文主要介绍C#使用StackExchange.Redis实现分布式锁的... 目录自定义分布式锁获取锁释放锁自动续期StackExchange.Redis分布式锁获取锁释放锁自动续期分布式

springboot使用Scheduling实现动态增删启停定时任务教程

《springboot使用Scheduling实现动态增删启停定时任务教程》:本文主要介绍springboot使用Scheduling实现动态增删启停定时任务教程,具有很好的参考价值,希望对大家有... 目录1、配置定时任务需要的线程池2、创建ScheduledFuture的包装类3、注册定时任务,增加、删

SpringBoot整合mybatisPlus实现批量插入并获取ID详解

《SpringBoot整合mybatisPlus实现批量插入并获取ID详解》这篇文章主要为大家详细介绍了SpringBoot如何整合mybatisPlus实现批量插入并获取ID,文中的示例代码讲解详细... 目录【1】saveBATch(一万条数据总耗时:2478ms)【2】集合方式foreach(一万条数