Apache Curator 클라이언트: Zookeeper 활용 가이드

세션 생성

Curator 클라이언트를 사용하여 Zookeeper 세션을 생성하는 방식은 다른 Zookeeper 클라이언트와는 차이가 있습니다. 주로 다음 두 단계로 진행됩니다:

  1. CuratorFrameworkFactory 팩토리 클래스의 정적 메서드를 통해 클라이언트 인스턴스를 생성합니다.
  2. 생성된 CuratorFramework 인스턴스의 start() 메서드를 호출하여 세션을 시작합니다.

CuratorFrameworkFactory에서 제공하는 주요 클라이언트 생성 메서드는 다음과 같습니다:

public static CuratorFramework newClient(String connectString, RetryPolicy retryPolicy)
public static CuratorFramework newClient(String connectString, int sessionTimeoutMs,
 int connectionTimeoutMs, RetryPolicy retryPolicy)

재시도 정책(RetryPolicy)은 Curator가 Zookeeper 연결 실패 시 어떻게 재시도할지 정의하는 인터페이스입니다. 이 인터페이스에는 allowRetry 메서드가 정의되어 있어 사용자가 직접 재시도 로직을 구현할 수 있습니다.

public boolean allowRetry(int retryCount, long elapsedTimeMs, RetrySleeper sleeper);

Curator에서 세션을 생성하는 기본적인 예제는 다음과 같습니다. ExponentialBackoffRetry는 Curator에서 기본으로 제공하는 재시도 정책 중 하나로, 재시도 간격이 점진적으로 증가합니다.

package com.example.curator.session;

import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.retry.ExponentialBackoffRetry;

public class ZkSessionCreator {
    public static void main(String[] args) throws InterruptedException {
        // 기본 슬립 타임 1000ms, 최대 재시도 횟수 3회로 설정
        ExponentialBackoffRetry retryStrategy = new ExponentialBackoffRetry(1000, 3);
        
        // Zookeeper 서버 주소와 재시도 정책으로 클라이언트 생성
        CuratorFramework zkClient = CuratorFrameworkFactory.newClient("127.0.0.1:2181", retryStrategy);
        
        // 클라이언트 세션 시작
        zkClient.start();
        System.out.println("Zookeeper 세션이 시작되었습니다.");
        
        // 프로그램 종료 방지 (실제 환경에서는 적절한 종료 로직 필요)
        Thread.sleep(Long.MAX_VALUE);
    }
}

ExponentialBackoffRetry 생성자는 다음과 같습니다:

public ExponentialBackoffRetry(int baseSleepTimeMs, int maxRetries)
public ExponentialBackoffRetry(int baseSleepTimeMs, int maxRetries, int maxSleepMs)

플루언트(Fluent) 스타일 API를 사용한 세션 생성

Curator API의 가장 큰 특징 중 하나는 플루언트(Fluent) 디자인 스타일을 따른다는 점입니다. 이는 메서드 체인을 통해 가독성 높은 코드를 작성할 수 있게 합니다.

package com.example.curator.session;

import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.retry.ExponentialBackoffRetry;

public class FluentSessionBuilder {
    public static void main(String[] args) {
        ExponentialBackoffRetry retryPolicy = new ExponentialBackoffRetry(1000, 3);
        
        // Builder 패턴을 사용하여 플루언트 스타일로 클라이언트 설정
        CuratorFramework zkClient = CuratorFrameworkFactory.builder()
                                 .connectString("127.0.0.1:2181")
                                 .sessionTimeoutMs(5000) // 세션 타임아웃 5초
                                 .retryPolicy(retryPolicy)
                                 .build();
        zkClient.start();
        System.out.println("플루언트 스타일로 Zookeeper 세션이 시작되었습니다.");
    }
}

격리된 네임스페이스를 포함한 세션 생성

다양한 Zookeeper 애플리케이션 간의 격리를 위해, 각 서비스에 독립적인 네임스페이스(즉, Zookeeper의 루트 경로)를 할당하는 경우가 많습니다. 예를 들어, 네임스페이스를 /my-app으로 지정하면, 이 클라이언트가 수행하는 모든 Zookeeper 데이터 노드 작업은 이 상대 경로를 기준으로 이루어집니다.

CuratorFramework clientWithNamespace = CuratorFrameworkFactory.builder()
    .connectString("127.0.0.1:2181")
    .sessionTimeoutMs(5000)
    .retryPolicy(new ExponentialBackoffRetry(1000, 3))
    .namespace("my-app") // 네임스페이스 설정
    .build();
