데보션앱 소개페이지 바로가기
로그인 선택

신고하기

CLOSE
신고사유 (대표 사유 1개)
상세내용 (선택)
0/200
  • 신고한 게시글은 더 이상 보이지 않습니다.
  • 이용약관과 운영정책에 따라 신고사유에 해당하는지 검토 후 조치됩니다.
  • 허위 신고인 경우, 신고자의 서비스 이용이 제한될 수 있으니 유의하시어 신중하게 신고해 주세요.
(이 회원이 작성한 모든 댓글과 커뮤니티 게시물이 보이지 않고, 알림도 오지 않습니다.)

미리보기

커뮤니티

      1,234

      badge 23.06.15

      글 등록

      카테고리를 선택해주세요.

      DEVOTEE를 활성화 시키면
      지금 작성한 커뮤니티 글에 대해 1개의 댓글을 달아줍니다.

      버튼을 누르면 글 수정 시 ChatGPT가 작성한 댓글이 수정됩니다.

      임시저장함에 저장되었습니다. 저장일시 : 2022.5.17 14:29:08

      임시저장함

      제목을 선택하시면 이어서 작성이 가능하며,
      최대 20건까지 저장합니다.
      컨텐츠 유형, 제목, 저장일시, 삭제로 이뤄진 임시저장 목록
      컨텐츠 유형 제목 저장일 삭제

      데보션 블로그 게재 요청

      CLOSE
      • *
      • *

      본인인증

      효율적인 데보션 서비스 이용 및
      고객님의 소중한 개인정보보호를 위해
      본인인증을 진행해주세요. 본인인증 미 진행 시 로그인이 제한됩니다.
      본인인증 실패

      본인인증 로그인에 실패하였습니다.
      회원이 아니시거나 본인인증 등록이
      완료되지 않은 사용자입니다.

      회원정보 연결

      Raft 알고리즘을 이용해 고가용 프로그램을 만들어보자!!

      marine 25.05.19
      3,246 4 1
      DEVOTEE 요약
      본 블로그는 검색인프라팀이 분산 시스템에서 데이터 일관성과 고가용성을 보장하기 위한 Raft 알고리즘의 적용 사례를 다룹니다. Raft는 이해가 쉽고 구현이 용이한 합의 알고리즘으로, 이를 기반으로 3노드 환경에서 동작하는 고가용성 우선순위 큐를 구현하였습니다. 이를 통해 노드 장애 상황에서도 데이터 일관성을 유지하는 중요성을 강조하며, Java 기반 라이브러리인 sofa-jraft를 활용한 구체적인 설계와 구현 방법을 설명합니다.

      개요

      검색인프라팀은 분산 시스템 구축과 운영에 대한 다양한 경험을 쌓아왔습니다.

      이전에는 SaaS형 Redis 클러스터 제공 서비스(RC)와 MinIO 및 Kubernetes를 활용한 사내 스토리지 서비스 구축 경험을 공유한 바 있습니다.

      이번에는 분산 시스템의 핵심 과제 중 하나인 ‘합의(Consensus)’ 문제를 해결하는 Raft 알고리즘과, 이를 활용한 고가용성 우선순위 큐 구현을 통해 실제 적용 방법을 살펴보고자 합니다.

      검색인프라팀의 주요 업무는 데이터 수집, 가공, 연동, 서비스 개발이며, 이 과정에서 데이터의 일관성과 고가용성을 보장하는 것은 필수적이지만 매우 까다로운 과제입니다.

      특히, 네트워크를 통해 수많은 데이터가 실시간으로 입·출력되고,

      마이크로서비스 아키텍처 기반의 여러 노드가 유기적으로 연결된 환경에서는 노드 장애가 발생하더라도 서비스가 중단되지 않고, 데이터 일관성을 유지하는 것이 핵심입니다.

      이러한 배경에서 Raft 알고리즘을 팀내 어플리케이션에 도입한 경험을 공유하고자 이 글을 작성하게 되었습니다.

      이번 글에서는 자바 기반 Raft 구현체인 sofa-jraft 라이브러리를 활용해 3노드 환경에서 운영되는 고가용성 우선순위 큐를 구현하고, 이를 통해 Raft 알고리즘의 개념과 실전 적용 방안을 살펴보겠습니다.


      분산 시스템과 합의 알고리즘의 필요성

      분산 시스템은 여러 개의 독립적인 컴퓨터가 네트워크를 통해 통신하며 하나의 시스템처럼 작동하는 구조를 말합니다.

      이러한 시스템에서는 노드 장애, 네트워크 지연, 메시지 손실 등 다양한 문제가 발생할 수 있으며, 이런 상황에서도 시스템은 일관된 상태를 유지해야 합니다.

      분산 시스템 구축시 해결 과제:

      • 일관성 유지: 모든 노드가 동일한 데이터와 상태를 가져야 함

      • 장애 허용성: 일부 노드가 실패해도 시스템이 계속 작동해야 함

      • 네트워크 파티션 대응: 네트워크 분할 상황에서도 안전성 보장 필요

      이러한 문제를 해결하기 위해 Paxos, Raft, PBFT 등 다양한 합의 알고리즘이 개발되었습니다.

      그 중에서도 Raft는 이해하기 쉽고 구현이 용이하다는 장점으로 많은 분산 시스템에서 채택되고 있습니다.

      Raft 알고리즘을 채택한 사례:

      • Kafka 3.3부터 기존 zookeeper 대신 kraft 사용

      • Etcd

      • Consul


      Raft 알고리즘 이해하기

      Raft는 분산 시스템에서 노드 간 합의를 이루기 위한 알고리즘으로, 스탠포드 대학의 Diego Ongaro와 John Ousterhout에 의해 개발되었습니다.

      Raft는 "이해하기 쉬운 합의 알고리즘"을 목표로 설계되었으며, 복잡한 Paxos 알고리즘의 대안으로 등장했습니다.

      Raft 논문: https://raft.github.io/raft.pdf

      Raft(뗏목)라는 이름의 유래

      뗏목

      위키피디아에 따르면, Raft는 "Reliable, Replicated, Redundant, Fault-Tolerant"의 약어에서 유래했다고 하지만, 이는 일부 해석일 뿐입니다.

      이 알고리즘을 만든 Diego Ongaro는 Google Groups에 직접 이름을 짓게 된 과정에 대해 남긴 기록이 있는데, 요약하면 다음과 같습니다.

      • 처음엔 특정한 줄임말을 염두에 두진 않았지만, "Reliable", "Replicated", "Redundant", "Fault-Tolerant" 같은 단어들을 떠올림

      • 로그(logs)를 떠올리며 이들로 무엇을 만들 수 있을지 고민(log는 '통나무'라는 의미도 있음)

      • 기존 합의 알고리즘인 Paxos라는 '섬'에서 벗어나는 방법을 고민하던 중 **‘뗏목(Raft)’**이라는 이름을 떠올리게 됨


      Raft의 핵심 구성 요소

      Raft는 세 가지 핵심 메커니즘으로 구성됩니다:

      알고리즘 동작 시각적 소개: https://thesecretlivesofdata.com/raft/

      1. 리더 선출 (Leader Election)

      모든 노드는 리더, 팔로워, 후보자 중 하나의 상태를 가집니다. 클러스터가 시작되면 모든 노드는 팔로워 상태로 시작하고, 일정 시간 동안 리더로부터 신호가 없으면 후보자가 되어 리더 선출 과정을 시작합니다.

      • 리더: 클라이언트 요청 처리 및 로그 복제 담당

      • 팔로워: 리더의 요청에 응답하고 로그를 복제

      • 후보자: 리더 선출 과정 중 투표를 요청하는 노드

      리더 선출은 랜덤한 타임아웃과 투표 과정을 통해 이루어지며, 과반수의 노드로부터 투표를 받은 후보자가 새 리더가 됩니다.

      2. 로그 복제 (Log Replication)

      클라이언트의 모든 요청은 리더가 받아서 처리합니다. 리더는 요청을 로그 항목으로 변환하여 자신의 로그에 추가한 후, 모든 팔로워에게 로그를 복제하도록 요청합니다.

      과반수(Quorum)의 노드가 로그를 복제하면 해당 항목이 '커밋'되고 상태 머신에 적용됩니다.

      로그 항목은 다음 정보를 포함합니다:

      • 명령(Command): 실행할 작업

      • 인덱스(Index): 로그 위치

      • 텀 번호(Term Number): 리더 선출 주기를 식별하는 번호

      3. 안전성 (Safety)

      Raft는 다음과 같은 안전성 속성을 보장합니다:

      • 리더 완전성(Leader Completeness): 커밋된 로그는 모든 미래 리더의 로그에도 포함됨

      • 상태 머신 안전성(State Machine Safety): 모든 노드는 커밋된 로그를 동일한 순서로 적용

      • 로그 일치 속성(Log Matching Property): 같은 인덱스와 텀을 가진 로그는 동일한 명령을 포함

      Raft의 장점

      Raft 알고리즘은 다음과 같은 주요 장점을 가집니다:

      • 이해하기 쉽고 구현이 간편함: 복잡한 Paxos에 비해 명확한 개념 구조

      • 장애 허용성(CFT): 전체 노드의 과반수가 정상 작동하는 한 시스템은 계속 작동

      • 즉각적인 블록 완결성: 합의가 이루어진 데이터는 즉시 확정되어 롤백되지 않음

      • 동적 멤버십 변경: 서비스 중단 없이 클러스터 구성 변경 가능


      Raft 활용한 고가용 분산 우선순위 큐 구현하기

      sofa-jraft 라이브러리 소개

      sofa-jraft는 알리바바에서 개발한 Java 기반 Raft 구현체로, Baidu의 Braft를 Java로 재구현하고 최적화한 라이브러리입니다.

      고성능 분산 시스템 구축에 적합하며, 다양한 기능과 최적화를 포함하고 있습니다.

      github: https://github.com/sofastack/sofa-jraft

      주요 특징

      • 완전한 Raft 구현: 리더 선출, 로그 복제, 스냅샷, 멤버십 관리 등 모든 Raft 기능 구현

      • 고성능: 완전 동시 복제, 복제 파이프라인 최적화 등을 통한 성능 향상

      • 선형화된 읽기: ReadIndex/LeaseRead 메커니즘을 통한 효율적인 읽기 작업 지원

      • MULTI-RAFT-GROUP: 복잡한 애플리케이션을 위한 다중 Raft 그룹 지원

      • 스냅샷 및 로그 압축: 효율적인 상태 관리 및 빠른 복구 지원

      핵심 컴포넌트

      sofa-jraft는 다음과 같은 주요 컴포넌트로 구성됩니다:

      1. Node: Raft 그룹의 구성원으로, 리더/팔로워/후보자 상태 관리

      2. StateMachine: 애플리케이션 로직을 구현하는 상태 머신

      3. LogManager: 로그 항목 관리 및 저장

      4. RPC: 노드 간 통신 처리

      5. RouteTable: 리더 노드 조회 및 관리

      분산 우선순위 큐 설계하기

      우선순위 큐는 각 요소가 우선순위를 가지며, 높은 우선순위의 요소가 먼저 처리되는 자료구조입니다.

      이 자료구조는 작업 스케줄링, 이벤트 처리, 네트워크 패킷 관리 등 다양한 분야에서 활용됩니다. 팀내에서는 수집 스케쥴 관리에 적용했습니다.

      우선순위 큐 기본 개념

      우선순위 큐는 일반적으로 힙(Heap) 자료구조를 사용하여 구현되며, 다음과 같은 핵심 연산을 제공합니다:

      • enqueue(element, priority): 우선순위와 함께 요소 추가 (O(log N))

      • dequeue(): 가장 높은 우선순위의 요소 제거 및 반환 (O(log N))

      • peek(): 가장 높은 우선순위의 요소 확인 (O(1))

      분산 환경에서의 설계 고려사항

      분산 우선순위 큐를 설계할 때는 다음과 같은 사항을 고려해야 합니다:

      1. 일관성 유지: 모든 노드가 동일한 큐 상태를 가져야 함

      2. 장애 허용: 일부 노드의 장애가 전체 시스템에 영향을 미치지 않아야 함

      3. 성능 최적화: 분산 환경에서의 오버헤드 최소화

      4. 확장성: 시스템 부하 증가에 따른 대응 가능성

      시스템 아키텍처

      우리의 분산 우선순위 큐는 3개의 노드로 구성된 Raft 클러스터 위에 구현되었습니다.

      각 노드는 sofa-jraft를 사용하여 Raft 프로토콜을 구현하고, 동일한 우선순위 큐 상태 머신을 복제합니다.

      우선순위큐 시스템 아키텍쳐

      시스템 구성:

      1. 클라이언트 계층: 큐 작업(enqueue, dequeue, peek)을 요청하는 인터페이스

      2. Raft 그룹: 3개 노드로 구성된 Raft 클러스터

      3. 상태 머신: 우선순위 큐 로직을 구현한 PriorityQueueStateMachine

      4. 스토리지 계층: 로그 및 스냅샷 저장

      sofa-jraft를 이용한 구현

      이제 sofa-jraft 라이브러리를 사용하여 분산 우선순위 큐를 구현하는 방법을 살펴보겠습니다.

      우선순위 큐 전체 구현 코드 github: https://github.com/kbbmanse/raft-priority-queue.git

      의존성 설정

      먼저 Maven 프로젝트에 sofa-jraft 의존성을 추가합니다.(현시점 가장 최신)

      <dependency>
        <groupId>com.alipay.sofa</groupId>
        <artifactId>jraft-core</artifactId>
        <version>1.3.15.bugfix</version>
      </dependency>

      우선순위 큐 상태 머신 구현

      Raft의 핵심은 상태 머신 복제입니다. 우선순위 큐 상태 머신은 다음과 같이 구현됩니다.

      public class PriorityQueueStateMachine extends StateMachineAdapter {
          // 우선순위 큐 구현을 위한 힙 자료구조
          private final PriorityQueue<Element> queue = new PriorityQueue<>((e1, e2) ->
                  Integer.compare(e2.getPriority(), e1.getPriority()));
      
          private final ReadWriteLock lock = new ReentrantReadWriteLock();
      
          /**
           * Raft 로그를 적용하여 상태 머신을 업데이트합니다.
           *
           * @param iter 로그 항목에 대한 반복자
           */
          @Override
          @SuppressWarnings("unchecked")
          public void onApply(Iterator iter) {
              lock.writeLock().lock(); // 쓰기 락 획득
              try {
                  while (iter.hasNext()) { // 로그 항목을 순회
                      PriorityQueueRequest op = deserialize(iter.getData().array()); // 로그 데이터에서 PriorityQueueRequest 역직렬화
                      try {
                          switch (op.getType()) { // 요청 타입에 따라 처리
                              case ENQUEUE: // 삽입 연산
                                  queue.add(op.getElement()); // 큐에 요소 삽입
                                  PriorityQueueClosure<Boolean> enqueueClosure = null; // 클로저 초기화
                                  if (iter.done() != null) { // 클로저가 존재하는 경우
                                      enqueueClosure = (PriorityQueueClosure<Boolean>) iter.done(); // 클로저 캐스팅
                                  }
                                  if (enqueueClosure != null) { // 클로저가 존재하는 경우
                                      enqueueClosure.setResponse(true); // 응답 설정
                                      enqueueClosure.run(Status.OK()); // 클로저 실행
                                  }
                                  LOG.info("Enqueued element: {}", op.getElement());
                                  break;
                              case DEQUEUE: // 삭제 연산
                                  PriorityQueueClosure<Element> dequeueClosure = null; // 클로저 초기화
                                  if (iter.done() != null) { // 클로저가 존재하는 경우
                                      dequeueClosure = (PriorityQueueClosure<Element>) iter.done(); // 클로저 캐스팅
                                  }
                                  if (queue.isEmpty()) { // 큐가 비어있는 경우
                                      if (dequeueClosure != null) { // 클로저가 존재하는 경우
                                          dequeueClosure.setResponse(null); // 응답 설정 (null)
                                          dequeueClosure.run(Status.OK()); // 클로저 실행
                                      }
                                      LOG.info("Dequeue operation on empty queue");
                                  } else { // 큐가 비어있지 않은 경우
                                      Element element = queue.poll(); // 큐에서 요소 삭제 (가장 우선순위가 높은 요소)
                                      if (dequeueClosure != null) { // 클로저가 존재하는 경우
                                          dequeueClosure.setResponse(element); // 응답 설정 (삭제된 요소)
                                          dequeueClosure.run(Status.OK()); // 클로저 실행
                                      }
                                      LOG.info("Dequeued element: {}", element);
                                  }
                                  break;
                              case PEEK: // 조회 연산
                                  PriorityQueueClosure<Element> peekClosure = null; // 클로저 초기화
                                  if (iter.done() != null) { // 클로저가 존재하는 경우
                                      peekClosure = (PriorityQueueClosure<Element>) iter.done(); // 클로저 캐스팅
                                  }
                                  if (queue.isEmpty()) { // 큐가 비어있는 경우
                                      if (peekClosure != null) { // 클로저가 존재하는 경우
                                          peekClosure.setResponse(null); // 응답 설정 (null)
                                          peekClosure.run(Status.OK()); // 클로저 실행
                                      }
                                      LOG.info("Peek operation on empty queue");
                                  } else { // 큐가 비어있지 않은 경우
                                      Element element = queue.peek(); // 큐에서 요소 조회 (가장 우선순위가 높은 요소)
                                      if (peekClosure != null) { // 클로저가 존재하는 경우
                                          peekClosure.setResponse(element); // 응답 설정 (조회된 요소)
                                          peekClosure.run(Status.OK()); // 클로저 실행
                                      }
                                      LOG.info("Peeked element: {}", element);
                                  }
                                  break;
                          }
                      } catch (Exception e) { // 예외 발생 시
                          LOG.error("Error in processing request", e);
                      } finally {
                          iter.next(); // 다음 로그 항목으로 이동
                      }
                  }
              } finally {
                  lock.writeLock().unlock(); // 쓰기 락 해제
              }
          }
          // 스냅샷 생성 및 로드 구현 생략...
      }

      onApply 메서드는 Raft 로그 항목이 커밋될 때 호출되며, 여기서 큐에 대한 모든 작업(enqueue, dequeue, peek)이 처리됩니다.

      큐 작업 정의

      큐 작업을 로그 항목으로 인코딩하기 위한 클래스를 정의합니다.

      public class PriorityQueueRequest implements Serializable {
          private final OperationType type;
          private final Element element;
      
          public PriorityQueueRequest(OperationType type, Element element) {
              this.type = type;
              this.element = element;
          }
      
          public OperationType getType() {
              return type;
          }
      
          public Element getElement() {
              return element;
          }
      }

      Raft 서버 구현

      Raft 서버를 초기화하고 실행하는 코드는 다음과 같습니다.

      public class PriorityQueueServer {
          private final RaftGroupService raftGroupService;
          private Node node;
          private final PriorityQueueStateMachine stateMachine;
          private final String dataPath;
          private final String groupId;
          private final PeerId serverId;
          private final RpcServer rpcServer;
          private final AtomicBoolean started = new AtomicBoolean(false);
      
          public PriorityQueueServer(String dataPath, String groupId, PeerId serverId, List<PeerId> peerIds) throws IOException {
              this.dataPath = dataPath;
              this.groupId = groupId;
              this.serverId = serverId;
              // RPC 서버 초기화
              this.rpcServer = RaftRpcServerFactory.createRaftRpcServer(serverId.getEndpoint());
              // Raft 옵션 설정
              RaftOptions raftOptions = new RaftOptions();
              // 상태 머신 생성
              this.stateMachine = new PriorityQueueStateMachine();
              // 노드 옵션 설정
              NodeOptions nodeOptions = new NodeOptions();
              nodeOptions.setFsm(this.stateMachine);
              nodeOptions.setLogUri(dataPath + File.separator + "log");
              nodeOptions.setRaftMetaUri(dataPath + File.separator + "raft_meta");
              nodeOptions.setSnapshotUri(dataPath + File.separator + "snapshot");
              nodeOptions.setInitialConf(new Configuration(peerIds));
              nodeOptions.setSnapshotIntervalSecs(30);
              // Raft 그룹 서비스 및 노드 초기화
              this.raftGroupService = new RaftGroupService(groupId, serverId, nodeOptions, rpcServer);
          }
          *// 메서드 생략...*
       }

      이 코드에서는 Raft 옵션을 설정하고, 상태 머신을 초기화한 후, Raft 그룹 서비스를 시작합니다.

      클라이언트 구현

      클라이언트는 큐 작업을 Raft 그룹에 요청하는 인터페이스를 제공합니다.

      public class PriorityQueueClient {
          private final CliClientServiceImpl cliClientService;
          private final String groupId;
          private PeerId leaderId;
          private final int rpcTimeoutMs = 5000;
      
          /**
           * 우선순위 큐에 대한 특정 연산(enqueue, dequeue, peek)을 수행하고 결과를 반환합니다.
           *
           * @param type    수행할 연산 타입 (OperationType enum)
           * @param element 연산에 필요한 요소 (enqueue의 경우 추가할 요소, dequeue/peek의 경우 null)
           * @param <T>     반환 타입 (enqueue는 boolean, dequeue/peek는 Element)
           * @return 연산 결과
           * @throws Exception 연산 실패 시 예외 발생
           */
          @SuppressWarnings("unchecked")
          private <T> T processOperation(OperationType type, Element element) throws Exception {
              if (!refreshLeader()) {
                  throw new IllegalStateException("Leader not available"); // 리더가 없으면 예외 발생
              }
      
              // 요청 객체를 직렬화
              final BytesValue msg = BytesValue.of(ByteString.copyFrom(serialize(new PriorityQueueRequest(type, element))));
              CompletableFuture<T> future = new CompletableFuture<>(); // 비동기 결과를 처리하기 위한 CompletableFuture 생성
      
              // RPC 호출 결과를 처리하기 위한 ClosureAdapter 생성
              RpcResponseClosureAdapter<BytesValue> closureAdapter = new RpcResponseClosureAdapter<>() {
                  @Override
                  public void run(Status status) {
                      if (status.isOk()) { // RPC 호출 성공 시
                          BytesValue response = getResponse(); // 응답 획득
                          byte[] responseData = response.getValue().toByteArray(); // 응답 데이터 획득
                          PriorityQueueResponse<?> queueResponse = (PriorityQueueResponse<?>) deserialize(responseData); // 응답 데이터 역직렬화
                          future.complete((T) queueResponse.getResult()); // future에 결과를 설정
                      } else { // RPC 호출 실패 시
                          future.completeExceptionally(new Throwable(status.getErrorMsg())); // future에 예외를 설정
                      }
                  }
              };
      
              // 리더에게 RPC 호출
              cliClientService.invokeWithDone(
                      leaderId.getEndpoint(), // 리더의 엔드포인트
                      msg,     // Message 타입으로 래핑된 객체
                      closureAdapter, // 결과를 처리할 ClosureAdapter
                      rpcTimeoutMs // RPC 타임아웃
              );
      
              // closure에서 결과 대기 및 반환
              return future.get(rpcTimeoutMs, TimeUnit.MILLISECONDS); // future에서 결과를 얻어 반환 (타임아웃 설정)
          }
      
          /**
           * 우선순위 큐에 요소를 삽입합니다.
           *
           * @param value    요소의 값
           * @param priority 요소의 우선순위
           * @return 삽입 성공 여부
           * @throws Exception 연산 실패 시 예외 발생
           */
          public boolean enqueue(String value, int priority) throws Exception {
              return processOperation(OperationType.ENQUEUE, new Element(value, priority)); // enqueue 연산 수행
          }
      
          /**
           * 우선순위 큐에서 가장 높은 우선순위를 가진 요소를 삭제하고 반환합니다.
           *
           * @return 삭제된 요소
           * @throws Exception 연산 실패 시 예외 발생
           */
          public Element dequeue() throws Exception {
              return processOperation(OperationType.DEQUEUE, null); // dequeue 연산 수행
          }
      
          /**
           * 우선순위 큐에서 가장 높은 우선순위를 가진 요소를 반환합니다 (삭제하지 않음).
           *
           * @return 큐의 head 에 있는 요소
           * @throws Exception 연산 실패 시 예외 발생
           */
          public Element peek() throws Exception {
              return processOperation(OperationType.PEEK, null); // peek 연산 수행
          }
      
          // 메서드 생략...
      }

      클라이언트는 refreshLeader() 메서드를 통해 현재 리더 노드를 찾고, 해당 노드에 작업을 요청합니다.

      데이터 흐름 및 처리 과정

      분산 우선순위 큐에서 요청이 처리되는 과정을 단계별로 살펴보겠습니다.

      • enqueue 작업 처리 흐름

        1. 클라이언트가 enqueue(element, priority) 요청

        2. 클라이언트는 RouteTable을 통해 리더 노드 식별

        3. 리더 노드에 ENQUEUE 타입의 작업을 포함한 로그 항목 전송

        4. 리더 노드는 로그 항목을 자신의 로그에 추가

        5. 리더는 AppendEntries RPC를 통해 팔로워들에게 로그 복제 요청

        6. 과반수 노드가 로그를 복제하면 리더는 로그 항목을 커밋

        7. 모든 노드에서 onApply() 메서드가 호출되어 큐에 요소 추가

        8. 리더는 작업 결과를 클라이언트에 반환

      • dequeue/peek 작업 처리 흐름

        dequeue 작업도 유사한 흐름으로 처리되지만, 큐에서 요소를 제거하는 작업이 포함됩니다. 이 과정을 통해 모든 노드의 우선순위 큐가 동일한 상태를 유지하게 됩니다.

        peek 작업은 dequeue와 유사한 흐름으로 처리되지만, 큐에서 요소를 제거하지 않습니다.

      장애 시나리오 테스트

      분산 시스템의 핵심 가치는 장애 상황에서의 대응 능력입니다. 다양한 장애 시나리오를 테스트하여 시스템의 견고함을 확인했습니다.

      • 팔로워 노드 장애

        하나의 팔로워 노드가 다운되는 시나리오에서는 시스템이 정상적으로 작동했습니다.

        Raft 알고리즘은 과반수의 노드가 정상 작동하는 한 합의를 계속 진행할 수 있기 때문에, 3노드 환경에서 1개 노드 장애는 시스템에 영향을 미치지 않았습니다.

        장애 노드가 복구되면 리더 노드로부터 로그를 복제받아 최신 상태로 자동 복구되었습니다.

      • 리더 노드 장애

        리더 노드가 다운되는 경우, 남은 두 팔로워 노드 중 하나가 후보자가 되어 새로운 리더로 선출되었습니다.

        리더 전환 과정에서 약 200~300ms의 서비스 지연이 발생했지만, 그 이후에는 정상적으로 작업을 처리했습니다.

        • 1. 리더 노드 장애 발생

        • 2. 팔로워 노드들의 선거 타임아웃 발생 (150~300ms)

        • 3. 후보자 상태로 전환 및 투표 요청

        • 4. 새 리더 선출 및 서비스 재개

        클라이언트는 RouteTable을 통해 새로운 리더를 자동으로 발견하고 요청을 전송할 수 있었습니다.

      • 네트워크 파티션

        네트워크 파티션이 발생하여 노드들이 서로 통신할 수 없는 상황에서는 Raft의 안전성 속성이 중요한 역할을 했습니다.

        소수 파티션에 있는 노드들은 새 리더를 선출할 수 없었고, 다수 파티션에서만 작업이 계속 처리되었습니다.

        파티션이 해소된 후에는 분리되었던 노드들이 새로운 리더로부터 로그를 복제받아 상태를 동기화했습니다.

      실제 운영 시 고려사항

      분산 우선순위 큐를 실제 운영 환경에 배포할 때 고려해야 할 몇 가지 중요한 사항이 있습니다.

      • 모니터링 및 알람 설정

        시스템 운영 상태를 모니터링하고 이상 징후를 감지할 수 있는 모니터링 시스템이 필요합니다.

        다음과 같은 지표를 모니터링하는 것이 좋습니다. (팀내 어플리케이션에서는 아래 지표들을 프로메테우스로 익스포트해서 모니터링 하고 있습니다)

        • 노드 상태 (리더/팔로워/후보자)

        • 로그 복제 지연 시간

        • 커밋 인덱스 및 적용 인덱스

        • 큐 크기 및 처리량

        • 리더 선출 횟수

      • 백업 및 복구 전략

        데이터 손실을 방지하기 위한 백업 및 복구 전략이 필요합니다.

        • 정기적인 스냅샷 생성 및 외부 저장소 백업

        • 로그 압축을 통한 디스크 공간 관리

        • 복구 시나리오 테스트 및 절차 문서화

      • 클러스터 확장

        시스템 부하가 증가할 경우 클러스터를 확장하는 방법을 고려해야 합니다.

        • 동적 멤버십 변경을 통한 노드 추가

        • 샤딩을 통한 수평적 확장

        • MULTI-RAFT-GROUP을 활용한 큐 파티셔닝


      결론 및 향후 과제

      이 글에서는 Raft 알고리즘을 간단히 살펴보고, sofa-jraft 라이브러리를 활용하여 고가용성을 갖춘 분산 우선순위 큐를 구현하는 방법을 소개했습니다.

      Raft 알고리즘은 분산 시스템 구축의 복잡성을 크게 줄여주며, sofa-jraft와 같은 안정적인 라이브러리를 활용하면 실제 프로덕션 환경에서도 활용 가능한 고가용성 시스템을 구현할 수 있습니다.

      향후에는 분산 우선순위 큐의 성능 최적화, 장애 복구 시나리오 테스트, 다양한 워크로드에 대한 확장성 검증 등 추가적인 연구와 개발이 필요합니다.

      또한, 실시간 모니터링 및 운영 자동화 도구와의 연동, 보안 강화, 그리고 다양한 언어 및 플랫폼 지원 등도 고려해볼 만한 과제입니다.

      분산 시스템 개발에 관심이 있는 분들께 이 글이 도움이 되길 바랍니다. 질문이나 의견이 있으시면 언제든 댓글로 남겨주세요.

      앞으로도 실전에서 활용할 수 있는 다양한 분산 시스템 구현 방법을 소개할 예정이니 많은 관심 부탁드립니다.

      댓글 0

      DEVOTEE를 활성화 시키면
      지금 작성한 댓글에 AI가 댓글을 달아줍니다.

      marine 님의 최신 블로그

      더보기
      동영상 기고하기