Post

Kafka 입문3

Kafka 입문3

카프카 리플리케이션

데이터 파이프라인의 메인 허브 역할을 하는 카프카 클러스터가 하드웨어의 문제나 점검 등으로 인해 정상적으로 동작하지 못한다거나, 카프카와 연결된 전체 데이터 파이프라인에 영향을 미친다면 매우 심각한 문제로 발생할 수 있기 때문에 안정성을 확보하기 위해 리플리케이션이 동작하게 된다.

동작 원리

  • 메시지들을 여러 개로 복제해서 카프카 클러스터 내 브로커들에 분산한다.
  • N개의 리플리케이션이 있는 경우 N-1까지의 브로커 장애가 발생해도 메시지 손실 없이 안정적으로 메시지를 주고받을 수 있다.

리더와 팔로워

리더

  • 리플리케이션 중 하나가 선정됨, 모든 읽기와 쓰기는 리더를 통해서만 가능
  • 프로듀서는 리더에게만 메시지를 전송하고, 컨슈머도 오직 리더로부터 메시지를 가져온다.

카프카 리더와 팔로워 이미지

팔로워

  • 리더에 문제가 발생하거나 이슈가 있을 경우를 대비해 새로운 리더가 될 준비를 한다.
  • 파티션의 리더가 새로운 메시지를 받았는지 확인하고, 새로운 메시지가 있다면 리더로부터 복제한다.

복제 유지와 커밋

ISR(InSyncReplica)

  • 리더와 팔로워를 묶는 논리적인 그룹
  • 해당 그룹 안에 속한 팔로워들만이 새로운 리더의 자격을 가질 수 있음
  • ISR 그룹에 속하지 못한 팔로워는 새로운 리더의 자격을 가질 수 없음

ISR 내의 팔로워들은 리더와의 데이터 일치를 유지하기 이해 지속적으로 리더의 데이터를 따라가게 되고, 리더는 ISR내 모든 팔로워가 메시지를 받을 때까지 기다린다.

하지만 팔로워가 네트워크 오류, 브로커 장애 등 여러 이유로 리더로부터 리플리케이션하지 못하는 경우도 발생할 수 있음

이때 팔로워는 리더와의 데이터 불일치 상태에 놓이게 되고, 이 팔로워에게 새로운 리더를 넘겨준다면 데이터의 정합성이나 메시지 손실 등의 문제가 발생함

그럼 리더와 팔로워 중 리플리케이션 동작을 잘하고 있는지 여부 등은 누가 판단하고, 어떤 기준으로 판단할까?

리더는 팔로워가 특정 주기의 시간만큼 복제 요청을 하지 않는다면, 해당 팔로워가 리플리케이션 동작에 문제가 발생했다고 판단해 ISR 그룹에서 추방함

ISR 내에서 모든 팔로워의 복제가 완료되면, 리더는 내부적으로 커밋되었다는 표시를 한다.

하이워터마크

  • 커밋 오프셋 위치
  • 커밋되었다는 것은 리프리케이션 팩터 수의 모든 리플리케이션이 전부 메시지를 저장했음을 의미
  • 컨슈머는 커밋된 메시지만 읽어갈 수 있다. => 메시지의 일관성을 유지하기 위해서

카프카 커밋 메시지 이미지

만약 커밋된 메시지만 읽지 않는 상황이라면??

카프카 커밋 메시지 않읽을 경우 이미지

컨슈머 A가 리더로부터 test message 1, test meesage 2를 읽었는데, 리더 브로커에 문제가 발생해 팔로워 중 하나가 새로운 리더가 된 경우

새로운 리더는 아직 이전 리더로부터 test message 2를 가지고 있지 않은 상황에서 컨슈머 B는 test message 1만 컨슘 할 수 있다.

따라서, 동일한 토픽의 파티션에서 컨슘했음에도 메시지가 일치하지 않은 현상이 발생할 수 있음.

그럼 커밋된 위치를 어떻게 알 수 있을까??

모든 브로커는 재시작될 때, 커밋된 메시지를 유지하기 위해 로컬 디스크의 replication-offset-chcekpoint라는 파일에 마지막 커밋 오프셋 위치를 저장함

