4편에서 Redis로 동시성 문제를 해결했다. 이번 편에서는 Kafka를 도입해서 주문 완료 이벤트를 비동기로 처리하는 구조를 만들었다.
주문이 완료되면 이메일 발송, 포인트 적립, 배송 시작 등 여러 작업이 필요하다. 이걸 동기로 처리하면 문제가 생긴다.
Kafka를 쓰면 주문 서비스는 메시지만 던지고 끝. 나머지 서비스들이 알아서 처리한다.
Topic — 메시지를 저장하는 공간. 우리 프로젝트에서는 order-complete.
Producer — 메시지를 Topic에 발행하는 쪽. 주문 서비스가 담당.
Consumer — Topic을 구독해서 메시지를 읽는 쪽. 이메일/포인트 서비스가 담당.
Consumer Group — 같은 Topic을 여러 서비스가 독립적으로 구독할 수 있음.
Offset — Consumer가 어디까지 읽었는지 기록. 장애 복구 시 이어서 읽기 가능.
implementation 'org.springframework.kafka:spring-kafka'
@Configuration
public class KafkaConfig {
@Bean
public ProducerFactory<String, String> producerFactory(
@Value("${spring.kafka.bootstrap-servers}") String bootstrapServers) {
Map<String, Object> config = new HashMap<>();
config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
return new DefaultKafkaProducerFactory<>(config);
}
@Bean
public KafkaTemplate<String, String> kafkaTemplate(ProducerFactory<String, String> producerFactory) {
return new KafkaTemplate<>(producerFactory);
}
@Bean
public ConsumerFactory<String, String> consumerFactory(
@Value("${spring.kafka.bootstrap-servers}") String bootstrapServers) {
Map<String, Object> config = new HashMap<>();
config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
config.put(ConsumerConfig.GROUP_ID_CONFIG, "order-group");
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
return new DefaultKafkaConsumerFactory<>(config);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
ConsumerFactory<String, String> consumerFactory) {
ConcurrentKafkaListenerContainerFactory<String, String> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory);
return factory;
}
}
@Component
@RequiredArgsConstructor
public class OrderEventProducer {
private final KafkaTemplate<String, String> kafkaTemplate;
public void sendOrderComplete(Long orderId) {
kafkaTemplate.send("order-complete", String.valueOf(orderId));
}
}
@Component
public class OrderEventConsumer {
private static final Logger log = LoggerFactory.getLogger(OrderEventConsumer.class);
@KafkaListener(topics = "order-complete", groupId = "order-group")
public void handleOrderComplete(String orderId) {
log.info("주문 완료 이벤트 수신 — orderId: {}", orderId);
}
}
orderRepository.save(order);
orderItemRepository.save(orderItem);
orderEventProducer.sendOrderComplete(order.getId()); // 주문 완료 이벤트 발행
spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.consumer.group-id=order-group
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.auto-offset-reset=earliest
주문 API 호출 후 Kafka Topic에 메시지가 쌓이는 것을 확인했다.
kafka-console-consumer --topic order-complete --bootstrap-server localhost:9092 --from-beginning
orderId가 순서대로 출력되는 것을 확인했다.
Kafka Consumer 전략은 크게 세 가지다.
| 전략 | 설명 | 특징 |
|---|---|---|
| At Most Once | offset 먼저 커밋 | 메시지 유실 가능 |
| At Least Once | 처리 후 offset 커밋 | 중복 가능, 멱등성 필요 |
| Exactly Once | 딱 한 번만 처리 | 구현 복잡, 성능 저하 |
우리 프로젝트는 At Least Once 전략을 채택했다. 중복 처리 방지는 orderId 기반 체크로 멱등성을 보장할 수 있다.
"At Least Once 전략을 사용하고, 중복 처리 방지를 위해 orderId 기반 멱등성을 보장했습니다."
이번 편에서 Kafka를 도입해서 주문 완료 이벤트를 비동기로 발행하는 구조를 완성했다.
다음 편에서는 성능 테스트와 README, GitHub 정리를 진행할 예정이다.