Robust And Versatile Open-source spring batch based backup solution

SangYeon Min·2025년 6월 30일

PROJECT-RAVO

목록 보기
1/2

RAVO는 Spring Batch와 ORM 기술을 활용하여 MySQL 데이터베이스의 데이터를 주기적으로 자동 백업하는 시스템을 구축하고 이를 Helm을 통해 배포 가능한 프로젝트를 구현하는 것을 목표로 한다.

특히, 대용량 데이터를 안정적이고 효율적으로 처리하기 위해 Chunk 기반 처리, 트랜잭션 관리, 체크포인트, 오류 복구 등을 도입하고자 한다.
이와 함께 자동화된 실시간 덤핑 및 전환 기능을 통해 서비스의 연속성을 보장하고자 하며 이는 장애 발생 시 Hot Standby DB Server로 신속하게 전환되어 다운타임을 최소화할 수 있을 것으로 기대된다.

또한, 자체적으로 배포 및 관리가 가능한 오픈소스 백업 솔루션을 도입하여 운영 부담을 경감하기 위해 개발자들이 익숙한 JPA 엔티티 기반으로 백업 및 복구 프로세스를 구성함으로써 사용 편의성 또한 향상시키는 것을 목표로 한다.

팀원 및 역할 분담

  • 민상연 (팀장) : K3s Cluster 구축 및 서비스 배포, RAVO-MANAGER 개발
    RAVO K3s 애플리케이션 및 Kafka, Debezium 등 배포 및 운영
  • 김규빈 (팀원) : Schedular Trigger, Cold Standby Backup process 개발
  • 노성준 (팀원) : Recovery API, dump data 기반 Recovery process 개발

프로토타입 구현 범위와 목적

RAVO 프로토타입은 “백업-동기화-모니터링” 핵심 흐름을 최소 기능 제품 수준으로 검증하는 데 목적을 둠
범위는 전체 시스템 아키텍처를 축소 및 단순화하여 리스크가 높은 기능만 우선 구현함, 실제 운영 환경과 유사한 단일 노드 K3s 클러스터에서 end-to-end 동작을 실험·측정할 수 있도록 한정함

  1. 주기적 백업 실행 및 자동화 (Cold Standby)

    • 스프링 배치(Spring Batch)를 활용한 주기적인 데이터 덤핑 기능 구현
    • 백업 작업 중 오류 발생 시 체크포인트 복구 및 재시도 로직 제공
  2. 데이터 복구 (Recovery)

    • 백업된 데이터 파일을 기반으로
      Active DB 또는 Standby DB로 복원 가능한 기능 제공
    • 복구 대상 시점(타임스탬프 또는 백업 버전)을 지정하여 복구 가능
    • 테이블 간 참조 무결성을 고려한 복구 순서 및 트랜잭션 처리 보장
  3. 실시간 데이터 백업 및 전환 (Hot Standby)

    • Active DB와 Hot Standby DB 간 실시간 동기화 기능 제공
    • 장애 발생 시, 자동으로 Hot Standby DB 서버로
      페일오버(전환)하여 다운타임 없이 서비스 연속성을 유지
    • 복구 후 Active DB의 정상 복구(Hot Standby DB와의 동기화) 프로세스 마련
  4. 모니터링 및 관리 인터페이스 (RAVO-MANAGER)

    • 시스템 상태, 백업 성공 여부, 데이터 연속성과 무결성을
      실시간으로 확인할 수 있는 대시보드 제공
    • 주기적인 데이터 연속성 및 무결성 검사 기능 포함

프로토타입 시스템 구성도

  • Active DB Server에 주기적으로 데이터를 요청하여
    Host Standby 상태의 별도 DB Server에 현재 데이터를 저장할 수 있도록 구현

    • 장애 발생 시 빠르게 Hot Standby DB Server로 Failover 가능
    • Active DB 복구 이후 Hot Standby DB Server와 Sync 가능
  • 이와 동시에 Cold Standby 상태의 Persistent Volume에 저장

    • Cold Standby 상태의 Storage를 통해 복구를 가능케 하거나
      DB가 공격 등으로 인해 손상되었을 때 특정 Snapshot으로 복구 가능하도록 설계
  • Cluster 자체 장애 시 Transaction 영향을 줄여 Down Time 최소화

  • Prod 환경은 단일 노드 K3s 클러스터로 구성되며 Prometheus-Grafana-Stack으로
    클러스터를 모니터링하고 nginx Proxy 설정을 통해 외부 접근을 제어

프로토타입 설계 내용

Spring Application 전체 구조

