msa

Kafka — 깊은 이야기

msa 로 돌아가기 / msa-design §6 참조

티켓팅 MSA에서 Kafka가 차지하는 자리, 왜 거기 있어야 하는지, 어떻게 동작하는지 정리. 우리가 직접 구현한 MyKafka(MVP) 가 좋은 reference.

📊 시각 자료: kafka-lifecycle-steps.html (단계별 인터랙티브 — 레코드/세그먼트/오프셋/zero-copy/MSA 5덱), kafka-lifecycle.html (정적 흐름도), kafka-summary.html (CRC·차이 상세).


§1. 왜 메시지 큐가 필요한가

동기 처리의 한계

이상적인 (아주 단순한) 예매 흐름:

Client → Booking API → [좌석 락 + DB INSERT + 결제] → 응답

문제: Booking API가 DB INSERT 끝날 때까지 응답을 못 함.

비동기 처리

Client → Booking API → [좌석 락(Redis) + Kafka publish] → 응답
                          ↓
                       Worker → DB INSERT (천천히)

메시지 큐가 주는 4가지

  1. 버퍼링 (back-pressure) — 생산 폭주를 소비 속도와 분리.
  2. 재시도/내구성 — Worker가 죽어도 메시지는 큐에 남아 다음 Worker가 처리.
  3. fan-out — 한 이벤트를 여러 소비자가 각자 다른 처리 (검색 인덱싱, 통계, 알림 등).
  4. 시간 분리(decoupling) — 생산자와 소비자가 같은 시점에 살아있을 필요 없음.

MSA 스케일에서 — 왜 대기업이 Kafka를 쓰는가

위는 단일 서비스(예매 API) 관점. 서비스가 수십~수백 개로 늘어나는 MSA에서는 문제가 질적으로 달라진다.

문제

  1. 강결합 (동기 호출 지옥) — A→B→C→D 직접 호출 체인. D가 느려지면 A까지 줄줄이 느려지고, 한 서비스 장애가 연쇄 장애(cascading failure) 로 전파. 새 소비자 붙이려면 호출하는 쪽 코드를 다 수정해야 함.
  2. 트래픽 스파이크 — 동기 처리는 가장 느린 지점(보통 DB)이 전체 처리량을 결정. 몰리면 커넥션/스레드 풀 고갈 → 5xx.

Kafka의 해결

문제 Kafka가 주는 것
폭주에 시스템이 죽음 비동기 버퍼 — API는 ms 응답, 워커가 자기 속도로 소비 (back-pressure)
서비스 강결합 decoupling — 발행만 하면 됨. 소비자가 죽어도 이벤트는 토픽에 남음
새 기능마다 생산자 수정 fan-out — 새 consumer group 추가에 기존 코드 0 수정
컨슈머 장애 시 유실 내구성 + replay — consume=삭제가 아니라 offset 전진. 되감아 재처리 가능
순서 vs 확장 partition — key로 순서 보장 + 병렬로 선형 확장

실전 쓰임새 — Kafka는 MSA의 중앙 신경계(central nervous system):

한 줄: 서비스가 많아질수록 동기 직접 호출은 결합·장애전파·스파이크에 취약해진다. Kafka는 그 사이에 재생 가능한 영속 로그 버퍼를 끼워 넣어 느슨하게 결합되고 탄력적이며 확장 가능한 시스템을 만든다. → 이벤트 드리븐 관점은 eventdriven 참조.


§2. Kafka의 핵심 개념 (5분 요약)

Topic / Partition / Offset

topic: reservation-events
├── partition-0: [m0, m1, m2, m3, ...]   ← append-only log
├── partition-1: [m0, m1, m2, ...]
└── partition-2: [m0, m1, m2, m3, m4, ...]

Producer

메시지를 발행. partition 선택 정책:

티켓팅 예: key = seatId 로 하면 같은 좌석 이벤트는 같은 partition에 모임 → 처리 순서 보장.

Consumer / Consumer Group

topic: reservation-events (partitions: 3)

group: "db-writer"        group: "search-indexer"
├── consumer-A → p0       ├── consumer-X → p0
├── consumer-B → p1       └── consumer-Y → p1, p2
└── consumer-C → p2

같은 메시지를 두 그룹이 각자 받음 (DB write 따로, 검색 인덱싱 따로).

Offset Commit

"여기까지 처리했어요"를 저장. consumer 재시작 시 거기서부터 이어 읽음.

ISR (In-Sync Replica)

→ MyKafka는 의도적으로 단일 노드라 이 부분 미구현. 학습 거리.


§3. 우리 MyKafka가 구현한 것 (자세히)

실제 코드 위치: ~/workspace/ticketing/MyKafka/src/main/kotlin/com/example/mykafka/ 토론용으로 코드 흐름을 따라가며 "왜 이렇게 짰는지" 풀어 씀.

3.0 패키지 구조 한 줄씩

mykafka/
├── Main.kt                          # 부팅 진입점, CLI 인자 파싱
├── server/
│   ├── BrokerServer.kt              # Netty TCP 서버 (boss/worker EventLoopGroup)
│   └── RequestRouter.kt             # Frame → ApiKey별 핸들러 분기
├── protocol/
│   ├── Frame.kt                     # (ApiKey + payload) 한 단위
│   ├── ApiKey.kt                    # enum: 5개 ApiKey
│   └── FrameCodec.kt                # length-prefix encode/decode (Netty handler)
├── log/
│   ├── Record.kt                    # 메모리상 record (offset, ts, key, value)
│   ├── RecordCodec.kt               # Record ↔ ByteArray + CRC32
│   ├── OffsetIndex.kt               # sparse index, mmap 기반
│   ├── Segment.kt                   # 단일 .log + .index = 한 segment
│   └── Log.kt                       # 한 파티션 (다수 segment 관리)
└── topic/
    ├── LogManager.kt                # 모든 토픽×파티션 관리 + partitioner
    └── OffsetStore.kt               # consumer offset 영속화 (자체 토픽 사용)