clientWithNamespace.start();
// 이제 이 clientWithNamespace는 모든 작업을 /my-app 경로 아래에서 수행합니다.
// 예를 들어, clientWithNamespace.create().forPath("/node")는 실제로는 /my-app/node를 생성합니다.

노드 작업 (CRUD)

Curator는 플루언트 스타일의 인터페이스를 제공하여 개발자가 다양한 유형의 노드 생성, 삭제, 읽기, 업데이트 작업을 유연하게 조합할 수 있도록 합니다.

노드 생성

  • 기본 노드 생성 (내용 없음)

    client.create().forPath("/test_path"); // 기본적으로 영구 노드(PERSISTENT) 생성, 내용 없음
  • 초기 내용을 포함한 노드 생성

    client.create().forPath("/data_path", "초기 데이터".getBytes());
  • 특정 모드로 노드 생성 (예: 임시 노드)

    client.create().withMode(CreateMode.EPHEMERAL).forPath("/ephemeral_node");
  • 부모 노드가 없는 경우 자동으로 생성

    client.create().creatingParentsIfNeeded().withMode(CreateMode.EPHEMERAL).forPath("/parent/child/temp_node");

노드 삭제

  • 기본 노드 삭제

    client.delete().forPath("/test_path");
  • 자식 노드를 포함하여 재귀적으로 삭제

    client.delete().deletingChildrenIfNeeded().forPath("/parent_node");
  • 버전을 지정하여 조건부 삭제 (CAS)

    client.delete().withVersion(1).forPath("/versioned_node"); // 버전이 1인 경우에만 삭제
  • 삭제 보장

    client.delete().guaranteed().forPath("/guaranteed_delete_node");

    guaranteed()는 클라이언트 세션이 유효한 동안 백그라운드에서 노드 삭제를 성공할 때까지 계속 시도하는 안전장치입니다.

데이터 읽기

  • 노드 데이터 읽기

    byte[] data = client.getData().forPath("/some_path");
  • 노드 데이터와 Stat 정보 함께 읽기

    import org.apache.zookeeper.data.Stat;
            // ...
            Stat nodeStat = new Stat();
            byte[] dataWithStat = client.getData().storingStatIn(nodeStat).forPath("/some_path");
            System.out.println("노드 데이터: " + new String(dataWithStat));
            System.out.println("노드 버전: " + nodeStat.getVersion());

데이터 업데이트

  • 노드 데이터 업데이트

    client.setData().forPath("/update_path", "새로운 데이터".getBytes());
  • 버전을 지정하여 조건부 업데이트 (CAS)

    int currentVersion = client.getData().storingStatIn(new Stat()).forPath("/cas_path").getVersion();
            client.setData().withVersion(currentVersion).forPath("/cas_path", "조건부 업데이트".getBytes());

    withVersion()은 Zookeeper의 CAS(Compare And Set) 기능을 사용하여, 지정된 버전과 현재 노드의 버전이 일치할 경우에만 데이터를 업데이트합니다. 버전 정보는 보통 기존 노드의 Stat 객체에서 가져옵니다.

비동기 API

Curator는 비동기 작업 처리를 위해 BackgroundCallback 인터페이스를 도입했습니다. 이 인터페이스는 비동기 호출 후 Zookeeper 서버로부터 결과가 반환될 때 호출됩니다.

BackgroundCallback 인터페이스에는 processResult 메서드 하나만 정의되어 있습니다. 이 메서드는 작업이 완료되면 비동기적으로 호출되며, CuratorEvent 객체를 통해 Zookeeper 서버가 클라이언트로 전송하는 다양한 이벤트 매개변수를 받을 수 있습니다. CuratorEvent에서 중요한 파라미터로는 이벤트 타입(CuratorEventType)과 응답 코드(int)가 있습니다.

비동기 API에서 주목할 만한 매개변수는 executor입니다. Zookeeper의 모든 비동기 알림 이벤트 처리는 EventThread라는 단일 스레드에서 직렬화되어 처리됩니다. EventThread의 이러한 직렬화 처리 메커니즘은 대부분의 경우 이벤트 처리 순서의 일관성을 보장하지만, 복잡한 처리 단위가 있을 경우 오랜 시간을 소모하여 다른 이벤트 처리를 지연시킬 수 있는 단점이 있습니다. Curator는 이러한 문제를 해결하기 위해 사용자에게 Executor 인스턴스를 전달할 수 있도록 허용합니다. 이를 통해 복잡한 이벤트 처리를 전용 스레드 풀에서 수행하여 EventThread의 부담을 줄일 수 있습니다.

package com.example.curator.async;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.framework.api.BackgroundCallback;
import org.apache.curator.framework.api.CuratorEvent;
import org.apache.curator.retry.ExponentialBackoffRetry;
import org.apache.zookeeper.CreateMode;