리더와 팔로워의 단게별 리플리케이션 동작

단계별 동작 예시

토픽 1개, 파티션 1개, 리플리케이션 팩터 3 (리더 1 + 팔로워 2)


[STEP 1] 프로듀서가 리더에게 메시지 1 전송

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
  Producer
     |
     | 메시지 1 전송
     ▼
 ┌────────────────────────┐
 │  Leader                │
 │  [ msg1 ]              │  ← offset 0
 └────────────────────────┘
 ┌────────────────────────┐
 │  팔로워 1                │
 │  [      ]              │  ← 아직 없음
 └────────────────────────┘
 ┌────────────────────────┐
 │  팔로워 2               │
 │  [      ]              │  ← 아직 없음
 └────────────────────────┘

[STEP 2] 팔로워들이 리더에게 메시지 1 복제 요청 (PULL)

1
2
3
4
5
6
7
8
9
10
11
12
13
 ┌────────────────────────┐
 │  Leader                │
 │  [ msg1 ]              │
 └───────────┬────────────┘
             │
    ┌────────┴────────┐
    │ fetch(offset 0) │ ← 팔로워들이 PULL 방식으로 요청
    │                 │
    ▼                 ▼
 ┌──────────┐    ┌──────────┐
 │ 팔로워 1  │    │ 팔로워 2  │
 │ [      ] │    │ [      ] │
 └──────────┘    └──────────┘

[STEP 3] 리더가 팔로워들에게 메시지 1 응답

1
2
3
4
5
6
7
8
9
10
11
12
13
 ┌────────────────────────┐
 │  Leader                │
 │  [ msg1 ]              │
 └───────────┬────────────┘
             │
    ┌────────┴────────┐
    │     msg1 전달   │
    │                 │
    ▼                 ▼
 ┌──────────┐    ┌──────────┐
 │ 팔로워 1  │    │ 팔로워 2  │
 │ [ msg1 ] │    │ [ msg1 ] │
 └──────────┘    └──────────┘

[STEP 4] 프로듀서가 리더에게 메시지 2 전송

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
  Producer
     |
     | 메시지 2 전송
     ▼
 ┌────────────────────────────────┐
 │  Leader                        │
 │  [ msg1 ][ msg2 ]              │  ← offset 1 추가
 └────────────────────────────────┘
 ┌────────────────────────────────┐
 │  팔로워 1                       │
 │  [ msg1 ]                      │  ← msg2 아직 없음
 └────────────────────────────────┘
 ┌────────────────────────────────┐
 │  팔로워 2                       │
 │  [ msg1 ]                      │  ← msg2 아직 없음
 └────────────────────────────────┘

[STEP 5] 팔로워들이 메시지 2 복제 요청 → 메시지 1 커밋 완료

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
 ┌──────────────────────────────────────┐
 │  Leader                              │
 │  [ msg1 ][ msg2 ]                    │
 └────────────────┬─────────────────────┘
                  │
     ┌────────────┴────────────┐
     │    fetch(offset 1)      │ ← offset 1 요청 = msg1은 이미 저장했다는 의미
     │                         │
     ▼                         ▼
 ┌──────────┐             ┌──────────┐
 │ 팔로워 1  │             │ 팔로워 2  │
 │ [ msg1 ] │             │ [ msg1 ] │
 └──────────┘             └──────────┘

 ISR 내 모든 팔로워가 msg1 복제 완료
       ↓
 리더가 msg1을 커밋 (하이워터마크 이동)

 ┌──────────────────────────────────────┐
 │  Leader                              │
 │  [ msg1 ][ msg2 ]                    │
 │      ↑                               │
 │  하이워터마크 (커밋 완료)               │
 └──────────────────────────────────────┘

msg2는 아직 팔로워들이 복제하지 못했으므로 커밋 전 상태이며, 컨슈머는 msg1까지만 읽을 수 있다.

  • RabbitMq와 같이 다른 메시징 시스템들의 경우 리더와 팔로워간에 메시지를 잘 받았는지 확인하기 위해 ACK 통신을 함. -> 하지만 카프카에서는 ACK 통신을 제거함으로써 리플리케이션 동작의 성능을 높였음

  • 리더가 PUSH 하는 방식이 아니라 팔로워들이 PULL 하는 방식으로 동작을 함. -> PULL 방식은 리플리케이션 동작에서 리더의 부하를 줄여주기 위함

