KIP-794에 실린 1시간짜리 실측에서 세 파티션의 배치 개수는 1683/1713/1711로 거의 같았는데, 느린 브로커가 리더인 파티션 0의 배치당 크기는 약 15.4KB, 나머지 둘은 약 4.5KB였다. 배치 수는 고른데 바이트는 세 배 넘게 차이가 난다. 즉 느린 브로커가 오히려 더 많은 데이터를 떠안고 있었다.
나는 "키 없는 레코드는 라운드로빈으로 고르게 퍼진다"고 알고 있었다. 소스(BuiltInPartitioner, RecordAccumulator)를 따라가 보니 현재 기본 동작은 라운드로빈도, KIP-480의 sticky도 아니었다. 이 글은 한 가지 질문만 따라간다. 파티션 전환을 "언제" 하느냐가 왜 부하 분포를 뒤집는가.
KIP-480 sticky 파티셔너는 배치를 작게 쪼개지 않으려고 한 파티션에 붙어 있다가 새 배치가 만들어질 때 다음 파티션으로 넘어갔다. 문제는 "새 배치 생성" 시점이 브로커가 얼마나 빨리 가져가느냐(drain)에 달려 있다는 점이다.
브로커 B가 잠깐 느려짐
→ B 파티션 배치가 전송되지 못하고 큐에서 대기
→ 대기하는 동안 그 배치가 계속 채워짐 (새 배치가 늦게 생김)
→ 전환이 늦어짐 = B에 머무는 시간이 김
→ B가 더 많이 받음 → 더 느려짐 (양의 되먹임)
빠른 브로커 (linger.ms=0)
→ 레코드 1개짜리 배치도 즉시 전송 → 바로 새 배치 → 즉시 전환
배치 개수로 균등해도 바이트로는 균등하지 않다. KIP-794는 이 상태를 "neither uniform nor sufficiently sticky"라고 요약한다.
KIP-794 이후에는 partitioner.class 기본값이 null이고, 키 없는 레코드는 send() 시점이 아니라 accumulator에 append하는 순간 BuiltInPartitioner가 파티션을 정한다. 전환 판단은 실제로 쓴 바이트 수로 한다.
// BuiltInPartitioner#updatePartitionInfo 발췌 (주석은 필자)
int producedBytes = info.producedBytes.addAndGet(appendedBytes);
if (producedBytes >= stickyBatchSize && enableSwitch // batch.size만큼 썼고 전환 허용
|| producedBytes >= stickyBatchSize * 2) { // 2배면 무조건 전환
stickyPartitionInfo.set(new StickyPartitionInfo(nextPartition(cluster)));
}
브로커가 빠르든 느리든 한 번 머물 때 받는 양은 대략 batch.size(기본 16384)로 같아지고, 지연과 분배량의 되먹임이 끊긴다. enableSwitch는 큐의 마지막 배치가 꽉 찼는지(allBatchesFull)로 정해진다. 반쯤 찬 배치를 남기고 떠나면 작은 배치가 늘기 때문이고, 대신 2 × batch.size에서는 강제로 넘어간다.
다음 파티션은 균등 랜덤이 아니다. RecordAccumulator#ready()가 돌 때마다 파티션별 배치 큐 길이를 모으고, (max+1 − 길이)로 뒤집어 누적 빈도표(CFT)를 만든 뒤 가중 랜덤으로 고른다.
큐가 길다 = 브로커가 못 가져가고 있다. 그래서 큐 길이를 뒤집은 값을 가중치로 쓰면 밀린 브로커일수록 덜 선택된다.
아래 코드는 이 계산을 그대로 옮긴 것이다. JShell에 붙여 넣으면 바로 돌아간다.
int[] q = {1, 5, 2}; // 파티션별 배치 큐 길이
int maxPlus1 = java.util.Arrays.stream(q).max().getAsInt() + 1;
int[] cft = new int[q.length];
cft[0] = maxPlus1 - q[0];
for (int i = 1; i < q.length; i++) cft[i] = maxPlus1 - q[i] + cft[i - 1];
int[] hit = new int[q.length];
for (int r = 0; r < cft[q.length - 1]; r++) { // r을 전 구간 순회해 확률 분포를 센다
int idx = Math.abs(java.util.Arrays.binarySearch(cft, 0, q.length, r) + 1);
hit[idx]++;
}
System.out.println(java.util.Arrays.toString(cft) + " " + java.util.Arrays.toString(hit));
// [5, 6, 10] [5, 1, 4] → p0 50%, p1 10%, p2 40%
큐 길이가 5로 가장 밀린 p1은 10%만 뽑힌다. +1 덕분에 가중치가 최소 1이라 굶지는 않는다. 큐 길이가 전부 같으면 표를 만들지 않고 균등 랜덤으로 돌아간다. KIP-480의 양의 되먹임이 음의 되먹임으로 바뀐 셈이다.
이 가중은 partitioner.adaptive.partitioning.enable(기본 true)로 켜고 끈다. 응답 없는 브로커를 아예 후보에서 빼려면 partitioner.availability.timeout.ms(기본 0, 비활성)를 쓴다.
소스를 보기 전에 내가 믿고 있던 것과 실제 동작을 표로 대조했다.
| 흔한 이해 | 실제 (KIP-794 이후 기본) |
|---|---|
| 키 없으면 라운드로빈 | batch.size 바이트만큼 한 파티션에 붙었다가 전환. 라운드로빈은 RoundRobinPartitioner를 직접 지정해야 쓰인다 |
| sticky 전환 = 새 배치 생성 시 | 그게 KIP-480의 불균등 원인이었고, 지금은 바이트 기준 |
| 다음 파티션은 균등 랜덤 | 큐 길이 역가중 랜덤, 큐가 전부 같을 때만 균등 |
UniformStickyPartitioner를 지정하면 최신 동작 | deprecated. partitioner.class를 비워 두고 키를 무시하려면 partitioner.ignore.keys=true |
마지막 줄은 실무에서 부딪히기 쉽다. 예전 가이드대로 partitioner.class에 deprecated 파티셔너를 박아 두면 append 시점 결정과 적응형 가중이라는 내장 경로를 오히려 우회하게 된다.
KIP-480은 "언제 전환하나"를 브로커 속도에 묶는 바람에 느린 브로커에 데이터를 몰아줬고, KIP-794는 전환을 바이트에 묶고 다음 파티션 선택을 큐 길이의 역수로 가중해 그 되먹임을 반대로 돌렸다. 다음에는 이 큐 길이를 실제로 줄이는 쪽, 즉 RecordAccumulator#drain의 노드별 순회와 BufferPool의 메모리 압박(max.block.ms)이 적응형 선택과 어떻게 맞물리는지 따라가 볼 생각이다.
BuiltInPartitioner.java, RecordAccumulator.java, KafkaProducer#partition, Sender.java, ProducerConfig.java