전체 ~470줄. 핵심은 log/topic/. server/protocol은 wire 처리.


3.1 Wire Protocol — Length-prefix Framing

모든 메시지는 길이 접두로 framed:

[ 4 bytes totalLength ][ 1 byte apiKey ][ payload (totalLength-1 bytes) ]

왜 length-prefix? TCP는 바이트 스트림. 메시지 경계가 없음.

FrameDecoder 의 정확한 동작 (FrameCodec.kt)

Netty ByteToMessageDecoder 를 상속. 한 번의 decode() 호출에서:

  1. readableBytes() < 4 → 길이 헤더도 못 읽음. 반환(대기)
  2. totalLength = input.readInt() → 그 만큼 안 모이면 resetReaderIndex() 후 반환(대기)
  3. apiKey = input.readByte() → 1바이트
  4. payload = input.readRetainedSlice(totalLength - 1) → 슬라이스 + retain
  5. out.add(Frame(apiKey, payload))

ByteToMessageDecoder 의 보석 같은 점은 부분 도착 자동 처리. 1MB짜리 frame이 TCP MSS 단위로 쪼개져 와도 알아서 누적해주고, 다 모이면 decode 호출.

→ 직접 짜면 partial-read 버그(읽기 시작 후 부족하면 어디까지 읽었는지 잃어버림) 잡기 어려움.


3.2 ApiKey 5개 — 각각의 wire format

protocol/ApiKey.kt:

enum class ApiKey(val code: Byte) {
    PRODUCE(0),
    FETCH(1),
    CREATE_TOPIC(2),
    COMMIT_OFFSET(3),
    FETCH_OFFSET(4);
}

실제 Kafka는 수십 개 (FetchSession, JoinGroup, Heartbeat, OffsetFetch v2…). MVP는 본질 5개만.

CREATE_TOPIC

req : [topicLen:2][topic][partitionCount:4]
resp: [status:1]   0=OK, 1=ALREADY_EXISTS, 2=INVALID

간단. broker는 디렉토리 만들고 Log 인스턴스 N개 생성.

PRODUCE (batch)

req : [topicLen:2][topic][partition:4][recordCount:4]
      for each record: [keyLen:4][key][valLen:4][value]
        - partition == -1  → broker가 partitioner로 결정 (첫 record의 key 사용)
        - keyLen == -1     → null key
resp: [errCode:1][partition:4][baseOffset:8][count:4]
        - 개별 offset = baseOffset, baseOffset+1, …, baseOffset+count-1

핵심 invariant: 한 batch의 N record는 모두 같은 (topic, partition). 실제 Kafka도 동일 — 클라이언트가 partition별로 미리 grouping해서 보낸다.

FETCH (pull)

req : [topicLen:2][topic][partition:4][offset:8][maxBytes:4]
resp: [errCode:1][recordCount:4][record1 bytes][record2 bytes]…

record는 self-delimiting (RecordCodec 포맷 그대로). client는 N번 decode.

진행 보장: record 1건이 maxBytes보다 커도 1건은 무조건 포함. 안 그러면 컨슈머 영원히 막힘. Kafka 동일 invariant.

COMMIT_OFFSET / FETCH_OFFSET

COMMIT_OFFSET:
  req : [groupLen:2][group][topicLen:2][topic][partition:4][offset:8]
  resp: [errCode:1]

FETCH_OFFSET:
  req : [groupLen:2][group][topicLen:2][topic][partition:4]
  resp: [errCode:1][offset:8]    offset=-1 → 미커밋

둘 다 내부 토픽 __consumer_offsets 에 저장. broker 재시작에도 살아남음. (§3.8 참조)


3.3 Record 포맷 — 한 메시지의 정확한 바이트

log/RecordCodec.kt:

[offset:8][timestamp:8][keyLen:4][key:keyLen][valueLen:4][value:valueLen][crc:4]

CRC 계산 (encode)

val crc = CRC32()
crc.update(buf.array(), 0, buf.position())   // 헤더+payload 전체
buf.putInt(crc.value.toInt())

CRC 검증 (decode)

decode 시 같은 범위로 다시 CRC 계산 → 저장된 값과 비교. 불일치면 error("CRC mismatch") → 호출자가 corruption으로 판단.

실 Kafka와 차이: MyKafka는 record별 CRC32. 실 Kafka(메시지 포맷 v2, 0.11+)는 batch별 CRC32C(Castagnoli, CPU HW 가속)이고, CRC 범위에서 baseOffset·leaderEpoch제외한다 → broker가 offset을 부여해도 CRC 재계산이 불필요 → zero-copy(§6) 성립의 전제. CRC는 "데이터가 저장 시점과 달라졌나"(disk bit rot 등 무결성)만 검출한다. 극히 드물게 깨진 데이터가 같은 CRC를 내는 충돌(~1/2³²)은 통과되어 silent corruption이 되지만, single-bit·짧은 burst 오류는 100% 검출. 악의적 변조는 막지 못함(그건 해시·서명의 영역).

부분 읽기 처리