본 프로토타입은 Spring Boot와 Spring Batch를 기반으로 설계 및 구현됐음. Spring Boot는 데이터베이스 연결 및 관련 설정을 유연하고 효율적으로 관리할 수 있게 지원하며, 시스템의 유지보수성과 확장성을 높이는 역할을 함.

Spring Batch는 대량 데이터 처리에 특화된 프레임워크로, 최상위 단위인 Job을 중심으로 데이터 처리 작업을 구성함. 각 Job은 여러 개의 Step으로 구성되며, Step 간 유기적이고 명확한 작업 흐름을 제공함. Job과 Step의 실행 관련 정보는 Job Repository에 기록돼, 작업 중 발생한 문제점을 신속히 파악하고 실행 시간과 결과 등을 효율적으로 모니터링할 수 있음.

프로토타입에서는 ColdstandbyBackupJobColdstandbyRecoveryJob 두 가지 Job을 구현했음. 추가로 Hotstandby Recovery 기능만 PoC 수준에서 부분적으로 구현된 상태임.


Coldstandby Backup 동작 흐름

  1. 스케줄러 동작

    • Spring Scheduler와 cron 표현식을 활용해 설정된 주기에 따라 Backup Job이 자동 실행됨.
    • 각 실행마다 고유한 타임스탬프를 파라미터로 전달해 재실행 가능성을 보장함.
  2. 백업 전략 선택

    • 애플리케이션은 데이터베이스 드라이버 정보를 기반으로 적절한 백업 전략 구현체를 선택하도록 설계됨.
    • 이를 통해 확장성과 유지보수성을 높이고 다양한 DBMS 환경에서 일관된 백업 전략을 적용할 수 있음.
  3. mysql dump backup 실행

    • Datasource 설정 정보(JDBC URL, username, password)를 기반으로 Active DB를 대상으로 mysqldump 명령어를 실행함.
    • --single-transaction, --quick, --add-drop-database 옵션을 사용해 인덱스 잠금을 최소화하고 대량 데이터 처리를 효율적으로 수행함.
  4. 백업 저장

    • 생성된 SQL 덤프 파일은 전용 저장소 디렉터리에 저장됨.
    • 파일명은 <DB_Name>_cold_YYYYMMDDHHmmss.sql 형태로 구성해 명확한 백업 시점을 식별할 수 있게 함.
  5. 백업 완료

    • 백업 완료 시 덤프 파일 위치 및 생성 시간이 로그로 기록돼 관리됨.

Coldstandby Recovery 동작 흐름

  1. 사용자 복구 요청

    • REST API 또는 웹 폼(HTML Form)을 통해 요청받음.
    • RESTful API:
      • GET /recovery: SQL 덤프 파일 목록 조회
      • POST /recovery: 요청 본문에 포함된 filename과 일치하는 덤프 파일로 복구 진행
  2. 복구 프로세스 호출

    • Controller가 요청을 처리해 Spring Batch Recovery Job을 실행함.
    • 복구할 파일 이름을 Job Parameter로 전달해 정확한 파일 복구를 보장함.
  3. SQL Dump 파일 조회

    • 파라미터에 해당하는 덤프 파일 존재 여부 확인
      • 파일이 존재하면 복구 프로세스 진행
      • 파일이 없으면 Exception 발생 후 사용자에게 오류 전달
  4. mysql dump recovery 실행

    • Datasource 설정 정보 기반으로 mysql CLI 명령어를 실행해 복구함.
    • 덤프 파일의 SQL 내용을 표준 입력으로 전달해 데이터를 복원함.
  5. 복구 완료

    • 복구 성공 시 데이터베이스 상태 및 처리 결과가 로그로 기록돼 운영자가 결과를 확인할 수 있음.

Hotstandby Recovery 동작 흐름

동작 흐름