리더에포크와 복구

리더에포크

  • 카프카의 파티션들이 복구 동작을 할 때 메시지의 일관성을 유지하기 위한 용도로 이용
  • 리더에포크는 컨트롤러에 의해 관리되는 32비트의 숫자로 표현됨
  • 리더에포크 정보는 리플리케이션 프로토콜에 의해 전파되고, 새로운 리더가 변경된 후 변경된 리더에 대한 정보는 팔로워에게 전달된다.

복구 과정 예시

파티션 1개, 리플리케이션 팩터 2 (B1: 리더, B2: 팔로워), min.insync.replicas=1

초기 상태 (장애 발생 전)

1
2
3
4
5
6
7
8
9
10
11
╔══════════════════════════════════════╗
║  B1 (Leader, Epoch 1)               ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=1      ║  ← msg2 미커밋
╚══════════════════════════════════════╝

╔══════════════════════════════════════╗
║  B2 (Follower)                      ║
║  offset:  0                         ║
║         [msg1]             HW=1     ║  ← msg2 복제 못 받은 상태
╚══════════════════════════════════════╝

① 리더에포크가 없다면? (문제 발생)

[STEP 1] B1 장애 발생

1
2
3
4
5
╔══════════════════════════════════════╗
║  B1 (Leader)  ❌ 장애 발생           ║
║  offset:  0        1                ║
║         [msg1]   [msg2]             ║  ← msg2 복제 전에 장애
╚══════════════════════════════════════╝

[STEP 2] B2가 뉴리더로 선출 후 B1 복구

1
2
3
4
5
6
7
8
9
10
11
╔══════════════════════════════════════╗
║  B2 (New Leader)                    ║
║  offset:  0                         ║
║         [msg1]             HW=1     ║
╚══════════════════════════════════════╝

╔══════════════════════════════════════╗
║  B1 (복구됨, 팔로워)                 ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=1      ║  ← 에포크 정보 없어 msg2를 잘라야 할지 모름
╚══════════════════════════════════════╝

[결과] 두 복제본의 데이터 불일치

1
2
3
4
5
6
7
8
9
10
11
12
13
14
╔══════════════════════════════════════╗
║  B1 (팔로워)                        ║
║  offset:  0        1                ║
║         [msg1]   [msg2]             ║
╚══════════════════════════════════════╝

╔══════════════════════════════════════╗
║  B2 (리더)                          ║
║  offset:  0                         ║
║         [msg1]                      ║
╚══════════════════════════════════════╝

  ⚠️ B1에는 msg2가 있고 B2에는 없음 → 복제본 간 데이터 불일치 발생!
      - 결국 HW(하이워터마크) 1까지 기준으로 msg1까지 잘라서 메시지 유실 발생

② 리더에포크가 있다면? (정상 복구)

[STEP 1] B2가 리더에포크를 포함한 fetch 요청 전송

1
2
3
4
5
6
7
8
╔══════════════════════════════════════╗
║  B1 (Leader, Epoch 1)               ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=1      ║
╚══════════════════════════════════════╝

  B2 → B1: "Epoch 1, offset 1부터 fetch 요청"
           (리더에포크 정보를 포함해 요청)

[STEP 2] B1이 msg2를 B2에게 전달 → 복제 완료

1
2
3
4
5
6
7
8
9
10
11
12
13
╔══════════════════════════════════════╗
║  B1 (Leader, Epoch 1)               ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=2      ║  ← ISR 전체 복제 완료, HW 상승
╚══════════════════════════════════════╝

  B1 → B2: msg2 전달

╔══════════════════════════════════════╗
║  B2 (Follower)                      ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=2      ║  ← msg2 복제 성공
╚══════════════════════════════════════╝

[STEP 3] B1 장애 → B2가 뉴리더로 선출

1
2
3
4
5
6
7
8
9
╔══════════════════════════════════════╗
║  B1 (Leader, Epoch 1)  ❌ 장애       ║
╚══════════════════════════════════════╝