public class AsyncNodeOperations {
    
    private static final String ZK_NODE_PATH = "/curator-async-node";
    
    // Zookeeper 클라이언트 인스턴스
    private static CuratorFramework asyncClient = CuratorFrameworkFactory.builder()
            .connectString("127.0.0.1:2181")
            .sessionTimeoutMs(5000)
            .retryPolicy(new ExponentialBackoffRetry(1000, 3))
            .build();
    
    // 비동기 콜백 처리를 위한 스레드 풀
    private static ExecutorService callbackThreadPool = Executors.newFixedThreadPool(2);
     
    public static void main(String[] args) throws Exception {
        asyncClient.start();
        System.out.println("메인 스레드: " + Thread.currentThread().getName());
        
        CountDownLatch latch = new CountDownLatch(1);

        asyncClient.create()
                   .creatingParentsIfNeeded()
                   .withMode(CreateMode.EPHEMERAL)
                   .inBackground(new BackgroundCallback() { // 비동기 콜백 정의
                       @Override
                       public void processResult(CuratorFramework client, CuratorEvent event) throws Exception {
                           System.out.println("---------- 비동기 콜백 실행 -----------");
                           System.out.println("결과 코드: " + event.getResultCode() + ", 이벤트 타입: " + event.getType());
                           System.out.println("콜백 처리 스레드: " + Thread.currentThread().getName());
                           latch.countDown();
                       }
                   }, callbackThreadPool) // 스레드 풀을 지정하여 콜백 실행
                   .forPath(ZK_NODE_PATH, "async data".getBytes());
        
        latch.await(); // 비동기 작업 완료 대기
        System.out.println("비동기 노드 생성 요청 완료.");
        Thread.sleep(Integer.MAX_VALUE); // 프로그램 종료 방지
    }
}

주요 활용 사례

1. 이벤트 리스닝 (캐시)

Zookeeper는 기본적으로 Watcher를 통해 이벤트 리스닝을 지원하지만, 사용자가 매번 워처를 등록해야 하는 번거로움이 있습니다. Curator는 이러한 불편함을 해소하기 위해 Cache 기능을 도입했습니다. Curator의 캐시는 Zookeeper 서버의 상태 변화를 로컬에 캐시된 뷰와 비교하여 효율적으로 이벤트를 감지하고, 워처의 반복적인 등록 작업을 자동으로 처리합니다.

Curator의 캐시는 크게 두 가지 유형으로 나뉩니다: 노드 자체를 감지하는 NodeCache와 자식 노드의 변화를 감지하는 PathChildrenCache입니다.

NodeCache

NodeCache는 지정된 Zookeeper 데이터 노드 자체의 변경(데이터 업데이트, 생성, 삭제)을 감지하는 데 사용됩니다.

NodeCache(CuratorFramework client, String path)
public NodeCache(CuratorFramework client, String path, boolean dataIsCompressed)

NodeCacheNodeCacheListener 콜백 인터페이스를 정의합니다. 데이터 노드의 내용이 변경되거나 노드의 존재 여부가 바뀔 때 이 메서드가 호출됩니다.

package com.example.curator.cache;

import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.framework.recipes.cache.NodeCache;
import org.apache.curator.framework.recipes.cache.NodeCacheListener;
import org.apache.curator.retry.ExponentialBackoffRetry;
import org.apache.zookeeper.CreateMode;

public class NodeCacheExample {
    
    private static final String NODE_CACHE_PATH = "/curator-cache/my-node";
    private static CuratorFramework zkClient = CuratorFrameworkFactory.builder()
            .connectString("127.0.0.1:2181")
            .connectionTimeoutMs(5000)
            .retryPolicy(new ExponentialBackoffRetry(3000, 3))
            .build();
    
