배포판마다 기본 포트가 다르다. 아파치 기본값과 Cloudera 는 9092, 옛 Hortonworks 배포본은 6667 을 쓴다. 보안 리스너를 쓰면 9093 (SASL_SSL) 처럼 따로 둔다. 실제 값은 브로커 설정에서 확인한다.
grep -E '^(listeners|advertised.listeners)' /opt/kafka/config/server.properties
ss -tlnp | grep -E '9092|9093|6667'
Kafka 4.x 는 KRaft 전용이며 모든 CLI 에서 --zookeeper 옵션이 사라졌다. --bootstrap-server 만 쓴다.
BS=broker01.example.com:9092
kafka-topics.sh --bootstrap-server $BS --list
kafka-topics.sh --bootstrap-server $BS --describe --topic order-events
kafka-topics.sh --bootstrap-server $BS --create \
--topic order-events --partitions 6 --replication-factor 3 \
--config retention.ms=604800000 --config cleanup.policy=delete
kafka-topics.sh --bootstrap-server $BS --alter --topic order-events --partitions 12
kafka-topics.sh --bootstrap-server $BS --delete --topic order-events
파티션은 늘릴 수만 있고 줄일 수 없다. 키 기반 파티셔닝을 쓰는 토픽은 파티션을 늘리는 순간 같은 키가 다른 파티션으로 가게 되어 순서 보장이 깨지므로, 늘리기 전에 소비 측 영향을 따져야 한다.
복제 계수는 --alter 로 바꿀 수 없다. 재배치 도구를 쓴다.
설정 변경은 kafka-configs.sh 로 한다.
kafka-configs.sh --bootstrap-server $BS --entity-type topics --entity-name order-events --describe
kafka-configs.sh --bootstrap-server $BS --entity-type topics --entity-name order-events \
--alter --add-config retention.ms=86400000
kafka-console-producer.sh --bootstrap-server $BS --topic order-events
kafka-console-producer.sh --bootstrap-server $BS --topic order-events \
--property parse.key=true --property key.separator=:
kafka-console-consumer.sh --bootstrap-server $BS --topic order-events --from-beginning
kafka-console-consumer.sh --bootstrap-server $BS --topic order-events \
--property print.key=true --property print.timestamp=true --max-messages 10
인증이 있는 클러스터에서는 클라이언트 설정 파일을 준다.
kafka-console-consumer.sh --bootstrap-server $BS --topic order-events \
--consumer.config /opt/kafka/config/client.properties
--from-beginning 없이 실행하면 실행 이후 들어오는 메시지만 본다. 데이터가 흐르는지 확인할 때는 이 차이를 의식해야 한다.
그룹 ID 는 컨슈머가 스스로 정해 보내는 값이다. 브로커나 토픽에 붙어 있는 값이 아니며, 토픽 하나를 여러 그룹이 각자의 진행 위치로 읽을 수 있다. 같은 그룹 안에서는 파티션이 나뉘어 배정되고, 그룹이 다르면 같은 메시지를 각각 받는다.
kafka-consumer-groups.sh --bootstrap-server $BS --list
kafka-consumer-groups.sh --bootstrap-server $BS --describe --group nifi-order-consumer
kafka-consumer-groups.sh --bootstrap-server $BS --describe --all-groups --state
--describe 출력의 열은 다음과 같다.
| 열 | 의미 |
|---|---|
| CURRENT-OFFSET | 그룹이 커밋한 위치 |
| LOG-END-OFFSET | 파티션의 마지막 위치 |
| LAG | 둘의 차이. 남은 건수 |
| CONSUMER-ID · HOST | 그 파티션을 맡은 컨슈머 |
CONSUMER-ID 가 비어 있으면 그 그룹에 동작 중인 컨슈머가 없는 것이다. 중단된 컨슈머의 남은 lag 을 볼 때 이 상태가 된다.
진행 위치를 옮기려면 그룹이 비활성 이어야 한다. 컨슈머를 모두 멈춘 뒤 실행한다.
kafka-consumer-groups.sh --bootstrap-server $BS --group nifi-order-consumer \
--topic order-events --reset-offsets --to-earliest --dry-run
kafka-consumer-groups.sh --bootstrap-server $BS --group nifi-order-consumer \
--topic order-events:3 --reset-offsets --to-offset 40920576 --execute
--dry-run 으로 먼저 확인하는 습관을 들인다. --to-earliest · --to-latest · --to-offset · --to-datetime · --shift-by 를 쓸 수 있다.
kafka-get-offsets.sh --bootstrap-server $BS --topic order-events
kafka-run-class.sh kafka.tools.GetOffsetShell --bootstrap-server $BS --topic order-events --time -1
--time -1 은 마지막 오프셋, -2 는 처음 오프셋이다. 두 번 찍어 차이를 보면 유입량을 어림할 수 있다.
컨슈머와 프로듀서의 사용량 자체는 CLI 로 직접 볼 수 없다. JMX 지표를 수집해 확인한다. kafka_exporter 로 컨슈머 그룹 lag 을, JMX Exporter 로 브로커의 초당 유입 바이트·메시지 수를 본다.
컨슈머 그룹의 커밋 오프셋을 저장하는 내부 토픽이다. 예전에는 ZooKeeper 에 저장하던 것을 토픽으로 옮긴 것으로, 기본 파티션 50 개에 cleanup.policy=compact 로 동작한다. 지우거나 직접 손대지 않는다. 브로커를 줄일 때는 이 토픽의 복제본 배치도 함께 옮겨야 한다.
NiFi 의 PublishKafka 로는 보내지는데 ConsumeKafka 로 받히지 않는 경우, 가장 흔한 원인은 두 가지다.
첫째, 새 컨슈머 그룹은 커밋된 오프셋이 없어 auto.offset.reset 을 따른다. 기본값 latest 면 그룹 등록 이후 들어오는 메시지만 받는다. 과거 메시지부터 받으려면 earliest 로 둔다.
둘째, Honor Transactions 설정이다. 이 값이 true 이면 컨슈머가 isolation.level=read_committed 로 동작해 커밋된 트랜잭션 메시지만 읽는다. 프로듀서가 트랜잭션을 쓰지 않았다면 그 메시지들은 그냥 일반 메시지이므로 읽히는 것이 정상이다. 다만 중간에 끝나지 않은 트랜잭션(LSO 를 붙잡는 미완료 트랜잭션)이 있으면 그 지점 이후가 보이지 않는다. 진단 삼아 false 로 바꿔 보아 갑자기 데이터가 들어오면 미완료 트랜잭션이 원인이다.
콘솔 컨슈머로 같은 조건을 재현해 보면 NiFi 문제인지 Kafka 문제인지 금방 가른다.
kafka-console-consumer.sh --bootstrap-server $BS --topic order-events \
--from-beginning --isolation-level read_committed --max-messages 5