╔══════════════════════════════════════╗
║  B2 (New Leader, Epoch 2)           ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=2      ║  ← msg2를 이미 보유
╚══════════════════════════════════════╝

[STEP 4] B1 복구 → 정상 동기화

1
2
3
4
5
6
7
8
9
10
11
12
13
╔══════════════════════════════════════╗
║  B1 (팔로워)                        ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=2      ║  ✓
╚══════════════════════════════════════╝

╔══════════════════════════════════════╗
║  B2 (리더)                          ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=2      ║  ✓
╚══════════════════════════════════════╝

  ✅ 메시지 손실 없이 로그 완전 일치!

리더에포크 덕분에 B2가 장애 전에 msg2를 정상적으로 복제하고 HW를 올릴 수 있었기 때문에, B1 장애 이후에도 데이터 손실 없이 복구가 가능하다.

여기까지 보면 그냥 복제 끝나고 죽었나 아니냐의 차이로 알 수 있다.

밑에 리더에포크의 진짜 역할인 상황이 존재한다.


③ 리더에포크 없이 불일치 후 msg4 수신 (문제 지속)

B1이 복구된 후 B2(뉴리더)가 msg3을 받으면서 두 브로커의 HW는 같아 보이지만, 실제 로그 내용이 다른 심각한 불일치 상태가 된다.

[STEP 1] B2(뉴리더)가 msg3 수신 및 커밋

1
2
3
4
5
6
7
8
9
  Producer
     │
     │  msg3 전송
     ▼
╔══════════════════════════════════════╗
║  B2 (New Leader)                    ║
║  offset:  0        1                ║
║         [msg1]   [msg3]   HW=2      ║  ← offset 1 = msg3, 커밋 완료
╚══════════════════════════════════════╝

B2가 msg3을 커밋하면서 HW=2로 올라가고, B1도 HW 업데이트를 받아 HW=2가 된다. 하지만 B1의 offset 1에는 여전히 msg2가 남아있다.

1
2
3
4
5
╔══════════════════════════════════════╗
║  B1 (팔로워)                        ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=2      ║  ← offset 1 = msg2 (다른 데이터!)
╚══════════════════════════════════════╝

[STEP 2] 겉으론 정상처럼 보이는 불일치 상태

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
╔══════════════════════════════════════╗
║  B1 (팔로워)                        ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=2      ║
╚══════════════════════════════════════╝

╔══════════════════════════════════════╗
║  B2 (리더)                          ║
║  offset:  0        1                ║
║         [msg1]   [msg3]   HW=2      ║
╚══════════════════════════════════════╝

  ⚠️ HW는 둘 다 2로 일치하지만 offset 1의 내용이 다름!
     B1: offset 1 = msg2
     B2: offset 1 = msg3

[STEP 3] msg4 수신 → 불일치 지속

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
  Producer
     │
     │  msg4 전송
     ▼
╔══════════════════════════════════════╗
║  B2 (리더)                          ║
║  offset:  0        1        2       ║
║         [msg1]   [msg3]   [msg4]    ║  HW=2
╚══════════════════════════════════════╝

  B2 → B1: msg4 복제

╔══════════════════════════════════════╗
║  B1 (팔로워)                        ║
║  offset:  0        1        2       ║
║         [msg1]   [msg2]   [msg4]    ║  HW=2
╚══════════════════════════════════════╝

  ❌ offset 1의 불일치가 해소되지 않은 채 msg4만 쌓임
     리더에포크 없이는 이 상태를 감지하거나 바로잡을 방법이 없다.

④ 리더에포크로 위 문제를 해결한다면?

[STEP 1] B1 복구 시 에포크 기반 절단점 확인

B1이 팔로워로 합류할 때 뉴리더(B2)에게 현재 에포크의 시작 오프셋을 먼저 질의한다.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
╔══════════════════════════════════════╗
║  B1 (복구됨)                        ║
║  offset:  0        1                ║
║         [msg1]   [msg2]   HW=1      ║
╚══════════════════════════════════════╝

  B1 → B2: "Epoch 2는 몇 번 오프셋부터 시작했나요?"
  B2 → B1: "Epoch 2는 offset 1부터 시작했습니다"

  B1은 offset 1부터 전부 잘라냄 → msg2 제거

