마이크로 서비스 간 End-To-End 통신 시 발생할 수 있는 서비스 간 의존성 증가, 서비스 간 커플링 발생 등 으로 인한 유지보수, 확장에 어려움이 발생할 수 있다.
이에 대한 해결책으로 중앙집중화, 이벤트 스트리밍 방식의 Kafka 가 많은 관심을 받고 있다.
Kafka는 Producer / Consumer 로 분리해서 메세지 큐(FIFO) 구조를 채택해서 데이터를 구분하는 단위인 토픽에 1개 이상의 FIFO 형태의 큐 구조인 파티션을 배치해서 Producer가 파티션에 전송한 메세지를 Consumer가 파티션에서 꺼내는 방식으로 동작한다.
위의 동작 방식으로 Kafka는 빅데이터 파이프라인 으로서 아래의 장점을 갖는다.
인프런 데브원영 님의 카프카 애플리케이션 강의 를 듣고 학습한 내용을 정리한 Github : https://github.com/onlydev7777/TIL/tree/main/kafka
해당 프로젝트에서는 사용자가 일기(EDIARY-DIARY) 작성 요청 시 사용자의 포인트(EDIARY-POINT)가 적립되는 비즈니스 로직에 Kafka 통신을 적용하였다.
일기 작성에 대한 비즈니스 로직에서 포인트 적립은 Kafka 통신으로 위임
일기 서비스(EDIARY-DIARY) 에 Kafka 프로듀서를 설정해서 'POINT-ADD' 토픽의 파티션에 메세지를 전송하도록 구현한다.
@Configuration
public class KafkaProducerConfig {
@Bean
public ProducerFactory<String, PointHistoryRequest> producerFactory() {
Map<String, Object> properties = new HashMap<>();
properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:7070");
properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
return new DefaultKafkaProducerFactory<>(properties);
}
@Bean
public KafkaTemplate<String, PointHistoryRequest> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
}
@Slf4j
@RequiredArgsConstructor
@Service
public class PointKafkaProducer {
private static final String TOPIC = "point-add";
private final KafkaTemplate<String, PointHistoryRequest> kafkaTemplate;
public void send(PointHistoryRequest request) {
kafkaTemplate.send(TOPIC, request);
log.info("{} 토픽 전송! memberId : {}", TOPIC, request.getMemberId());
}
}
@PostMapping
public ResponseEntity<ApiResult<DiaryResponse>> saveDiary(@RequestBody DiaryRequest request) {
DiaryResponse savedDiaryResponse = mapper.toResponse(service.save(mapper.toDto(request)));
PointHistoryRequest pointHistoryRequest = PointHistoryRequest.builder()
.memberId(request.getMemberId())
.score(10)
.details("일기쓰기")
.build();
log.info("kafkaProducer.send 호출 : memberId : {}", pointHistoryRequest.getMemberId());
kafkaProducer.send(pointHistoryRequest);
return ResponseEntity.ok(ApiResult.OK(savedDiaryResponse));
}
포인트 서비스(EDIARY-POINT) 에 Kafka 컨슈머를 설정해서 'POINT-ADD' 토픽의 파티션에 메세지를 수신해서 비즈니스 로직을 처리하도록 구현한다.
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Import(KafkaListenerConfigurationSelector.class)
public @interface EnableKafka {
}
@EnableKafka
@Configuration
public class KafkaConsumerConfig {
@Bean
public ConsumerFactory<String, PointHistoryRequest> consumerFactory() {
Map<String, Object> properties = new HashMap<>();
properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:7070");
properties.put(ConsumerConfig.GROUP_ID_CONFIG, "point-group");
properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
properties.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
return new DefaultKafkaConsumerFactory<>(properties, new StringDeserializer(), new JsonDeserializer<>(PointHistoryRequest.class, false));
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, PointHistoryRequest> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, PointHistoryRequest> kafkaListenerContainerFactory = new ConcurrentKafkaListenerContainerFactory<>();
kafkaListenerContainerFactory.setConsumerFactory(consumerFactory());
return kafkaListenerContainerFactory;
}
}
@Slf4j
@RequiredArgsConstructor
@Service
public class PointKafkaConsumer {
private final PointService service;
private final PointHistoryMapper historyMapper;
@KafkaListener(topics = "point-add", groupId = "point-group")
public void add(PointHistoryRequest request) {
PointHistoryDto dto = historyMapper.toDto(request);
PointDto pointDto = service.getOrSave(request.getMemberId());
dto.setPointId(pointDto.getId());
PointHistoryDto savedDto = service.add(dto);
log.info("PointKafkaConsumer.add > memberId : {}, pointId : {}, pointHistoryId : {}", pointDto.getMemberId(), pointDto.getId(), savedDto.getId());
}
}




front-end("front-msa" 브랜치) : https://github.com/onlydev7777/emotion-diary-react
back-end : https://github.com/onlydev7777/emotion-diary-msa
inflearn-msa : https://github.com/onlydev7777/springboot-msa-3.0/tree/master
inflearn-kafka : https://github.com/onlydev7777/TIL/tree/main/kafka