    public static void main(String[] args) throws Exception {
        zkClient.start();

        // 노드 생성 (NodeCache가 감지하기 전에 미리 생성)
        zkClient.create()
                .creatingParentsIfNeeded()
                .withMode(CreateMode.PERSISTENT) // 임시 노드가 아닌 영구 노드로 생성
                .forPath(NODE_CACHE_PATH, "initial data".getBytes());
        
        final NodeCache nodeCache = new NodeCache(zkClient, NODE_CACHE_PATH, false);
        nodeCache.start(true); // true로 설정하여 캐시 시작 시 Zookeeper에서 초기 데이터 로드

        nodeCache.getListenable().addListener(new NodeCacheListener() {
            @Override
            public void nodeChanged() throws Exception {
                // 노드가 존재하면 현재 데이터를 출력, 없으면 "노드 삭제됨" 출력
                if (nodeCache.getCurrentData() != null) {
                    System.out.println("노드 데이터 업데이트: 새 데이터 = " + 
                                       new String(nodeCache.getCurrentData().getData()));
                } else {
                    System.out.println("노드가 삭제되었습니다.");
                }
            }
        });
        
        System.out.println("노드 데이터를 'update 1'로 변경합니다.");
        zkClient.setData().forPath(NODE_CACHE_PATH, "update 1".getBytes());
        Thread.sleep(2000);

        System.out.println("노드 데이터를 'update 2'로 변경합니다.");
        zkClient.setData().forPath(NODE_CACHE_PATH, "update 2".getBytes());
        Thread.sleep(2000);
        
        System.out.println("노드를 삭제합니다.");
        zkClient.delete().deletingChildrenIfNeeded().forPath(NODE_CACHE_PATH);
        Thread.sleep(2000);
        
        System.out.println("종료 대기 중...");
        Thread.sleep(Integer.MAX_VALUE);
    }
}

NodeCache.start(true) 호출 시 true 매개변수는 캐시가 처음 시작될 때 Zookeeper로부터 해당 노드의 현재 데이터를 읽어와 캐시에 저장하도록 지시합니다. NodeCache는 노드 데이터 변경뿐만 아니라 노드 자체의 존재 여부(생성 또는 삭제)도 감지할 수 있습니다.

PathChildrenCache

PathChildrenCache는 지정된 Zookeeper 데이터 노드의 자식 노드 변경 사항을 감지하는 데 사용됩니다. 즉, 자식 노드의 추가, 삭제, 데이터 변경 이벤트를 모니터링합니다.

PathChildrenCachePathChildrenCacheListener 콜백 인터페이스를 정의합니다. 지정된 노드의 자식 노드에서 변화가 발생하면 이 메서드가 호출됩니다.

package com.example.curator.cache;

import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.framework.recipes.cache.PathChildrenCache;
import org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent;
import org.apache.curator.framework.recipes.cache.PathChildrenCacheListener;
import org.apache.curator.framework.recipes.cache.PathChildrenCache.StartMode;
import org.apache.curator.retry.ExponentialBackoffRetry;
import org.apache.zookeeper.CreateMode;

public class PathChildrenCacheExample {
    
    private static final String PARENT_NODE_PATH = "/curator-parent-cache";
    private static CuratorFramework zkClient = CuratorFrameworkFactory.builder()
            .connectString("127.0.0.1:2181")
            .connectionTimeoutMs(5000)
            .retryPolicy(new ExponentialBackoffRetry(3000, 3))
            .build();
    
    public static void main(String[] args) throws Exception {
        zkClient.start();

        // 부모 노드가 없으면 미리 생성
        zkClient.create().creatingParentsIfNeeded().forPath(PARENT_NODE_PATH);
        
        PathChildrenCache childrenCache = new PathChildrenCache(zkClient, PARENT_NODE_PATH, true);
        childrenCache.start(StartMode.POST_INITIALIZED_EVENT); // 초기화 완료 후 이벤트 발생

        childrenCache.getListenable().addListener(new PathChildrenCacheListener() {
            @Override
            public void childEvent(CuratorFramework client, PathChildrenCacheEvent event) throws Exception {
                switch (event.getType()) {
                    case CHILD_ADDED:
                        System.out.println("자식 노드 추가: " + event.getData().getPath() + 
                                           ", 데이터: " + new String(event.getData().getData()));
                        break;
                    case CHILD_UPDATED:
                        System.out.println("자식 노드 업데이트: " + event.getData().getPath() + 
                                           ", 새 데이터: " + new String(event.getData().getData()));
                        break;
                    case CHILD_REMOVED:
                        System.out.println("자식 노드 삭제: " + event.getData().getPath());
                        break;
                    default:    
                        System.out.println("다른 이벤트: " + event.getType());
                        break;
                }
            }
        });
        
        System.out.println("자식 노드 '/child1'를 생성합니다.");
        zkClient.create().withMode(CreateMode.EPHEMERAL).forPath(PARENT_NODE_PATH + "/child1", "data1".getBytes());
        Thread.sleep(2000);
        
        System.out.println("자식 노드 '/child2'를 생성합니다.");
        zkClient.create().withMode(CreateMode.EPHEMERAL).forPath(PARENT_NODE_PATH + "/child2", "data2".getBytes());
        Thread.sleep(2000);
        
        System.out.println("자식 노드 '/child1'의 데이터를 업데이트합니다.");
        zkClient.setData().forPath(PARENT_NODE_PATH + "/child1", "updated_data1".getBytes());
        Thread.sleep(2000);
        
        System.out.println("자식 노드 '/child2'를 삭제합니다.");
        zkClient.delete().forPath(PARENT_NODE_PATH + "/child2");
        Thread.sleep(2000);
        
        System.out.println("종료 대기 중...");
        Thread.sleep(Integer.MAX_VALUE); // 프로그램 종료 방지
    }
}