╔══════════════════════════════════════╗
║  B1 (절단 후)                       ║
║  offset:  0                         ║
║         [msg1]             HW=1     ║
╚══════════════════════════════════════╝

[STEP 2] B2가 msg3 수신 → B1이 정상 복제

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
  Producer
     │
     │  msg3 전송
     ▼
╔══════════════════════════════════════╗
║  B2 (New Leader, Epoch 2)           ║
║  offset:  0        1                ║
║         [msg1]   [msg3]   HW=2      ║
╚══════════════════════════════════════╝

  B2 → B1: msg3 복제

╔══════════════════════════════════════╗
║  B1 (팔로워)                        ║
║  offset:  0        1                ║
║         [msg1]   [msg3]   HW=2      ║  ← offset 1 = msg3 (일치!)
╚══════════════════════════════════════╝

[STEP 3] msg4 수신 → 정상 복제

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
  Producer
     │
     │  msg4 전송
     ▼
╔══════════════════════════════════════╗
║  B2 (리더)                          ║
║  offset:  0        1        2       ║
║         [msg1]   [msg3]   [msg4]    ║  HW=3
╚══════════════════════════════════════╝

  B2 → B1: msg4 복제

╔══════════════════════════════════════╗
║  B1 (팔로워)                        ║
║  offset:  0        1        2       ║
║         [msg1]   [msg3]   [msg4]    ║  HW=3
╚══════════════════════════════════════╝

  ✅ 두 브로커의 로그가 완전히 일치, 데이터 정합성 유지!

리더에포크를 통해 B1이 복구 시점에 정확한 절단 위치를 파악함으로써, 잘못된 msg2를 제거하고 올바른 msg3부터 동기화할 수 있었다.

즉 리더에포크가 하는 진짜 역할은 “장애가 언제 났는지”와 무관하게, 각 메시지가 “몇 번째 에포크에서 쓰인 것인지”를 태깅해서, 팔로워가 재합류할 때 “새 에포크가 시작된 정확한 offset”을 물어보고 그 지점부터 정확히 잘라낼 수 있게 해주는 것이다. HW라는 부정확한 지표 대신 에포크라는 정확한 지표로 절단점을 잡는 것이다.

그럼 msg2는 아예 유실 되는 것인가?

acks=all + min.insync.replicas≥2 인 경우

  • 리더가 메시지를 받아도 ISR 전체(팔로워 포함)에 복제가 끝나야 producer에게 성공 응답을 줌
  • 글의 예시(msg2)처럼 복제되기 전에 리더가 죽으면, producer는애초에 성공 응답을 못 받은 상태 → producer 입장에선 “실패”로 인식하고 재전송(retry)함
  • 그래서 이 경우는 진짜 “유실”이 아니라, 애초에 커밋 안 된 걸로 취급되고 재전송으로 복구되는 흐름. 소비자 입장에서도 커밋 전 메시지는 애초에 안 보였으니 일관성 문제 없음

컨트롤러

  • 카프카 클러스터 중 하나의 부로커가 컨트롤러 역할을 함
  • 파티션의 ISR 리스트 중에서 리더를 선출함
  • 컨트롤러는 브로커가 실패하는 것을 예의주시 하고 있으며, 만약 브로커의 실패가감지되면 즉시 ISR 리스트 중 하나를 새로운 파티션 리더로 선출함

리더 선출 과정 예시

파티션 1개, 리플리케이션 팩터 2, 브로커: B1(리더) · B3(팔로워), 컨트롤러: B3

초기 상태

1
2
3
4
5
6
7
8
9
10
11
12
  ┌─────────────────────────────────────────────────────┐
  │  Kafka Cluster                                      │
  │                                                     │
  │  ╔══════════════════╗     ╔══════════════════════╗  │
  │  ║  Broker 1        ║     ║  Broker 3            ║  │
  │  ║  partition 0     ║     ║  partition 0         ║  │
  │  ║  [Leader]        ║     ║  [Follower]          ║  │
  │  ║                  ║     ║  [Controller] ★      ║  │
  │  ╚══════════════════╝     ╚══════════════════════╝  │
  └─────────────────────────────────────────────────────┘

  ISR: [B1, B3]   /   partition 0 Leader: B1