decode는 buf 안에 record가 다 안 모였으면 null 반환 + buf position을 시작 위치로 복원. 호출자는 "truncated tail"로 판단해 거기서 멈춤.


3.4 OffsetIndex — Sparse + Mmap

log/OffsetIndex.kt. 한 segment에 하나씩 딸려 있음.

엔트리 포맷

[relativeOffset:4][filePosition:4]    ← 8바이트 entry

Sparse 정책

모든 record를 인덱싱하지 않음. Segment.append() 안에서:

if (lastIndexedPos < 0 || pos - lastIndexedPos >= indexIntervalBytes) {
    index.append((offset - baseOffset).toInt(), pos.toInt())
    lastIndexedPos = pos
}

Lookup 알고리즘 (이진 탐색)

fun lookup(targetRelativeOffset: Int): IntArray {
    var lo = 0; var hi = entryCount - 1; var bestIdx = -1
    while (lo <= hi) {
        val mid = (lo + hi) ushr 1
        val midOff = mmap.getInt(mid * ENTRY_SIZE)
        if (midOff <= targetRelativeOffset) {
            bestIdx = mid; lo = mid + 1
        } else hi = mid - 1
    }
    // ... bestIdx 에서 (rel, pos) 반환
}

target 이하 가장 큰 엔트리 찾기. 거기서부터 짧은 sequential scan으로 정확한 offset 도달.

Mmap 사용 이유

mmap = channel.map(FileChannel.MapMode.READ_WRITE, 0, fileSize)

Pre-allocate

파일 사이즈를 시작 시 maxEntries * 8 만큼 미리 truncate. 추후 append할 때 파일 크기 변경 syscall 없음.

자기 완결성

인덱스가 깨져도 .log 만 멀쩡하면 시작 시 완전 재구축. 검증에서 모든 .index를 0으로 깨도 정상 동작.

fun rebuildIndex(truncateTail: Boolean): Long {
    index.reset()
    scan { rec, startPos -> if (조건) index.append(...) }
    // 마지막 segment면 부분 쓰기된 trailing 자르기
}

3.5 Segment — 단일 (.log + .index)

log/Segment.kt. 파일명 규약:

00000000000000000000.log    ← %020d.format(baseOffset)
00000000000000000000.index
00000000000000001000.log    ← 1000 offset부터 새 segment
00000000000000001000.index

파일명 = baseOffset → 디렉토리 자체가 인덱스. 임의 offset이 어느 파일에 있는지 파일명만으로 이진 탐색 가능.

append (단건)

fun append(offset: Long, bytes: ByteArray) {
    val pos = sizeBytes
    channel.position(pos)
    channel.write(ByteBuffer.wrap(bytes))
    sizeBytes += bytes.size
    // sparse 인덱스 갱신 (§3.4)
}

호출자(=Log)가 락을 잡았다고 가정. FileChannel.write() 한 번 = 1 syscall.

appendBatch — 한 번의 write로 N record

fun appendBatch(baseOffsetOfBatch: Long, recordBytes: List<ByteArray>) {
    val totalLen = recordBytes.sumOf { it.size }
    val combined = ByteBuffer.allocate(totalLen)
    for (b in recordBytes) combined.put(b)
    combined.flip()
    channel.write(combined)    // ← 1 syscall로 모두 기록
    // 인덱스 엔트리는 record 단위로 갱신
}

핵심: N개 record를 메모리에서 합쳐 한 번에 write. syscall 수 = 1.

readFrom — pull-based 핵심

fun readFrom(startOffset: Long, maxBytes: Int): Pair<List<Record>, Int> {
    val rel = (startOffset - baseOffset).toInt()
    val (_, startPos) = index.lookup(rel).let { it[0] to it[1] }
    // startPos 부터 sizeBytes 끝까지 한 번에 읽어서 메모리 buf
    // buf 에서 RecordCodec.decode 반복
    // startOffset 이전은 skip, maxBytes 누적 초과하면 break
    // 단, 첫 record는 항상 포함 (진행 보장)
}

scan — recovery용 전체 순회

fun scan(onRecord: (Record, Long) -> Unit): Long {
    // .log 처음부터 끝까지 record 단위로 decode
    // CRC 깨진 record가 나오면 그 직전 byte 위치 반환
    return lastGoodPos
}

rebuildIndex — 시작 시 복구

  1. index reset
  2. .log 전체 scan → sparse 규칙대로 index 재구성
  3. 활성 segment면 마지막 정상 위치까지만 truncate (부분 쓰기된 꼬리 자르기)
  4. nextOffset = 마지막 record offset + 1

이게 "꼬리 자르기 복구" 의 정체. 전원 나가 부분 쓰기된 trailing record는 신뢰하지 않고 잘라낸다 (실 Kafka도 동일).


3.6 Log — 한 파티션 = 다수 segment 관리

log/Log.kt. 파티션 하나당 인스턴스 하나.

class Log(
    private val dir: Path,
    private val maxSegmentBytes: Long = 1MB,      // 기본 1MB (학습용)
    private val indexIntervalBytes: Int = 4KB,
) {
    private val lock = ReentrantLock()             // 파티션당 한 락
    private val segments: MutableList<Segment> = mutableListOf()
    @Volatile private var nextOffset: Long = 0
}

시작 시

  1. loadSegments() — 디렉토리에서 *.log 파일들 발견 → baseOffset 기준 정렬
  2. recoverAll() — 각 segment rebuildIndex(truncateTail = isActive)
  3. nextOffset = 활성 segment의 nextOffset

append (단건)

