[토이프로젝트][MSA] 감정일기장-8 : Kafka 통신

onlydev7777·2024년 9월 24일

☀️ 개요

마이크로 서비스 간 End-To-End 통신 시 발생할 수 있는 서비스 간 의존성 증가, 서비스 간 커플링 발생 등 으로 인한 유지보수, 확장에 어려움이 발생할 수 있다.

이에 대한 해결책으로 중앙집중화, 이벤트 스트리밍 방식의 Kafka 가 많은 관심을 받고 있다.

Kafka는 Producer / Consumer 로 분리해서 메세지 큐(FIFO) 구조를 채택해서 데이터를 구분하는 단위인 토픽에 1개 이상의 FIFO 형태의 큐 구조인 파티션을 배치해서 Producer가 파티션에 전송한 메세지를 Consumer가 파티션에서 꺼내는 방식으로 동작한다.

위의 동작 방식으로 Kafka는 빅데이터 파이프라인 으로서 아래의 장점을 갖는다.

  1. 높은 처리량
    • 많은 양의 데이터를 묶어 배치로 처리
    • 파티션 단위를 통해 동일 목적의 데이터를 여러 파티션에 분배, 데이터 병렬 처리 가능
    • 파티션, 컨슈머 개수를 늘려서 데이터 처리량을 높일 수 있음
  2. 확장성
    • 카프카 클러스터의 브로커를 최소한의 개수로 운영
    • 데이터가 많아지면 브로커 개수를 자동으로 늘려서 Scale-Out
    • 데이터가 적어지면 브로커 개수를 자동으로 줄여서 Scale-In
    • 무중단으로 Scale-In/Out 가능
  3. 영속성
    • 데이터를 메모리에 저장하지 않고 파일시스템에 저장
    • OS 레벨에서 파일 I/O 처리
    • 파일 I/O 성능 향상을 위해 페이지 캐시 영역에 메모리 생성
    • 한번 읽은 파일은 페이지 케시 메모리에 저장 시켜서 다시 사용
    • 브로커(카프카 서버)가 장애 발생 하더라도 파일 시스템에 데이터가 존재하므로 안정적인 서비스 운영 가능
  4. 고가용성
    • 일부 서버에 장애가 발생하더라도 무중단 처리
    • 리더 브로커 - 팔로워 브로커 간 데이터 복제
    • 리더 브로커 장애 발생 하더라도 팔로워 브로커에 데이터 복제 되어 있으므로 지속적인 서비스 운영 가능

► 참고

인프런 데브원영 님의 카프카 애플리케이션 강의 를 듣고 학습한 내용을 정리한 Github : https://github.com/onlydev7777/TIL/tree/main/kafka

해당 프로젝트에서는 사용자가 일기(EDIARY-DIARY) 작성 요청 시 사용자의 포인트(EDIARY-POINT)가 적립되는 비즈니스 로직에 Kafka 통신을 적용하였다.

일기 작성에 대한 비즈니스 로직에서 포인트 적립은 Kafka 통신으로 위임

1️⃣ Kafka 프로듀서 애플리케이션 구현

일기 서비스(EDIARY-DIARY) 에 Kafka 프로듀서를 설정해서 'POINT-ADD' 토픽의 파티션에 메세지를 전송하도록 구현한다.

1. KafkaProducerConfig

  • Kafka 프로듀서 설정 Configuration
  • ProducerFactory 빈 등록
    • Kafka 서버 URL 설정
    • Kafka 서버의 파티션 메세지 전송 시 Key/Value 값 Serializer 설정
  • KafkaTemplate<K, V> 빈 등록
    • Kafka 서버로 파티션 전송을 수행하는 역할
  @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());
    }
  }

2. PointKafkaProducer

  • Kafka 파티션으로 메세지 전송을 수행하는 서비스 컴포넌트
  @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());
    }
  }

3. DiaryController

  • 일기 작성 API 요청 수행 컨트롤러
  • 일기 작성 트랜잭션이 종료 후 Kafka 프로듀서를 통해 Kafka 통신
    @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));
    }

2️⃣ Kafka 컨슈머 애플리케이션 구현

포인트 서비스(EDIARY-POINT) 에 Kafka 컨슈머를 설정해서 'POINT-ADD' 토픽의 파티션에 메세지를 수신해서 비즈니스 로직을 처리하도록 구현한다.

1. EnableKafka

  • KafkaListenerConfigurationSelector 를 import 해서 KafkaBootstrapConfiguration 에서 KafkaListenerAnnotationBeanPostProcessor 빈을 수행시킨다.
  • KafkaListenerAnnotationBeanPostProcessor 는 빈을 후처리 하면서 @KafkaListener 가 붙은 메서드를 Kafka 컨슈머로서 동작할 수 있도록 설정
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Import(KafkaListenerConfigurationSelector.class)
public @interface EnableKafka {
}

2. KafkaConsumerConfig

  • Kafka 컨슈머 설정 Configuration
  • ConsumerFactory 빈 등록
    • Kafka 서버 URL 설정
    • Kafka 서버의 파티션에 담긴 메세지 수신 시 Key/Value 값 Deserializer 설정
    • Kafka 컨슈머 그룹 설정
  • ConcurrentKafkaListenerContainerFactory 빈 등록
    • 동시성(비동기)을 지원하는 리스너 생성
    • ConsumerFactory 빈을 기반으로 메세지를 수신할 리스너 컨테이너 생성
    • 생성된 리스너 컨테이너는 @KafkaListener 가 선언된 메서드로 메세지 전달
    • 즉, 비동기 적으로 처리하는 리스너 컨테이너를 생성해서 리스너 컨테이너가 Kafka 메세지를 Deserialize 한 후에 @KafkaListener 가 선언된 메서드에 수신된 메세지를 전달, 호출 하도록 설정
  @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;
    }
  }

3. PointKafkaConsumer

  • Kafka Consumer 서비스 컴포넌트
  • @KafkaListener 어노테이션이 선언된 add 메서드 선언
  • ConcurrentKafkaListenerContainerFactory 빈으로 부터 생성된 리스너 컨테이너가 Kafka 서버의 point-add 토픽의 파티션에 메세지가 수신되면 add 메서드 호출
  • add 메서드는 point-add 토픽에 대한 비즈니스 로직 정의
  @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());
    }
  }

3️⃣ 수행결과

1. 수행 로그

  1. DiaryController.saveDiary
  2. KafkaProducer.send
    DiaryController-KafkaProducer
  1. KafkaConsumer.add
    KafkaConsumer

2. Kafka-Console-Consumer

Kafka-Console-Consumer

3. DB 확인

DB

★ Github

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

profile
https://github.com/onlydev7777

0개의 댓글