학습 목표

시나리오

서울 리전(primary)에 운영 Kafka 클러스터가 있습니다. 규제 요건으로 다른 리전(secondary)에 재해 복구용 사본을 유지해야 합니다. RPO는 5분, RTO는 30분입니다.

추가 요구사항이 있습니다. 분석 팀이 운영 클러스터에 직접 붙는 것을 막고 싶어서, 분석용 컨슈머는 secondary에서만 읽게 하려 합니다. 그리고 장애 시 운영 컨슈머를 secondary로 전환할 때 처음부터 다시 읽지 않아야 합니다.

MirrorMaker 2로 active/passive 구성을 만들고, 오프셋 동기화까지 켭니다.

아키텍처

MirrorMaker 2는 세 종류의 커넥터로 구성됩니다. 전용 클러스터 모드(connect-mirror-maker.sh)로 띄우면 이들이 자동으로 생성됩니다.

MirrorMaker 2의 커넥터 3종
커넥터옮기는 것주요 설정
MirrorSourceConnector 토픽 데이터, 토픽 설정, ACL, 파티션 구조 topics(기본 .*), topics.exclude, sync.topic.configs.enabled(기본 true), sync.topic.acls.enabled(기본 true), replication.factor(기본 2)
MirrorCheckpointConnector 컨슈머 그룹 오프셋 매핑(체크포인트) groups(기본 .*), emit.checkpoints.enabled(기본 true), sync.group.offsets.enabled(기본 false)
MirrorHeartbeatConnector 하트비트 레코드 (복제 지연 측정용) heartbeats.replication.enabled(기본 true)

토픽 이름이 바뀝니다

기본 정책인 DefaultReplicationPolicy는 복제된 토픽 이름에 소스 클러스터 별칭을 접두어로 붙입니다.

토픽 이름 변환
primary 클러스터의   orders
        ↓  primary->secondary 복제
secondary 클러스터의  primary.orders

오프셋은 그대로 일치하지 않습니다

같은 레코드가 primary에서 오프셋 1000이었다고 secondary에서도 1000인 것은 보장되지 않습니다. 복제가 시작된 시점, 리텐션으로 삭제된 구간, 트랜잭션 마커 때문에 어긋납니다.

그래서 MirrorCheckpointConnector가 {source}.checkpoints.internal 토픽에 "소스 오프셋 X ↔ 타깃 오프셋 Y" 매핑을 기록합니다. sync.group.offsets.enabled=true로 켜면 이 매핑을 이용해 타깃 클러스터의 __consumer_offsets에 변환된 오프셋을 직접 커밋해 줍니다.

사전 요구사항

전체 코드

디렉터리 구조
mm2-replication/
├── docker-compose.secondary.yml   # 목적지 클러스터 (단일 노드 KRaft)
├── docker-compose.mm2.yml         # MirrorMaker 2 전용 클러스터
└── mm2.properties                 # 복제 흐름 정의 (핵심)

목적지 클러스터

mm2-replication/docker-compose.secondary.yml
# 복제 목적지 클러스터. 단일 노드 combined 모드입니다.
# 내부 토픽 RF 를 1로 내려야 단일 노드에서 뜹니다(기본값은 3).
---
name: kafka-secondary

services:
  kafka-secondary:
    image: apache/kafka:4.3.1
    hostname: kafka-secondary
    container_name: kafka-secondary
    ports:
      # primary(29092/39092/49092)와 겹치지 않는 포트를 씁니다.
      - '19092:9092'
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: 'broker,controller'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka-secondary:9093'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT'
      KAFKA_LISTENERS: 'PLAINTEXT://:29092,CONTROLLER://:9093,PLAINTEXT_HOST://:9092'
      KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
      KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka-secondary:29092,PLAINTEXT_HOST://localhost:19092'
      # primary 와 다른 클러스터이므로 클러스터 ID 도 달라야 합니다.
      CLUSTER_ID: 'Zt9m2QeTSySn0-hK1p4vLg'
      # 단일 노드이므로 내부 토픽 RF 를 1로 내립니다.
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_SHARE_COORDINATOR_STATE_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_SHARE_COORDINATOR_STATE_TOPIC_MIN_ISR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      KAFKA_LOG_DIRS: '/var/lib/kafka/data'
      KAFKA_HEAP_OPTS: '-Xmx512m -Xms512m'
    volumes:
      - secondary-data:/var/lib/kafka/data
    networks:
      # primary 와 같은 네트워크에 두어 MirrorMaker 가 양쪽에 붙을 수 있게 합니다.
      - kafka-lab_default
    healthcheck:
      test: ['CMD-SHELL', '/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:29092 >/dev/null 2>&1']
      interval: 10s
      timeout: 10s
      retries: 12
      start_period: 30s