[STEP 1] B1(리더) 강제 종료

1
2
3
4
5
6
7
  ╔══════════════════╗     ╔══════════════════════╗
  ║  Broker 1  ❌    ║     ║  Broker 3            ║
  ║  강제 종료        ║     ║  [Follower]          ║
  ╚══════════════════╝     ║  [Controller] ★      ║
                           ╚══════════════════════╝

  B1이 ZooKeeper에 등록한 세션이 만료됨

[STEP 2] 컨트롤러(B3)가 B1 장애 감지

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
  ZooKeeper
  ┌─────────────────────────────┐
  │  /brokers/ids/1  ← 세션 만료 │  → 컨트롤러(B3)에 알림
  └─────────────────────────────┘
              │
              ▼
  ╔══════════════════════╗
  ║  Broker 3            ║
  ║  [Controller] ★      ║
  ║                      ║
  ║  "B1 장애 감지!"     ║
  ║  B1이 리더인 파티션  ║
  ║  확인 중...          ║
  ╚══════════════════════╝

  → partition 0의 리더가 B1임을 확인
  → ISR 목록 확인: [B1, B3]
  → B1 제외 후 남은 ISR: [B3]

[STEP 3] 컨트롤러가 B3를 새 리더로 선출 및 메타데이터 갱신

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
  ╔══════════════════════╗
  ║  Broker 3            ║
  ║  [Controller] ★      ║
  ║                      ║
  ║  partition 0         ║
  ║  새 리더 = B3 선출   ║
  ╚══════════════════════╝
              │
              │  LeaderAndIsr 요청 전송
              ▼
  ZooKeeper
  ┌───────────────────────────────────────┐
  │  /brokers/topics/.../partition 0      │
  │  leader: B3  (B1 → B3 로 갱신)        │
  │  ISR: [B3]                            │
  └───────────────────────────────────────┘

[STEP 4] B3가 partition 0의 새 리더로 승격

1
2
3
4
5
6
7
8
9
10
11
12
13
14
  ┌─────────────────────────────────────────────────────┐
  │  Kafka Cluster                                      │
  │                                                     │
  │  ╔══════════════════╗     ╔══════════════════════╗  │
  │  ║  Broker 1  ❌    ║     ║  Broker 3            ║  │
  │  ║  (오프라인)       ║     ║  partition 0         ║  │
  │  ╚══════════════════╝     ║  [New Leader] ★      ║  │
  │                           ║  [Controller] ★      ║  │
  │                           ╚══════════════════════╝  │
  └─────────────────────────────────────────────────────┘

  ISR: [B3]   /   partition 0 Leader: B3

  ✅ 프로듀서·컨슈머는 새 메타데이터를 받아 B3로 요청을 전환

제어된 종료 과정 (SIGTERM)

강제 종료(kill -9, SIGKILL)와 달리 kafka-server-stop.sh는 SIGTERM을 전달해 브로커가 스스로 정리할 시간을 줍니다.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
[STEP 1] kafka-server-stop.sh가 브로커에 SIGTERM 전달

  kafka-server-stop.sh ──→ Broker 1 프로세스에 SIGTERM 전송


[STEP 2] JVM이 SIGTERM을 받아 등록된 셧다운 훅 실행

  Broker 1 (JVM)
  ┌─────────────────────────────────────┐
  │  SIGTERM 수신                       │
  │       ↓                             │
  │  Shutdown Hook 실행                 │
  │  (카프카가 JVM 시작 시 등록해둔 훅)  │
  └─────────────────────────────────────┘


[STEP 3] 셧다운 훅이 컨트롤러에게 제어된 종료 요청

  Broker 1 ──→ Controller(B3): ControlledShutdownRequest 전송


[STEP 4] 컨트롤러가 B1의 파티션 리더를 B3로 이전 후 ZooKeeper 갱신

  Controller(B3) ──→ B3: "partition 0 리더 받아주세요"
  B3 리더 승격 완료

  ZooKeeper
  ┌───────────────────────────────────────┐
  │  /brokers/topics/.../partition 0      │
  │  leader: B3  (B1 → B3 로 갱신)        │
  │  ISR: [B3]                            │
  ├───────────────────────────────────────┤
  │  /brokers/ids/1  ← B1 세션 유지 중    │  ← 아직 살아있음
  └───────────────────────────────────────┘


