ZooKeeper Overview
ZooKeeper is an Apache Hadoop subproject that provides a hierarchical directory service. It functions as a distributed, open-source coordination service for distributed applications.
Primary Capabilities: Configuration Management, Distributed Locking, Cluster Management
ZooKeeper Data Model
ZooKeeper maintains a tree-structured namespace similar to Unix filesystem hierarchies. Each unit in this structure is called a ZNode, which stores both its own data and metadata. ZNodes can contain child nodes and support storage of up to 1MB of data per node.
Node Types
| Type | Description | Flag |
|---|---|---|
| PERSISTENT | Persistent node, survives client disconnection | default |
| EPHEMERAL | Temporary node, auto-deleted on client disconnnect | -e |
| PERSISTENT_SEQUENTIAL | Persistent node with auto-incrementing suffix | -s |
| EPHEMERAL_SEQUENTIAL | Temporary node with auto-incrementing suffix | -es |
Server Commands
# Start ZooKeeper service
./zkServer.sh start
# Check service status
./zkServer.sh status
# Stop service
./zkServer.sh stop
# Restart service
./zkServer.sh restart
Client Commands
# Connect to server
./zkCli.sh -server 192.168.1.100:2181
# Disconnect
quit
# Show command help
help
# List directory contents
ls /path
# Create node with value
create /nodepath value
# Retrieve node value
get /nodepath
# Update node value
set /nodepath value
# Delete single node
delete /nodepath
# Delete node with children
deleteall /nodepath
Temporary and Sequential Nodes
# Create ephemeral node
create -e /temppath value
# Create sequential node
create -s /seqpath value
# Display detailed node information
ls -s /nodepath
Node Stat Fields
- czxid: Transaction ID when node was created
- ctime: Creation timestamp
- mzxid: Last modification transaction ID
- mtime: Last modification timestamp
- pzxid: Last child list modification transaction ID
- cversion: Child node version count
- dataversion: Data version number
- aclversion: Permission version number
- ephemeralOwner: Owner transaction ID for temporary nodes (0 for persistent)
- dataLength: Size of stored data
- numChildren: Number of direct children
Java Client Integration with Curator
Curator is Apache ZooKeeper's official Java client library, designed to simplify ZooKeeper usage. It evolved from a Netflix project and became an Apache top-level project.
Available Clients:
- Native ZooKeeper API
- ZkClient
- Curator (recommended)
Establishing Connection
public class ZkConnectionTest {
private CuratorFramework client;
@Before
public void setup() {
RetryPolicy retryPolicy = new ExponentialBackoffRetry(3000, 10);
client = CuratorFrameworkFactory.builder()
.connectString("192.168.1.100:2181")
.sessionTimeoutMs(60 * 1000)
.connectionTimeoutMs(15 * 1000)
.retryPolicy(retryPolicy)
.namespace("production")
.build();
client.start();
}
}
Node Operations
public class NodeOpsTest {
@Test
public void createBasicNode() throws Exception {
String path = client.create().forPath("/service", "data".getBytes());
System.out.println(path);
}
@Test
public void createWithData() throws Exception {
String path = client.create().forPath("/config", "settings".getBytes());
System.out.println(path);
}
@Test
public void createEphemeralNode() throws Exception {
String path = client.create().withMode(CreateMode.EPHEMERAL).forPath("/session");
System.out.println(path);
}
@Test
public void createNestedNodes() throws Exception {
String path = client.create()
.creatingParentsIfNeeded()
.forPath("/department/team/member");
System.out.println(path);
}
}
Query Operations
public class QueryOpsTest {
@Test
public void fetchNodeData() throws Exception {
byte[] payload = client.getData().forPath("/config");
System.out.println(new String(payload));
}
@Test
public void listChildren() throws Exception {
List<String> children = client.getChildren().forPath("/");
System.out.println(children);
}
@Test
public void fetchNodeWithStats() throws Exception {
Stat metadata = new Stat();
client.getData().storingStatIn(metadata).forPath("/config");
System.out.println("Version: " + metadata.getVersion());
System.out.println("Created: " + metadata.getCtime());
}
}
Update Operations
public class UpdateOpsTest {
@Test
public void simpleUpdate() throws Exception {
client.setData().forPath("/config", "newvalue".getBytes());
}
@Test
public void versionedUpdate() throws Exception {
Stat status = new Stat();
client.getData().storingStatIn(status).forPath("/config");
int currentVersion = status.getVersion();
client.setData()
.withVersion(currentVersion)
.forPath("/config", "updated".getBytes());
}
}
Delete Operations
public class DeleteOpsTest {
@Test
public void removeSingleNode() throws Exception {
client.delete().forPath("/temp");
}
@Test
public void removeWithChildren() throws Exception {
client.delete().deletingChildrenIfNeeded().forPath("/branch");
}
@Test
public void guaranteedDelete() throws Exception {
client.delete().guaranteed().forPath("/critical");
}
@Test
public void asyncDelete() throws Exception {
client.delete().guaranteed().inBackground(new BackgroundCallback() {
@Override
public void processResult(CuratorFramework client, CuratorEvent event) {
System.out.println("Deletion completed: " + event.getPath());
}
}).forPath("/temp");
}
}
Watch Mechanism Overview
ZooKeeper's watch feature enables clients to register watchers on specific nodes, receiving notifications when events occur. This forms the foundation for publish-subscribe functionality in distributed systems.
Limitation of Native API: Requires manual re-registration after each triggered event.
Curator Cache Solutions:
| Cache Type | Scope |
|---|---|
| NodeCache | Monitors a single specific node |
| PathChildrenCache | Monitors all children of a node |
| TreeCache | Monitors entire subtree including the root |
NodeCache Implementation
public class SingleNodeWatcherTest {
@Test
public void monitorNodeChanges() throws Exception {
final NodeCache cache = new NodeCache(client, "/config");
cache.getListenable().addListener(new NodeCacheListener() {
@Override
public void nodeChanged() throws Exception {
System.out.println("Node updated");
byte[] updated = cache.getCurrentData().getData();
System.out.println("New data: " + new String(updated));
}
});
cache.start(true);
Thread.sleep(Long.MAX_VALUE);
}
}
PathChildrenCache Implementation
public class ChildNodeWatcherTest {
@Test
public void monitorChildChanges() throws Exception {
PathChildrenCache watcher = new PathChildrenCache(client, "/services", true);
watcher.getListenable().addListener(new PathChildrenCacheListener() {
@Override
public void childEvent(CuratorFramework client, PathChildrenCacheEvent event) {
System.out.println("Child event: " + event.getType());
if (event.getType() == PathChildrenCacheEvent.Type.CHILD_UPDATED) {
byte[] payload = event.getData().getData();
System.out.println("Changed: " + new String(payload));
}
}
});
watcher.start();
Thread.sleep(Long.MAX_VALUE);
}
}
TreeCache Implementation
public class SubtreeWatcherTest {
@Test
public void monitorEntireTree() throws Exception {
TreeCache watcher = new TreeCache(client, "/root");
watcher.getListenable().addListener(new TreeCacheListener() {
@Override
public void childEvent(CuratorFramework client, TreeCacheEvent event) {
System.out.println("Tree event: " + event.getType());
}
});
watcher.start();
Thread.sleep(Long.MAX_VALUE);
}
}
Distributed Locking with ZooKeeper
Problem Statement
In distributed environments with multiple JVMs, traditional threading locks (synchronized, Lock) cannot coordinate across machines. ZooKeeper provides a mechanism for cross-machine process synchronization.
Lock Acquisition Algorithm
- Client creates an ephemeral sequential node under the lock znode
- Client retrieves all child nodes and checks if its own node has the lowest sequence number
- If lowest: lock acquired → after use, delete the node
- If not lowest: register a watcher on the predecessor node and wait
- When predecessor is deleted, re-check sequence numbers
public class DistributedLock {
private static final String LOCK_PATH = "/distributed-lock";
private final CuratorFramework zkClient;
public void acquireLock() throws Exception {
String nodePath = zkClient.create()
.withMode(CreateMode.EPHEMERAL_SEQUENTIAL)
.forPath(LOCK_PATH + "/lock-", "".getBytes());
while (true) {
List<String> children = zkClient.getChildren().forPath(LOCK_PATH);
Collections.sort(children);
if (nodePath.endsWith(children.get(0))) {
return;
}
String predecessor = LOCK_PATH + "/" + children.get(
children.indexOf(getNodeName(nodePath)) - 1);
CountDownLatch latch = new CountDownLatch(1);
PathChildrenCache cache = new PathChildrenCache(zkClient, predecessor, true);
cache.getListenable().addListener((client, event) -> latch.countDown());
cache.start();
latch.await();
}
}
}
Leader Election in ZooKeeper Cluster
Election Criteria
Server ID: Each server in the cluster has a unique ID. Higher IDs carry greater weight in the election algorithm.
Zxid (Transaction ID): Represents the highest data version held by a server. Servers with more recent data have higher priority. A server needs majority votes to become leader.
Cluster Failure Scenarios
Scenario 1: With 3-node cluster, if 2 followers fail, the remaining follower cannot continue operations since it lacks majority support (less than 50% of cluster size).
Scenario 2: When the leader fails, surviving servers initiate election and produce a new leader.
Scenario 3: New servers joining an active cluster do not disrupt the current leader's operation.
Server Roles
Leader
- Processes all write requests (transaction handling)
- Coordinates internal cluster operations
Follower
- Handles read requests and forwards write requests to leader
- Participates in leader election voting
Observer
- Handles read requests and forwards writes to leader
- Does not participate in election (reduces voting overhead)
# Observer configuration in zoo.cfg
peerType=observer
server.1=host1:2888:3888:observer
server.2=host2:2888:3888:observer
server.3=host3:2888:3888