fun append(key: ByteArray?, value: ByteArray, timestamp: Long): Long =
    lock.withLock {
        maybeRoll()                                   // 활성 segment 크기 검사
        val offset = nextOffset
        val bytes = RecordCodec.encode(offset, timestamp, key, value)
        activeSegment().append(offset, bytes)
        nextOffset = offset + 1
        offset
    }

Broker가 offset 부여. 클라이언트가 못 정함 — 두 producer가 같은 partition에 쓸 때 충돌 방지.

maybeRoll — segment 분할 트리거

private fun maybeRoll() {
    val active = activeSegment()
    if (active.sizeBytes < maxSegmentBytes) return
    val newBase = nextOffset
    segments.add(newSegment(newBase))                 // 새 파일 생성
}

왜 분할?

find / read — 이진 탐색으로 segment 결정

private fun segmentIndexFor(offset: Long): Int? {
    // segments는 baseOffset 오름차순 정렬 → 이진 탐색
    // target 이하의 가장 큰 baseOffset 가진 segment 반환
}

fun read(startOffset: Long, maxBytes: Int): List<Record> {
    var idx = segmentIndexFor(startOffset)
    while (idx <= segments.lastIndex) {
        val (recs, used) = segments[idx].readFrom(off, budget)
        if (recs.isEmpty()) break
        out.addAll(recs); bytesUsed += used; off = recs.last().offset + 1
        idx++                                          // ← cross-segment 자동
    }
}

Cross-segment FETCH: 시작 offset이 segment-A에 있고 다음 record가 segment-B에 있어도 자동으로 넘어감. 컨슈머는 segment 경계를 의식할 필요 없음.


3.7 LogManager — Topic × Partition × Partitioner

topic/LogManager.kt. 한 브로커가 책임지는 모든 토픽 관리.

디렉토리 = 단일 진실 출처

<rootDir>/
  <topicName>/
    0/   ← partition 0
    1/   ← partition 1

별도 메타 파일 없음. 디렉토리 트리 자체가 토픽 목록 + 파티션 수 + segment 목록. 실 Kafka도 거의 동일 (<topic>-<partition>/ 단일 디렉토리지만 의미는 같음).

시작 시 loadExisting()

디렉토리 스캔 → 토픽별로 partition 디렉토리 발견 → 각각 Log 인스턴스 생성.

val ids = partDirs.map { it.fileName.toString().toInt() }
if (ids != (0 until ids.size).toList()) {
    log.warn("skipping topic={} — non-contiguous partition ids", name)
}

파티션 ID는 0..n-1 연속이라는 단순 가정. 학습용.

Partitioner — key 있으면 sticky hash

fun pickPartition(topic: String, key: ByteArray?): Int? {
    val t = topics[topic] ?: return null
    val n = t.partitions.size
    return if (key == null) t.nextRoundRobin(n)
           else Math.floorMod(key.contentHashCode(), n)
}

파티션마다 별도 Log → 락 충돌 0

Log 안에서만 락. 파티션 A 쓰기와 파티션 B 쓰기는 서로 안 기다림. "파티션 = 병렬성의 단위" 정체성을 코드로 표현.


