실무 예제 · 11
MirrorMaker 2 클러스터 간 복제
MirrorMaker 2는 Kafka Connect 프레임워크 위에서 도는 클러스터 간(inter-cluster) 복제 도구입니다. 클러스터 안의 복제(replication factor)와는 완전히 다른 것입니다. 토픽·설정·ACL·컨슈머 그룹 오프셋까지 옮길 수 있는 대신, 토픽 이름이 바뀌고 오프셋이 그대로 일치하지 않는다는 성질을 이해하지 않으면 재해 복구 시점에 당황합니다.
학습 목표
- 복제 흐름(
{source}->{target}) 문법으로 active/passive와 active/active를 구성할 수 있습니다. DefaultReplicationPolicy가 붙이는 소스 별칭 접두어가 순환 복제를 어떻게 막는지 설명할 수 있습니다.- 컨슈머 그룹 오프셋 동기화(
sync.group.offsets.enabled)의 기본값과 한계를 알 수 있습니다. - MirrorMaker 2가 만드는 내부 토픽들과 각각의 역할을 구분할 수 있습니다.
시나리오
서울 리전(primary)에 운영 Kafka 클러스터가 있습니다.
규제 요건으로 다른 리전(secondary)에 재해 복구용 사본을
유지해야 합니다. RPO는 5분, RTO는 30분입니다.
추가 요구사항이 있습니다. 분석 팀이 운영 클러스터에 직접 붙는 것을 막고 싶어서,
분석용 컨슈머는 secondary에서만 읽게 하려 합니다.
그리고 장애 시 운영 컨슈머를 secondary로 전환할 때
처음부터 다시 읽지 않아야 합니다.
MirrorMaker 2로 active/passive 구성을 만들고, 오프셋 동기화까지 켭니다.
아키텍처
MirrorMaker 2는 세 종류의 커넥터로 구성됩니다.
전용 클러스터 모드(connect-mirror-maker.sh)로 띄우면 이들이 자동으로 생성됩니다.
| 커넥터 | 옮기는 것 | 주요 설정 |
|---|---|---|
| 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에
변환된 오프셋을 직접 커밋해 줍니다.
사전 요구사항
- 예제 1의 클러스터를
primary로 씁니다. - 두 번째 클러스터(
secondary)가 필요합니다. 아래 compose가 단일 노드로 띄웁니다 — 복제 대상만 확인하는 목적이므로 3노드가 아니어도 됩니다. - Apache Kafka 4.3.1 (MirrorMaker 2는 배포판에 포함됩니다)
전체 코드
mm2-replication/
├── docker-compose.secondary.yml # 목적지 클러스터 (단일 노드 KRaft)
├── docker-compose.mm2.yml # MirrorMaker 2 전용 클러스터
└── mm2.properties # 복제 흐름 정의 (핵심)
목적지 클러스터
# 복제 목적지 클러스터. 단일 노드 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
복제 흐름 정의 — 이 예제의 핵심
# =============================================================================
# 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 실행 컨테이너
---
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. 복제된 토픽 이름과 내부 토픽
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
| 토픽 | 역할 |
|---|---|
primary.orders | 복제된 데이터 토픽. DefaultReplicationPolicy가 primary. 접두어를 붙였습니다 |
primary.checkpoints.internal | 컨슈머 그룹 오프셋 매핑(소스 오프셋 ↔ 타깃 오프셋) |
mm2-offset-syncs.secondary.internal | 오프셋 변환의 기준점이 되는 싱크 레코드 |
primary.heartbeats / heartbeats | 복제 지연 측정용 하트비트 |
mm2-configs / mm2-offsets / mm2-status .*.internal | Connect 프레임워크의 내부 토픽 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. 컨슈머 그룹 오프셋이 동기화되었는가
이것이 이 예제의 핵심 검증입니다. 재해 복구 시 재처리량을 결정합니다.
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.enabled를 false(기본값)로 두면
이 그룹이 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)을 측정
# 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. 페일오버 시뮬레이션
# 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. 기존 클러스터는 preparing → enabled 2단계 롤링이 필요합니다 |
| 모니터링 | 수동 확인 | 하트비트 기반 end-to-end 지연을 RPO와 비교해 알림. MirrorMaker의 Connect 태스크 상태도 함께 봐야 합니다(예제 10) |
| 오프셋 동기화 | 켬 | active/passive에서만 켜세요. 타깃에 같은 그룹 ID의 활성 컨슈머가 있으면 충돌합니다 |
| DR 훈련 | 없음 | 정기적으로 페일오버를 실제로 수행하세요. 접두어 때문에 애플리케이션 설정 변경이 필요하다는 사실을 장애 당일에 알면 RTO를 못 지킵니다 |
자주 하는 실수
관련 케이스 스터디
이어서 볼 곳
공식 문서 출처
- Geo-Replication (Cross-Cluster Data Mirroring) — MirrorMaker 2가 Connect 기반이라는 점, 복제 흐름 문법
{source}->{target}, active/active·active/passive·aggregation·fan-out 패턴,clusters와{cluster}.bootstrap.servers,{cluster}.{config_name}/{source}.consumer.*/{target}.producer.*/{cluster}.admin.*형식,tasks.max를 최소 2 이상으로 권장, 클러스터 간 복제와 클러스터 내 복제가 다르다는 명시, EOS는 3.5.0부터이며exactly.once.source.support와dedicated.mode.enable.internal.rest가 필요하다는 서술 - MirrorMaker Configs —
topics=.*,topics.exclude=mm2.*\.internal,.*\.replica,__.*,replication.factor=2,sync.topic.configs.enabled=true,sync.topic.acls.enabled=true,refresh.topics.interval.seconds=600,offset.lag.max=100,offset-syncs.topic.replication.factor=3,groups=.*,groups.exclude=console-consumer-.*,connect-.*,__.*,emit.checkpoints.enabled=true,emit.checkpoints.interval.seconds=60,sync.group.offsets.enabled=false,sync.group.offsets.interval.seconds=60,checkpoints.topic.replication.factor=3 - MirrorClientConfig (Apache Kafka 4.3) —
replication.policy.class의 기본값이DefaultReplicationPolicy라는 것 - Connect Worker Configs —
offset.storage.replication.factor등 내부 토픽 설정(MirrorMaker 전용 모드에서도 적용됩니다)