[STEP 5] 컨트롤러가 종료 승인 응답

  Controller(B3) ──→ Broker 1: ControlledShutdownResponse (성공)


[STEP 6] Broker 1 정상 종료 및 ZooKeeper 세션 해제

  Broker 1: 네트워크 연결 닫기 → 로그 flush → 프로세스 종료

  ZooKeeper
  ┌───────────────────────────────────────┐
  │  /brokers/ids/1  ← 삭제됨             │  ← 종료 후 세션 만료
  └───────────────────────────────────────┘

제어된 종료 vs 갑작스러운 종료 (다운타임 비교)

갑작스러운 종료 (SIGKILL)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
╔══════════════════════════════════════════════════════════════╗
║  [1] B1 즉시 종료  →  프로듀서·컨슈머 요청 실패 시작         ║
╠══════════════════════════════════════════════════════════════╣
║  [2] ZooKeeper가 B1 세션 만료를 감지할 때까지 대기           ║
║      (기본값: 수 초 ~ 수십 초)                               ║
║                                                              ║
║      ⏳ 이 구간 동안 프로듀서·컨슈머는 계속 실패             ║
╠══════════════════════════════════════════════════════════════╣
║  [3] 컨트롤러가 세션 만료를 감지                             ║
║      → ISR에서 새 리더(B3) 선출                              ║
║      → ZooKeeper 메타데이터 갱신                             ║
╠══════════════════════════════════════════════════════════════╣
║  [4] 프로듀서·컨슈머가 새 리더(B3)로 전환                    ║
║                                                              ║
║  총 다운타임 = 세션 타임아웃 대기 + 리더 선출 시간           ║
╚══════════════════════════════════════════════════════════════╝

제어된 종료 (SIGTERM)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
╔══════════════════════════════════════════════════════════════╗
║  [1] SIGTERM 수신 → 셧다운 훅 실행                          ║
╠══════════════════════════════════════════════════════════════╣
║  [2] B1이 살아있는 상태에서 컨트롤러에게 리더 이전 요청      ║
║      → B3가 새 리더로 승격                                   ║
║      → ZooKeeper 메타데이터 갱신                             ║
║                                                              ║
║      ✅ B1이 아직 살아있으므로 요청 실패 없음                ║
╠══════════════════════════════════════════════════════════════╣
║  [3] 리더 이전 완료 후 B1 종료                               ║
║      → 로그 flush → 프로세스 종료                            ║
║      → ZooKeeper 세션 즉시 해제                              ║
╠══════════════════════════════════════════════════════════════╣
║  [4] 프로듀서·컨슈머는 이미 B3로 연결된 상태                 ║
║                                                              ║
║  총 다운타임 ≈ 거의 없음 (리더 이전 시간만큼)                ║
╚══════════════════════════════════════════════════════════════╝
구분갑작스러운 종료제어된 종료
ZooKeeper 세션타임아웃까지 대기종료 시 즉시 해제
리더 이전 시점종료 후 (감지 후 선출)종료 전 (살아있는 동안)
다운타임세션 타임아웃 + 선출 시간리더 이전 시간만큼
데이터 유실 위험있음 (로그 미flush)낮음 (flush 후 종료)

로그 (로그 세그먼트)

  • 카프카의 토픽으로 들어오는 메시지는 세그먼트(로그 세그먼트) 파일에 저장됨
  • 메시지는 정해진 형식에 맞추어 순차적으로 로그 세그먼트 파일에 저정함
    • 메시지 뿐만 아니라 메시지의 키, 밸류, 오프셋, 메시지 크기 같은 정보가 함께 저장됨
  • 로그 세그먼트 파일들은 브로커의 로컬 디스크에 보관됨
  • 로그 세거먼트의 최대 크기는 1GB가 기본값
    • 1GB보다 크면 해당 파일은 닫고, 새로운 세그먼트를 생성하는 방식으로 진행됨

