【问题标题】:Broadcasting using the protocol Zab in ZooKeeper在 ZooKeeper 中使用 Zab 协议进行广播
【发布时间】:2014-02-20 22:16:56
【问题描述】:

早安,

我是 ZooKeeper 及其协议的新手,我对它的广播协议 Zab 很感兴趣。

能否提供一个简单的使用 Zookeeper 的 Zab 协议的 java 代码?我一直在搜索,但我没有成功找到显示如何使用 Zab 的代码。

实际上我需要的很简单,我有一个 MapReduce 代码,我希望所有映射器在成功找到更好的 X 值(即更大的值)时更新一个变量(比如说 X)。在这种情况下,领导者必须比较旧值和新值,然后将实际的最佳值广播给所有映射器。我怎么能在Java中做这样的事情?

提前致谢, 问候

【问题讨论】:

    标签: hadoop mapreduce apache-zookeeper


    【解决方案1】:

    您不需要使用 Zab 协议。相反,您可以按照以下步骤操作:

    你在 Zookeeper 上有一个 Znode 说 /bigvalue。所有映射器在启动时都会读取存储在其中的值。他们还监视 Znode 上的数据变化。每当映射器获得更好的值时,它都会用更好的值更新 Znode。所有映射器都会收到数据更改事件的通知,他们会读取新的最佳值并重新建立数据更改的监视。这样,它们与最新的最佳值同步,并且可以在有更好的值时更新最新的最佳值。

    实际上 zkclient 是一个非常好的库,可以与 Zookeeper 一起使用,它隐藏了很多复杂性(https://github.com/sgroschupf/zkclient)。下面是一个示例,演示了如何观察 Znode“/bigvalue”的任何数据更改。

    package geet.org;
    
    import java.io.UnsupportedEncodingException;
    import org.I0Itec.zkclient.IZkDataListener;
    import org.I0Itec.zkclient.ZkClient;
    import org.I0Itec.zkclient.exception.ZkMarshallingError;
    import org.I0Itec.zkclient.exception.ZkNodeExistsException;
    import org.I0Itec.zkclient.serialize.ZkSerializer;
    import org.apache.zookeeper.data.Stat;
    
    public class ZkExample implements IZkDataListener, ZkSerializer {
        public static void main(String[] args) {
            String znode = "/bigvalue";
            ZkExample ins = new ZkExample();
            ZkClient cl = new ZkClient("127.0.0.1", 30000, 30000,
                    ins);
            try {
                cl.createPersistent(znode);
            } catch (ZkNodeExistsException e) {
                System.out.println(e.getMessage());
            }
            // Change the data for fun
            Stat stat = new Stat();
            String data =  cl.readData(znode, stat);
            System.out.println("Current data " + data + "version = " + stat.getVersion());
            cl.writeData(znode, "My new data ", stat.getVersion());
    
            cl.subscribeDataChanges(znode, ins);
            try {
                Thread.sleep(36000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    
        @Override
        public void handleDataChange(String dataPath, Object data) throws Exception {
            System.out.println("Detected data change");
            System.out.println("New data for " + dataPath + " " + (String)data);
        }
    
        @Override
        public void handleDataDeleted(String dataPath) throws Exception {
            System.out.println("Data deleted " + dataPath);
        }
    
        @Override
        public byte[] serialize(Object data) throws ZkMarshallingError {
            if (data instanceof String){
                try {
                    return ((String) data).getBytes("UTF-8");
                } catch (UnsupportedEncodingException e) {
                    e.printStackTrace();
                }
            }
            return null;
        }
    
        @Override
        public Object deserialize(byte[] bytes) throws ZkMarshallingError {
            try {
                return new String(bytes, "UTF-8");
            } catch (UnsupportedEncodingException e) {
                e.printStackTrace();
            }
            return null;
        }
    }
    

    【讨论】:

    • 感谢您的快速回答,这是否意味着映射器将自动获得数据更改事件的通知?而我所要做的就是读取值并监视 Znode 上的数据变化?另一个问题:我是ZooKeeper的初学者,在创建Znode/bigvalue时,会保存在哪里?所有映射器如何访问它?它像 HDFS 吗?抱歉问了这么多问题
    • 正确。他们会自动收到通知。
    • 好的,谢谢。对于 ZooKeeper 的安装,我是否需要在我拥有的所有 Hadoop 集群的机器上安装它?正如您的提议一样,我所需要的只是一个 Znode ! ZooKeeper 将这些 Znode 保存在哪里?我仍然很难理解!
    • Znode 是 Zookeeper 数据树中的一个节点,不是机器。如果高可用性很重要,您可能希望拥有一组 Zookeeper 节点。
    • 好的好的。所以 Znode = 数据树中的一个节点,ZooKeeper 节点 = machine。非常感谢您的所有回答。
    猜你喜欢
    • 1970-01-01
    • 2017-05-16
    • 1970-01-01
    • 1970-01-01
    • 2017-01-18
    • 2020-06-09
    • 2010-10-10
    • 1970-01-01
    • 2016-07-30
    相关资源
    最近更新 更多