Kafka ์ํคํ
์ฒ, KRaft & Docker Compose ํ๊ฒฝ ๊ตฌ์ถ, Topic ๊ด๋ฆฌ, Console Producer/Consumer, Consumer Group & Offset ๋ฆฌ์
, ํต์ฌ ํ๋ผ๋ฏธํฐ ํ๋, Spring Kafka ์ฐ๋ ๋ฐ ์ฅ์ ํธ๋ฌ๋ธ์ํ
์๋ฒฝ ๊ฐ์ด๋
## 1. Kafka ํต์ฌ ์ํคํ
์ฒ & ๋์ ์๋ฆฌ (Architecture)
```mermaid
graph LR
subgraph Producers
P1[Producer App 1]
P2[Producer App 2]
end
subgraph Kafka Cluster [Kafka Cluster / KRaft Controller]
subgraph Topic: orders [Topic: orders]
P_0[Partition 0<br/>Leader: Broker 1<br/>ISR: 1, 2]
P_1[Partition 1<br/>Leader: Broker 2<br/>ISR: 2, 3]
P_2[Partition 2<br/>Leader: Broker 3<br/>ISR: 3, 1]
end
end
subgraph Consumer Group: order-service [Consumer Group: order-service]
C1[Consumer 1<br/>reads P0]
C2[Consumer 2<br/>reads P1]
C3[Consumer 3<br/>reads P2]
end
P1 -->|Key=101 -> P0| P_0
P2 -->|Key=102 -> P1| P_1
P1 -->|Key=103 -> P2| P_2
P_0 --> C1
P_1 --> C2
P_2 --> C3
```
```text
โ ํต์ฌ ์ฉ์ด ์ ๋ฆฌ:
โข Topic: ๋ฉ์์ง๋ฅผ ๊ตฌ๋ถํ๋ ๋
ผ๋ฆฌ์ ์ฑ๋ (RDBMS์ ํ
์ด๋ธ ๊ฐ๋
)
โข Partition: ํ ํฝ์ ๋ถํ ํ์ฌ ๋ณ๋ ฌ ์ฒ๋ฆฌ ๋ฐ ์ํ ํ์ฅ์ ๊ฐ๋ฅํ๊ฒ ํ๋ ๋ฌผ๋ฆฌ์ ํ ๋จ์ (์์ ๋ณด์ฅ ๋จ์)
โข Offset: ํํฐ์
๋ด์์ ๊ฐ ๋ฉ์์ง๊ฐ ๋ถ์ฌ๋ฐ๋ ๊ณ ์ ํ ์์ฐจ์ ๋ฒํธ (0๋ถํฐ ๋จ์กฐ ์ฆ๊ฐ, ๋ถ๋ณ)
โข Broker: Kafka ์๋ฒ ์ธ์คํด์ค (๋ฉ์์ง ์์ , ๋์คํฌ ๊ธฐ๋ก, ์คํ์
๊ด๋ฆฌ, ์ฅ์ ๋ณต๊ตฌ ๋ด๋น)
โข Replication Factor: ํํฐ์
์ ๋ณต์ ๋ณธ ์ (Leader 1๊ฐ + Follower N-1๊ฐ, ๊ณ ๊ฐ์ฉ์ฑ ๋ณด์ฅ)
โข ISR (In-Sync Replicas): Leader์ ๋๊ธฐํ ์ํ๋ฅผ ์ ์งํ๊ณ ์๋ ๋ณต์ ๋ณธ ๊ทธ๋ฃน
โข Consumer Group: ๋์ผํ ํ ํฝ์ ๋ถ์ฐ/๋ณ๋ ฌ ์ฒ๋ฆฌํ๊ธฐ ์ํด ํ๋ ฅํ๋ ์ปจ์๋จธ ์ธ์คํด์ค๋ค์ ์งํฉ
โข KRaft (Kafka Raft): Kafka 3.x+๋ถํฐ ZooKeeper ์์ด ์์ฒด Raft ์๊ณ ๋ฆฌ์ฆ์ผ๋ก ๋ฉํ๋ฐ์ดํฐ๋ฅผ ๊ด๋ฆฌํ๋ ์ฐจ์ธ๋ ์ํคํ
์ฒ
```
---
## 2. Docker & Docker Compose ๋น ๋ฅธ ์คํ (KRaft Mode)
```bash
# ๐ณ [1] Docker ๋จ์ผ ๋ธ๋ก์ปค ์ฆ์ ์คํ (ZooKeeper ์๋ ์ต์ KRaft ๋ชจ๋)
docker run -d \
--name kafka-standalone \
-p 9092:9092 \
-e KAFKA_NODE_ID=1 \
-e KAFKA_PROCESS_ROLES=broker,controller \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
-e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT \
-e KAFKA_CONTROLLER_QUORUM_VOTERS=1@localhost:9093 \
-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
-e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 \
-e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 \
-e KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS=0 \
-e KAFKA_NUM_PARTITIONS=3 \
apache/kafka:latest
# ์ปจํ
์ด๋ ์ ์ ์
docker exec -it kafka-standalone /bin/bash
```
```yaml
# ๐ณ [2] docker-compose.yml (Kafka KRaft ๋ธ๋ก์ปค + Kafka UI ์น ๋์๋ณด๋)
version: '3.8'
services:
kafka:
image: apache/kafka:latest
container_name: kafka-kraft
ports:
- "9092:9092"
environment:
# ๋
ธ๋ ์ญํ ๋ฐ ID (๋ธ๋ก์ปค + ์ปจํธ๋กค๋ฌ ํตํฉ ๋ชจ๋)
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
# ๋ฆฌ์ค๋ ์ค์ (์ธ๋ถ ์ ์: 9092, ๋ด๋ถ ์ปจํธ๋กค๋ฌ: 9093)
KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
# ๊ฐ๋ฐ์ฉ ๋จ์ผ ๋
ธ๋ ์ค์
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_LOG_DIRS: /tmp/kraft-combined-logs
volumes:
- kafka_data:/tmp/kraft-combined-logs
# Kafka UI: ์น ๋ธ๋ผ์ฐ์ (http://localhost:8080)๋ก ํ ํฝ, ๋ฉ์์ง, ์ปจ์๋จธ ๊ทธ๋ฃน ํ์ธ
kafka-ui:
image: provectuslabs/kafka-ui:latest
container_name: kafka-ui
ports:
- "8080:8080"
environment:
KAFKA_CLUSTERS_0_NAME: local-cluster
KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9092
depends_on:
- kafka
volumes:
kafka_data:
```
---
## 3. ํ ํฝ ๊ด๋ฆฌ CLI (`kafka-topics.sh`)
```bash
# โ ๏ธ Kafka ๊ธฐ๋ณธ CLI ์คํ ์์น: $KAFKA_HOME/bin/ ๋๋ docker exec -it kafka-kraft ...
# [1] ํ ํฝ ์์ฑ (ํํฐ์
3๊ฐ, ๋ณต์ ๊ณ์ 1๊ฐ)
kafka-topics.sh --bootstrap-server localhost:9092 \
--create \
--topic payment-events \
--partitions 3 \
--replication-factor 1
# [2] ํ ํฝ ์์ฑ ์ ๋ณด์กด ์ ์ฑ
(Retention) ๋์ ์ค์ (์: 7์ผ = 604800000ms)
kafka-topics.sh --bootstrap-server localhost:9092 \
--create \
--topic order-events \
--partitions 3 \
--replication-factor 1 \
--config retention.ms=604800000 \
--config segment.bytes=1073741824
# [3] ์ ์ฒด ํ ํฝ ๋ชฉ๋ก ํ์ธ
kafka-topics.sh --bootstrap-server localhost:9092 --list
# [4] ํน์ ํ ํฝ ์์ธ ์ ๋ณด ์กฐํ (ํํฐ์
๋ณ Leader, Replicas, ISR ํ์ธ)
kafka-topics.sh --bootstrap-server localhost:9092 \
--describe \
--topic payment-events
# [5] ํํฐ์
๊ฐ์ ์ฆ๊ฐ (โ ๏ธ ํํฐ์
์๋ ๋๋ฆด ์๋ง ์๊ณ ์ค์ผ ์ ์์!)
kafka-topics.sh --bootstrap-server localhost:9092 \
--alter \
--topic payment-events \
--partitions 5
# [6] ํ ํฝ ์ญ์ (delete.topic.enable=true ํ์)
kafka-topics.sh --bootstrap-server localhost:9092 \
--delete \
--topic payment-events
```
---
## 4. ํ ํฝ ๋์ ์ค์ ๋ณ๊ฒฝ (`kafka-configs.sh`)
```bash
# [1] ํ ํฝ์ ํ์ฌ ์ค์ ๊ฐ ์ ์ฒด ์กฐํ
kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics \
--entity-name order-events \
--describe
# [2] ๋ฉ์์ง ๋ณด์กด ์๊ฐ(retention.ms) ๋์ ๋ณ๊ฒฝ (์: 24์๊ฐ = 86400000ms)
kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics \
--entity-name order-events \
--alter \
--add-config retention.ms=86400000
# [3] ์ต๋ ๋ฉ์์ง ํฌ๊ธฐ ๋ณ๊ฒฝ (๊ธฐ๋ณธ 1MB -> 10MB๋ก ํ์ฅ)
kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics \
--entity-name order-events \
--alter \
--add-config max.message.bytes=10485760
# [4] ์ ๋ฆฌ ์ ์ฑ
(Cleanup Policy) ๋ณ๊ฒฝ: ์ญ์ (delete) vs ์์ถ(compact)
# compact: ์ต์ Key ๊ฐ๋ง ์ ์งํ๋ ๋ก๊ทธ ์์ถ ๋ฐฉ์
kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics \
--entity-name user-profile \
--alter \
--add-config cleanup.policy=compact
# [5] ๋์ ์ผ๋ก ์ถ๊ฐํ๋ ํ ํฝ ์ค์ ์ญ์ (๊ธฐ๋ณธ๊ฐ ๋ณต์)
kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics \
--entity-name order-events \
--alter \
--delete-config retention.ms
```
---
## 5. ๋ฉ์์ง ๋ฐํ: ์ฝ์ ํ๋ก๋์ (`kafka-console-producer.sh`)
```bash
# [1] ๊ธฐ๋ณธ ๋ฉ์์ง ๋ฐํ (ํ ์ค ์
๋ ฅํ ๋๋ง๋ค Value๋ก ๋ฐํ, Ctrl+C๋ก ์ข
๋ฃ)
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic order-events
# [2] Key-Value ํํ ๋ฉ์์ง ๋ฐํ (๋์ผ Key๋ ๋์ผ Partition์ผ๋ก ๋ผ์ฐํ
๋จ)
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic order-events \
--property "parse.key=true" \
--property "key.separator=:"
# ์
๋ ฅ ์์:
# order-101:{"itemId": "item-A", "price": 15000}
# order-102:{"itemId": "item-B", "price": 28000}
# [3] ํ์ผ ๋ด์ฉ ๋๋ ํ์ค ์
๋ ฅ์ ํตํ ๋๋ ๋ฉ์์ง ์ ์ก
cat sample_events.json | kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic order-events
# [4] Producer ACKs ์ต์
์ ๋ถ์ฌํ์ฌ ์์ ํ๊ฒ ๋ฐํ
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic order-events \
--request-required-acks all
```
---
## 6. ๋ฉ์์ง ์๋น: ์ฝ์ ์ปจ์๋จธ (`kafka-console-consumer.sh`)
```bash
# [1] ์ค์๊ฐ ๋ฉ์์ง ๊ตฌ๋
(๋ช
๋ น์ด ์คํ ์ดํ ์ธ์
๋๋ ๋ฉ์์ง๋ง ์๋น)
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic order-events
# [2] ์ฒ์ ์คํ์
(0๋ฒ)๋ถํฐ ๊ณผ๊ฑฐ ๋ชจ๋ ๋ฉ์์ง ์ฝ๊ธฐ (--from-beginning)
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic order-events \
--from-beginning
# [3] ๋ฉ์์ง Key, Timestamp, Partition, Offset์ ํจ๊ป ์ถ๋ ฅ (๋๋ฒ๊น
ํ์)
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic order-events \
--from-beginning \
--property print.key=true \
--property print.timestamp=true \
--property print.partition=true \
--property print.offset=true \
--property key.separator=" | "
# [4] ํน์ ์ปจ์๋จธ ๊ทธ๋ฃน(Consumer Group)์ ์ง์ ํ์ฌ ๋ฉ์์ง ์๋น (์คํ์
์ปค๋ฐ๋จ)
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic order-events \
--group order-worker-group
# [5] ํน์ ํํฐ์
์ ์ง์ ์คํ์
๋ถํฐ ์ฝ๊ธฐ
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic order-events \
--partition 0 \
--offset 100 \
--max-messages 50
```
---
## 7. ์ปจ์๋จธ ๊ทธ๋ฃน ๋ฐ ์คํ์
๊ด๋ฆฌ (`kafka-consumer-groups.sh`)
```bash
# [1] ๋ฑ๋ก๋ ๋ชจ๋ ์ปจ์๋จธ ๊ทธ๋ฃน ๋ชฉ๋ก ํ์ธ
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
# [2] ํน์ ์ปจ์๋จธ ๊ทธ๋ฃน์ ์๋น ์งํ ์ํ ๋ฐ Lag(์ง์ฐ) ์ ๊ฒ
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe \
--group order-worker-group
# ์ถ๋ ฅ ๊ฒฐ๊ณผ ์ปฌ๋ผ ํด์:
# TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST
# order-events 0 1452 1460 8 consumer-1... /127.0.0.1
# -> LAG: Log-End-Offset - Current-Offset = ์ปจ์๋จธ๊ฐ ์์ง ์ฒ๋ฆฌํ์ง ๋ชปํ๊ณ ์์ธ ๋ฉ์์ง ์!
# โ ๏ธ [3] ์คํ์
๋ฆฌ์
(๋ฐ๋์ ํด๋น ์ปจ์๋จธ ๊ทธ๋ฃน ํ๋ก์ธ์ค๊ฐ ์ค์ง๋ ์ํ์ฌ์ผ ํจ)
# 3-1. ๊ฐ์ฅ ์ฒ์ ์คํ์
์ผ๋ก ๋ฆฌ์
(--to-earliest)
# --dry-run: ์ค์ ๋ฐ์ ์ ์๋ฎฌ๋ ์ด์
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-worker-group \
--topic order-events \
--reset-offsets \
--to-earliest \
--dry-run
# 3-2. ์ค์ ์คํ์
๋ฆฌ์
์คํ (--execute)
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-worker-group \
--topic order-events \
--reset-offsets \
--to-earliest \
--execute
# 3-3. ์ต์ ์คํ์
์ผ๋ก ๋ฆฌ์
(๊ณผ๊ฑฐ ๋ฉ์์ง ๊ฑด๋๋ฐ๊ธฐ)
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-worker-group \
--topic order-events \
--reset-offsets \
--to-latest \
--execute
# 3-4. ํน์ ์ผ์ ๊ธฐ์ค์ผ๋ก ๋ฆฌ์
(--to-datetime)
# ํ์: YYYY-MM-DDTHH:mm:ss.sss
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-worker-group \
--topic order-events \
--reset-offsets \
--to-datetime 2026-09-30T00:00:00.000 \
--execute
# 3-5. ํ์ฌ ์คํ์
์์ N๊ฑด ๋ค๋ก ๋๋๋ฆฌ๊ธฐ (--shift-by)
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-worker-group \
--topic order-events \
--reset-offsets \
--shift-by -100 \
--execute
```
---
## 8. ํ๋ก๋์(Producer) ํต์ฌ ํ๋ ํ๋ผ๋ฏธํฐ
```properties
# 1. ๋ฉ์์ง ์ ์ค ๋ฐฉ์ง (์ ๋ขฐ์ฑ vs ์ฒ๋ฆฌ๋)
# acks=0: ์๋ต ๋๊ธฐ ์์ (์ต๋ ์๋, ์ ์ค ๊ฐ๋ฅ)
# acks=1: Leader ํํฐ์
๋์คํฌ ๊ธฐ๋ก ํ์ธ ํ ์๋ต (๊ธฐ๋ณธ๊ฐ)
# acks=all ๋๋ -1: ๋ชจ๋ ISR ๋ณต์ ๋ณธ ๊ธฐ๋ก ์๋ฃ ํ ์๋ต (๊ฐ๋ ฅ ๊ถ์ฅ, ๋ฐ์ดํฐ ์ ์ค 0)
acks=all
# 2. ๋ฉฑ๋ฑ์ฑ ํ๋ก๋์ (Idempotence) - ๋คํธ์ํฌ ์ฌ์๋ ์ ์ค๋ณต ๋ฉ์์ง ์ ์ฅ ๋ฐฉ์ง
enable.idempotence=true
# ๋ฉฑ๋ฑ์ฑ ํ์ฑํ ์ ํ์ ์ฐ๊ณ ์ต์
:
# acks=all
# retries > 0
# max.in.flight.requests.per.connection <= 5
# 3. ์ฌ์๋ ๋ฐ ํ์์์
retries=2147483647
delivery.timeout.ms=120000
retry.backoff.ms=100
# 4. ์ผ๊ด ์ ์ก(Batching)์ผ๋ก ์ฒ๋ฆฌ๋ ๊ทน๋ํ
# linger.ms: ๋ฒํผ๊ฐ ๋ ์ฐจ๋ ์ต๋ N๋ฐ๋ฆฌ์ด ๋์ ๋๊ธฐ ํ ๋ชจ์์ ์ ์ก (๊ธฐ๋ณธ 0 -> 5~20ms ๊ถ์ฅ)
linger.ms=20
# batch.size: ๋จ์ผ ๋ฐฐ์น ์ต๋ ๋ฐ์ดํธ ํฌ๊ธฐ (๊ธฐ๋ณธ 16KB -> ๋์ฉ๋ ํธ๋ํฝ ์ 64KB ๊ถ์ฅ)
batch.size=65536
# 5. ๋ฉ์์ง ์์ถ ์๊ณ ๋ฆฌ์ฆ (๋คํธ์ํฌ ๋์ญํญ ์ ์ฝ)
# snappy: CPU ์ฌ์ฉ๋ ๋ฎ๊ณ ์์ถ/ํด์ ์๋ ๋งค์ฐ ๋น ๋ฆ (์ค๋ฌด ํ์ค ๊ถ์ฅ)
# zstd: ๋์ ์์ถ๋ฅ (๋ก๊ทธ, ํ
์คํธ ๋ฐ์ดํฐ์ ์ ํฉ)
compression.type=snappy
# 6. ๋ฒํผ ๋ฉ๋ชจ๋ฆฌ
buffer.memory=33554432 # 32MB ํ๋ก๋์ ๋ฒํผ
max.block.ms=60000
```
---
## 9. ์ปจ์๋จธ(Consumer) ํต์ฌ ํ๋ ํ๋ผ๋ฏธํฐ
```properties
# 1. ์ด๊ธฐ ์คํ์
๊ฒฐ์ ์ ๋ต (์ปค๋ฐ๋ ์คํ์
์ด ์์ ๋)
# earliest: ๊ฐ์ฅ ์ค๋๋ ๋ฉ์์ง๋ถํฐ ์ฝ๊ธฐ
# latest: ๊ฐ์ฅ ์ต์ ๋ฉ์์ง๋ถํฐ ์ฝ๊ธฐ (๊ธฐ๋ณธ๊ฐ)
# none: ์ปค๋ฐ ๊ธฐ๋ก ์์ผ๋ฉด ์์ธ(NoOffsetForPartitionException) ๋ฐ์
auto.offset.reset=earliest
# 2. ์คํ์
์๋ ์ปค๋ฐ (Auto Commit) vs ์๋ ์ปค๋ฐ (Manual Commit)
# ์ค๋ฌด ๊ถ์ฅ: false ์ค์ ํ ๋น์ฆ๋์ค ๋ก์ง ์ ์ ์๋ฃ ์์ ์ ์๋ ์ปค๋ฐ (At-Least-Once ๋ณด์ฅ)
enable.auto.commit=false
auto.commit.interval.ms=5000
# 3. ํด๋ง(Poll) ๋ ์ฝ๋ ์ ๋ฐ ์ธํฐ๋ฒ ํ๋ (๋ฆฌ๋ฐธ๋ฐ์ฑ ๋ฐฉ์ง ํต์ฌ)
# max.poll.records: ๋จ์ผ poll() ํธ์ถ ์ ๊ฐ์ ธ์ฌ ์ต๋ ๋ ์ฝ๋ ์ (๊ธฐ๋ณธ 500)
# ์ฒ๋ฆฌ ์๊ฐ์ด ๊ธด ์์
์ธ ๊ฒฝ์ฐ 50~100๊ฐ๋ก ์ค์ฌ ๋ฆฌ๋ฐธ๋ฐ์ฑ ํ์์์ ๋ฐฉ์ง
max.poll.records=100
# max.poll.interval.ms: poll() ํธ์ถ ๊ฐ ์ต๋ ํ์ฉ ์๊ฐ (๊ธฐ๋ณธ 300,000ms = 5๋ถ)
# ์ด ์๊ฐ ๋ด์ ๋ค์ poll()์ด ํธ์ถ๋์ง ์์ผ๋ฉด ๋ธ๋ก์ปค๋ ์ปจ์๋จธ๊ฐ ์ฃฝ์ ๊ฒ์ผ๋ก ๊ฐ์ฃผํ๊ณ Rebalance ์ ๋ฐ
max.poll.interval.ms=300000
# 4. ํํธ๋นํธ ๋ฐ ์ธ์
ํ์์์
# session.timeout.ms: ๋ธ๋ก์ปค๊ฐ ์ปจ์๋จธ์ ์ฅ์ ๋ฅผ ๊ฐ์งํ๋ ์๊ฐ (๊ธฐ๋ณธ 45์ด)
session.timeout.ms=45000
# heartbeat.interval.ms: session.timeout.ms์ 1/3 ์์ค ๊ถ์ฅ (๊ธฐ๋ณธ 3์ด)
heartbeat.interval.ms=3000
# 5. ํํฐ์
ํ ๋น ์ ๋ต (Partition Assignment Strategy)
# CooperativeStickyAssignor: ๋ฆฌ๋ฐธ๋ฐ์ฑ ๋ฐ์ ์ ์ ์ฒด ์ค๋จ(Stop-the-world) ์์ด ์ ์ง์ ์ฌํ ๋น
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
```
---
## 10. ๋ธ๋ก์ปค(Broker) `server.properties` ํ์ ์ค์
```properties
# 1. ๋ธ๋ก์ปค ๊ณ ์ ID ๋ฐ ๋ก๊ทธ ๋๋ ํ ๋ฆฌ
broker.id=1
log.dirs=/var/lib/kafka/data
# 2. ํ ํฝ ๊ธฐ๋ณธ ํํฐ์
๋ฐ ๋ณต์ ๊ณ์
num.partitions=3
default.replication.factor=3
min.insync.replicas=2 # acks=all๊ณผ ๊ฒฐํฉํ์ฌ ์ต์ 2๋ ์ด์ ์ฐ๊ธฐ ์ฑ๊ณตํด์ผ ์น์ธ
# 3. ํ ํฝ ์๋ ์์ฑ ๋ฐฉ์ง (์ค๋ฌด ํ์: ์คํ๋ก ์ธํ ํ ํฝ ๋๋ฆฝ ๋ฐฉ์ง)
auto.create.topics.enable=false
# 4. ์ ๋ขฐ์ฑ ์๋ ๋ฆฌ๋ ์ ์ถ ๋ฐฉ์ง (ISR์ ์๋ ๋ณต์ ๋ณธ์ด ๋ฆฌ๋ ๋๋ ๊ฒ ์ฐจ๋จ -> ์ ์ค ๋ฐฉ์ง)
unclean.leader.election.enable=false
# 5. ๋ก๊ทธ ๋ณด์กด ๋ฐ ์ธ๊ทธ๋จผํธ ์ ์ฑ
log.retention.hours=168 # 7์ผ ๋ณด์กด
log.retention.bytes=107374182400 # ๋ธ๋ก์ปค๋น ํ ํฝ ์ต๋ 100GB ๋ณด์กด
log.segment.bytes=1073741824 # ์ธ๊ทธ๋จผํธ ํ์ผ 1GB ๋จ์๋ก ๋กค๋ง
log.cleanup.policy=delete # ์ด๊ณผ ์ ์ญ์
# 6. ๋คํธ์ํฌ & I/O ์ค๋ ๋ ์ (CPU ์ฝ์ด ์ ๊ธฐ๋ฐ ํ๋)
num.network.threads=8 # ์์ฒญ ์์ /์๋ต ์ค๋ ๋
num.io.threads=16 # ๋์คํฌ ์ฝ๊ธฐ/์ฐ๊ธฐ ์ค๋ ๋
```
---
## 11. Spring Boot ์ฐ๋: `application.yml` ์ค์
```yaml
spring:
kafka:
# ์นดํ์นด ๋ธ๋ก์ปค ํด๋ฌ์คํฐ ์ฃผ์ ๋ชฉ๋ก
bootstrap-servers: localhost:9092
# ํ๋ก๋์ ์ ์ญ ์ค์
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all
retries: 3
properties:
enable.idempotence: true
linger.ms: 20
compression.type: snappy
# ์ปจ์๋จธ ์ ์ญ ์ค์
consumer:
group-id: order-service-group
auto-offset-reset: earliest
enable-auto-commit: false # ์๋ ์ปค๋ฐ ๋ชจ๋ ์ฌ์ฉ
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "com.example.*"
max.poll.records: 100
# ๋ฆฌ์ค๋ ์ปจํ
์ด๋ ์ค์ (์๋ Ack ์ง์)
listener:
ack-mode: MANUAL_IMMEDIATE # ์๋ ์ฆ์ ์ปค๋ฐ
concurrency: 3 # ์ปจ์๋จธ ์ค๋ ๋ ๋์์ฑ ์
```
---
## 12. Spring Boot ํ๋ก๋์ & ์ปจ์๋จธ ๊ตฌํ ์ฝ๋
```java
// [1] Kafka Producer ์๋น์ค (๋น๋๊ธฐ ์ฝ๋ฐฑ ์ ์ก)
@Service
@RequiredArgsConstructor
@Slf4j
public class OrderProducer {
private final KafkaTemplate<String, OrderEvent> kafkaTemplate;
public void sendOrderEvent(String orderId, OrderEvent event) {
// ํ ํฝ๋ช
, ๋ฉ์์ง ํค(orderId ํํฐ์
๋ ๋ณด์ฅ), ํ์ด๋ก๋
CompletableFuture<SendResult<String, OrderEvent>> future =
kafkaTemplate.send("order-events", orderId, event);
future.whenComplete((result, ex) -> {
if (ex == null) {
RecordMetadata meta = result.getRecordMetadata();
log.info("๋ฐํ ์ฑ๊ณต: topic={}, partition={}, offset={}",
meta.topic(), meta.partition(), meta.offset());
} else {
log.error("๋ฐํ ์คํจ: orderId={}, error={}", orderId, ex.getMessage());
}
});
}
}
```
```java
// [2] Kafka Consumer ๋ฆฌ์ค๋ (์๋ ์ปค๋ฐ + DLT ์ฌ์๋ ์ง์)
@Component
@Slf4j
public class OrderConsumer {
@KafkaListener(
topics = "order-events",
groupId = "order-service-group",
containerFactory = "kafkaListenerContainerFactory"
)
public void consumeOrder(
@Payload OrderEvent event,
@Header(KafkaHeaders.RECEIVED_KEY) String key,
@Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
@Header(KafkaHeaders.OFFSET) long offset,
Acknowledgment ack
) {
try {
log.info("์์ : key={}, partition={}, offset={}, payload={}",
key, partition, offset, event);
// ๋น์ฆ๋์ค ๋ก์ง ์ํ
processOrder(event);
// ์ฒ๋ฆฌ ์๋ฃ ํ ์คํ์
์๋ ์ปค๋ฐ
ack.acknowledge();
} catch (Exception e) {
log.error("๋ฉ์์ง ์ฒ๋ฆฌ ์คํจ, ์ฌ์๋ ๋๊ธฐ: {}", e.getMessage());
// ์์ธ๋ฅผ ๋์ง๋ฉด ErrorHandler์ ์ํด DLT(Dead Letter Topic)๋ก ์ด๋ ๊ฐ๋ฅ
throw e;
}
}
private void processOrder(OrderEvent event) {
// ์ฃผ๋ฌธ ์ฒ๋ฆฌ ๋ก์ง
}
}
```
---
## 13. ์ค๋ฌด ์ฅ์ ํธ๋ฌ๋ธ์ํ
& ์ฑ๋ฅ ์ต์ ํ (Troubleshooting)
```text
โ 1. Consumer Lag(์ง์ฐ)์ด ์ง์์ ์ผ๋ก ์ฆ๊ฐํ ๋
โ ์์ธ ๋ถ์:
- ์ปจ์๋จธ์ ๋น์ฆ๋์ค ๋ก์ง ์ฒ๋ฆฌ ์๋(DB ์ฐ๊ธฐ, ์ธ๋ถ API ํธ์ถ ์ง์ฐ)๊ฐ ํ๋ก๋์ ์ธ์
๋๋ณด๋ค ๋๋ฆผ
- ํน์ ํํฐ์
์ Key ๋ถ๊ท ํ์ผ๋ก ํซ์คํ(Hotspot) ๋ฐ์
โก ํด๊ฒฐ ์กฐ์น:
- ํ ํฝ ํํฐ์
์ ์ฆ๊ฐ + ๋์ผ ์ปจ์๋จธ ๊ทธ๋ฃน์ ์ธ์คํด์ค/์ค๋ ๋ ์ฆ์ค (1:1 ๋งคํ)
- DB Batch Insert ๋๋ ๋น๋๊ธฐ ๋ฉํฐ์ค๋ ๋ ํ์ดํ๋ผ์ธ ์ฒ๋ฆฌ
- Key๊ฐ ํน์ ๊ฐ์ ์น์ฐ์น์ง ์๋๋ก Salt ์ถ๊ฐ ๋๋ ํด์ ๋ถ์ฐ
โ 2. Rebalance Storm (๋น๋ฒํ ์ปจ์๋จธ ๋ฆฌ๋ฐธ๋ฐ์ฑ ๋ฐ์)
โ ์์ธ ๋ถ์:
- ์ปจ์๋จธ๊ฐ 1ํ poll()ํด์จ ๋ฉ์์ง๋ค์ ์ฒ๋ฆฌํ๋ ๋ฐ ๊ฑธ๋ฆฌ๋ ์๊ฐ์ด `max.poll.interval.ms`(๊ธฐ๋ณธ 5๋ถ)๋ฅผ ์ด๊ณผ
- ๋ธ๋ก์ปค๊ฐ ์ปจ์๋จธ๊ฐ ๋ค์ด๋ ๊ฒ์ผ๋ก ํ์ ํ๊ณ ์ ์ฒด ์ปจ์๋จธ ํํฐ์
์ฌํ ๋น ๋ฐ๋ณต
โก ํด๊ฒฐ ์กฐ์น:
- `max.poll.records` ํฌ๊ธฐ๋ฅผ 500 -> 50~100์ผ๋ก ์ค์ฌ 1ํ ์ฒ๋ฆฌ ์๊ฐ ๋จ์ถ
- `max.poll.interval.ms` ๋๊ธฐ ์๊ฐ์ 10~15๋ถ์ผ๋ก ์ถฉ๋ถํ ์ฐ์ฅ
- ํํฐ์
ํ ๋น ์ ๋ต์ `CooperativeStickyAssignor`๋ก ์ ํํ์ฌ ์ ๋ฉด ์ค๋จ(Stop-the-world) ๋ฐฉ์ง
โ 3. ๋ฉ์์ง ์ ์ค ๋ฐฉ์ง(Zero Data Loss) 3๋ ๊ณจ๋ ๋ฃฐ
โ Producer: `acks=all`, `retries=int_max`, `enable.idempotence=true`
โก Broker: `min.insync.replicas=2` (๋ณต์ ๊ณ์ 3 ๊ธฐ์ค), `unclean.leader.election.enable=false`
โข Consumer: `enable.auto.commit=false` (๋ก์ง ์ฑ๊ณต ํ ๋ช
์์ `ack.acknowledge()`)
โ 4. OS ๋ฐ ๋์คํฌ I/O ์ฑ๋ฅ ํ๋
โ Zero-Copy & Page Cache:
- Kafka๋ JVM ํ ๋์ OS ์ปค๋์ PageCache๋ฅผ ์ ๊ทน ํ์ฉํจ (sendfile ์์คํ
์ฝ)
- ๋ธ๋ก์ปค JVM ํ ํฌ๊ธฐ๋ 6~8GB ์ ๋๋ก ์๊ฒ ๋๊ณ , ๋จ์ ์์คํ
๋ฉ๋ชจ๋ฆฌ๋ฅผ OS Page Cache์ ํ ๋นํ ๊ฒ!
โก ๋์คํฌ XFS ํ์ผ ์์คํ
์ฌ์ฉ ๊ถ์ฅ (ext4 ๋๋น ๋์ฉ๋ ์ธ๊ทธ๋จผํธ ํ๋ฌ์ ์ฑ๋ฅ ์ฐ์)
```
์๊ฒฌ ๋ฐ ์ง๋ฌธ
0์์ง ๋ฑ๋ก๋ ์๊ฒฌ์ด ์์ต๋๋ค. ์ฒซ ๋ฒ์งธ ๋๊ธ์ ๋จ๊ฒจ๋ณด์ธ์!
๋๊ธ ์์
๋๊ธ ์ญ์