volumes:
  secondary-data:

networks:
  kafka-lab_default:
    external: true

복제 흐름 정의 — 이 예제의 핵심

mm2-replication/mm2.properties
# =============================================================================
# MirrorMaker 2 — primary → secondary (active/passive)
# 실행: connect-mirror-maker.sh mm2.properties
# =============================================================================

# --- 클러스터 별칭 정의 -------------------------------------------------------
# 이 별칭이 복제된 토픽 이름의 접두어가 되고, 내부 토픽 이름에도 들어갑니다.
# 한번 정하면 바꾸기 어렵습니다(토픽이 전부 새로 생깁니다).
clusters = primary, secondary

# 각 별칭의 bootstrap.servers. 컨테이너 네트워크 주소를 씁니다.
primary.bootstrap.servers = kafka-1:19092,kafka-2:19092,kafka-3:19092
secondary.bootstrap.servers = kafka-secondary:29092

# --- 복제 흐름 ---------------------------------------------------------------
# {source}->{target}.enabled 로 방향을 켭니다.
# active/active 를 원하면 secondary->primary.enabled = true 를 함께 켭니다.
primary->secondary.enabled = true
# 반대 방향은 끕니다(active/passive).
secondary->primary.enabled = false

# 복제 대상 토픽. 기본값은 .* (전부)입니다.
# 정규식이므로 필요한 것만 명시하는 편이 안전합니다.
primary->secondary.topics = orders, orders\..*, payments, users\..*

# 제외할 토픽. 기본값은 mm2.*\.internal,.*\.replica,__.*
# 기본값을 덮어쓸 때 __.* (내부 토픽) 제외를 빼먹으면
# __consumer_offsets 까지 복제하려 해 문제가 생깁니다.
primary->secondary.topics.exclude = mm2.*\.internal, .*\.replica, __.*, .*\.DLT, .*\.DLQ

# 복제 대상 컨슈머 그룹. 기본값 .*
# 기본 제외는 console-consumer-.*,connect-.*,__.*
primary->secondary.groups = orders-.*, payment-.*
primary->secondary.groups.exclude = console-consumer-.*, connect-.*, __.*

# --- Connect 워커 설정 (MirrorMaker 는 Connect 위에서 돕니다) ------------------
# 기본값 1 이지만 공식 문서는 최소 2 이상(가능하면 더 높게)을 권장합니다.
# 태스크가 여러 프로세스로 분산되어야 처리량이 나옵니다.
tasks.max = 4

# --- 타깃에 만들어질 토픽의 복제 계수 -----------------------------------------
# MirrorSourceConnector 의 replication.factor 기본값은 2 입니다.
# 이 예제의 secondary 는 단일 노드이므로 1로 내려야 토픽 생성이 성공합니다.
# 실제 DR 클러스터라면 3으로 두세요.
primary->secondary.replication.factor = 1

# MirrorMaker 내부 토픽들의 복제 계수. 단일 노드 secondary 기준으로 1입니다.
# 기본값은 모두 3 입니다.
primary->secondary.checkpoints.topic.replication.factor = 1
primary->secondary.heartbeats.topic.replication.factor = 1
primary->secondary.offset-syncs.topic.replication.factor = 1
# Connect 자체의 내부 토픽 RF (전용 모드에서도 필요합니다)
secondary.offset.storage.replication.factor = 1
secondary.config.storage.replication.factor = 1
secondary.status.storage.replication.factor = 1
primary.offset.storage.replication.factor = 3
primary.config.storage.replication.factor = 3
primary.status.storage.replication.factor = 3

