ZooKeeper 데이터 저장 구조
Apache ZooKeeper는 파일 시스템과 유사한 계층적 트리 구조의 ZNode를 기반으로 데이터를 저장합니다. 이는 마치 인메모리 데이터베이스처럼 동작하며, ZNode는 경로, 데이터, 접근 권한 정보 등을 포함합니다. ZooKeeper는 이러한 모든 정보를 주기적으로 디스크에 저장하여 영속성을 확보합니다.
내부적으로는 DataTree 객체를 사용하여 전체 데이터 저장소를 관리하며, 주로 노드 경로와 그 내용을 저장합니다. 이 DataTree는 ConcurrentHashMap과 같은 동시성 컬렉션을 기반으로 구현되어 있습니다.
데이터 저장 메커니즘
트랜잭션 로그 (Transaction Log)
모든 노드 생성, 수정, 삭제와 같은 트랜잭션 작업은 트랜잭션 로그로 기록됩니다. ZooKeeper 설정 파일(zoo.cfg)에 명시된 datadir 경로에 저장됩니다. 이러한 로그의 디스크 I/O 성능은 ZooKeeper 자체의 처리 성능에 직접적인 영향을 미치므로, 실제 운영 환경에서는 트랜잭션 로그를 전용 디스크에 분리하여 저장하는 것이 일반적입니다.
스냅샷 로그 (Snapshot Log)
스냅샷 로그는 ZooKeeper의 특정 시점 전체 데이터(트리 구조의 모든 ZNode 정보)를 기록한 파일입니다. 데이터 백업과 유사한 개념으로, 트랜잭션 로그와 함께 datadir 경로에 저장됩니다. 이를 통해 시스템 재시작 시 빠른 데이터 복구를 가능하게 합니다.
런타임 로그 (Runtime Log)
ZooKeeper 서버의 운영 상태나 발생한 이벤트는 bin/zookeeper.out 파일에 기록됩니다. 이 로그를 통해 ZooKeeper 인스턴스의 동작 상황을 모니터링하고 문제 발생 시 진단할 수 있습니다.
Java API를 이용한 ZooKeeper 기본 사용
ZooKeeper는 커맨드 라인 인터페이스(zkCli.sh)를 통해 조작할 수도 있지만, Java 애플리케이션에서는 공식 Java API를 사용하여 프로그래밍 방식으로 상호작용합니다.
의존성 추가
Maven 프로젝트의 경우, pom.xml 파일에 다음 의존성을 추가합니다.
<dependency>
<groupId>org.apache.zookeeper</groupId>
<artifactId>zookeeper</artifactId>
<version>3.8.0</version> <!-- 최신 안정 버전으로 업데이트하세요 -->
</dependency>
ZooKeeper 연결 설정
다음은 ZooKeeper 앙상블(클러스터)에 연결하는 기본적인 Java 코드 예시입니다. 초기 연결 상태를 확인하고 세션을 닫는 과정을 보여줍니다.
import org.apache.zookeeper.ZooKeeper;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
public class ZkBasicConnector {
private static final String ZK_SERVER_LIST = "192.168.136.128:2181,192.168.136.129:2181,192.168.136.130:2181";
private static final int SESSION_TIMEOUT_MS = 4000;
public static void main(String[] args) {
try {
// ZooKeeper 객체 생성 (Watcher는 null로 설정)
ZooKeeper zkClient = new ZooKeeper(ZK_SERVER_LIST, SESSION_TIMEOUT_MS, null);
System.out.println("초기 연결 상태: " + zkClient.getState());
// 연결이 완료될 때까지 잠시 대기
TimeUnit.SECONDS.sleep(2);
System.out.println("2초 후 연결 상태: " + zkClient.getState());
// 연결 종료
zkClient.close();
System.out.println("ZooKeeper 연결 종료.");
} catch (IOException | InterruptedException e) {
System.err.println("연결 중 오류 발생: " + e.getMessage());
e.printStackTrace();
}
}
}
위 코드를 실행하면 초기에는 CONNECTING 상태였다가, 잠시 후 CONNECTED 상태로 변경되는 것을 볼 수 있습니다. 그러나 Thread.sleep() 방식은 연결의 정확한 완료 시점을 보장하지 못합니다.
Watcher를 활용한 연결 성공 보장
정확히 연결 성공 시점을 확인하고 다음 작업을 진행하기 위해서는 Watcher 메커니즘을 사용합니다. CountDownLatch와 함께 SyncConnected 이벤트를 기다림으로써 안정적인 연결을 확보할 수 있습니다.
import org.apache.zookeeper.Watcher;
import org.apache.zookeeper.ZooKeeper;
import java.io.IOException;
import java.util.concurrent.CountDownLatch;
public class ZkReliableConnector {
private static final String ZK_SERVER_LIST = "192.168.136.128:2181,192.168.136.129:2181,192.168.136.130:2181";
private static final int SESSION_TIMEOUT_MS = 4000;
private static CountDownLatch connectionLatch = new CountDownLatch(1);
public static void main(String[] args) {
try {
ZooKeeper zkClient = new ZooKeeper(ZK_SERVER_LIST, SESSION_TIMEOUT_MS, event -> {
// 연결 이벤트 처리 Watcher
if (Watcher.Event.KeeperState.SyncConnected == event.getState()) {
System.out.println("ZooKeeper 연결 성공!");
connectionLatch.countDown(); // 연결 성공 시 래치 감소
} else {
System.out.println("ZooKeeper 연결 상태 변경: " + event.getState());
}
});
// 연결이 성공할 때까지 대기
connectionLatch.await();
System.out.println("최종 연결 상태: " + zkClient.getState()); // CONNECTED 출력
// 예제 로직 수행 후 연결 종료
// ... (CRUD 작업 등)
zkClient.close();
} catch (IOException | InterruptedException e) {
System.err.println("연결 중 예외 발생: " + e.getMessage());
e.printStackTrace();
}
}
}
ZNode 생성, 조회, 수정, 삭제 (CRUD)
ZooKeeper Java API를 사용하여 ZNode를 생성(Create), 조회(Read), 수정(Update), 삭제(Delete)하는 예시입니다. Stat 객체를 통해 ZNode의 메타데이터를 관리하며, 특히 setData 시에는 버전 정보를 활용하여 낙관적 락(optimistic locking)을 구현할 수 있습니다.
import org.apache.zookeeper.*;
import org.apache.zookeeper.data.Stat;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
public class ZkCrudOperations {
private static final String ZK_SERVER_LIST = "192.168.136.128:2181,192.168.136.129:2181,192.168.136.130:2181";
private static final int SESSION_TIMEOUT_MS = 4000;
private static CountDownLatch connectionLatch = new CountDownLatch(1);
private static ZooKeeper zkClient;
public static void main(String[] args) {
try {
zkClient = new ZooKeeper(ZK_SERVER_LIST, SESSION_TIMEOUT_MS, event -> {
if (Watcher.Event.KeeperState.SyncConnected == event.getState()) {
connectionLatch.countDown();
}
});
connectionLatch.await();
System.out.println("ZooKeeper 연결 완료. 상태: " + zkClient.getState());
String testNodePath = "/my-app-config";
String initialData = "initial-value";
String updatedData = "new-value-updated";
// 1. ZNode 생성 (PERSISTENT 모드)
// Acl: OPEN_ACL_UNSAFE는 모든 권한을 허용
zkClient.create(testNodePath, initialData.getBytes(StandardCharsets.UTF_8),
ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
System.out.println("ZNode '" + testNodePath + "' 생성됨. 초기 데이터: " + initialData);
// 2. ZNode 데이터 조회
Stat nodeStat = new Stat();
byte[] dataBytes = zkClient.getData(testNodePath, null, nodeStat);
System.out.println("ZNode '" + testNodePath + "' 데이터 조회: " + new String(dataBytes, StandardCharsets.UTF_8));
System.out.println("현재 버전: " + nodeStat.getVersion());
TimeUnit.SECONDS.sleep(1); // 작업 간격
// 3. ZNode 데이터 수정 (현재 버전 사용)
zkClient.setData(testNodePath, updatedData.getBytes(StandardCharsets.UTF_8), nodeStat.getVersion());
System.out.println("ZNode '" + testNodePath + "' 데이터 수정됨. 새 데이터: " + updatedData);
// 수정된 데이터 및 버전 재조회
dataBytes = zkClient.getData(testNodePath, null, nodeStat);
System.out.println("ZNode '" + testNodePath + "' 수정 후 데이터: " + new String(dataBytes, StandardCharsets.UTF_8));
System.out.println("새로운 버전: " + nodeStat.getVersion());
TimeUnit.SECONDS.sleep(1); // 작업 간격
// 4. ZNode 삭제 (현재 버전 사용)
zkClient.delete(testNodePath, nodeStat.getVersion());
System.out.println("ZNode '" + testNodePath + "' 삭제됨.");
// 삭제 확인 (NodeExistsException 발생 예상)
System.out.println("삭제 확인: " + (zkClient.exists(testNodePath, false) == null ? "ZNode 없음" : "ZNode 존재"));
} catch (IOException | InterruptedException | KeeperException e) {
System.err.println("ZooKeeper 작업 중 예외 발생: " + e.getMessage());
e.printStackTrace();
} finally {
if (zkClient != null) {
try {
zkClient.close();
} catch (InterruptedException e) {
System.err.println("ZooKeeper 클라이언트 종료 중 예외 발생: " + e.getMessage());
}
}
}
}
}
ZooKeeper Watcher 메커니즘 심층 분석
Watcher 메커니즘은 ZooKeeper의 핵심 기능 중 하나로, 분산 환경에서 발생하는 데이터 변경을 클라이언트에게 비동기적으로 통지하는 발행-구독(Publish-Subscribe) 모델을 제공합니다. 클라이언트는 특정 ZNode에 Watcher를 등록하여 해당 노드의 데이터나 자식 노드 변경 이벤트를 구독할 수 있습니다.
Watcher의 특징
- 일회성 트리거(One-time Trigger): Watcher는 한 번만 이벤트를 통지합니다. 이벤트가 발생하여 클라이언트에게 통지되면 해당 Watcher는 즉시 소멸됩니다. 지속적인 모니터링을 위해서는 클라이언트가 이벤트를 받은 후 Watcher를 재등록해야 합니다.
- 이벤트 전송 보장: ZooKeeper는 데이터 변경이 발생하면 해당 Watcher를 등록한 모든 클라이언트에게 이벤트를 전송합니다.
- 비동기성: 이벤트 통지는 비동기적으로 처리됩니다. 클라이언트는 이벤트가 발생했을 때 즉시 응답을 받는 것이 아니라, 별도의 스레드에서 이벤트가 처리됩니다.
Watcher 등록 방식
다음 ZooKeeper API 메서드를 호출할 때 Watcher를 등록할 수 있습니다.
getData(path, watcher, stat): 특정 ZNode의 데이터 변경을 감지합니다.exists(path, watcher): 특정 ZNode의 생성, 삭제, 데이터 변경을 감지합니다.getChildren(path, watcher): 특정 ZNode의 자식 노드 생성 또는 삭제를 감지합니다.
메서드의 watcher 인자에 true를 전달하면 클라이언트 인스턴스에 설정된 기본 Watcher를 사용하고, Watcher 인터페이스의 구현체를 직접 전달하여 특정 이벤트에 대한 커스텀 로직을 정의할 수 있습니다.
Watcher 트리거 조건
Watcher는 다음과 같은 트랜잭션 타입 작업에 의해 트리거됩니다.
create(path, data, ...): ZNode 생성delete(path, version): ZNode 삭제setData(path, data, version): ZNode 데이터 변경
영구적인 Watcher 등록 예시
다음 코드는 exists 메서드를 사용하여 ZNode의 변경을 감지하고, 이벤트 발생 후 Watcher를 재등록하여 지속적으로 모니터링하는 방법을 보여줍니다.
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.CreateMode;
import org.apache.zookeeper.data.Stat;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
public class ZkPersistentWatcherDemo {
private static final String ZK_SERVER_LIST = "192.168.136.128:2181,192.168.136.129:2181,192.168.136.130:2181";
private static final int SESSION_TIMEOUT_MS = 4000;
private static CountDownLatch connLatch = new CountDownLatch(1);
private static ZooKeeper zkClient;
private static String targetNode = "/watched-node";
public static void main(String[] args) {
try {
zkClient = new ZooKeeper(ZK_SERVER_LIST, SESSION_TIMEOUT_MS, event -> {
System.out.println("글로벌 Watcher 이벤트: " + event.getType() + " (" + event.getState() + ") -> " + event.getPath());
if (Watcher.Event.KeeperState.SyncConnected == event.getState()) {
connLatch.countDown();
}
});
connLatch.await();
System.out.println("ZooKeeper 연결 완료.");
// 테스트 노드 미리 삭제 (이전 실행 잔여 방지)
if (zkClient.exists(targetNode, false) != null) {
zkClient.delete(targetNode, -1);
}
// Watcher를 특정 노드에 등록
// exists()를 통해 노드의 존재 여부, 데이터 변경, 삭제 이벤트를 감시
Stat initialStat = zkClient.exists(targetNode, new Watcher() {
@Override
public void process(WatchedEvent event) {
System.out.println("노드 전용 Watcher 이벤트: " + event.getType() + " -> " + event.getPath());
try {
// 이벤트 발생 후 Watcher를 재등록하여 지속적인 감시
zkClient.exists(event.getPath(), this); // 'this'는 현재 Watcher 객체
} catch (KeeperException | InterruptedException e) {
System.err.println("Watcher 재등록 중 오류: " + e.getMessage());
}
}
});
// 노드가 존재하지 않으면 생성, 존재하면 데이터 업데이트
if (initialStat == null) {
zkClient.create(targetNode, "data-v1".getBytes(StandardCharsets.UTF_8),
ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
System.out.println("ZNode '" + targetNode + "' 생성 (data-v1)");
} else {
zkClient.setData(targetNode, "data-updated".getBytes(StandardCharsets.UTF_8), -1);
System.out.println("ZNode '" + targetNode + "' 업데이트 (data-updated)");
}
TimeUnit.SECONDS.sleep(2); // 이벤트 처리 대기
// 데이터 변경으로 Watcher 트리거
Stat currentStat = zkClient.exists(targetNode, false); // Watcher를 다시 등록하지 않고 Stat만 가져옴
if (currentStat != null) {
zkClient.setData(targetNode, "data-v2".getBytes(StandardCharsets.UTF_8), currentStat.getVersion());
System.out.println("ZNode '" + targetNode + "' 데이터 변경 (data-v2)");
}
TimeUnit.SECONDS.sleep(2); // 이벤트 처리 대기
// 노드 삭제로 Watcher 트리거
if (zkClient.exists(targetNode, false) != null) {
zkClient.delete(targetNode, -1); // -1은 모든 버전에 대해 삭제 허용
System.out.println("ZNode '" + targetNode + "' 삭제");
}
// 메인 스레드 블로킹하여 Watcher 이벤트 수신 대기
System.out.println("실행을 종료하려면 Enter 키를 누르세요...");
System.in.read();
} catch (IOException | InterruptedException | KeeperException e) {
System.err.println("예외 발생: " + e.getMessage());
e.printStackTrace();
} finally {
if (zkClient != null) {
try {
zkClient.close();
} catch (InterruptedException e) {
System.err.println("클라이언트 종료 중 예외 발생: " + e.getMessage());
}
}
}
}
}
Watcher 이벤트 유형 및 트리거 조건
다음 표는 Watcher 이벤트 유형과 그 트리거 조건을 요약한 것입니다.
이벤트 유형 (Event.EventType) |
설명 | 주요 트리거 | Watcher 등록 메서드 |
|---|---|---|---|
None(-1) |
클라이언트의 연결 상태가 변경될 때 발생합니다. (예: SyncConnected, Disconnected) |
연결/세션 상태 변화 | 클라이언트 생성 시 글로벌 Watcher |
NodeCreated(1) |
등록된 경로에 ZNode가 새로 생성될 때 발생합니다. | create() |
exists(), getChildren() |
NodeDeleted(2) |
등록된 ZNode가 삭제될 때 발생합니다. | delete() |
exists(), getData(), getChildren() |
NodeDataChanged(3) |
등록된 ZNode의 데이터가 변경될 때 발생합니다. | setData() |
exists(), getData() |
NodeChildrenChanged(4) |
등록된 ZNode의 자식 노드가 생성되거나 삭제될 때 발생합니다. 자식 노드의 데이터 변경은 포함되지 않습니다. | 자식 노드 create() 또는 delete() |
getChildren() |
Watcher 구현 원리
클라이언트가 ZooKeeper 서버에 Watcher를 등록하면, 서버는 이 Watcher 정보를 내부적으로 관리합니다. ZNode에 대한 변경 요청(create, delete, setData 등)이 들어오면, 서버는 해당 ZNode에 등록된 Watcher가 있는지 확인하고, 있다면 이벤트를 생성하여 클라이언트에게 전송합니다.
ZooKeeper 클라이언트 측에서는 ZooKeeper 클래스의 exists(), getData(), getChildren() 등의 메서드를 통해 Watcher가 등록됩니다. 예를 들어, exists() 메서드를 살펴보면 내부적으로 WatchRegistration 객체를 생성하여 Watcher와 경로 정보를 캡슐화합니다. 이 정보는 서버로 전송되어 서버의 Watcher Manager에 저장됩니다.
서버는 NettyServerCnxn과 같은 클래스를 통해 클라이언트 요청을 수신하고 처리합니다. 요청 처리 과정에서 데이터 변경이 발생하면, 관련된 Watcher 이벤트를 OutgoingQueue에 추가한 후 클라이언트에게 비동기적으로 전송합니다. 이러한 과정은 파이프라인(Pipeline) 또는 체인 오브 리스폰서빌리티(Chain of Responsibility) 패턴과 유사하게 요청을 여러 단계로 나누어 처리하며, 대부분의 작업은 멀티스레딩을 통해 비동기적으로 수행되어 성능과 확장성을 확보합니다.
Curator 클라이언트 활용: ZK API의 고급 추상화
기본 ZooKeeper Java API는 유연하지만, 연결 관리, Watcher 재등록, 복잡한 분산 레시피(분산 락, 리더 선출 등) 구현에는 상당한 노력이 필요합니다. 이러한 불편함을 해소하기 위해 Netflix에서 개발하여 오픈 소스로 공개한 것이 Apache Curator입니다.
Curator는 ZooKeeper의 원시 API 위에 구축된 고수준의 추상화 계층을 제공합니다. 이는 단지 API를 래핑하는 것을 넘어, 분산 환경에서 자주 사용되는 패턴(예: 분산 락, 리더 선출, 캐시)을 미리 구현하여 사용자가 간편하게 활용할 수 있도록 돕습니다. Curator는 Fluent API 스타일을 채택하여 코드를 더욱 읽기 쉽고 간결하게 만듭니다.
Curator 의존성 추가
Maven 프로젝트에 Curator를 추가하려면 다음 의존성을 pom.xml에 포함합니다.
<dependency>
<groupId>org.apache.curator</groupId>
<artifactId>curator-framework</artifactId>
<version>5.5.0</version> <!-- 최신 안정 버전으로 업데이트하세요 -->
</dependency>
<dependency>
<groupId>org.apache.curator</groupId>
<artifactId>curator-recipes</artifactId>
<version>5.5.0</version> <!-- 레시피(분산 락 등) 사용 시 필요 -->
</dependency>
Curator를 이용한 ZNode CRUD 예시
Curator는 연결 설정부터 ZNode 조작까지 훨씬 간결한 API를 제공합니다. 특히 creatingParentContainersIfNeeded()와 같은 메서드를 통해 부모 노드가 없어도 한 번에 노드를 생성할 수 있습니다.
import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.retry.ExponentialBackoffRetry;
import org.apache.zookeeper.CreateMode;
import org.apache.zookeeper.data.Stat;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.TimeUnit;
public class CuratorCrudExample {
private static final String ZK_SERVER_LIST = "192.168.136.128:2181,192.168.136.129:2181,192.168.136.130:2181";
private static final int SESSION_TIMEOUT_MS = 4000;
private static final String NAMESPACE = "curator-app"; // 특정 비즈니스 로직을 위한 네임스페이스
public static void main(String[] args) {
CuratorFramework curatorClient = null;
try {
// CuratorFramework 인스턴스 생성 및 연결 설정
curatorClient = CuratorFrameworkFactory.builder()
.connectString(ZK_SERVER_LIST)
.sessionTimeoutMs(SESSION_TIMEOUT_MS)
.retryPolicy(new ExponentialBackoffRetry(1000, 3)) // 재시도 정책: 1초 간격, 최대 3회 재시도
.namespace(NAMESPACE) // 모든 경로 앞에 /curator-app이 자동으로 붙음
.build();
curatorClient.start(); // 클라이언트 시작 (ZooKeeper 연결)
System.out.println("Curator ZooKeeper 클라이언트 연결 완료.");
String nodePath = "/serviceA/config";
String initialContent = "version_1.0";
String updatedContent = "version_1.1_released";
// 1. ZNode 생성 (부모 컨테이너가 없으면 자동으로 생성)
curatorClient.create()
.creatingParentContainersIfNeeded() // 부모 노드가 없으면 생성
.withMode(CreateMode.PERSISTENT) // 영구 노드
.forPath(nodePath, initialContent.getBytes(StandardCharsets.UTF_8));
System.out.println("ZNode '" + NAMESPACE + nodePath + "' 생성됨. 데이터: " + initialContent);
// 2. ZNode 데이터 조회
Stat stat = new Stat();
byte[] fetchedData = curatorClient.getData()
.storingStatIn(stat) // Stat 정보를 저장
.forPath(nodePath);
System.out.println("ZNode '" + NAMESPACE + nodePath + "' 데이터 조회: " + new String(fetchedData, StandardCharsets.UTF_8));
System.out.println("현재 버전: " + stat.getVersion());
TimeUnit.SECONDS.sleep(1);
// 3. ZNode 데이터 수정 (낙관적 락을 위한 버전 지정)
curatorClient.setData()
.withVersion(stat.getVersion()) // 현재 버전과 일치할 때만 수정
.forPath(nodePath, updatedContent.getBytes(StandardCharsets.UTF_8));
System.out.println("ZNode '" + NAMESPACE + nodePath + "' 데이터 수정됨. 새 데이터: " + updatedContent);
// 수정 후 데이터 재조회
fetchedData = curatorClient.getData()
.storingStatIn(stat)
.forPath(nodePath);
System.out.println("수정 후 데이터: " + new String(fetchedData, StandardCharsets.UTF_8) + ", 새 버전: " + stat.getVersion());
TimeUnit.SECONDS.sleep(1);
// 4. ZNode 삭제 (자식 노드가 있어도 함께 삭제)
curatorClient.delete()
.deletingChildrenIfNeeded() // 자식 노드가 있으면 함께 삭제
.withVersion(stat.getVersion()) // 현재 버전과 일치할 때만 삭제
.forPath(nodePath);
System.out.println("ZNode '" + NAMESPACE + nodePath + "' 삭제됨.");
} catch (Exception e) {
System.err.println("Curator 작업 중 예외 발생: " + e.getMessage());
e.printStackTrace();
} finally {
if (curatorClient != null) {
curatorClient.close(); // 클라이언트 종료
System.out.println("Curator ZooKeeper 클라이언트 연결 종료.");
}
}
}
}
Curator의 고급 Watcher (Cache) 메커니즘
Curator는 원시 ZooKeeper API의 일회성 Watcher 한계를 극복하기 위해 지속적인 이벤트 감지를 위한 캐시(Cache) 메커니즘을 제공합니다. 이는 클라이언트 측에서 Watcher 재등록 로직을 직접 구현할 필요 없이, ZooKeeper 데이터를 로컬 캐시에 동기화하고 변경 사항을 편리하게 감지할 수 있도록 합니다.
주요 Cache 유형은 다음과 같습니다:
PathChildrenCache: 특정 경로의 자식 노드 추가, 삭제, 데이터 변경 이벤트를 감지합니다. (노드 자체의 데이터 변경은 감지하지 않음)NodeCache: 특정 단일 ZNode의 데이터 변경 및 노드 생성/삭제 이벤트를 감지합니다.TreeCache:PathChildrenCache와NodeCache의 기능을 결합하여, 특정 경로 및 모든 하위 자식 노드들의 생성, 삭제, 데이터 변경 이벤트를 감지하는 가장 포괄적인 캐시입니다.
Curator Watcher (Cache) 예시
import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.framework.recipes.cache.*;
import org.apache.curator.retry.ExponentialBackoffRetry;
import org.apache.zookeeper.CreateMode;
import org.apache.zookeeper.KeeperException;
import org.apache.zookeeper.data.Stat;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.TimeUnit;
public class CuratorWatcherExample {
private static final String ZK_SERVER_LIST = "192.168.136.128:2181,192.168.136.129:2181,192.168.136.130:2181";
private static final int SESSION_TIMEOUT_MS = 4000;
private static final String NAMESPACE = "curator-watch";
private static final String BASE_NODE_PATH = "/monitor-root";
public static void main(String[] args) throws Exception {
CuratorFramework curatorClient = CuratorFrameworkFactory.builder()
.connectString(ZK_SERVER_LIST)
.sessionTimeoutMs(SESSION_TIMEOUT_MS)
.retryPolicy(new ExponentialBackoffRetry(1000, 3))
.namespace(NAMESPACE)
.build();
curatorClient.start();
System.out.println("Curator 클라이언트 연결 완료.");
// 테스트를 위한 초기 노드 상태 정리
try {
if (curatorClient.checkExists().forPath(BASE_NODE_PATH) != null) {
curatorClient.delete().deletingChildrenIfNeeded().forPath(BASE_NODE_PATH);
System.out.println(BASE_NODE_PATH + " 노드 및 자식 노드 삭제 완료.");
}
} catch (KeeperException.NoNodeException e) {
// 노드가 없으면 무시
}
// 다양한 캐시 리스너 등록
addNodeCacheListener(curatorClient, BASE_NODE_PATH);
addPathChildrenCacheListener(curatorClient, BASE_NODE_PATH);
addTreeCacheListener(curatorClient, BASE_NODE_PATH);
// --- 노드 조작으로 이벤트 트리거 ---
System.out.println("\n--- 노드 조작 시작 ---");
// 1. 기본 노드 생성
curatorClient.create().creatingParentContainersIfNeeded().withMode(CreateMode.PERSISTENT)
.forPath(BASE_NODE_PATH, "root-data-v1".getBytes(StandardCharsets.UTF_8));
System.out.println("생성: " + NAMESPACE + BASE_NODE_PATH + " (root-data-v1)");
TimeUnit.SECONDS.sleep(2);
// 2. 기본 노드 데이터 변경
curatorClient.setData().forPath(BASE_NODE_PATH, "root-data-v2".getBytes(StandardCharsets.UTF_8));
System.out.println("변경: " + NAMESPACE + BASE_NODE_PATH + " (root-data-v2)");
TimeUnit.SECONDS.sleep(2);
// 3. 자식 노드 생성
String childNodePath = BASE_NODE_PATH + "/child-node-1";
curatorClient.create().withMode(CreateMode.PERSISTENT)
.forPath(childNodePath, "child-data-v1".getBytes(StandardCharsets.UTF_8));
System.out.println("생성: " + NAMESPACE + childNodePath + " (child-data-v1)");
TimeUnit.SECONDS.sleep(2);
// 4. 자식 노드 데이터 변경
curatorClient.setData().forPath(childNodePath, "child-data-v2".getBytes(StandardCharsets.UTF_8));
System.out.println("변경: " + NAMESPACE + childNodePath + " (child-data-v2)");
TimeUnit.SECONDS.sleep(2);
// 5. 손자 노드 생성
String grandChildNodePath = childNodePath + "/grandchild-1";
curatorClient.create().withMode(CreateMode.PERSISTENT)
.forPath(grandChildNodePath, "grandchild-data-v1".getBytes(StandardCharsets.UTF_8));
System.out.println("생성: " + NAMESPACE + grandChildNodePath + " (grandchild-data-v1)");
TimeUnit.SECONDS.sleep(2);
// 6. 손자 노드 삭제
curatorClient.delete().forPath(grandChildNodePath);
System.out.println("삭제: " + NAMESPACE + grandChildNodePath);
TimeUnit.SECONDS.sleep(2);
// 7. 자식 노드 삭제
curatorClient.delete().forPath(childNodePath);
System.out.println("삭제: " + NAMESPACE + childNodePath);
TimeUnit.SECONDS.sleep(2);
// 8. 기본 노드 삭제
curatorClient.delete().forPath(BASE_NODE_PATH);
System.out.println("삭제: " + NAMESPACE + BASE_NODE_PATH);
TimeUnit.SECONDS.sleep(2);
System.out.println("\n--- 노드 조작 완료 ---");
// 프로그램이 종료되지 않도록 대기
System.out.println("Watcher 이벤트를 계속 수신하려면 Enter 키를 누르세요...");
System.in.read();
curatorClient.close();
System.out.println("Curator 클라이언트 종료.");
}
/**
* NodeCache: 특정 단일 ZNode의 데이터 변경 및 노드 생성/삭제 이벤트 감지
* (대상 노드 자체의 변경만 감지)
*/
public static void addNodeCacheListener(CuratorFramework client, String path) throws Exception {
NodeCache nodeCache = new NodeCache(client, path);
nodeCache.getListenable().addListener(() -> {
ChildData data = nodeCache.getCurrentData();
if (data != null) {
System.out.println("[NodeCache] 이벤트 발생 - 경로: " + data.getPath() + ", 데이터: " + new String(data.getData(), StandardCharsets.UTF_8));
} else {
System.out.println("[NodeCache] 이벤트 발생 - 노드 삭제됨: " + path);
}
});
nodeCache.start(true); // true는 캐시 초기화 시 노드 데이터를 읽어옴
System.out.println("[NodeCache] " + path + " 에 리스너 등록.");
}
/**
* PathChildrenCache: 특정 경로의 자식 노드 추가, 삭제, 데이터 변경 이벤트 감지
* (대상 노드 자체의 변경은 감지하지 않음)
*/
public static void addPathChildrenCacheListener(CuratorFramework client, String path) throws Exception {
// cacheData: true는 이벤트와 함께 노드 데이터를 캐시
PathChildrenCache pathChildrenCache = new PathChildrenCache(client, path, true);
pathChildrenCache.getListenable().addListener((curatorFramework, event) -> {
String eventPath = (event.getData() != null) ? event.getData().getPath() : "N/A";
String eventData = (event.getData() != null) ? new String(event.getData().getData(), StandardCharsets.UTF_8) : "N/A";
System.out.println("[PathChildrenCache] 이벤트 유형: " + event.getType() + ", 경로: " + eventPath + ", 데이터: " + eventData);
});
pathChildrenCache.start(PathChildrenCache.StartMode.BUILD_INITIAL_CACHE); // 초기 캐시 빌드
System.out.println("[PathChildrenCache] " + path + " 하위 자식 노드 리스너 등록.");
}
/**
* TreeCache: NodeCache와 PathChildrenCache의 기능을 결합하여,
* 특정 경로 및 모든 하위 자식 노드들의 생성, 삭제, 데이터 변경 이벤트를 감지하는 가장 포괄적인 캐시.
*/
public static void addTreeCacheListener(CuratorFramework client, String path) throws Exception {
TreeCache treeCache = TreeCache.newBuilder(client, path).setCacheData(true).build();
treeCache.getListenable().addListener((curatorFramework, event) -> {
String eventPath = (event.getData() != null) ? event.getData().getPath() : "N/A";
String eventData = (event.getData() != null) ? new String(event.getData().getData(), StandardCharsets.UTF_8) : "N/A";
System.out.println("[TreeCache] 이벤트 유형: " + event.getType() + ", 경로: " + eventPath + ", 데이터: " + eventData);
});
treeCache.start();
System.out.println("[TreeCache] " + path + " 및 모든 하위 노드 리스너 등록.");
}
}