
Kafk의 Topic과 Partition

Topic
카프카 토픽은 구체화된 이벤트 스트림을 뜻한다. 토픽은 연관된 이벤트들을 묶어 영속화한다. 이때, 이벤트는 불변(immutable)이기 때문에 토픽에 한번 추가되면 이후로는 수정 불가능하다.
토픽은 Producer와 Consumer 역할을 분리하는 기점이다.
Producer는 카프카 토픽에 메시지를 저장(Push)하고, Consumer는 카프카 토픽에 저장된 메시지를 읽어 온다(Pull).

Partition
카프카의 토픽들은 여러 파티션으로 나눠진다. 토픽이 카프카의 논리적 개념이라면, 파티션은 토픽에 속한 레코드를 실제 저장소에 저장하는 가장 작은 단위다. 각각의 파티션은 Append-Only 방식으로 기록되는 하나의 로그 파일이다.

Offset
파티션의 레코드는 각각 Offset이라 불리는, 파티션 내 고유한 레코드의 순서를 의미하는 식별자 정보를 가진다.
하지만 카프카는 일반적으로 메시지의 순서를 보장하지 않는다. 파티션 내에서는 순서를 보장하더라도, 파티션 간에는 순서가 보장되지 않기 때문이다.
Offset 정보는 카프카에 의해서 관리되고, 값이 계속 증가하며, 불변하는 숫자 정보다. 레코드가 파티션에 쓰일 대 항상 기록의 맨 뒤에 쓰이면서 다음 순서의 Offset 값을 갖게 된다.
Partition에서 Records 읽기
일반적인 pub/sub 모델과 달리, Kafka는 메시지를 Consumer에 전달(Push)하지 않고, 파티션으로부터 메시지를 읽어(pull) 가야 한다.

메시지의 offset은 Consumer 측에서 커서처럼 동작한다. 추적한 메시지 offset을 통해 이미 소비한 메시지를 저장한다.
각 파티션에서 마지막으로 소비된 메시지의 offset을 기억함으로써 consumer는 어떤 시점에 파티션의 어떤 위치에서 재시작할 수 있다. 이를 통해 consumer는 장애를 복구한 뒤 메시지 소비를 재시작할 수 있다.