# --- 토픽/ACL 설정 동기화 ------------------------------------------------------
# 기본값 true. 소스 토픽의 설정(retention 등)을 타깃에도 반영합니다.
primary->secondary.sync.topic.configs.enabled = true
primary->secondary.sync.topic.configs.interval.seconds = 600
# 기본값 true. ACL 도 함께 옮깁니다. 보안 모델이 다르면 끄세요.
primary->secondary.sync.topic.acls.enabled = false

# 새 토픽·파티션 자동 감지. 기본값 true, 주기 600초.
primary->secondary.refresh.topics.enabled = true
primary->secondary.refresh.topics.interval.seconds = 60

# --- 컨슈머 그룹 오프셋 --------------------------------------------------------
# 체크포인트 발행. 기본값 true, 주기 60초.
# 이것만으로는 타깃의 __consumer_offsets 가 갱신되지 않습니다.
primary->secondary.emit.checkpoints.enabled = true
primary->secondary.emit.checkpoints.interval.seconds = 30

# ★ 기본값이 false 입니다. 켜야 타깃의 컨슈머 그룹 오프셋이 실제로 갱신됩니다.
#   장애 시 컨슈머를 secondary 로 전환할 때 "처음부터 다시 읽기" 를 막는 설정입니다.
#   주의: 타깃에서 같은 그룹 ID 의 컨슈머가 활성 상태이면 충돌합니다.
#         active/passive 에서만 켜세요.
primary->secondary.sync.group.offsets.enabled = true
primary->secondary.sync.group.offsets.interval.seconds = 30

# 오프셋 싱크 레코드를 발행하는 최소 간격(레코드 수). 기본값 100.
# 작게 하면 오프셋 변환이 정밀해지지만 내부 토픽 쓰기가 늘어납니다.
primary->secondary.offset.lag.max = 100

# --- 클라이언트 세부 설정 -----------------------------------------------------
# {source}.consumer.* / {target}.producer.* / {cluster}.admin.* 형식으로 넘깁니다.
# 소스에 트랜잭션 프로듀서가 있다면 read_committed 로 읽어야
# abort 된 레코드를 복제하지 않습니다.
primary.consumer.isolation.level = read_committed
primary.consumer.auto.offset.reset = earliest
# 타깃 쓰기의 내구성. 복제본이 유실되면 DR 의 의미가 없습니다.
secondary.producer.acks = all
secondary.producer.enable.idempotence = true
secondary.producer.compression.type = lz4
secondary.producer.linger.ms = 50

# --- 복제 정책 ---------------------------------------------------------------
# 기본값은 DefaultReplicationPolicy 이며 "{source}." 접두어를 붙입니다.
# IdentityReplicationPolicy 는 접두어를 붙이지 않지만
# active/active 에서는 순환 복제가 발생하므로 단방향에서만 안전합니다.
replication.policy.class = org.apache.kafka.connect.mirror.DefaultReplicationPolicy
# 접두어와 토픽명 사이의 구분자. 기본값은 "." 입니다.
replication.policy.separator = .

# --- EOS (선택) ---------------------------------------------------------------
# MirrorMaker 전용 클러스터의 exactly-once 는 3.5.0부터 지원됩니다.
# 켜려면 두 설정이 함께 필요하고, 기존 클러스터는 preparing → enabled 2단계
# 롤링이 필요합니다.
# secondary.exactly.once.source.support = enabled
# dedicated.mode.enable.internal.rest = true
# listeners = http://0.0.0.0:8080

MirrorMaker 실행 컨테이너

mm2-replication/docker-compose.mm2.yml
---
name: kafka-mm2

services:
  mirrormaker:
    image: apache/kafka:4.3.1
    hostname: mirrormaker
    container_name: mirrormaker
    restart: unless-stopped
    # 전용 클러스터 모드로 실행합니다.
    # 이 스크립트가 MirrorSource/Checkpoint/Heartbeat 커넥터를 자동으로 만듭니다.
    command:
      - /opt/kafka/bin/connect-mirror-maker.sh
      - /etc/mm2/mm2.properties
    volumes:
      - ./mm2.properties:/etc/mm2/mm2.properties:ro
    environment:
      KAFKA_HEAP_OPTS: '-Xmx512m -Xms512m'
    networks:
      - kafka-lab_default