로그 세그먼트를 관리하는 방법은 로그 세그먼트 삭제와 컴팩션이 존재함

로그 세그먼트 삭제

  • 브로커의 설정 파일인 server.properties에서 log.cleanup.policy가 delete로 명시되어야 함
  • 삭제 기준은 시간(log.retention.ms)과 크기(log.retention.bytes) 두 가지가 있음
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
파티션의 로그 세그먼트 (디스크에 저장된 파일들)

  ┌──────────────────────────────────────────────────────────┐
  │  segment-00000.log     offset 0  ~ 999   (3일 전 생성)   │  ← 보존 기간 초과
  │  segment-01000.log     offset 1000~1999  (2일 전 생성)   │  ← 보존 기간 초과
  │  segment-02000.log     offset 2000~2999  (6시간 전 생성) │  ← 보존 중
  │  segment-03000.log     offset 3000~      (현재 활성)     │  ← 활성 세그먼트
  └──────────────────────────────────────────────────────────┘

  log.retention.hours=24 (보존 기간: 24시간) 설정 시


[STEP 1] 브로커가 주기적으로 각 세그먼트의 최종 수정 시간을 확인 후 삭제 대상 표시

  ┌──────────────────────────────────────────────────────────┐
  │  segment-00000.log  ❌ 삭제 대상 (3일 전)                │
  │  segment-01000.log  ❌ 삭제 대상 (2일 전)                │
  │  segment-02000.log  ✅ 보존 (6시간 전)                   │
  │  segment-03000.log  ✅ 활성 세그먼트 (삭제 불가)          │
  └──────────────────────────────────────────────────────────┘


[STEP 2] 삭제 대상 세그먼트를 .deleted 확장자로 변경 후 실제 삭제

  segment-00000.log  →  segment-00000.log.deleted  →  파일 삭제
  segment-01000.log  →  segment-01000.log.deleted  →  파일 삭제


[결과]

  ┌──────────────────────────────────────────────────────────┐
  │  segment-02000.log     offset 2000~2999                  │
  │  segment-03000.log     offset 3000~  (활성)              │
  └──────────────────────────────────────────────────────────┘

로그 세그먼트 컴팩션

  • 로그를 삭제하지 않고 컴팩션하여 보관할 수 있음
  • 기본적으로 로컬 디스크에 저장되어 있는 세그먼트를 대상으로 실행되는데, 현재 활성화된 세그먼트는 제외하고 나머지 세그먼트들을 대상으로 컴팩션이 실행됨

카프카 로그 컴팩션 과정 이미지

  • 장애 복구 시 젠체 로그를 복구하지 않고, 메시지의 키를 기준으로 최신의 상태만 복구함
  • 따라서 전체 로그를 복구할 때보다 복구 시간을 줄일 수 있다는 장점이 있음 => 빠른 장애 복구

하지만 모든 토픽에 로그 컴팩션을 적용하는 것은 좋지 않음

키값을 기준으로 최종값만 필요한 워크로드에 적용하는 것이 바람직함

옵션 이름옵션값적용 범위설명
cleanup.policycompact토픽의 옵션으로 적용토픽 레벨에서 로그 컴팩션을 설정할 때 적용하는 옵션
log.cleanup.policycompact브로커의 설정 파일에 적용브로커의 레벨에서 로그 컴팩션을 설정할 때 적용하는 옵션
log.cleaner.min.compaction.lag.ms0브로커의 설정 파일에 적용메시지가 기록된 후 컴팩션하기 전 경과되어야 할 최소 시간을 지정함. 만약 이 옵션을 설정하지 않으면, 마지막 세그먼트를 제외하고 모든 세그먼트를 컴팩션할 수 있음
log.cleaner.max.compaction.lag.ms9223372036854775807브로커의 설정 파일에 적용메시지가 기록된 후 컴팩션하기 전 경과되어야 할 최대 시간을 지정함
log.cleaner.min.cleanable.ration0.5브로커의 설정 파일에 적용로그에서 압축이 되지 않은 부분을 더티라고 표현함. 전체 로그 대비 더티의 비율이 50%가 넘으면 로그 컴팩션이 실행됨
This post is licensed under CC BY 4.0 by the author.

Trending Tags