파티션은 하나 혹은 그 이상의 consumer들로부터 소비될 수 있고, 각각 서로 다른 offset을 통해 메시지를 읽을 수 있다.
카프카의 consumer group은 동일 토픽을 소모하는 consumer를 그룹으로 묶는 개념이다. 동일한 consumer group에 있는 consumer들은 동일한 group-id를 부여받는다.
Spring Boot / Kotlin에 Kafka 적용하기
문제 설정
현재 프로젝트는 백엔드의 주문 서비스에서 주문 생성 (checkout) / 취소 (canceled) / 완료 (sold) 처리 시 알림 서비스를 직접 호출하여 알림을 생성하는 동기화 fire-and-forget 방식이다.
BFF.checkout()
→ orderClient.getMyCart() (REST, sync)
→ orderClient.checkout() (REST, sync)
→ storeClient.findStore() (REST, sync — to get storeOwnerId + storeName)
→ storeClient.listProducts() (REST, sync — to get productNames)
→ notificationClient.createNotification() (REST, coroutine, circuit-breaker, fire-and-forget)
BFF.markSold()
→ orderClient.markSold() (REST, sync)
→ storeClient.findStore() (REST, sync — storeName)
→ storeClient.listProducts() (REST, sync — productNames)
→ notificationClient.createNotification() (REST, coroutine, circuit-breaker, fire-and-forget)
이 구조는 주문 서비스가 알림 서비스에 종속되는 구조로 아래와 같은 문제가 발생한다 :
- 알림 서비스가 다운되면 주문 서비스에서 발생하는 모든 이벤트가 소실된다.
- 알림 서비스 다운에 즉시 대응하기 위해 circuit breaker가 요구된다.
Kafka를 도입하면 해당 프로세스는 아래 방색으로 변경할 수 있다.
BFF.checkout()
→ orderClient.getMyCart() (REST, sync — unchanged)
→ storeClient.findStore() (REST, sync — now feeds CheckoutRequest, not notification)
→ storeClient.listProducts() (REST, sync — productNames for CheckoutRequest items)
→ orderClient.checkout(enrichedRequest) (REST, sync)
← OrderService publishes OrderEvent to Kafka (durable)
BFF.markSold()
→ orderClient.markSold() (REST, sync — that's it, no extra calls)
← OrderService publishes OrderEvent to Kafka
[Decoupled, async]
Kafka → NotificationService consumer → persist notification
Kafka → StoreService consumer → increment product popularity
- Kafka에서 이벤트 보유 기한(retention)을 설정하여 주문 서비스가 다운되더라도 로그가 7일까지 저장된다.
- BFF는 더 이상 주문 이벤트 각각에 대해 알림 서비스를 호출할 필요가 없어진다. (비동기화)
이벤트 다이어그램
Client BFF OrderService Kafka NotificationService StoreService
│ │ │ │ │ │
├─checkout──>│ │ │ │ │
│ ├─getMyCart─────────>│ │ │ │
│ ├─findStore(store)──>│(StoreService) │ │ │
│ ├─listProducts──────>│(StoreService) │ │ │
│ ├─checkout(enriched)>│ │ │ │
│ │ │─publish──────>│ │ │
│ │<───────────────────┤ NEW_ORDER │ │ │
│<───────────┤ │ │─consume──────────>│ │
│ 200 OK │ │ │ │ createFromEvent │
│ │ │ │ │ │
├─markSold──>│ │ │ │ │
│ ├─markSold──────────>│ │ │ │
│ │ │─publish──────>│ │ │
│ │<───────────────────┤ ORDER_SOLD │─consume──────────>│ │
│<───────────┤ │ │─consume────────────────────────────>│
│ 200 OK │ │ │ │ incrementPopularity
현재 프로젝트에서 이벤트 Producer는 주문 서비스(OrderService)다.
서비스 트랜잭션(@Transaction) 내에서 이벤트를 커밋하여 이벤트 유실을 방지하기 위함이다.
Kafka 토픽 설계 및 설정
토픽 설계
| 토픽 | 파티션 수 | 키 | 보유 기간 | Consumers |
| baemin.order.events | 3 | storeId | 7일 | 주문 / 가게 서비스 |
Kafka Helm 템플릿 작성
<!--helm/kafka/Chart.yaml-->
apiVersion: v2
name: kafka
description: Single-node Apache Kafka (KRaft) for Baemin
type: application
version: 0.2.0
appVersion: "3.7.1"
<!--helm/kafka/values.yaml-->
image:
repository: apache/kafka
tag: "3.7.1"
service:
port: 9092
kafka:
nodeId: 1
offsetsTopicReplicationFactor: 1
transactionStateLogReplicationFactor: 1
transactionStateLogMinIsr: 1
autoCreateTopicsEnable: true
- node.id : KRaft 모드에 필요한 고유 ID.
- offsets.topic.replication.factor : 카프카는 consumer group offsets(__consumer_offsets)를 내부 토픽으로 저장함. 설정된 수 만큼 복제 팩터가 설정되지 않으면 내부 토픽 생성이 제한됨.
- transaction.state.log.replication.factor : 다른 내부 토픽(__transaction_state)으로, transactional producer 상태를 저장함. 설정된 수 만큼 복제 팩터가 설정되지 않으면 내부 토픽 생성이 제한됨.
- transaction.state.log.min.isr : 트랜잭션 토픽에 대한 쓰기로 인정하기 위해 필요한 최소 동기화 복제본(In-Sync Replicas) 수
- auto.create.topics.enable : 서버에 대한 토픽 자동 생성 허용 (NewTopic bean in KafkaConfig.kt to create topics)
※ KRaft (Kafka Raft)
아파치 카프카의 분산 시스템을 관리하기 위해 도입된 메커니즘이다.
KRaft 이전의 카프카는 클러스터 메타데이터를 관리하기 위해 아파치 주키퍼를 사용했었다. 때문에 브로커는 모든 토픽과 파티션에 대한 메타데이터를 주키퍼에서 읽어야 했다. 메타데이터 업데이트는 주키퍼에 동기 방식으로, 브로커에는 비동기 방식으로 일어났고, 이 과정에서 주키퍼와 브로커 간 메타데이터 불일치도 발생할 수 있었다. 특히 토픽과 파티션이 많은 대규모 카프카 클러스터에서는 병목 현상이 발생했다.
또한 주키퍼와 카프카는 서로 완전히 다른 애플리케이션으로, 구성 파일, 환경, 서비스 데몬을 각각 가지고 있어 관리자가 직접 동시에 운용해야 했다. 모니터링을 적용하는 방법과 각 앱이 보이는 주요 metrics도 달랐기 때문에 관리적인 측면에서도 어려움이 존재했다.
KRaft는 카프카와 결합하여 이 운영 복잡성을 줄이고, 카프카의 전반적인 신뢰성과 관리 용이성을 개선하는데 기여했다.
주문 서비스 (Producer) 설정
Producer Helm 템플릿 작성
<!--helm/order-service/values.yaml-->
config:
kafkaBootstrapServers: "baemin-kafka:9092"
<!--helm/order-service/templates/deployment.yaml-->
- name: KAFKA_BOOTSTRAP_SERVERS
value: {{ .Values.config.kafkaBootstrapServers | quote }}
KafkaConfig 구성
// build.gradle.kts
implementation("org.springframework.boot:spring-boot-starter-kafka")
---
// application.yaml (Producer)
kafka:
bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JacksonJsonSerializer
properties:
spring.json.add.type.headers: false
---
// KafkaConfig.kt
package order.config
import common.event.OrderEvent
import org.apache.kafka.clients.admin.NewTopic
import org.springframework.boot.kafka.autoconfigure.KafkaProperties
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration
import org.springframework.kafka.config.TopicBuilder
import org.springframework.kafka.core.DefaultKafkaProducerFactory
import org.springframework.kafka.core.KafkaTemplate
import org.springframework.kafka.core.ProducerFactory
@Configuration
class KafkaConfig(private val kafkaProperties: KafkaProperties) {
@Bean
fun orderEventTopic(): NewTopic =
TopicBuilder.name("baemin.order.events").partitions(3).replicas(1).build()
@Bean
fun orderEventProducerFactory(): ProducerFactory<String, OrderEvent> =
DefaultKafkaProducerFactory(kafkaProperties.buildProducerProperties())
@Bean
fun orderEventKafkaTemplate(pf: ProducerFactory<String, OrderEvent>): KafkaTemplate<String, OrderEvent> =
KafkaTemplate(pf)
}
- orderEventTopic() : "baemin.order.events" 토픽 생성
- orderEventProducerFactory() : buildProducerProperties()를 통해 application.yml에 정의된 카프카 설정 읽기
- orderEventKafkaTemplate() : 카프카 이벤트 생성에 사용할 ProducerFactory 정의
EventRelay 구성
package order.config
import common.event.OrderEvent
import org.slf4j.LoggerFactory
import org.springframework.kafka.core.KafkaTemplate
import org.springframework.stereotype.Component
import org.springframework.transaction.event.TransactionPhase
import org.springframework.transaction.event.TransactionalEventListener
private val log = LoggerFactory.getLogger(OrderEventRelay::class.java)
@Component
class OrderEventRelay(private val kafkaTemplate: KafkaTemplate<String, OrderEvent>) {
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
fun onOrderEvent(event: OrderEvent) {
kafkaTemplate.send("baemin.order.events", event.storeId.toString(), event)
.whenComplete { result, ex ->
if (ex != null) {
log.error("Failed to publish {} for order {}: {}", event.eventType, event.orderId, ex.message, ex)
} else {
log.debug(
"Published {} for order {} (partition={}, offset={})",
event.eventType, event.orderId,
result.recordMetadata.partition(), result.recordMetadata.offset()
)
}
}
}
}
- 이벤트 발행 시, Service 레이어에서 직접 호출하지 않고 ApplicationEventPublisher를 통해 발행하기 위함.
- 이 설정을 통해 이벤트는 DB 트랜잭션 커밋 이후에만 카프카로 전송되는 것을 보장한다.
주문 서비스 Event Publisher 구현
// checkout()
eventPublisher.publishEvent(
OrderEvent (
eventType = OrderEventType.NEW_ORDER,
...
)
)
// markSold()
eventPublisher.publishEvent(
OrderEvent (
eventType = OrderEventType.ORDER_SOLD,
...
)
)
// markCanceled()
eventPublisher.publishEvent(
OrderEvent (
eventType = OrderEventType.ORDER_CANCELLED,
...
)
)
알림 서비스 (Consumer) 설정
KafkaConfig 설정
// build.gradle.kts
implementation("org.springframework.boot:spring-boot-starter-kafka")
---
// application.yaml (Consumer)
kafka:
bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
consumer:
group-id: notification-service
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
properties:
spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JacksonJsonDeserializer
spring.json.trusted.packages: "common.event"
spring.json.value.default.type: "common.event.OrderEvent"
---
// KafkaConsumerConfig.kt
package notification.config
import org.slf4j.LoggerFactory
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration
import org.springframework.kafka.listener.DefaultErrorHandler
import org.springframework.util.backoff.FixedBackOff
private val log = LoggerFactory.getLogger(KafkaConsumerConfig::class.java)
@Configuration
class KafkaConsumerConfig {
@Bean
fun kafkaErrorHandler(): DefaultErrorHandler =
DefaultErrorHandler(
{ record, ex ->
log.error(
"Giving up on Kafka record after retries: topic={} partition={} offset={}: {}",
record.topic(), record.partition(), record.offset(), ex.message, ex
)
},
FixedBackOff(1000L, 2L)
)
}
알림 서비스 Event Consumer 구현
package notification.service.user
import common.event.OrderEvent
import common.event.OrderEventType
import notification.entity.user.NotificationType
import org.springframework.kafka.annotation.KafkaListener
import org.springframework.stereotype.Component
private data class NotifContext(
val recipientId: Long,
val type: NotificationType,
val title: String,
val content: String
)
@Component
class OrderEventConsumer(private val notificationService: NotificationService) {
@KafkaListener(topics = ["baemin.order.events"], groupId = "notification-service")
fun consume(event: OrderEvent) {
val ctx = when (event.eventType) {
OrderEventType.NEW_ORDER -> NotifContext(
recipientId = event.storeOwnerId,
type = NotificationType.NEW_ORDER,
title = "새 주문 접수",
content = "새 주문이 접수되었습니다."
)
OrderEventType.ORDER_SOLD -> NotifContext(
recipientId = event.userId,
type = NotificationType.ORDER_SOLD,
title = "주문 완료",
content = "주문이 완료되었습니다."
)
OrderEventType.ORDER_CANCELLED -> NotifContext(
recipientId = event.userId,
type = NotificationType.ORDER_CANCELED, // NotificationType enum uses single-L CANCELED
title = "주문 취소",
content = "주문이 취소되었습니다."
)
}
notificationService.createFromEvent(
recipientId = ctx.recipientId,
type = ctx.type,
title = ctx.title,
content = ctx.content,
storeId = event.storeId,
storeName = event.storeName.ifBlank { null },
items = event.items,
occurredAt = event.occurredAt
)
}
}
결과 확인
kafka-topics.sh : 토픽 관리 및 점검
# list all topics
kafka-topics.sh --bootstrap-server localhost:9092 --list
__consumer_offsets
baemin.order.events
baemin.order.events-dlt
# describe a topic (partitions, replication, config)
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic baemin.order.events
Topic: baemin.order.events TopicId: zyNIsrlaTgKw8ryaH0lN_w PartitionCount: 3 ReplicationFactor: 1 Configs:
Topic: baemin.order.events Partition: 0 Leader: 1 Replicas: 1 Isr: 1
Topic: baemin.order.events Partition: 1 Leader: 1 Replicas: 1 Isr: 1
Topic: baemin.order.events Partition: 2 Leader: 1 Replicas: 1 Isr: 1
# describe the DLT
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic baemin.order.events-dlt
Topic: baemin.order.events-dlt TopicId: ctr-ROlQQn-cecSzpRcoSA PartitionCount: 3 ReplicationFactor: 1 Configs:
Topic: baemin.order.events-dlt Partition: 0 Leader: 1 Replicas: 1 Isr: 1
Topic: baemin.order.events-dlt Partition: 1 Leader: 1 Replicas: 1 Isr: 1
Topic: baemin.order.events-dlt Partition: 2 Leader: 1 Replicas: 1 Isr: 1
kafka-console-consumer(producer}.sh : 토픽에 메시지를 넣고(Push) 읽어오기(Pull)
# watch live events as they arrive
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic baemin.order.events
# publish a message manually
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic baemin.order.events
kafka-consumer-groups.sh : Consumer group 상태 확인
# list all consumer groups
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
# show committed offsets and lag for both groups
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group store-service
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group notification-service
주문 처리 결과
psql -h 192.168.160.104 -U notif_svc -d notificationdb \
-c "SELECT id, user_id, type, title, created_at FROM notifications ORDER BY created_at DESC LIMIT 10;"
id | user_id | type | title | created_at
----+---------+----------------+--------------+---------------
25 | 1 | ORDER_SOLD | 주문 완료 | 1789389272921
24 | 2 | NEW_ORDER | 새 주문 접수 | 1789389038533
23 | 1 | ORDER_SOLD | 주문 완료 | 1789132092990
22 | 1 | ORDER_SOLD | 주문 완료 | 1789132087816
21 | 5 | ORDER_CANCELED | 주문 취소 | 1789132086473
20 | 2 | NEW_ORDER | 새 주문 접수 | 1789132051490
19 | 2 | NEW_ORDER | 새 주문 접수 | 1789132037118
18 | 2 | NEW_ORDER | 새 주문 접수 | 1788960053482
17 | 1 | ORDER_SOLD | 주문 완료 | 1788695021955
16 | 1 | ORDER_SOLD | 주문 완료 | 1788694560116
(10 rows)
...
psql -h 192.168.160.102 -U store_svc -d storedb \
-c "SELECT id, name, popularity FROM products ORDER BY popularity DESC;"
id | name | popularity
----+---------------+------------
4 | 양념 | 7
3 | 후라이드 | 4
1 | 페퍼로니 피자 | 4
5 | 간장 | 3
2 | 시카고 피자 | 3
6 | 고르곤졸라 | 2
(6 rows)
참고 출처
- [Backend] Message Queue (ft. Kafka, RabbitMQ, ZeroMQ)
- What are partitions? -- Red Hat Developer
- [Kafka] Kafka 의 Topic 과 Partition
- Understanding Kafka Topics and Partitions
- Apache Kafka의 새로운 협의 프로토콜인 KRaft에 대해(1)
- Spring Boot에서 Kafka 제대로 쓰기: Producer부터 DLT까지
- https://rudaks.tistory.com/entry/Spring-Boot-%EB%B2%88%EC%97%AD-Apache-Kafka-Support
'Backend > Spring' 카테고리의 다른 글
| [Kotlin + Spring] plugin.spring과 plugin.jpa가 존재하는 이유 (0) | 2026.03.24 |
|---|---|
| [Spring] Spring Data JPA Audit, Kotlin으로 구현하기 (0) | 2026.03.16 |
| [Spring] Spring 이해하기(2) - POJO 프로그래밍을 돕는 IoC/DI, AOP, PSA (1) | 2026.02.17 |
| [Spring] Spring 이해하기(1) - POJO, Java Bean, Spring Bean (0) | 2026.02.17 |
| [Spring Data JPA] @Lob + ByteArray 조합으로 인한 bytea / oid 혼동 문제 (1) | 2026.02.12 |