networks:
  kafka-lab_default:
    external: true

실행 방법

순서대로 실행
# 0. primary 클러스터(예제 1)가 떠 있어야 합니다.
cd kafka-lab && docker compose ps

# 1. secondary 클러스터 기동
cd ../mm2-replication
docker compose -f docker-compose.secondary.yml up -d

# 2. secondary 가 준비되었는지 확인
docker exec kafka-secondary /opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server localhost:29092 --list

# 3. primary 에 데이터를 넣습니다(복제 대상이 있어야 확인이 됩니다).
cd ../kafka-lab
for i in $(seq 1 100); do
  echo "ORD-$i:{\"orderId\":\"ORD-$i\",\"amount\":$((1000 * i))}"
done | ./kcli kafka-console-producer.sh --topic orders \
        --property parse.key=true --property key.separator=:

# 4. primary 에서 컨슈머 그룹을 하나 만들어 오프셋을 남깁니다.
#    (오프셋 동기화를 확인하기 위해)
./kcli kafka-console-consumer.sh --topic orders --from-beginning \
  --group orders-dr-test --timeout-ms 8000 --max-messages 50 >/dev/null 2>&1 || true
./kcli kafka-consumer-groups.sh --describe --group orders-dr-test

# 5. MirrorMaker 2 기동
cd ../mm2-replication
docker compose -f docker-compose.mm2.yml up -d
docker logs -f mirrormaker | head -50

검증 방법

1. 복제된 토픽 이름과 내부 토픽

secondary의 토픽 목록
docker exec kafka-secondary /opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server localhost:29092 --list
기대 출력
heartbeats
mm2-configs.primary.internal
mm2-offset-syncs.secondary.internal
mm2-offsets.primary.internal
mm2-status.primary.internal
primary.checkpoints.internal
primary.heartbeats
primary.orders
MirrorMaker 2가 만드는 토픽
토픽역할
primary.orders복제된 데이터 토픽. DefaultReplicationPolicyprimary. 접두어를 붙였습니다
primary.checkpoints.internal컨슈머 그룹 오프셋 매핑(소스 오프셋 ↔ 타깃 오프셋)
mm2-offset-syncs.secondary.internal오프셋 변환의 기준점이 되는 싱크 레코드
primary.heartbeats / heartbeats복제 지연 측정용 하트비트
mm2-configs / mm2-offsets / mm2-status .*.internalConnect 프레임워크의 내부 토픽 3개(설정·오프셋·상태)

2. 데이터가 실제로 복제되었는가

건수 비교
echo -n 'primary   orders        : '
docker exec kafka-1 /opt/kafka/bin/kafka-get-offsets.sh \
  --bootstrap-server kafka-1:19092 --topic orders --time -1 \
  | awk -F: '{s+=$3} END {print s}'

echo -n 'secondary primary.orders: '
docker exec kafka-secondary /opt/kafka/bin/kafka-get-offsets.sh \
  --bootstrap-server localhost:29092 --topic primary.orders --time -1 \
  | awk -F: '{s+=$3} END {print s}'

두 값이 정확히 같지 않을 수 있습니다. 복제가 진행 중이면 secondary가 작고, primary의 리텐션으로 삭제된 구간이 있으면 오프셋 번호 자체가 어긋납니다. 중요한 것은 레코드 내용이 같은지입니다.

내용 비교
docker exec kafka-secondary /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server localhost:29092 --topic primary.orders \
  --from-beginning --timeout-ms 10000 --property print.key=true 2>/dev/null | head -5

3. 컨슈머 그룹 오프셋이 동기화되었는가

이것이 이 예제의 핵심 검증입니다. 재해 복구 시 재처리량을 결정합니다.

secondary에서 같은 그룹 ID를 조회
docker exec kafka-secondary /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server localhost:29092 --describe --group orders-dr-test
기대 출력 — 토픽 이름이 primary.orders이고 오프셋이 변환되어 있습니다
GROUP           TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
orders-dr-test  primary.orders  0          8               17              9
orders-dr-test  primary.orders  1          9               16              7
orders-dr-test  primary.orders  2          7               17              10

