기본개념 · 11장
운영 기초
클러스터를 세우는 것과 운영하는 것은 다른 일입니다.
이 장은 네 개의 대주제 — 보안 · 모니터링 · 성능 튜닝 · 파티션 재할당 — 와
4.2에서 production-ready가 된 Share Groups를 다룹니다.
특히 listeners · advertised.listeners · listener.security.protocol.map
세 설정의 관계는 운영자가 가장 많이 틀리는 지점이므로 별도 절로 정리했습니다.
CCAAK의 중심 내용이면서 CCDAK의 Application Observability(13%)와도 겹칩니다.
학습 목표
- 암호화 · 인증 · 인가 · 감사 네 계층을 구분하고
security.protocol네 값의 의미를 정확히 설명할 수 있습니다. listeners/advertised.listeners/listener.security.protocol.map의 역할을 구분하고, KRaft의 컨트롤러 리스너 요구사항을 설명할 수 있습니다.- ACL의 principal × resource × operation 모델과
super.users·allow.everyone.if.no.acl.found의 영향을 설명할 수 있습니다. - 반드시 감시해야 할 브로커 JMX 메트릭과 각각의 정상 범위를 말할 수 있습니다.
- consumer lag을 측정하는 세 가지 방법과 각각의 함정을 구분할 수 있습니다.
kafka-reassign-partitions.sh의 3단계 워크플로와 throttle 해제의 중요성을 설명할 수 있습니다.- Share Groups가 컨슈머 그룹과 무엇이 다른지, ack 타입 4종의 의미를 설명할 수 있습니다.
보안 — 네 개의 계층
Kafka의 보안은 서로 독립적인 네 계층으로 이루어집니다. 공식 문서는 지원되는 보안 조치를 연결 인증(SSL 또는 SASL), 전송 데이터 암호화(SSL), 클라이언트 읽기·쓰기 인가(authorization), 그리고 인가의 플러그인화와 외부 인가 서비스 연동으로 정리합니다. 문서는 보안이 선택 사항이며, 인증·비인증·암호화·비암호화 클라이언트를 섞어 쓸 수도 있다는 점도 명시합니다.
security.protocol 네 값
| 값 | 암호화 | 인증 | 설명 |
|---|---|---|---|
PLAINTEXT |
없음 | 없음 | 공식 문서 표현대로 "어떤 보안도 제공하지 않고 추가 설정도 필요하지 않은" 프로토콜. 브로커 listeners의 기본값(PLAINTEXT://:9092)입니다 |
SSL |
TLS | 선택적 (mTLS) | 전송 암호화 + 클라이언트 인증서 기반 인증(선택). 브로커의 ssl.client.auth가 required일 때만 클라이언트 인증이 강제됩니다 (기본값은 none) |
SASL_PLAINTEXT |
없음 | SASL | 인증은 하지만 전송은 평문입니다. PLAIN 메커니즘과 조합하면 비밀번호가 네트워크에 평문으로 흐릅니다 |
SASL_SSL |
TLS | SASL | SASL로 인증 + TLS로 암호화. 프로덕션 권장 조합입니다 |
SASL 메커니즘
| 메커니즘 | 자격증명 | 자격증명 저장 위치 | 언제 고르나 |
|---|---|---|---|
GSSAPI (Kerberos)0.9.0.0 |
Kerberos 티켓 / keytab | Kerberos KDC (외부) | 사내에 이미 Kerberos·Active Directory가 있는 경우. sasl.enabled.mechanisms의 기본값이자 유일한 기본 활성 메커니즘입니다 |
PLAIN0.10.0.0 |
사용자명 / 비밀번호 | 브로커의 JAAS 설정 파일 (user_{name} 속성) |
가장 단순합니다. 대신 사용자를 추가할 때마다 브로커 설정을 고치고 재시작해야 하고, 반드시 SASL_SSL과 함께 써야 합니다 |
SCRAM-SHA-256 / SCRAM-SHA-5120.10.2.0 |
사용자명 / 비밀번호 (챌린지-응답) | Kafka 클러스터 메타데이터. kafka-configs.sh로 런타임에 추가·삭제 |
외부 인증 시스템 없이 사용자 관리를 하고 싶을 때의 기본 선택. 브로커 재시작 없이 사용자 추가가 가능합니다 |
OAUTHBEARER2.0 |
OAuth 2.0 Bearer 토큰 (JWT) | 외부 인증 서버 (JWKS 엔드포인트로 검증) | 이미 OAuth/OIDC 기반 인증 체계가 있는 경우. 토큰 만료·갱신이 자동화됩니다 |
# 사용자 alice 에게 SCRAM-SHA-256 자격증명 부여 — 브로커 재시작이 필요 없습니다
bin/kafka-configs.sh --bootstrap-server localhost:9092 --alter \
--add-config 'SCRAM-SHA-256=[iterations=8192,password=alice-secret]' \
--entity-type users --entity-name alice \
--command-config client.properties
# 조회
bin/kafka-configs.sh --bootstrap-server localhost:9092 --describe \
--entity-type users --entity-name alice --command-config client.properties
# 삭제
bin/kafka-configs.sh --bootstrap-server localhost:9092 --alter \
--delete-config 'SCRAM-SHA-256' \
--entity-type users --entity-name alice --command-config client.properties
# ── server.properties (브로커) ─────────────────────────────
listeners=CLIENT://:9092,BROKER://:9093,CONTROLLER://:9094
inter.broker.listener.name=BROKER
controller.listener.names=CONTROLLER
listener.security.protocol.map=CLIENT:SASL_SSL,BROKER:SASL_SSL,CONTROLLER:SASL_SSL
# 기본값은 GSSAPI 하나뿐이므로 반드시 명시합니다
sasl.enabled.mechanisms=SCRAM-SHA-256
sasl.mechanism.inter.broker.protocol=SCRAM-SHA-256
sasl.mechanism.controller.protocol=SCRAM-SHA-256
ssl.keystore.location=/etc/kafka/ssl/broker.keystore.jks
ssl.keystore.password=changeit
ssl.truststore.location=/etc/kafka/ssl/broker.truststore.jks
ssl.truststore.password=changeit
# 인가를 켭니다. KRaft 에서는 StandardAuthorizer 를 씁니다
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
# 구분자는 세미콜론입니다 (SSL 사용자명에 콤마가 들어갈 수 있으므로)
super.users=User:admin
# ── client.properties (클라이언트) ─────────────────────────
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-256
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \
username="alice" \
password="alice-secret";
ssl.truststore.location=/etc/kafka/ssl/client.truststore.jks
ssl.truststore.password=changeit
listeners · advertised.listeners · listener.security.protocol.map
이 세 설정을 헷갈리는 것이 Kafka 운영에서 가장 흔한 실수입니다.
브로커는 시작되고 kafka-topics.sh --list도 동작하는데
애플리케이션만 붙지 못하는 상황의 원인이 거의 항상 여기 있습니다.
| 설정 | 기본값 | 질문에 답한다 | 형식 |
|---|---|---|---|
listeners |
PLAINTEXT://:9092 |
"브로커가 어느 주소·포트에 바인딩해 소켓을 열 것인가" | {리스너이름}://{호스트}:{포트} 콤마 구분. 호스트를 비우면 기본 인터페이스, 0.0.0.0은 모든 인터페이스 |
advertised.listeners |
null |
"클라이언트와 다른 브로커에게 알려 줄 접속 주소는 무엇인가" | 같은 형식. null이면 listeners 값이 그대로 광고됩니다 |
listener.security.protocol.map |
SASL_SSL:SASL_SSL, |
"각 리스너 이름이 어떤 보안 프로토콜을 쓰는가" | {리스너이름}:{보안프로토콜} 콤마 구분 |
관계를 한 문장으로 정리하면 이렇습니다 —
listeners는 바인딩, advertised.listeners는 광고,
listener.security.protocol.map은 이름과 프로토콜의 매핑입니다.
리스너 이름을 왜 따로 두는가
리스너 이름이 보안 프로토콜 이름과 같다면 listener.security.protocol.map을 생략할 수 있습니다.
listeners=SSL://localhost:9092,PLAINTEXT://localhost:9093이 그런 예입니다.
하지만 공식 문서는 리스너의 용도가 명확해지므로 명시적인 이름을 권장합니다.
그리고 같은 보안 프로토콜을 두 개 이상의 포트에서 쓰려면 이름을 따로 붙이는 것이 필수입니다 —
내부 트래픽과 외부 트래픽을 분리하면서 양쪽 다 SSL을 쓰는 경우가 대표적입니다.
# 두 리스너가 같은 SSL 프로토콜을 쓰므로 이름을 따로 붙여야 합니다
listeners=INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:9094
listener.security.protocol.map=INTERNAL:SSL,EXTERNAL:SASL_SSL
# 바인딩 주소와 광고 주소가 다릅니다.
# 컨테이너 안에서는 0.0.0.0 에 바인딩하지만,
# 클라이언트에게는 실제로 도달 가능한 이름을 알려 줘야 합니다
advertised.listeners=INTERNAL://kafka-1.internal:9092,EXTERNAL://kafka.example.com:9094
# 브로커 간 복제는 내부 리스너로
inter.broker.listener.name=INTERNAL
KRaft에서의 리스너 규칙
KRaft에서는 process.roles에 broker가 있으면 브로커,
controller가 있으면 컨트롤러이고, 역할에 따라 리스너 구성이 달라집니다.
inter.broker.listener.name이 가리키는 리스너는 브로커 간 요청 전용입니다. 주 용도는 파티션 복제입니다. 지정하지 않으면security.inter.broker.protocol(기본PLAINTEXT)이 결정합니다. 두 설정을 동시에 지정하면 오류입니다.- 컨트롤러는
controller.listener.names로 정의된 별도 리스너를 써야 하며, 브로커 간 리스너와 같은 값일 수 없습니다. controller역할이 없는 순수 브로커도 컨트롤러 리스너를 정의해야 합니다. 컨트롤러에 요청을 보내야 하기 때문입니다. 단 그 리스너를listeners에 넣지는 않습니다 — 브로커는 컨트롤러 리스너를 노출하지 않습니다.- combined 모드(
broker,controller)에서는 컨트롤러 리스너를listeners에 포함해야 합니다. - 컨트롤러 리스너가 여럿이면 목록의 첫 번째가 아웃바운드 요청에 쓰입니다.
순수 브로커인데 컨트롤러 리스너를 listeners에 넣었습니다. 브로커는 컨트롤러 리스너를 노출하지 않습니다.
process.roles=broker
listeners=BROKER://localhost:9092,CONTROLLER://localhost:9093
inter.broker.listener.name=BROKER
controller.listener.names=CONTROLLER
컨트롤러 리스너는 정의만 하고 바인딩하지 않습니다. 포트는 controller.quorum.*가 결정합니다.
process.roles=broker
listeners=BROKER://localhost:9092
inter.broker.listener.name=BROKER
controller.quorum.bootstrap.servers=localhost:9093
controller.listener.names=CONTROLLER
listener.security.protocol.map=BROKER:SASL_SSL,CONTROLLER:SASL_SSL
ACL — 인가
ACL은 플러그인 가능한 인가 프레임워크이며 authorizer.class.name으로 지정합니다.
KRaft 클러스터에서는 모든 노드(브로커·컨트롤러·combined)에
org.apache.kafka.metadata.authorizer.StandardAuthorizer를 설정합니다.
이 구현은 ACL을 클러스터 메타데이터(KRaft 메타데이터 로그)에 저장합니다.
ACL의 일반 형식은 공식 문서 표현 그대로입니다 — "Principal {P}가 Host {H}에서 ResourcePattern {RP}에 매칭되는 Resource {R}에 대해 Operation {O}를 [허용|거부]한다".
super.users와 allow.everyone.if.no.acl.found가 그 판단을 어떻게 바꾸는지
| 오퍼레이션 | 주요 용도 |
|---|---|
Read | 토픽 소비, 그룹 참여 |
Write | 토픽 발행 |
Create | 토픽 생성 |
Delete | 토픽·그룹 삭제 |
Alter | 리소스 변경 (파티션 추가 등) |
Describe | 메타데이터 조회 |
ClusterAction | 클러스터 수준 내부 동작 (브로커 간 요청 등) |
DescribeConfigs | 설정 조회 |
AlterConfigs | 설정 변경 |
IdempotentWrite | 멱등 프로듀서 (2.8부터 deprecated) |
CreateTokens | 위임 토큰 생성 |
DescribeTokens | 위임 토큰 조회 |
All | 모든 오퍼레이션 |
| 리소스 | 무엇을 보호하는가 | 오류 |
|---|---|---|
Topic | 토픽 읽기·쓰기 등 | TOPIC_AUTHORIZATION_FAILED (29) |
Group | 컨슈머 그룹 참여 등 | GROUP_AUTHORIZATION_FAILED (30) |
Cluster | 클러스터 전체에 영향을 주는 동작 | CLUSTER_AUTHORIZATION_FAILED (31) |
TransactionalId | 트랜잭션 관련 동작 | TRANSACTIONAL_ID_AUTHORIZATION_FAILED (53) |
DelegationToken | 위임 토큰 | — |
User | 다른 사용자의 토큰 생성·조회 권한 | — |
# 특정 호스트에서만 Read/Write 허용
bin/kafka-acls.sh --bootstrap-server localhost:9092 --add \
--allow-principal User:Bob --allow-principal User:Alice \
--allow-host 198.51.100.0 --allow-host 198.51.100.1 \
--operation Read --operation Write --topic Test-topic
# 모두 허용하되 특정 사용자·호스트만 거부 (Deny 가 Allow 보다 우선합니다)
bin/kafka-acls.sh --bootstrap-server localhost:9092 --add \
--allow-principal User:'*' --allow-host '*' \
--deny-principal User:BadBob --deny-host 198.51.100.3 \
--operation Read --topic Test-topic
# 프로듀서 편의 옵션 — 필요한 오퍼레이션을 한 번에 부여합니다
bin/kafka-acls.sh --bootstrap-server localhost:9092 --add \
--allow-principal User:Bob --producer --topic Test-topic
# 컨슈머 편의 옵션 — 토픽 Read 와 그룹 Read 가 함께 필요합니다
bin/kafka-acls.sh --bootstrap-server localhost:9092 --add \
--allow-principal User:Bob --consumer --topic Test-topic --group Group-1
# 접두어 기반 (prefixed) 패턴 — 토픽이 늘어날 때 유용합니다
bin/kafka-acls.sh --bootstrap-server localhost:9092 --add \
--allow-principal User:Jane --producer \
--topic Test- --resource-pattern-type prefixed
# 목록 조회
bin/kafka-acls.sh --bootstrap-server localhost:9092 --list --topic Test-topic
모니터링
Kafka는 JMX로 메트릭을 노출합니다. 문제는 개수가 수백 개라는 점입니다. 아래는 공식 문서가 "반드시 감시해야 한다"고 정리한 서버 메트릭과 그 정상 범위입니다.
| 메트릭 | MBean | 정상 범위 | 벗어나면 |
|---|---|---|---|
| Under-replicated partitions | kafka.server:type=ReplicaManager, |
0 | ISR 크기가 전체 레플리카 수보다 작습니다. 브로커 장애, 느린 팔로워, 네트워크 문제를 의심합니다. 가장 먼저 봐야 할 메트릭입니다 |
| Under-min-ISR partitions | kafka.server:type=ReplicaManager, |
0 | ISR이 min.insync.replicas 미달입니다. acks=all 쓰기가 즉시 거부됩니다 |
| Offline partitions | kafka.controller:type=KafkaController, |
0 | 리더가 없는 파티션입니다. 읽기도 쓰기도 불가합니다. 즉시 대응 대상 |
| Active controller | kafka.controller:type=KafkaController, |
클러스터 전체에서 정확히 하나만 1 | 합계가 0이면 컨트롤러 없음, 2 이상이면 split-brain 의심입니다 |
| 리더 교체 조짐 | kafka.server:type=ReplicaManager, · IsrExpandsPerSeckafka.controller:type=ControllerStats, |
평시 0 | 평상시에 오르면 브로커가 반복적으로 이탈·복귀하고 있습니다.
ZooKeeper 시절의 LeaderElectionRateAndTimeMs는 KRaft에서 제거되었으므로 쓰지 마세요 |
| Unclean leader election | kafka.controller:type=ControllerStats, |
0 | 데이터 유실이 발생했다는 뜻입니다. ISR 밖 레플리카가 리더가 되었습니다 |
| ISR shrink rate | kafka.server:type=ReplicaManager, |
0 (브로커 장애 중 제외) | 공식 문서 표현: 브로커가 죽으면 일부 파티션의 ISR이 줄고, 복귀해 따라잡으면 다시 늘어납니다. 그 외에는 shrink·expand 모두 0이 기대값입니다 |
| Request handler idle | kafka.server:type=KafkaRequestHandlerPool, |
0~1 사이. 이상적으로 0.3 초과 | 0.3 아래로 떨어지면 I/O 스레드가 포화입니다. num.io.threads를 검토합니다 |
| Network processor idle | kafka.network:type=SocketServer, |
0~1 사이. 이상적으로 0.3 초과 | 네트워크 스레드 포화입니다. num.network.threads를 검토합니다 |
| Request total time | kafka.network:type=RequestMetrics, |
워크로드에 따라 다름 | 공식 문서에 따르면 queue · local · remote · response send 시간으로 분해됩니다. 어느 구간이 늘었는지가 원인을 가릅니다 |
| Offline log directory | kafka.log:type=LogManager, |
0 | 디스크 장애입니다 |
클라이언트 메트릭
| 메트릭 | MBean | 의미 |
|---|---|---|
records-lag-max |
kafka.consumer:type=consumer-fetch-manager-metrics,client-id={client-id} |
컨슈머가 프로듀서보다 뒤처진 레코드 수. 공식 문서가 명시하듯 브로커가 아니라 컨슈머가 발행하는 메트릭입니다 |
record-error-rate / record-retry-rate |
kafka.producer:type=producer-topic-metrics |
전송 실패·재시도율. 재시도율이 높으면 브로커 상태나 타임아웃 설정을 봅니다 |
request-latency-avg |
kafka.producer:type=producer-metrics |
프로듀서 요청 왕복 시간 |
buffer-available-bytes |
kafka.producer:type=producer-metrics |
남은 버퍼. 0에 가까워지면 send()가 max.block.ms만큼 블록됩니다 |
commit-latency-avg |
kafka.consumer:type=consumer-coordinator-metrics,client-id={client-id} |
커밋 요청 평균 시간 |
consumer lag을 측정하는 세 가지 방법
| 방법 | 어떻게 | 함정 |
|---|---|---|
① CLI (kafka-consumer-groups.sh) |
--describe --group {그룹}으로 CURRENT-OFFSET · LOG-END-OFFSET · LAG을 봅니다 |
커밋된 오프셋 기준입니다. 처리는 진행되는데 커밋 주기가 길면 lag이 실제보다 크게 보입니다. 반대로 처리 전에 커밋하는 코드면 실제보다 작게 보입니다. 또 스냅샷이라 추세를 못 봅니다 |
② 컨슈머 JMX (records-lag-max) |
컨슈머 프로세스가 직접 발행하는 메트릭을 수집합니다 | 컨슈머가 죽으면 메트릭도 사라집니다. "lag이 0"이 아니라 "컨슈머가 없음"인데 알림이 안 울립니다. 또 그 컨슈머가 할당받은 파티션만 반영합니다 |
| ③ AdminClient로 직접 계산 | listConsumerGroupOffsets()와 listOffsets(LATEST)의 차를 외부 수집기가 계산합니다 |
컨슈머가 죽어도 계산되므로 ①·②의 사각지대를 메웁니다. 다만 폴링 주기만큼 늦고, 조회 권한(Describe)이 필요합니다 |
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group my-group
# 출력 형식 (공식 문서 예시)
# TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
# topic3 0 241019 395308 154289 consumer2-… /… consumer2
# 그룹 상태 요약 — 멤버 수와 코디네이터
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group my-group --state
# 어느 멤버가 어느 파티션을 갖고 있는지
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group my-group --members --verbose
JMX에서 대시보드까지
Kafka는 JMX를 통해서만 메트릭을 노출하므로, 일반적인 파이프라인은 다음과 같습니다.
- 브로커에서 JMX 활성화 —
JMX_PORT환경 변수를 지정해 기동합니다. - JMX → Prometheus 변환 — JMX exporter를 자바 에이전트로 붙이거나 별도 프로세스로 띄워 HTTP 엔드포인트로 노출합니다. (exporter는 Apache Kafka 배포판에 포함되지 않는 외부 프로젝트입니다.)
- Prometheus 수집 — 브로커·컨트롤러·Connect 워커·Streams 애플리케이션을 각각 스크레이프합니다.
- Grafana 시각화 + 알림 — 아래 임계값을 룰로 등록합니다.
Kafka에는 클라이언트 메트릭을 브로커 쪽 플러그인으로 수집하는 경로(KIP-714)도 있고, 4.0부터 Kafka Streams 런타임 자체의 메트릭까지 이 경로로 수집할 수 있습니다(KIP-1076). 클라이언트에 exporter를 붙이기 어려운 환경에서 검토할 만합니다.
| 대상 | 경고 | 긴급 | 근거 |
|---|---|---|---|
OfflinePartitionsCount |
— | > 0 (즉시) | 공식 문서 정상값 0. 읽기·쓰기 불가 상태입니다 |
UnderMinIsrPartitionCount |
— | > 0 (즉시) | 공식 문서 정상값 0. acks=all 쓰기가 거부됩니다 |
UnderReplicatedPartitions |
> 0이 수 분 지속 | 계속 증가 | 공식 문서 정상값 0. 순간적인 값은 롤링 재시작 중 정상일 수 있습니다 |
ActiveControllerCount 합 |
— | ≠ 1 | 공식 문서: 클러스터에서 정확히 하나만 1 |
UncleanLeaderElectionsPerSec |
— | > 0 | 공식 문서 정상값 0. 유실 발생 신호입니다 |
RequestHandlerAvgIdlePercent |
< 0.3 | < 0.1 | 공식 문서: "이상적으로 0.3 초과" |
NetworkProcessorAvgIdlePercent |
< 0.3 | < 0.1 | 공식 문서: "이상적으로 0.3 초과" |
| consumer lag | 서비스별 SLO에서 역산. 절대값보다 "증가 추세"에 알림을 거는 편이 유용합니다 | 공식 문서에 권장 임계값이 없습니다 — 워크로드에 따라 다릅니다 | |
| 컨슈머 그룹 멤버 수 | 기대 인스턴스 수보다 적을 때 | lag 알림의 사각지대를 메웁니다 | |
성능 튜닝
처리량과 지연의 트레이드오프
| 설정 | 기본값 | 처리량에 유리 | 지연에 유리 |
|---|---|---|---|
linger.ms producer |
5 | 올림 (배치를 더 모음) | 0으로 내림 |
batch.size producer |
16384 | 올림 | 내림 (또는 그대로) |
compression.type producer |
none | lz4 / zstd로 네트워크·디스크 절약 |
none 또는 CPU가 싼 lz4 |
acks producer |
all | 1로 낮추면 빨라지지만 내구성을 잃습니다 |
같음 |
fetch.min.bytes consumer |
1 | 올림 (한 번에 더 많이 받음) | 1 유지 |
max.poll.records consumer |
500 | 올림 | 내림 (max.poll.interval.ms 초과 위험도 함께 줄어듭니다) |
num.io.threads broker |
8 | RequestHandlerAvgIdlePercent가 낮으면 올림 |
같음 |
num.network.threads broker |
3 | NetworkProcessorAvgIdlePercent가 낮으면 올림. 컨트롤러 리스너를 제외한 각 리스너가 자기 스레드 풀을 갖습니다 |
같음 |
num.replica.fetchers broker |
1 | 올림 — 복제 I/O 병렬성이 늘지만 CPU·메모리를 더 씁니다. 브로커당 총 fetcher 수는 이 값 × 브로커 수입니다 | — |
queued.max.requests broker |
500 | 큐를 늘려 버스트를 흡수 | 내리면 큐 대기 시간이 줄지만 네트워크 스레드가 더 자주 블록됩니다 |
socket.send.buffer.bytes / socket.receive.buffer.bytes broker |
102400 (100KiB) |
고지연·고대역 링크(DC 간)에서 올림. -1이면 OS 기본값 |
— |
num.recovery.threads.per.data.dir broker |
2 | — | 올리면 unclean shutdown 후 복구가 빨라지지만 I/O를 더 씁니다 (4.0에서 1→2로 변경) |
OS 레벨
공식 문서는 "OS 수준 튜닝이 많이 필요하지는 않지만 세 가지가 중요할 수 있다"고 말합니다.
| 항목 | 권고 | 근거 |
|---|---|---|
| 파일 디스크립터 한계 | 브로커 프로세스에 최소 100000개를 출발점으로 허용 | Kafka는 로그 세그먼트와 연결에 fd를 씁니다. 필요량은 대략 (파티션 수) × (파티션 크기 / 세그먼트 크기) + 연결 수입니다. mmap()은 close() 후에도 참조를 남긴다는 점도 문서가 지적합니다 |
| 최대 소켓 버퍼 크기 | DC 간 고성능 전송이 필요하면 올립니다 | 브로커의 socket.*.buffer.bytes가 OS 상한을 넘을 수 없습니다 |
vm.max_map_count |
파티션 수를 늘릴 때 반드시 함께 확인 | 여러 Linux에서 기본값이 대략 65535입니다. 로그 세그먼트 하나가 index·timeindex 두 개의 map area를 쓰므로, 파티션당 최소 2개입니다. 문서의 예시: 브로커에 파티션 50000개를 만들면 100000개의 map area가 필요해 기본값에서는 OutOfMemoryError (Map failed)로 브로커가 죽을 수 있습니다 |
JVM
공식 문서가 제시하는 인자 예시와 실제 운영 수치는 다음과 같습니다.
-Xmx6g -Xms6g -XX:MetaspaceSize=96m -XX:+UseG1GC
-XX:MaxGCPauseMillis=20 -XX:InitiatingHeapOccupancyPercent=35 -XX:G1HeapRegionSize=16M
-XX:MinMetaspaceFreeRatio=50 -XX:MaxMetaspaceFreeRatio=80 -XX:+ExplicitGCInvokesConcurrent
용량 산정 워크시트
정확한 사이징은 워크로드 측정으로만 가능하지만, 출발점을 잡는 데 쓰는 계산은 다음과 같습니다.
공식 문서가 제시하는 유일한 메모리 근거는 "활성 reader/writer를 버퍼링하기 위해
30초 분량, 즉 write_throughput × 30"이라는 대략적 추정입니다.
| 항목 | 계산 | 주의 |
|---|---|---|
| 디스크 용량 | 일 유입량 × 보관 일수 × RF × (1 − 압축률) + 여유 | 삭제는 세그먼트 단위로 일어나므로 실제 사용량은 계산값보다 큽니다(7장). 여유는 최소 30% 이상 두세요 |
| 파티션 수 (소비 기준) | 목표 처리량 ÷ 컨슈머 인스턴스 하나의 처리량 | 파티션 수가 컨슈머 병렬성의 상한입니다. 줄일 수 없습니다 |
| 파티션 수 (브로커 한계) | 브로커당 파티션 수 × 2 ≤ vm.max_map_count (세그먼트가 1개일 때의 최소) |
세그먼트가 여러 개면 훨씬 커집니다. 파일 디스크립터도 함께 확인하세요 |
| 브로커 수 | max(디스크 요구 ÷ 브로커당 디스크, RF, 목표 처리량 ÷ 브로커당 처리량) | RF보다 적은 브로커로는 그 RF를 만들 수 없습니다. 롤링 재시작 중에도 min.insync.replicas를 만족해야 하므로 여유가 필요합니다 |
| 컨트롤러 수 | 3 또는 5 | 쿼럼이므로 홀수여야 하고, 3이면 1대, 5면 2대 장애를 견딥니다. 상세는 3장 KRaft |
| 메모리 | JVM 힙(예시 6GB) + 나머지는 페이지 캐시로 남김 | 공식 문서의 대략적 추정: write_throughput × 30초 |
파티션 재할당
브로커를 추가·제거하거나 데이터 분포를 고칠 때 쓰는 도구가
kafka-reassign-partitions.sh입니다.
공식 문서는 이 도구가 데이터 분포를 자동으로 분석해 균형을 맞춰 주지는 않으며,
어떤 토픽·파티션을 옮길지는 관리자가 판단해야 한다고 명시합니다.
| 모드 | 무엇을 하는가 |
|---|---|
--generate | 토픽 목록과 대상 브로커 목록을 받아 후보 재할당 계획을 만듭니다. 실제로 옮기지는 않습니다 |
--execute | --reassignment-json-file의 계획대로 재할당을 시작합니다 |
--verify | 직전 --execute의 진행 상태를 확인합니다 (완료 / 실패 / 진행 중). 완료됐으면 throttle을 해제합니다 |
--cancel | 진행 중인 재할당을 취소합니다 |
--list | 현재 활성 재할당을 모두 나열합니다 |
# ── 0. 옮길 토픽 목록을 JSON 으로 ─────────────────────────
cat > topics-to-move.json <<'JSON'
{
"topics": [
{ "topic": "foo1" },
{ "topic": "foo2" }
],
"version": 1
}
JSON
# ── 1. --generate : 후보 계획 생성 ────────────────────────
# 출력의 "Current partition replica assignment" 를 반드시 파일로 저장하세요.
# 롤백할 때 이것이 --reassignment-json-file 이 됩니다.
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--topics-to-move-json-file topics-to-move.json \
--broker-list "5,6" --generate
# 제안된 계획을 저장
# → expand-cluster-reassignment.json
# ── 2. --execute : throttle 과 함께 시작 ──────────────────
# throttle 단위는 B/s 입니다. 운영 트래픽을 지키려면 반드시 지정하세요.
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--reassignment-json-file expand-cluster-reassignment.json --execute \
--throttle 50000000
# 진행 중 throttle 을 올리려면 --additional 로 같은 파일을 다시 실행합니다
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--additional --execute \
--reassignment-json-file expand-cluster-reassignment.json \
--throttle 700000000
# ── 3. --verify : 완료 확인 + throttle 해제 ───────────────
# --execute 에 쓴 것과 같은 JSON 파일을 써야 합니다.
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--reassignment-json-file expand-cluster-reassignment.json --verify
브로커를 완전히 제거하려면 재할당으로 그 브로커의 모든 레플리카를 옮긴 뒤 브로커를 정지하고 클러스터에서 등록을 해제합니다.
bin/kafka-cluster.sh unregister --bootstrap-server localhost:9092 --id 1
Share Groups — Queues for Kafka (KIP-932)
컨슈머 그룹의 파티션 배정 모델은 순서 보장과 확장성을 함께 주는 강력한 조합이지만, 공식 문서의 표현대로 "전통적인 메시징 워크로드 일부는 컨슈머 그룹에 잘 맞지 않습니다." share group은 컨슈머 그룹과 나란히 존재하는 또 다른 그룹 종류이며, 파티션과 레코드를 더 세밀하게 나눠 쓰고 싶을 때의 대안입니다. share group의 컨슈머는 share consumer라고 부르고 별도의 프로그래밍 인터페이스를 씁니다.
컨슈머 그룹과 무엇이 다른가
| 컨슈머 그룹 | share group | |
|---|---|---|
| 파티션 소유권 | 한 파티션은 그룹 안에서 정확히 한 컨슈머가 소비 | 한 파티션을 여러 컨슈머가 함께 소비합니다 (배타적 배정이 없습니다) |
| 컨슈머 수 상한 | 파티션 수를 넘으면 유휴 | 파티션 수를 넘어설 수 있습니다 |
| 진행 위치 단위 | 파티션당 오프셋 하나 | 레코드 단위 개별 확인(ack). 다만 효율을 위해 배치 처리에 최적화되어 있습니다 |
| 전달 시도 추적 | 없음 (재처리는 오프셋 되돌리기) | 전달 시도 횟수를 집계합니다 → 처리 불가 레코드의 자동 처리가 가능합니다 |
| 순서 보장 | 파티션 단위 순서 보장 | 여러 컨슈머가 같은 파티션을 나눠 처리하므로 파티션 단위 순서를 기대할 수 없습니다 |
| 적합한 워크로드 | 순서 있는 스트림 처리, 상태 있는 집계 | 레코드를 하나씩 처리하는 작업 큐. 처리 시간 편차가 크거나, 컨슈머를 파티션 수와 무관하게 늘리고 싶을 때 |
획득 락과 ack 타입
share consumer가 레코드를 fetch하면, 그 레코드는 시간 제한이 있는 획득 락(acquisition lock)과 함께 이 컨슈머에게 전달됩니다. 락이 걸려 있는 동안 그 레코드는 같은 share group의 다른 컨슈머에게 보이지 않습니다. 락은 기간이 지나면 자동으로 풀리고 레코드는 다시 전달 대상이 됩니다. 브로커는 이렇게 컨슈머가 죽어도 전달이 계속 진행되도록 보장합니다.
| 행동 | ack 타입 | 결과 |
|---|---|---|
| 처리 성공을 확인 | ACCEPT | 레코드가 처리 완료로 기록됩니다 |
| 해제(release) | RELEASE | 처리하지 못했으므로 다시 전달 대상으로 돌립니다 |
| 거부(reject) | REJECT | 처리 불가 레코드로 표시하고 더 이상 전달하지 않습니다 |
| 갱신(renew) | RENEW | 아직 처리 중이므로 획득 락 기간을 연장합니다 |
| 아무것도 하지 않음 | — | 락 기간이 지나면 자동 해제되어 다시 전달됩니다 |
| 설정 | 기본값 | 설명 |
|---|---|---|
share.record.lock.duration.ms |
30000 (30초) | 획득 락 유지 기간. 최소 1000 |
share.delivery.count.limit |
5 | 한 레코드의 최대 전달 시도 횟수. 최소 2. 이 횟수를 넘으면 더 전달하지 않습니다 |
share.partition.max.record.locks |
2000 | share-partition당 락 상한. 도달하면 락이 자연히 만료될 때까지 fetch가 일시적으로 아무 레코드도 반환하지 않습니다. 브로커 쪽 대응 설정은 group.share.partition.max.record.locks(기본 2000, 유효 범위 100~10000)입니다 |
share.auto.offset.reset |
latest |
share-partition 시작 오프셋 초기화 전략. earliest / latest / by_duration:PnDTnHnMn.nS (ISO 8601) |
share.isolation.level |
read_uncommitted |
read_committed면 커밋된 트랜잭션 레코드만 전달합니다 |
share.session.timeout.ms |
45000 (45초) | share group 프로토콜에서의 클라이언트 장애 감지 시간 |
share.heartbeat.interval.ms |
5000 (5초) | 멤버에게 주어지는 하트비트 간격 |
share.renew.acknowledge.enable |
true |
RENEW ack 타입 사용 가능 여부 |
group.share.max.size broker |
200 | 하나의 share group이 수용할 수 있는 최대 멤버 수 |
# feature 활성화 (4.2 이상으로 부트스트랩한 클러스터는 이미 켜져 있습니다)
bin/kafka-features.sh --bootstrap-server localhost:9092 \
upgrade --feature share.version=1
# 모든 그룹 종류를 한 번에 보기 (Consumer / Share / Streams)
bin/kafka-groups.sh --bootstrap-server localhost:9092 --list
# GROUP TYPE PROTOCOL
# my-consumer-group Consumer consumer
# my-share-group Share share
# share group 목록
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --list
# 진행 상황 — START-OFFSET 은 "전달 평가 대상 in-flight 레코드 중 가장 이른 오프셋"입니다.
# 그 뒤 레코드 중 일부는 이미 전달이 끝났을 수 있습니다.
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 \
--describe --group my-share-group
# 활성 멤버 / 상태 요약
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 \
--describe --group my-share-group --members
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 \
--describe --group my-share-group --state
# 오프셋 리셋 (활성 멤버가 없어야 합니다)
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 \
--reset-offsets --group my-share-group --topic topic1 --to-latest --execute
# 그룹 삭제 — 활성 멤버가 없는 share group 만 삭제됩니다
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 \
--delete --group my-share-group
# 그룹별 설정 override
bin/kafka-configs.sh --bootstrap-server localhost:9092 --alter \
--entity-type groups --entity-name my-share-group \
--add-config share.record.lock.duration.ms=60000
제약 사항
- 파티션 단위 순서 보장이 없습니다. 여러 컨슈머가 같은 파티션을 나눠 처리하므로 구조적으로 불가능합니다.
- 레코드당 상태를 브로커가 관리하므로, share-partition당 in-flight 레코드 수에 상한이 있습니다(
share.partition.max.record.locks기본 2000). 상한에 닿으면 fetch가 일시적으로 빈 결과를 반환합니다. - 전달 시도 횟수 상한이 있습니다(
share.delivery.count.limit기본 5). 넘으면 그 레코드는 더 전달되지 않습니다. - 내부 토픽이 하나 더 생깁니다(
__share_group_state). 브로커 3대 미만이면 사전 설정이 필요합니다. - Kafka Streams와 함께 쓰는 경로가 아닙니다. Streams는 컨슈머 그룹 또는 streams group을 씁니다(10장).
- 4.2.0에는 KIP-932 share group 경로에 심각한 데드락 버그(KAFKA-20505)가 있었고 4.2.1에서 수정되었습니다.
시험 포인트 정리
확인 문제
네 가지 문항 유형(단일 선택 · 복수 선택 · 연결형 · 순서 배열)이 섞여 있습니다. 리스너 3종의 역할 구분과 메트릭 정상 범위는 반드시 맞히고 넘어가세요.
이어서 볼 곳
공식 문서 출처
이 장의 설정 기본값·메트릭 정상 범위·CLI 옵션은 모두 아래에서 확인했습니다 (Apache Kafka 4.3 문서 기준).
- Security Overview — 지원 보안 조치 4가지, SASL 메커니즘별 도입 버전
- Listener Configuration —
listeners형식,listener.security.protocol.map, 보안 프로토콜 4종, KRaft의 컨트롤러 리스너 규칙과 예시 설정 - Authentication using SASL —
sasl.enabled.mechanisms·sasl.mechanism.inter.broker.protocol, SCRAM 사용자 생성 명령, JAAS 설정 - Authorization and ACLs —
StandardAuthorizer, ACL 일반 형식, 오퍼레이션 13종·리소스 6종과 오류 코드,super.users·allow.everyone.if.no.acl.found, KRaft principal forwarding,kafka-acls.sh예시 - Monitoring — 서버 필수 메트릭과 정상 범위(URP 0, ActiveControllerCount 1, idle percent > 0.3 등),
records-lag-max가 컨슈머 발행임 - Hardware and OS — 파일 디스크립터 100000,
vm.max_map_count65535와 map area 계산, XFS vs EXT4 비교,noatime, 플러시 정책 권고 - Java Version — Java 17·21·25 완전 지원, JVM 인자 예시, LinkedIn 클러스터 수치
- Expanding your cluster / Partition Reassignment —
--generate/--execute/--verify, 계획 JSON 형식, RF 유지, 롤백 계획 저장 - Limiting Bandwidth Usage during Data Migration —
--throttle,--additional, throttle 해제 경고,leader/follower.replication.throttled.* - Managing Groups —
kafka-groups.sh,kafka-consumer-groups.sh --describe출력 형식,kafka-share-groups.sh전체 옵션, 내부 토픽 파티션 수 변경 금지 경고 - The Share Consumer — share group과 consumer group의 근본 차이 4가지, 획득 락 기본 30초, 5가지 처리 행동,
group.share.partition.max.record.locks - Group Configs —
share.record.lock.duration.ms30000,share.delivery.count.limit5,share.partition.max.record.locks2000,share.auto.offset.reset,share.isolation.level,share.renew.acknowledge.enable - Broker Configs —
num.io.threads8,num.network.threads3,num.replica.fetchers1,queued.max.requests500,socket.*.buffer.bytes102400,num.recovery.threads.per.data.dir2,replica.lag.time.max.ms30000,ssl.client.authnone - Upgrading Apache Kafka — 4.2 Share Groups production-ready,
__share_group_state와 3브로커 요건, KIP-1240 그룹 설정 추가, 4.2.1의 KAFKA-20505 수정, 4.0의linger.ms0→5 - AcknowledgeType (Apache Kafka 4.3 소스) —
ACCEPT(1) ·RELEASE(2) ·REJECT(3) ·RENEW(4) - ReassignPartitionsCommandOptions (Apache Kafka 4.3 소스) —
--cancel·--list·--preserve-throttles등 실제 옵션 목록