Active DB 다운 감지 → Standby DB로 failover → 이후 Standby DB로 들어오는 쓰기 요청이 Active DB 복구 시 순차 반영돼야 함

  1. 변경 이벤트 캡처 (Debezium → Kafka)

    1. Debezium이 Standby DB의 integrity_data 테이블 변경(insert/update/delete) 이벤트를 CDC 방식으로 캡처
    2. 캡처된 이벤트를 JSON 형태로 Kafka 토픽에 발행
  2. Kafka 소비자(StandbyConsumer)에서 이벤트 수신

    1. @KafkaListener가 활성화돼 ravo_db.ravo_db.integrity_data 토픽 수신
    2. 레코드 수신 시 DebeziumToSqlConverter.convertDebeziumToSQL()에 JSON 전달
  3. JSON → SQL 변환 (DebeziumToSqlConverter)

    1. JSON 파싱
    2. payload.op 값("c", "u", "d")에 따라 SQL 생성:
      • "c": buildInsertSQL(after)
      • "u": buildUpdateSQL(before, after)
      • "d": buildDeleteSQL(before)
  4. SQL 큐에 등록 (InMemorySqlQueueService.enqueue)

    • ConcurrentLinkedQueue<String>에 SQL 문자열 저장
  5. 주기적 배치 실행 (flushQueue)

    1. @Scheduled(fixedDelay = 1000)로 1초마다 flushQueue() 호출
    2. Active DB 헬스체크 후 큐에서 SQL을 꺼내 실행
  6. Active DB 헬스체크 구현 (JdbcActiveDbHealthChecker)

    • jdbcTemplate.queryForObject("select 1", Integer.class) 실행
    • 정상 응답 시 true, 예외 발생 시 false 반환

<요약>

[Active DB 다운]
↓ 페일오버
[Standby DB + Debezium] → Kafka(토픽: integrity_data)
↓ consume
[StandbyConsumer] → [DebeziumToSqlConverter]
→ [InMemorySqlQueueService.queue]
↓ (매 1초)
[flushQueue] --헬스체크--> [Active DB 복구 전: 대기]
└→ [Active DB 복구 후: queue → execute → poll]

프로토타입 시험 시나리오

단계고위험 상황 가정실제 테스트 방안RAVO 제공 기능성공 기준
Cold Standby 백업 엔진중간고사 48시간 전 랜섬웨어 감염Cold-Standby 경과 추적02:30 자동 풀-Dump
백업 성공률 모니터링
백업 성공률 ≥99%
스토리지 비용 −40%
Cold Standby Recovery API24시간 내 DB 복구 의무 규정 존재기존 Active DB 데이터 삭제 후
RAVO-SPRING Recovery 요청
Point-in-Time 복구 호출
무결성 검사
RTO < 30분 (Recovery Time Objective)
데이터 손실 0건
Hot Standby 활성화수강신청 D-Day 전원·냉각 동시 장애mysql-active replicas=0 조정30초 내 자동 Hot Standby 승격
K3s SVC EP 재포인팅
Failover ≤30초
인력 투입 0명
Standby Active 재동기화장애 복구 후 Active DB 재투입Standby 활성화 후 MySQL DML 명령
mysql-active replicas=1 조정
K3s SVC EP 재포인팅
Probe 기반 RW 전환 검사
재동기화 < 10분
인력 투입 0명

단계 1. Cold Standby 백업 엔진

Coldstandby Backup Job은 application.yml 파일에 설정된 cron 스케줄러 시간에 따라 자동으로 트리거링됨
스케줄러 시간에 맞춰 Backup Job 실행되는 것 확인함
실행 완료되면 작업 결과에 따라 Exit Status 반환됨
Exit Status가 COMPLETE일 경우 정상 완료 의미함, FAILED일 경우 정책에 따라 작업 재시도하거나 실패 정보 알림 이메일 사용자에게 발송 등 추가 처리 가능함
프로토타입에서는 이 추가 처리 구현이 되지 않은 상태임

백업 작업 성공적으로 완료되면 application.yml에서 설정한 backup.output-dir 경로에 덤프 SQL 파일 생성됨
파일명은 백업이 실행된 시간을 기준으로 지정됨
또한, 백업 파일 압축 시 원래 파일 크기의 최대 30%까지 줄일 수 있으나 프로토타입에 압축 기능 아직 구현되지 않은 상태임

백업 성공률 ≥ 99% / 스토리지 비용 -40% 목표 달성

단계 2. Cold Standby 백업 엔진

Coldstnadby Recovery는 Active DB 테이블 데이터를 전부 삭제한 후 Recovery 로직을 수행하여 기존 테이블 정보 정상 복구 확인하는 방식으로 테스트함

이후 백업 파일을 선택하여 복구 진행

RTO (Recovery Time Objective) < 30분 / 데이터 손실 0건 목표 달성

단계 3. Hot Standby 활성화

apiVersion: v1
kind: ServiceAccount
metadata:
  name: mysql-failover-watcher
  namespace: default
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
  name: mysql-failover-role
  namespace: default