CURRENT-OFFSET이 0이 아니면 동기화가 동작한 것입니다. sync.group.offsets.enabledfalse(기본값)로 두면 이 그룹이 secondary에 아예 존재하지 않습니다 — 장애 시 전환하면 auto.offset.reset에 따라 처음부터 전부 재처리하거나 전부 건너뜁니다.

체크포인트 토픽을 직접 들여다보기
docker exec kafka-secondary /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server localhost:29092 --topic primary.checkpoints.internal \
  --from-beginning --timeout-ms 8000 \
  --formatter org.apache.kafka.connect.mirror.formatters.CheckpointFormatter \
  2>/dev/null | head -5

4. 복제 지연(replication lag)을 측정

MirrorMaker의 소스 측 컨슈머 그룹으로 확인
# MirrorSourceConnector 는 primary 에서 Connect 컨슈머 그룹으로 읽습니다.
docker exec kafka-1 /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server kafka-1:19092 --list | grep -i mirror

# 그 그룹의 lag 이 곧 복제 지연입니다.
docker exec kafka-1 /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server kafka-1:19092 --describe \
  --group connect-MirrorSourceConnector 2>/dev/null || true

더 정확한 방법은 하트비트 토픽을 쓰는 것입니다. MirrorHeartbeatConnector가 일정 주기로 하트비트를 발행하므로, secondary의 primary.heartbeats에 도착한 마지막 하트비트의 타임스탬프와 현재 시각의 차이가 end-to-end 복제 지연입니다. 이 값을 예제 10의 Prometheus에 넣고 RPO(5분)를 넘으면 알림을 걸면 됩니다.

5. 페일오버 시뮬레이션

secondary에서 소비를 이어받습니다
# 1) MirrorMaker 를 멈춥니다(primary 가 죽은 상황을 흉내냅니다).
docker compose -f docker-compose.mm2.yml stop mirrormaker

# 2) 같은 그룹 ID 로 secondary 에서 소비를 시작합니다.
#    동기화된 오프셋부터 이어서 읽습니다.
docker exec -i kafka-secondary /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server localhost:29092 \
  --topic primary.orders --group orders-dr-test \
  --timeout-ms 10000 --property print.key=true 2>/dev/null | head -10

프로덕션 고려사항

로컬 예제와 프로덕션의 차이
항목이 예제프로덕션
MirrorMaker 프로세스 수 1개 2개 이상. 단일 프로세스는 그 자체가 단일 장애점이고, tasks.max를 올려도 한 프로세스에서만 돕니다
tasks.max 4 복제할 토픽-파티션 수와 하드웨어에 맞춰 조정합니다. 공식 문서는 최소 2 이상, 가능하면 더 높게 권장합니다
타깃 토픽 RF 1 (단일 노드) 3. replication.factor 기본값이 2이므로 명시하지 않으면 2가 됩니다 — DR 클러스터로는 부족합니다
네트워크 같은 Docker 네트워크 리전 간이면 대역폭과 지연이 병목입니다. compression.type을 반드시 켜고, 쿼터로 복제 트래픽을 제한합니다
보안 PLAINTEXT 양쪽 클러스터에 {cluster}.security.protocol·sasl.*·ssl.*을 각각 설정합니다. 리전 간 평문 전송은 허용되지 않습니다
EOS 3.5.0+에서 {target}.exactly.once.source.support=enabled + dedicated.mode.enable.internal.rest=true. 기존 클러스터는 preparingenabled 2단계 롤링이 필요합니다
모니터링 수동 확인 하트비트 기반 end-to-end 지연을 RPO와 비교해 알림. MirrorMaker의 Connect 태스크 상태도 함께 봐야 합니다(예제 10)
오프셋 동기화 active/passive에서만 켜세요. 타깃에 같은 그룹 ID의 활성 컨슈머가 있으면 충돌합니다
DR 훈련 없음 정기적으로 페일오버를 실제로 수행하세요. 접두어 때문에 애플리케이션 설정 변경이 필요하다는 사실을 장애 당일에 알면 RTO를 못 지킵니다

자주 하는 실수

이어서 볼 곳

공식 문서 출처