3.8 OffsetStore — "모든 게 로그" 의 자기 dogfooding {#38}

topic/OffsetStore.kt. consumer group offset 영속화.

어디에 저장?

새 저장소 만들지 않음. broker가 이미 가진 메커니즘(Log)을 그대로 사용:

시작 시 cache 재구축

init {
    logManager.createTopic("__consumer_offsets", 1)
    val partition = logManager.getPartition("__consumer_offsets", 0)!!
    for (rec in partition.readAll()) {
        val k = decodeKey(rec.key!!)
        cache[k] = ByteBuffer.wrap(rec.value).long
    }
    log.info("loaded $count commits (${cache.size} unique keys)")
}

시작 시 내부 토픽을 처음부터 끝까지 scan → ConcurrentHashMap cache 재구축.

COMMIT

fun commit(group, topic, partition, offset) {
    val keyBytes = encodeKey(group, topic, partition)
    val valueBytes = ByteBuffer.allocate(8).putLong(offset).array()
    logManager.getPartition("__consumer_offsets", 0)!!
        .append(keyBytes, valueBytes)             // 그냥 또 다른 append
    cache[Key(group, topic, partition)] = offset
}

append + cache 업데이트. 끝.

FETCH

fun fetch(group, topic, partition): Long =
    cache[Key(group, topic, partition)] ?: -1L

cache lookup. 디스크 접근 0.

왜 이게 우아한가


3.9 RequestRouter — Frame을 핸들러로 분기

server/RequestRouter.kt. Netty SimpleChannelInboundHandler<Frame>.

override fun channelRead0(ctx, msg: Frame) {
    try {
        when (msg.apiKey) {
            ApiKey.PRODUCE       -> handleProduce(ctx, msg.payload)
            ApiKey.FETCH         -> handleFetch(ctx, msg.payload)
            ApiKey.CREATE_TOPIC  -> handleCreateTopic(ctx, msg.payload)
            ApiKey.COMMIT_OFFSET -> handleCommitOffset(ctx, msg.payload)
            ApiKey.FETCH_OFFSET  -> handleFetchOffset(ctx, msg.payload)
        }
    } finally {
        msg.payload.release()                       // Netty ByteBuf release 책임
    }
}

handleProduce 흐름

  1. payload에서 topic, partition, recordCount 읽기
  2. record N개 디코드 (각각 keyLen/key/valLen/value)
  3. partition == -1이면 logManager.pickPartition(topic, records[0].first)
  4. logManager.getPartition(topic, partition)Log 객체
  5. log.appendBatch(records) → baseOffset 반환
  6. 응답 [errCode:1][partition:4][baseOffset:8][count:4]

handleFetch 흐름

  1. topic, partition, offset, maxBytes 읽기
  2. log.read(offset, maxBytes)List<Record>
  3. 각 record를 RecordCodec.encode로 직렬화 → 합쳐서 응답

→ 응답 시 또 직렬화하는 건 zero-copy 미구현의 흔적. 실 Kafka는 sendfile() 또는 FileRegion 으로 OS가 디스크→소켓 직접 전송. byte copy 생략. → 다음 학습 거리.


3.10 검증된 동작 (DESIGN.md §8)

시나리오 결과
PRODUCE 라운드트립, monotonic offset
재시작 후 nextOffset 정확 복구
부분 쓰기 → 꼬리 자르기 복구
Segment 롤링 (파일명 = baseOffset)
Multi-segment recovery
Sparse mmap index, 임의 offset 빠르게 find
모든 .index 0으로 깨도 .log에서 재구축
같은 key sticky partition
30 keys → 3 partitions 10/10/10 분산
null key 엄격 round-robin
토픽/파티션 디렉토리 자동 복구
Batch PRODUCE all-or-none
1000건 single vs 1×1000 batch → 61x speedup
Cross-segment FETCH
maxBytes=1 진행 보장 (1건은 무조건)
Consumer loop (25 round, 50건 수신, self-offset)
__consumer_offsets 자동 생성
5 commits → unique key 2개, 최신 우선
재시작 → 컨슈머 group이 정확히 마지막 commit에서 이어 시작

61x speedup 의 정확한 숫자

"Kafka가 수십만 msgs/sec를 내는 비밀"의 정체. 데이터 양은 같다.


3.11 직접 짜본 결과 알게 된 통찰

코드 짜기 전에는 "그렇구나" 였지만 직접 짜보고서 와닿은 것들:

  1. "모든 게 로그" 는 정말 강력하다. OffsetStore 가 별도 저장소 없이 자기 자신을 dogfooding 하는 게 추상적인 슬로건이 아니라 실제 코드 절약. 새 코드는 key/value 인코딩 + cache 뿐.

  2. mmap이 마법이 아니다. 인덱스가 작아서 효과가 있는 거지, 큰 파일 mmap 하면 페이지 fault 폭발. 작은 데이터에 자주 접근하는 케이스에만 적합.

  3. CRC가 늦게 잡힌다. Record decode 시 CRC 검증을 하지만, 부분 쓰기는 decode 자체가 null 반환으로 잡힘. CRC는 disk corruption (드물지만 치명적) 용. 두 보호막이 다른 위협을 막음.

  4. 락은 작게. Log.append 안에만 락. partition 간에는 락 X. partition 수를 늘리면 throughput이 정확히 비례.

  5. 클라이언트가 offset 들고 다니는 비대칭의 깔끔함. broker가 컨슈머 상태를 안 갖는다 = 컨슈머 수에 비례한 broker 상태 폭발 없음. 컨슈머 N억 명도 broker는 그대로.

  6. 파일명 = baseOffset 이라는 한 줄의 미학. 메타 파일 없이 디렉토리만으로 segment 위치 결정. find data -name "*.log" 가 곧 toc.


§4. 의도적으로 안 넣은 것

영역 MVP에서 생략 다음 단계 학습 포인트
Replication 단일 노드 Leader/Follower + ISR + acks=all
Controller 메타 단순 관리 KRaft (Raft 기반 controller, Zookeeper 제거)
Log compaction sparse만 구현 같은 key 옛 record 삭제 — __consumer_offsets 정리에 즉시 활용
Exactly-once at-least-once producer idempotence + transactions
Zero-copy byte copy sendfile() / Netty FileRegion
Group coordination 단일 group string rebalance, member tracking, sticky assignment
Durability tuning OS flush 의존 channel.force() 정책 (per-request / per-N / time-based)

각 항목 모두 학습 거리. 부재가 한계로 드러나면 그때 구현.

클라이언트 SDK도 미구현

서버(BrokerServer/RequestRouter)는 있지만 외부에서 PRODUCE/FETCH 보내는 Kotlin 클라이언트 라이브러리는 없음. ticket-command가 MyKafka에 붙으려면 다음 단계에서 만들 것:

MyKafka/src/main/kotlin/com/example/mykafka/client/
├── Producer.kt   ← Netty 기반 PRODUCE 보내기
└── Consumer.kt   ← FETCH 폴링 루프 + 자동 offset commit

§5. 티켓팅에 어떻게 적용할 것인가

5.1 Worker 패턴 (다이어그램의 핵심)

Booking API ──PRODUCE──> Kafka(reservation-events)
                            │
                            └──FETCH──> Worker ──INSERT──> RDB

발행 시점

좌석 락(Redis) + 결제 검증 통과 직후. DB INSERT 전.

// command service의 의사 코드
@Transactional  // Redis 락만 잡고 짧게 끝
fun reserve(req: ReserveRequest): ReservationToken {
    val seat = redis.acquireSeatLock(req.seatId)  // 짧음
    val token = ReservationToken(uuid())
    kafka.publish("reservation-events", req.seatId, ReservationCreated(token, req.userId, req.seatId))
    return token  // 응답 즉시
}

소비 시점

Worker는 자기 속도대로 consume → DB INSERT.

// worker service 의사 코드
consumer.subscribe("reservation-events")
while (true) {
    val records = consumer.poll(timeoutMs = 100)
    for (r in records) {
        db.insertReservation(r.value)
        consumer.commitOffset(r.offset)
    }
}

5.2 왜 key=seatId 인가

5.3 Outbox 패턴 (학습 거리)

문제: command service가

1. DB에 reservation 행 INSERT
2. Kafka에 ReservationCreated publish

1만 성공하고 2 실패하면? 또는 반대?

outbox 테이블 도입:

1. (한 트랜잭션에) DB에 reservation 행 INSERT + outbox 행 INSERT
2. 별도 outbox-relay 가 outbox 폴링 → Kafka publish → outbox 행 삭제

DB 트랜잭션 하나로 원자성 보장. Kafka 발행은 최소 한 번(at-least-once)이 보장됨.

→ 정확히는 우리 다이어그램은 Booking API가 DB INSERT를 안 함. Worker가 함. 그래서 outbox는 다른 상황(예: 예매 취소 시 환불 알림 발행)에서 학습 거리.


§6. Kafka가 빠른 이유 (정리)

흔히 "Kafka는 빠르다"고만 말함. 왜 빠른가:

  1. 순차 디스크 쓰기 (append-only log) — 랜덤 access 대비 HDD/SSD 모두 압도적으로 빠름. 메모리에 가까운 속도.
  2. Page cache 활용 — OS 페이지 캐시가 hot data를 메모리에 유지. JVM 힙 부담 X.
  3. Zero-copy (sendfile()) — consumer fetch 시 kernel-space → socket 직접 전송. user-space 복사 생략.
  4. Batch 처리 — producer가 메시지를 모아 한 번에 전송. 네트워크 RTT 분할상환.
  5. Compression (선택) — batch를 압축해서 전송. gzip/snappy/lz4/zstd.
  6. Sparse index — 인덱스 자체가 작아 메모리에 다 올라감. 인덱스 lookup도 빠름.
  7. Partition별 병렬화 — partition 수만큼 producer/consumer 병렬.

Zero-copy 자세히 (전통 경로 vs sendfile)

1. 기존의 파일 전송 방식 (Traditional Approach)

일반적으로 서버에서 파일을 읽어 소켓으로 보낼 때는 유저 공간(Application)과 커널 공간(OS)을 서너 번 오가며, CPU가 데이터를 복사하느라 바쁩니다.

이동 경로 (4번의 복사, 4번의 컨텍스트 스위칭)

  1. Disk ➔ 페이지 캐시 (Kernel): 디스크에서 데이터를 읽어 OS 영역의 페이지 캐시(Read Buffer)로 가져옵니다. (DMA 카피)

  2. 페이지 캐시 (Kernel) ➔ 어플리케이션 버퍼 (User): JVM이나 카프카 소스코드 엔진이 쓰기 위해 데이터를 유저 메모리로 가져옵니다. (CPU 카피 - 1과 2를 오가며 Context Switch 발생)

  3. 어플리케이션 버퍼 (User) ➔ 소켓 버퍼 (Kernel): 네트워크로 보내기 위해 OS의 소켓(Socket) 버퍼로 데이터를 보냅니다. (CPU 카피 - 다시 Context Switch 발생)

  4. 소켓 버퍼 (Kernel) ➔ NIC 버퍼 (Hardware): 최종적으로 네트워크 카드(NIC)로 데이터를 보냅니다. (DMA 카피)

💡 문제점: 카프카 같은 메시지 브로커는 파일 내용을 수정하지 않고 '읽어서 그대로 보내기만' 하는 경우가 대부분입니다. 그런데 중간에 어플리케이션 버퍼(User Space)를 거치느라 CPU 자원과 메모리가 낭비되는 비효율이 발생합니다.

2. 카프카의 제로 카피 방식 (Zero Copy)

카프카는 자바의 FileChannel.transferTo() 메서드를 사용하며, 이는 내부적으로 리눅스의 sendfile 시스템 콜을 호출합니다. 핵심은 "어플리케이션 버퍼를 완전히 패싱한다"는 점입니다.

이동 경로 (2번의 복사, 2번의 컨텍스트 스위칭)

  1. Disk ➔ 페이지 캐시 (Kernel): 디스크에서 데이터를 읽어 페이지 캐시로 가져옵니다. (DMA 카피)

  2. 페이지 캐시 (Kernel) ➔ NIC 버퍼 (Hardware): 중간 과정을 전부 생략하고, 페이지 캐시의 데이터를 네트워크 카드로 직접 쏩니다. (DMA 카피)

(참고: 리눅스 커널 버전에 따라 소켓 버퍼에 데이터의 위치와 길이 같은 '최소한의 메타데이터'만 복사하는 하드웨어 내장 스캐터-개더(Scatter-Gather) DMA 방식을 사용하여 효율을 극대화합니다.)

왜 Kafka가 이게 되나 — broker가 데이터를 변형 없이 저장 포맷 그대로 흘려보내기 때문. CRC가 offset을 안 덮어서(§3.3) offset 부여해도 재계산 불필요, 압축도 그대로 전달 → user-space로 끌어올릴 이유가 없음.

한계 — TLS 암호화를 켜면 user-space 암호화 때문에 복사가 부활 → zero-copy 깨짐. 재압축·포맷 변환도 동일. 보안 vs 성능 트레이드오프.

→ MyKafka는 fetch 응답을 재직렬화(byte copy)하므로 미구현(§3.9). §7-[7]에서 Netty FileRegion으로 도입 예정.

우리 MyKafka가 검증한 것


§7. 학습 진화 단계 (Kafka 측면)

[현재] MyKafka MVP (단일 노드, client SDK 부재)
   ↓ client SDK 구현
[1] Producer / Consumer 라이브러리 → ticket-command 가 publish
   ↓ Worker 도입
[2] Worker service 신규 → consume → DB INSERT (비동기 영속화)
   ↓ 부하 측정으로 효과 검증
[3] partition 수 늘리며 처리량 변화 측정
   ↓ 신뢰성 보강
[4] Replication 구현 (Leader/Follower + ISR + acks=all)
   ↓ 운영성
[5] Log compaction (consumer_offsets 정리)
   ↓ 정확성
[6] Producer idempotence + Transactions (exactly-once)
   ↓ 성능
[7] Zero-copy (`FileRegion` Netty)

각 단계마다 무엇을 측정해서 효과를 보여줄 것인지 미리 정해두는 게 학습 핵심.


§8. 자주 나오는 오해

오해 진실
"Kafka는 메시지 큐다" 정확히는 distributed commit log. 큐처럼 쓰지만 "consume = 삭제"가 아니라 "consume = offset 전진". 같은 메시지를 N번 다시 읽을 수 있음.
"Kafka는 빠르니까 어디나 쓰자" 1KB 미만 짧은 메시지에는 오히려 RabbitMQ가 latency 낮을 수 있음. Kafka는 throughput 최적화.
"exactly-once 보장하니까 안전" exactly-once는 Kafka 내부 + producer side 한정. consumer 측 DB INSERT까지는 보장 X. 우리가 처리(idempotent consumer) 해야 함.
"partition 늘리면 무조건 빠름" partition 수 = consumer 병렬도 상한. consumer 수 < partition 수면 idle consumer 생김. 그리고 partition은 한번 늘리면 줄이기 어려움.
"offset commit은 자동" auto-commit은 데이터 손실 위험. 처리 완료 후 manual commit이 일반적.

§9. 토론 prompts 일부

자세한 건 discussion §3 참조.


§10. 코드 reading map

MyKafka 코드 읽을 때 권장 순서. 한 호흡에 한 단계.

[1] protocol/Frame.kt + FrameCodec.kt           "메시지 경계 어떻게 잡나?"
    └─ length-prefix framing의 5줄짜리 본질

[2] log/Record.kt + RecordCodec.kt              "한 메시지의 정확한 바이트는?"
    └─ CRC 위치와 검증 흐름

[3] log/OffsetIndex.kt                          "왜 sparse + mmap인가?"
    └─ lookup 이진 탐색 + sequential scan 의 결합

[4] log/Segment.kt — append/readFrom/scan       "한 segment의 라이프사이클"
    └─ rebuildIndex(truncateTail) 가 꼬리 자르기

[5] log/Log.kt — maybeRoll/find/read            "여러 segment 어떻게 관리?"
    └─ segmentIndexFor 이진 탐색 + cross-segment

[6] topic/LogManager.kt — pickPartition          "partition 어떻게 정해?"
    └─ floorMod(hash, n) vs round-robin

[7] topic/OffsetStore.kt                         "오프셋도 그냥 또 다른 로그"
    └─ 자기 dogfooding의 우아함

[8] server/RequestRouter.kt                      "ApiKey 분기 끝"
    └─ handleProduce / handleFetch 흐름 따라가기

[9] DESIGN.md §6 8단계                           "왜 그 순서로 만들었나"
    └─ Step1 framing → Step2 append-only → ... → Step8 offset 영속화

[1][2] 가 wire/포맷. [3][5] 가 디스크 추상. [6]~[8] 이 broker. 각 층이 다음 층의 invariant를 깔끔히 가정해서 위에 쌓는 구조.

토론하면 좋은 코드 포인트

Segment.appendBatch — 왜 syscall 수가 throughput을 지배하나?

N record를 메모리에서 합쳐 한 번에 write.

: syscall 한 번엔 고정 비용이 붙는다 — user→kernel 모드 전환, CPU 캐시·TLB 오염, 커널 진입 경로. 데이터 양과 무관한 이 고정 오버헤드 × 호출 횟수가 병목이다. 1000건을 1000번 write하면 모드전환 1000번 + 락 1000번 (+네트워크면 RTT 1000번). batch로 합치면 메모리 복사(싸다)는 N번이지만 syscall·락·RTT는 1번. 디스크도 큰 덩어리 순차쓰기가 페이지 캐시에 유리. 측정: 84ms→1.4ms(61배) — 데이터 양은 동일, 줄어든 건 오직 호출 횟수. 이게 분할상환(amortization) 의 정체.

Log.read — 컨슈머가 segment 경계를 의식 안 해도 되는 게 왜 중요한가?

cross-segment 자동 처리. offset이 segment-A에서 끝나고 B로 넘어가도 알아서 이어 읽음.

: 내부 표현과 외부 계약의 분리(캡슐화). 컨슈머는 "offset부터 maxBytes"만 알 뿐 segment가 1개인지 1000개인지, 경계가 어딘지 모른다. 만약 경계를 알아야 했다면 — ① segment 크기·롤링 정책이 클라이언트 API로 새어나가 서버가 정책을 못 바꿈, ② 클라이언트가 경계마다 재요청 로직 필요(복잡+버그), ③ retention으로 segment가 삭제되면 클라이언트 상태가 깨짐. 자동 처리 덕에 서버는 segment 크기 조절·compaction·retention을 자유롭게 하고, 클라이언트는 단순하게 유지된다.

OffsetIndex.lookup — sparse + sequential scan trade-off의 의미는?

이진 탐색으로 "그 이하 가장 가까운" 엔트리를 찾고, 거기서부터 짧게 scan.

: 공간 vs 시간 trade-off, 그리고 "메모리에 다 올라가는 것"의 승리. 모든 record를 인덱싱(dense)하면 lookup은 정확하지만 인덱스가 데이터만큼 커져 메모리에 못 올림. sparse(4KB 간격)는 인덱스가 수백 배 작아 통째로 mmap/page cache에 상주 → 이진 탐색이 디스크 안 타고 메모리에서 끝난다. 대신 정확한 위치가 아닌 "근처"를 주지만, 거기서부터의 sequential scan은 indexInterval(≤4KB) 이내로 bounded이고 순차 읽기라 빠르다. 즉 "거의 다 와서 마지막 몇 걸음만 걷기" — 작은 scan 비용을 내주고 인덱스를 메모리에 통째로 올리는 이득을 산다.

OffsetStore.commit — 두 줄로 durable consumer state가 생긴 마법은?

append(__consumer_offsets) + cache.put, 끝.

: "모든 게 로그"의 dogfooding. consumer offset도 결국 key=(group,topic,partition) → value=offset의 시계열 사실일 뿐. 이미 broker가 가진 Log(append-only + 재시작 recovery + durability)를 그대로 재사용하면 새로 만들 건 key/value 인코딩 + cache뿐이다. append 한 줄 = durability(디스크 영속·복구 공짜), cache.put 한 줄 = 읽기 성능. 별도 DB·저장엔진·복구로직 0. 같은 추상이 사용자 데이터와 메타데이터를 동시에 담당하니 코드·검증·운영 표면이 한 곳. 같은 key 최신값만 의미 있다는 compaction 의미까지 공짜로 따라온다. 마법의 정체 = 이미 있는 강력한 추상을 한 번 더 쓴 것.

Log.appendlock.withLock — partition 늘리면 throughput이 비례하는 이유는?

파티션 하나당 Log 인스턴스 하나, 그 안에 ReentrantLock 하나.

: 락 경계가 곧 병렬성의 경계. 락이 partition 단위로 독립이라 partition A 쓰기와 B 쓰기는 서로 다른 락 → 안 기다린다(no contention). 동시에 쓸 수 있는 writer 수 = partition 수. 경합이 없으면 CPU 코어·디스크 대역폭이 허용하는 한 선형 확장한다. 만약 broker 전역 락 하나였다면 partition을 아무리 늘려도 모든 쓰기가 직렬화되어 throughput이 고정(Amdahl의 직렬 구간). "partition = 병렬성의 단위"라는 정체성이 코드의 락 경계로 그대로 드러난 것. 단, '비례'의 전제: ① 자원(코어·디스크 IO)에 여유가 있어야 하고 — 물리 상한에 닿으면 더는 안 늘어남, ② key가 고르게 분산돼야 한다 — 한 partition에 쓰기가 몰리면 그 락에서 경합이 생겨 hot partition이 병목. 그래서 partition 수와 key 설계는 같이 고민해야 함(§9, §5.2 참조).

1. 카프카가 "극단적 패싱(Zero Copy)"을 하는 이유

카프카의 목적은 오직 하나, "데이터를 안 깨뜨리고 가장 빠르게 전달하는 것"입니다.

카프카 브로커(Broker)는 프로듀서가 보낸 메시지가 무슨 내용인지(회원 가입 정보인지, 결제 내역인지) 전혀 궁금해하지 않습니다. 그냥 '바이트 덩어리(Byte Array)'로 취급합니다.

2. 스프링은 왜 그걸 안 할까? (정확히는 못 할까?)

스프링의 목적은 "데이터를 열어보고, 쪼개고, 조립해서 비즈니스 로직을 수행하는 것"입니다.

우리가 만드는 일반적인 스프링 부트(Spring Boot) 웹 애플리케이션의 동작 과정을 생각해 보면 왜 안 쓰는지 알 수 있습니다.

  1. 데이터 변환 및 파싱: DB에서 바이너리 데이터를 긁어오면, MyBatis나 JPA가 이를 Java 객체(Entity, DTO)로 변환(역직렬화)해야 합니다.

  2. 비즈니스 로직 수행: if (user.isVip()) { discount(); } 같은 조건문을 체크하고 데이터를 수정합니다.

  3. 보안 및 권한 검증: 스프링 시큐리티가 돌면서 이 유저가 이 데이터를 볼 자격이 있는지 코드로 검사해야 합니다.

  4. 결과물 생성: 최종 객체를 다시 Jackson 라이브러리를 통해 JSON 문자열로 변환(직렬화)해서 클라이언트에 줍니다.

💡 핵심: 스프링은 데이터를 열어서 수정하고 가공해야 하는 요리사입니다. 재료를 만지려면 당연히 요리대(JVM 유저 메모리 공간) 위로 데이터를 올려놓아야 합니다. 데이터를 보지도 않고 네트워크 카드로 바로 밀어 넣어버리는 Zero Copy는 애초에 불가능한 구조인 것이죠.

3. 반전: 스프링도 '필요할 때'는 Zero Copy를 쓴다!

그렇다면 스프링은 평생 Zero Copy를 안 쓸까요? 아닙니다. 스프링도 카프카처럼 "데이터를 안 열어보고 그대로 전송만 해도 되는 상황"에서는 적극적으로 제로 카피를 씁니다.