2. 분산 마스터 선출

클러스터 환경에서 특정 복잡한 작업을 하나의 노드에서만 처리해야 할 때가 있습니다. 이러한 유형의 분산 문제를 "마스터 선출"이라고 부릅니다. Zookeeper를 활용하면 마스터 선출 기능을 쉽게 구현할 수 있습니다.

기본 아이디어는 다음과 같습니다:

  1. 하나의 루트 노드(예: /master-election)를 선택합니다.
  2. 클러스터의 여러 머신이 동시에 이 루트 노드 아래에 자식 노드(예: /master-election/lock)를 생성하려고 시도합니다.
  3. Zookeeper의 고유한 특성 덕분에, 최종적으로 단 한 대의 머신만 이 자식 노드를 성공적으로 생성할 수 있습니다. 노드 생성에 성공한 머신이 마스터 역할을 수행합니다.

Curator는 이러한 노드 생성, 이벤트 리스닝, 자동 선출 로직을 추상화하여 개발자가 간단한 API 호출만으로 마스터 선출을 구현할 수 있도록 합니다.

package com.example.curator.master;

import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.framework.recipes.leader.LeaderSelector;
import org.apache.curator.framework.recipes.leader.LeaderSelectorListenerAdapter;
import org.apache.curator.retry.ExponentialBackoffRetry;

public class MasterElectionExample {

    private static final String ELECTION_ROOT_PATH = "/curator-master-election";

    private static CuratorFramework zkClient = CuratorFrameworkFactory.builder()
            .connectString("127.0.0.1:2181")
            .connectionTimeoutMs(5000)
            .retryPolicy(new ExponentialBackoffRetry(3000, 3))
            .build();
    
    public static void main(String[] args) throws InterruptedException {
        zkClient.start();
        
        // LeaderSelector 인스턴스 생성: 마스터 선출 로직을 캡슐화
        LeaderSelector leaderSelector = new LeaderSelector(zkClient, ELECTION_ROOT_PATH, 
            new LeaderSelectorListenerAdapter() { // 마스터가 되었을 때 호출될 리스너
                @Override
                public void takeLeadership(CuratorFramework client) throws Exception {
                    // 이 메서드가 호출되면 현재 인스턴스가 마스터가 된 것입니다.
                    System.out.println(Thread.currentThread().getName() + " -> 마스터 역할을 맡았습니다.");
                    try {
                        // 마스터로서 수행할 작업 로직
                        Thread.sleep(5000); // 5초 동안 마스터 역할 유지
                    } finally {
                        System.out.println(Thread.currentThread().getName() + " -> 마스터 작업을 완료하고 권한을 반환합니다.");
                        // takeLeadership 메서드가 종료되면 마스터 권한을 포기하고 다음 선출에 참여합니다.
                    }
                }
            });
        
        leaderSelector.autoRequeue(); // 마스터 권한을 반환한 후 자동으로 다음 선출에 재참여
        leaderSelector.start(); // 마스터 선출 시작
        
        System.out.println("마스터 선출 대기 중...");
        Thread.sleep(Integer.MAX_VALUE); // 프로그램 종료 방지
    }
}

위 코드에서 LeaderSelector 인스턴스는 Zookeeper 서버와의 상호작용을 포함한 모든 마스터 선출 관련 로직을 캡슐화합니다. ELECTION_ROOT_PATH는 이번 마스터 선출이 이루어지는 루트 노드를 나타냅니다.

LeaderSelector 인스턴스를 생성할 때 LeaderSelectorListenerAdapter 리스너를 함께 전달합니다. 이 리스너는 개발자가 구현해야 하며, Curator는 이 인스턴스가 성공적으로 마스터 권한을 획득했을 때 takeLeadership 메서드를 콜백합니다.

takeLeadership 메서드 실행이 완료되면, Curator는 즉시 마스터 권한을 포기하고 새로운 마스터 선출 라운드를 시작합니다 (leaderSelector.autoRequeue() 설정 시). 이를 통해 마스터 노드의 교체와 고가용성이 보장됩니다.

태그: Curator ZooKeeper 분산 시스템 java 재시도 정책

9월 22일 16:14에 게시됨