rules:
  - apiGroups: [""]
    resources: ["services","endpoints","pods"]
    verbs: ["get","patch","list","watch"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: mysql-failover-binding
  namespace: default
roleRef:
  apiGroup: rbac.authorization.k8s.io
  kind: Role
  name: mysql-failover-role
subjects:
  - kind: ServiceAccount
    name: mysql-failover-watcher
---
apiVersion: apps/v1
kind: Deployment
metadata:
  name: mysql-failover-watcher
spec:
  replicas: 1
  selector:
    matchLabels:
      app: mysql-failover-watcher
  template:
    metadata:
      labels:
        app: mysql-failover-watcher
    spec:
      hostNetwork: true
      serviceAccountName: mysql-failover-watcher
      volumes:
        - name: state-volume
          emptyDir: {}
      containers:
        # kubectl-watcher 컨테이너: 서비스 셀렉터 패치
        - name: kubectl-watcher
          image: bitnami/kubectl:latest   # kubectl 내장 이미지를 사용
          volumeMounts:
            - name: state-volume
              mountPath: /state
          command:
            - /bin/sh
            - -c
            - |
              PREV_STATE=""
              while true; do
                # Active Deployment Pod 중 Ready 상태 개수 조회
                EP_ACTIVE=$(kubectl get pods -l app=mysql-active \
                  -o jsonpath='{.items[?(@.status.containerStatuses[0].ready==true)]}' \
                  | wc -w)

                # EP_ACTIVE가 0이면 standby, 그 외엔 active
                if [ "$EP_ACTIVE" -lt 1 ]; then
                  DESIRED_STATE="standby"
                else
                  DESIRED_STATE="active"
                fi

                # 상태가 바뀌었을 때만 Service selector 패치
                if [ "$DESIRED_STATE" != "$PREV_STATE" ]; then
                  kubectl patch svc mysql-active-service \
                    --type=json \
                    -p "[{\"op\":\"replace\",\"path\":\"/spec/selector/app\",\"value\":\"mysql-${DESIRED_STATE}\"}]"
                  
                  # 상태를 /state/state 파일에 기록
                  echo "${DESIRED_STATE}" > /state/state
                fi

                # 로그 출력
                echo "$(date '+%Y-%m-%d %H:%M:%S') [FailoverWatcher] State→${DESIRED_STATE} (prev=${PREV_STATE}) ActiveReady=${EP_ACTIVE}"

                PREV_STATE="$DESIRED_STATE"
                sleep 15
              done
        # conntrack-watcher 테이너 : conntrack 테이블 플러시
        # # Standby -> Active 전환시 아래 Sidecar Container를 통해 Gracefully MySQL TCP 연결 해제
        - name: conntrack-watcher
          image: nicolaka/netshoot:latest
          securityContext:
            capabilities:
              add: ["NET_ADMIN"]
          volumeMounts:
            - name: state-volume
              mountPath: /state
          command:
            - /bin/sh
            - -c
            - |
              PREV_FILE_STATE=""
              while true; do
                if [ -f /state/state ]; then
                  CUR_STATE=$(cat /state/state)
                  # 상태가 바뀌었으면 conntrack 삭제
                  if [ "$CUR_STATE" != "$PREV_FILE_STATE" ]; then
                    conntrack -D -p tcp --dport 32306
                    echo "$(date '+%Y-%m-%d %H:%M:%S') Flushed conntrack for 32306 (state=${CUR_STATE})"
                  fi
                  PREV_FILE_STATE="$CUR_STATE"
                fi
                sleep 5
              done
  • default 네임스페이스에 mysql-failover-watcher ServiceAccount 생성

  • Role 및 RoleBinding 통해 services, endpoints, pods 자원 get, patch, list, watch 권한 부여

  • Deployment는 hostNetwork 모드 사용

  • state-volume 빈 디렉터리를 두 개의 사이드카 컨테이너가 공유

  • bitnami/kubectl

    • 15초 간격 mysql-active 레이블 파드 준비 상태 조회
    • 준비된 파드 없으면 standby, 있으면 active 판단
    • Service selector.app 값을 mysql-standby/ mysql-active로 교체
    • 변경 상태 /state/state 파일 기록
    • 상태 변화와 준비된 파드 수 표준 출력 로깅
  • nicolaka/netshoot

    • NET_ADMIN 권한 conntrack 유틸리티 실행
    • /state/state 파일 5초 주기 모니터링
    • 상태 변경 시 MySQL 포트(32306) conntrack 테이블 항목 삭제
    • TCP 연결 해제 후 새로운 연결 원활 수립
    • 자동 페일오버 이후 안정적 DB 연결 유지

초기 DB 상태
replicas=0으로 조정 (Active DB Down)
Standby DB로 자동 Failover
replicas=1로 조정 (Active DB Up)
Endpoint Recover (Application 수정 없이 TCP 연결 해제)

Failover ≤ 30초 / 인력 투입 0명 목표 달성

단계 4. Standby → Active 재동기화

Debezium은 Database에서 발생하는 변경사항을 추적할 수 있는 일종의 Apache Kafka Connect의 Source Connector

각각의 connector은 해당 데이터베이스의 CDC(Change Data Capture)와 관련된 기능을 활용해서 변경된 데이터에 대한 정보를 가져옴. 또한 성공적으로 commit이 발생한 데이터에 대해서만 변경사항이 전파되기 때문에 실패한 트랜잭션은 고려할 필요가 없음

더 나아가 Debezium은 변경사항을 디스크에 저장하기 때문에 데이터의 변경사항을 전달받아야 하는 애플리케이션이 다운되더라도 문제없음

# bitnami 최신 Kafka 차트와 Kafka 최신 버전 (3.x 이상)에서는 KRaft 모드 사용 권장 

# KRaft 모드는 ZooKeeper 없이 내장된 Raft 합의 기반 Quorum Controller로 메타데이터 관리를
# Kafka 자체에 통합하여 아키텍처를 단순화하고 배포·관리 효율을 높인다
# 별도의 ZooKeeper 클러스터 없이 Kafka 브로커들이 스스로 메타데이터(토픽·파티션 정보, 리더 선출 등)를 관리

# KRaft 모드는 기존 ZooKeeper 클러스터와의 직접 마이그레이션이 불가능하고
# 컨트롤러 롤링 업데이트, 버전 업/다운그레이드를 지원하지 않으며
# 일부 인증, 스토리지 기능도 미지원되어 운영 유연성과 안정성이 제한
kraft:
  enabled: false

# autoDiscovery Init-Container 를 쓰려면 RBAC 와 ServiceAccount 토큰 권한이 필요
# RBAC 리소스 생성
rbac:
  create: true

# 하지만 broker의 리소스를 아껴 전체 클러스터의 리소스 제한에 맞추기 위해
# zookeeper 모드를 통해 Kafka 클러스터 배포
# KRaft 모드에선 Zookeeper 비활성화
zookeeper:
  enabled: true
  replicaCount: 1
  resources:
    requests:
      cpu: 100m
      memory: 256Mi
    limits:
      cpu: 250m
      memory: 512Mi

### [TroubleShooting: SASL Error]
# broker 정의 상단에서 진행해야 자동 주입 방지 가능
# 인증 방식을 기본 SASL → plaintext로 변경
auth:
  clientProtocol: plaintext
  interBrokerProtocol: plaintext
  sasl:
    # 빈 배열로 두면 SASL 완전 비활성화
    enabledMechanisms: []
  tls:
    type: ""

### [TroubleShooting: SASL Error]
# 브로커는 토픽 생성 요청을 받은 뒤, 내부 ZK_BROKER(SASL)나 CLIENT 리스너로 컨트롤러에게 파티션 할당을 요청해야 함
# 그 리스너들(SASL 9095, CLIENT 9092)은 외부에 노출되지 않아(=포워딩/NGINX 처리되지 않아) 연결이 끊기고, 결국 할당 대기 타임아웃 발생
# Broker Pod 의 환경변수로 Kafka 설정 옵션을 넘기는 방식

# 28.x 이후 Bitnami Kafka 차트는 extraEnvVars: 블록을 브로커별(server‑0,1,2) ConfigMap 이 아닌 제너릭 값으로 취급
# 그러나 템플릿 내부에서 listeners.* 기본값을 나중에 다시 써 넣기 때문에 우리가 넣은 덮어쓰기 줄이 앞쪽에 들어가고, 뒤에 오는 디폴트가 다시 9094·9095 리스너를 살려 버림림
# 중복을 전부 덮어쓰려면 extraConfig 를 써야 함함

# EXTERNAL 리스너가 없어도 외부 접속은 충분히 가능
# 리스너 이름의 경우 브로커 내부에서 구분용으로 붙인 라벨일 뿐 외부‑내부를 가르는 기능이 존재X
# Bitnami 차트 기본 SASL 리스너를 완전히 덮어쓰는 환경변수

### [Trouble Shooting: Timed out waiting for a node assignment]
# Kafka는 메타데이터 기반 다단계 연결
# 1. 클라이언트가 --bootstrap-server 로 지정한 호스트:포트에 접속 → 브로커로부터 토픽·파티션·리더 정보 등 메타데이터 수신
# 2. 메타데이터에 담긴 advertised.listeners(호스트:포트) 정보를 이용해, 실제 데이터를 주고받을 브로커에 새로운 TCP 연결을 시도
# 3. 이후 프로듀서/컨슈머는 이 새 연결을 통해 입출력을 수행
# advertised.listeners에 외부에서 접근할 수 없는 내부 IP(또는 NAT IP)를 두면
# 1단계(bootstrap)는 통과해도 2단계(메타데이터 재접속)에서 클라이언트가 잘못된 주소로 접속을 시도해 타임아웃
extraConfig: |-
  listener.security.protocol.map=CLIENT:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT

externalAccess:
  enabled: true
  autoDiscovery:
    # 브로커가 자신의 외부 IP 자동 광고 끄기
    # NAT 외부 호스트의 External IP를 직접 지정하기 위함
    enabled: false
  broker:
    service:
      type: NodePort
      # 기존 bitnami 차트 설정에 맞게 EXTERNAL 포트 (9095) 명시
      ports:
        external: 9095
      # 브로커 1개
      nodePorts: [30092]
      # [TroubleShooting: k3s 의 externalIPs hairpin/reflection 미지원 → Pod 로 가는 리턴 패킷이 날아가지 못해 세션이 끊어짐]
      # 이 옵션을 끄면 Service는 externalIPs: 가 아닌 일반 NodePort 경로로만 동작
      useHostIPs: false
      # 외부 클라이언트가 연결할 Public IP
      domain: <Host Public IP>

# [__consumer_offsets 토픽 Replica 3 ERR]
# internal topic replication factor 조정
configuration: |-
  offsets.topic.replication.factor=1
  transaction.state.log.replication.factor=1

broker:
  # Pod 에 토큰 마운트
  automountServiceAccountToken: true
  # Bitnami Kafka 차트에서 브로커가 실제로는 9092 만 LISTEN 하는데
  # startupProbe·readinessProbe 기본값의 문제로 Startup Probe 설정 적용이 불가
  # 아무런 로그 없이 broker-0 Pod가 정지되고 재시작
  # Bitnami Kafka 템플릿이 실제로 참조하는 필드는 두 곳이기에 두 곳 모두 변경

  # Pod 클라이언트/내부/외부 리스너 포트 모두 정의
  containerPorts:
    client:   9092
    internal: 9094
    external: 9095
  # 서비스/광고/외부 접근용
  # Pod 클라이언트/내부/외부 리스너 포트 모두 정의
  ports:
    client:   9092
    internal: 9094
    external: 9095
  # 브로커 1 개
  replicaCount: 1
  resources:
    requests:
      cpu: 300m
      memory: 1024Mi
    limits:
      cpu: 600m
      memory: 1536Mi
  # JVM Heap 옵션을 리소스 제한 보다 낮게 설정
  heapOpts: "-Xmx1280m -Xms1280m"
  startupProbe:
    enabled: true
    failureThreshold: 120
    periodSeconds: 10
    # probe 에서 broker ContainerPort 번호를 직접 지정
    tcpSocket:
      port: client
  readinessProbe:
    periodSeconds: 15
    tcpSocket:
      port: client
# 기본값인 false 상태면 컨트롤러 설정 요구
defaultInitContainers:
  autoDiscovery:
    enabled: true

# 호스트 경로 방식의 간편한 저장소
persistence:
  size: 5Gi

# initContainer 권한을 위해 단순 명시 필요
controller:
  replicaCount: 0
  automountServiceAccountToken: true

# controller:
#   replicaCount: 0
#   resources:
#     requests:
#       cpu: 100m
#       memory: 512Mi
#     limits:
#       cpu: 200m
#       memory: 768Mi
#   # JVM Heap 옵션을 리소스 제한 보다 낮게 설정
#   heapOpts: "-Xmx512m -Xms512m"
#   # Kafka 컨트롤러(Quorum Controller)처럼 초기 구동에 시간이 오래 걸리는 애플리케이션이
#   # 완전히 시작될 때까지 Kubernetes의 liveness나 readiness 프로브에 의해 잘못 재시작되거나 서비스 트래픽에서 제외되는 것을 방지

#   # 저사양 디스크 + 512 Mi heap이면 Controller 세팅에 1~3분, 간혹 5분 이상 소요
#   # 따라서 Startup Probe를 최대 20분으로 설정
#   startupProbe:
#     enabled: true
#     failureThreshold: 120   # 120×10s = 20 min
#     periodSeconds: 10
#   readinessProbe:
#     periodSeconds: 15
# DB 접속 자격증명용 Secret
apiVersion: v1
kind: Secret
metadata:
  name: debezium-secret
type: Opaque
stringData:
  MYSQL_USER: debezium
  MYSQL_PASSWORD: dbz!pass
---
# Debezium Server MicroProfile Config
apiVersion: v1
kind: ConfigMap
metadata:
  name: debezium-config
data:
  application.properties: |
    ### Sink → Kafka (기본 설정)
    debezium.sink.type=kafka
    debezium.sink.kafka.bootstrap.servers=kafka-broker-headless.kafka.svc.cluster.local:9092
    debezium.sink.kafka.producer.bootstrap.servers=${debezium.sink.kafka.bootstrap.servers}
    debezium.sink.kafka.producer.key.serializer=org.apache.kafka.common.serialization.StringSerializer
    debezium.sink.kafka.producer.value.serializer=org.apache.kafka.common.serialization.StringSerializer

    ### Source → MySQL(standby)
    debezium.source.connector.class=io.debezium.connector.mysql.MySqlConnector
    debezium.source.database.hostname=mysql-standby-service
    debezium.source.database.port=3306
    debezium.source.database.user=${MYSQL_USER}
    debezium.source.database.password=${MYSQL_PASSWORD}
    debezium.source.database.server.id=223344
    debezium.source.database.server.name=ravo_db
    debezium.source.database.include.list=ravo_db
    debezium.source.table.include.list=ravo_db.integrity_data
    debezium.source.topic.prefix=ravo_db

    ### Kafka Connect Worker가 요구하는 최상위 설정들
    offset.storage=org.apache.kafka.connect.storage.KafkaOffsetBackingStore
    offset.storage.bootstrap.servers=${debezium.sink.kafka.bootstrap.servers}
    # 최상위 debezium.source.bootstrap.servers 설정 필수적으로 추가
    debezium.source.bootstrap.servers=${debezium.sink.kafka.bootstrap.servers}
    offset.storage.topic=debezium-offsets
    offset.storage.partitions=1
    offset.storage.replication.factor=1

    ### Offset 저장소 (Connector 레벨 설정, 가독성용으로 유지)
    debezium.source.offset.storage=org.apache.kafka.connect.storage.KafkaOffsetBackingStore
    debezium.source.offset.storage.bootstrap.servers=${debezium.sink.kafka.bootstrap.servers}
    debezium.source.offset.storage.topic=debezium-offsets
    debezium.source.offset.storage.partitions=1
    debezium.source.offset.storage.replication.factor=1

    ### Schema-history 저장소 (Kafka)
    debezium.source.schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
    debezium.source.schema.history.internal.kafka.bootstrap.servers=${debezium.sink.kafka.bootstrap.servers}
    debezium.source.schema.history.internal.kafka.topic=debezium-schema-history

---
apiVersion: apps/v1
kind: Deployment
metadata:
  name: debezium-server
spec:
  replicas: 1
  selector:
    matchLabels: { app: debezium-server }
  template:
    metadata:
      labels: { app: debezium-server }
    spec:
      containers:
        - name: dbz
          #  Quay 레지스트리 :contentReference[oaicite:0]{index=0} 
          image: quay.io/debezium/server:3.1.1.Final
          envFrom:
            - secretRef: { name: debezium-secret }
          volumeMounts:
            # Config 파일을 Debezium이 찾는 위치(/debezium/config)로 폴더 통째로 주입
            - name: cfg
              mountPath: /debezium/config
          resources:
            requests:
              memory: "350Mi"
              cpu: "100m"
            limits:
              memory: "500Mi"
              cpu: "300m"
          ports:
            - containerPort: 8080   # health / metrics
          # probe 에서 broker ContainerPort 번호를 직접 지정
          # startupProbe를 길게 설정하여 초기 debezium 서버 배포 안정화
          startupProbe:
            httpGet: { path: /q/health/live,  port: 8080 }
            failureThreshold: 600
            periodSeconds: 10
          livenessProbe:
            httpGet: { path: /q/health/live,  port: 8080 }
            initialDelaySeconds: 20
            periodSeconds: 15
          readinessProbe:
            httpGet: { path: /q/health/ready, port: 8080 }
            initialDelaySeconds: 20
            periodSeconds: 15
      volumes:
         - name: cfg
           configMap:
             name: debezium-config
---
# 내부 접근용 Service
apiVersion: v1
kind: Service
metadata:
  name: debezium-server
spec:
  type: ClusterIP
  selector:
    app: debezium-server
  ports:
    - name: http
      port: 8080
      targetPort: 8080

Debezium을 배포하기 위한 매니페스트는 MySQL 접속 자격 증명을 담은 Opaque 타입 Secret을 정의하여 Debezium Server 컨테이너가 MYSQL_USERMYSQL_PASSWORD를 환경변수로 안전하게 주입받을 수 있도록 구성함

또한 ConfigMap에는 Debezium이 Kafka 싱크, MySQL 스탠바이 소스, 오프셋 저장소 및 스키마 히스토리 토픽 설정을 일괄 로드할 수 있는 application.properties가 포함됨

Deployment는 replicas: 1로 단일 인스턴스를 배포하며 quay.io/debezium/server:3.1.1.Final 이미지를 사용함. envFrom을 통해 Secret을 참조하고 volumeMount로 ConfigMap을 /debezium/config에 마운트해 모든 설정을 자동으로 적용함. 리소스 요청 및 제한을 통해 CPU와 메모리 사용량을 제어하며 startupProbe, livenessProbe, readinessProbe/q/health/live/q/health/ready 엔드포인트를 모니터링함

ClusterIP 타입 Service는 Debezium Server의 HTTP 포트를 노출하여 클러스터 내부 애플리케이션이 메트릭과 상태 확인 엔드포인트에 접근할 수 있도록 지원
이를 통해 이벤트 캡처와 데이터 파이프라인의 안정적인 운영을 보장하며, 모든 구성 요소를 하나의 YAML 파일로 중앙 집중 관리할 수 있음

Debezium Kafka Topic 및 정상 작동 확인
임의 테스트 데이터 DML 연산 후 Topic 내 Message 생성 확인

Standby -> Active 재동기화를 위한 Application 로직

  1. Kafka message Listener
    1) kafka 리스너 등록

    • @KafkaListener애노테이션을 통해 미리 정의한 토픽으로부터 메시지를 구독

    2) 메시지 수신

    • 토픽에 새 레코드가 들어오면 listen() 메서드가 호출되고, ConsumerRecord<String, String> record로 JSON 문자열을 받음

    3) SQL 변환

    • 받은 JSON 메시지를 DebeziumToSqlConverter.convertDebeziumToSQL(record.value())로 전달해 INSERT/UPDATE/DELETE용 SQL 문자열로 변환

    4) 큐에 적재

    • sqlQueueService.enqueue(sql)를 호출해 내부의 스레드 세이프 큐(InMemorySqlQueueService)에 SQL을 저장
  2. DebeziumToSqlConverter

    1) 입력 받은 JSON 파싱 (convertDebeziumToSQL)

    • ObjectMapper로 JSON 문자열을 JsonNode트리로 변환
    • 최상위 노드에서 "payload"를 꺼내고 그 안의 "op"(operation) 값을 확인

    2) 변경 타입 분기 (switch (op))

    • "c"(create) : buildInsertSQL(payload.get("after")) 호출
    • "u"(update) : buildUpdateSQL(payload.get("before")) 호출
    • "d"(delete) : buildDeleteSQL(payload.get("before")) 호출
    • 그 외 : Unsupported operation type 예외 발생

    3) INSERT SQL 생성 (buildInsertSQL)

    • after노드에서 id, data, checked_at(밀리초) 값 추출
    • checked_at을 formatMillisToDateTime으로 "yyyy-MM-dd HH:mm:ss" 문자열로 변환

    4) UPDATE SQL 생성 (buildUpdateSQL)

    • before에서 id, after에서 data, checked_at 추출 및 포맷

    5) DELETE SQL 생성 (buildDeleteSQL)

    • before에서 id 추출
  3. Enqueue SQL into thread-safe in-memory Queue & flush Queue when Active DB is healthy
    1) SQL 저장 (enqueue)

    • Debezium → Kafka → StandbyConsumer를 거쳐 전달된 SQL 문자열을
    • 스레드-세이프한 ConcurrentLinkedQueue(queue)에 queue.add(sql)로

Active/Standby DB Active DB 준비 및 연결 상태

Active DB replicas=0 설정 이후 헬스체크 실패 상태

Standby DB Debezium Produced 데이터 추가 및 메시지 수신

Active DB replicas=1 Kafka Partition 이후 재연결 및 내 데이터 Flush

재동기화 < 10 / 0 분 인력 투입 명 목표 달성

Github 링크

RAVO-MANAGERRAVO-BackendRAVO-INFRA
https://github.com/TEAM-RAVO-KU/RAVO-MANAGERhttps://github.com/TEAM-RAVO-KU/RAVO-Backendhttps://github.com/TEAM-RAVO-KU/RAVO-INFRA

Youtube 시연 영상

RAVO Failover 테스트RAVO Standby Active Sync 테스트
https://www.youtube.com/watch?v=TLTcue46aTchttps://www.youtube.com/watch?v=Csq_bNbw6nk

0개의 댓글