데보션앱 소개페이지 바로가기
로그인 선택

신고하기

CLOSE
신고사유 (대표 사유 1개)
상세내용 (선택)
0/200
  • 신고한 게시글은 더 이상 보이지 않습니다.
  • 이용약관과 운영정책에 따라 신고사유에 해당하는지 검토 후 조치됩니다.
  • 허위 신고인 경우, 신고자의 서비스 이용이 제한될 수 있으니 유의하시어 신중하게 신고해 주세요.
(이 회원이 작성한 모든 댓글과 커뮤니티 게시물이 보이지 않고, 알림도 오지 않습니다.)

미리보기

커뮤니티

      1,234

      badge 23.06.15

      글 등록

      카테고리를 선택해주세요.

      DEVOTEE를 활성화 시키면
      지금 작성한 커뮤니티 글에 대해 1개의 댓글을 달아줍니다.

      버튼을 누르면 글 수정 시 ChatGPT가 작성한 댓글이 수정됩니다.

      임시저장함에 저장되었습니다. 저장일시 : 2022.5.17 14:29:08

      임시저장함

      제목을 선택하시면 이어서 작성이 가능하며,
      최대 20건까지 저장합니다.
      컨텐츠 유형, 제목, 저장일시, 삭제로 이뤄진 임시저장 목록
      컨텐츠 유형 제목 저장일 삭제

      데보션 블로그 게재 요청

      CLOSE
      • *
      • *

      본인인증

      효율적인 데보션 서비스 이용 및
      고객님의 소중한 개인정보보호를 위해
      본인인증을 진행해주세요. 본인인증 미 진행 시 로그인이 제한됩니다.
      본인인증 실패

      본인인증 로그인에 실패하였습니다.
      회원이 아니시거나 본인인증 등록이
      완료되지 않은 사용자입니다.

      회원정보 연결

      OpenLab - Kotlin 6회차, Kafka

      kchabin 24.07.21
      618 4 1
      DEVOTEE 요약
      이번 블로그는 오픈랩 코틀린 스터디 마지막 6회차에서 배운 Kafka 내용을 정리하고 회고하려고 합니다. Kafka는 로그 기반의 분산 스트리밍 플랫폼으로, 대용량의 실시간 데이터 처리가 가능하고, 이를 SpringBoot 프로젝트에서 사용하는 방법도 학습했습니다. 스터디 참여를 통해 Spring 공부에 방향을 잡을 수 있었고, 많은 새로운 기술들을 배우며 성장할 수 있었습니다.
      DEVOTEE 추천 블로그

      안녕하세요, kchabin이란 이름으로 활동 중인 강다빈입니다.😎

      저번엔 오픈랩 코틀린 스터디 5회차 후기글을 작성했었는데요, 이번엔 마지막 6회차 스터디에서 배운 Kafka 내용 정리와 회고를 해보고자 합니다.

      20240709_195948.jpg


      Kafka

      image.png

      Kafka는 로그 기반의 분산 스트리밍 플랫폼입니다.

      Kafka와 RabbitMQ의 차이점

      둘 다 스트림 처리에 사용할 수 있는 메시지 대기열 시스템입니다. Kafka는 대용량의 실시간 데이터를 처리하고 저장하는데 중점을 두고, RabbitMQ는 메시지를 송신자에서 수신자로 전달하는데 중점을 둡니다.

      (AWS:Kafka와 RabbitMQ의 차이점은 무엇인가요?)

      image.png

      이미지 출처 : Part 1: RabbitMQ for beginners - What is RabbitMQ?

      RabbitMQ로 전달된 메시지는 queue로 전달되지 않고, 전화교환원 같은 exchange로 전달됩니다.

      여기서 key에 따라서 어떤 큐로 넣을지 달라지고, 이 exchange가 SOPF가 될 수 있어 이를 안 두는 방식으로 쓰기 위해 Kafka를 사용한다는 걸 지난 스터디에서 간단히 배웠습니다.


      Kafka는 RabbitMQ와 메시지 큐 전달 방식이 다릅니다.

      스크린샷 2024-07-21 오후 6.56.08.png

      • Broker : producer가 consumer에게 데이터를 스트리밍하게 해주는 kafka의 서버.

      • Topic = 카테고리로, 메시지 송수신용 통로 역할을 합니다. broker 내부에 존재합니다.

      • zookeeper는 kafka 메타데이터와 클러스터 상태를 관리합니다.

      • Partition : Topic 내부엔 여러 개의 partition이 존재합니다. kafka producer가 write 동작을 할 때, producer가 생성한 레코드는 broker -> topic -> partition 순으로 이동하여 저장됩니다.

        • partition 개수는 얼마만큼 병렬 처리를 할 것인가에 따라서 달라집니다.

      kafka의 주요 기능들을 좀 더 알아보겠습니다.

      Offset

      image.png

      Consumer가 읽은 곳에 offset이 가 있게 됩니다. 읽어가는 위치를 찍어놓고, 메시지가 들어올 때마다 producer offeset은 최신 위치를 가리키게 됩니다.

      스크린샷 2024-07-21 오후 7.12.21.png

      • kafka는 offset을 통해서 레코드를 읽게 됩니다.

      • READ : 레코드를 파티션으로부터 읽어오기

      • PROCESS : 특정 처리 수행

      • COMMIT : 처리 완료 후 커밋을 브로커로 보내서 처리 완료 알림

      Partitioner

      image.png

      partition 내부는 큐같아서, 순서를 보장하지만 partition끼리는 순서가 보장되지 않습니다.

      그림처럼, partitioner가 producer가 생성한 레코드를 지정된 토픽의 특정 파티션으로 전달할 때 key의 유무에 따라 방식이 달라집니다.

      Key가 있다

      해시 값으로 결정

      Key가 없다

      라운드로빈으로 전달 -> 순서 보장 x

      hash값에 따라 파티셔닝을 할 때, 동일한 해시 결과가 나오는 키를 가진 레코드는 동일한 파티션에만 들어갑니다.

      이 때문에 해시값을 사용하는 방식은 파티션 내 레코드의 순서가 보장됩니다.


      Kafka with SpringBoot

      Kafka를 스프링부트 프로젝트에서 사용하려면 어떻게 해야하는지도 간단히 배웠습니다.

      dependecies:

      • Spring web

      • Lombok

      • Spring for Apache Kafka

      스터디에선 docker-compose.yaml 파일을 만들어 zookeeper, kafka broker를 설정했습니다.

      의존성 설정 -> 애플리케이션 설정 -> kafka producer 설정 -> 메시지 객체 생성 -> ProducerController 작성 -> kafka consumer 설정 -> kafka listener 설정 등이 이뤄져야 합니다.

      의존성 설정

       implementation ('org.springframework.kafka:spring-kafka:2.8.1')

      Kafka 접속 설정

      # Kafka endpoint for endpoint
       kafka.bootstrap-servers=localhost:29092,localhost:29093,localhost:29094
      • bootstrap-servers : broker를 위한 endpoint 경로

      KafkaAdmin

      Spring Kafka에서 제공하는 클래스로, Topic을 생성 및 수정, 삭제하는 역할을 합니다.

      기본적으로 localhost:9092에 접속하도록 설정되어있어 bootstrap-servers로 접속하도록 수정해야 합니다.

      @Configuration
      public class KafkaAdminConfig {
      
          @Value("${kafka.bootstrap-servers}")
          private String bootstrapServer;
      
          @Bean
          public KafkaAdmin kafkaAdmin() {
              Map<String, Object> configs = new HashMap<>();
              configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer);
      
              return new KafkaAdmin(configs);
          }
      }

      Topic 생성

       @Autowired
       private KafkaAdmin kafkaAdmin;
       
      //토픽 생성 메서드
       private NewTopic defaultTopic() {
           return TopicBuilder.name(DEFAULT_TOPIC)
                      .partitions(2)
                      .replicas(2)
                      .build();
          }
      
      private NewTopic topicWithKey() {
          return TopicBuilder.name(TOPIC_WITH_KEY)
                      .partitions(2)
                      .replicas(2)
                      .build();
          }
      • defaultTopic() : key가 없는 토픽을 생성하는 메서드.

      • topicWithKey() : 키가 있는 토픽 객체를 생성하는 메서드.

      Producer 생성

      @Configuration
      public class KafkaProducerConfig {
      
          @Value("${kafka.bootstrap-servers}")
          private String bootstrapServer;
      
          private ProducerFactory<String, Object> producerFactory() {
              Map<String, Object> configProps = new HashMap<>();
              configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer);
              configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
              configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
      
              return new DefaultKafkaProducerFactory<>(configProps);
          }
      
          @Bean
          public KafkaTemplate<String, Object> kafkaProducerTemplate() {
              return new KafkaTemplate<>(producerFactory());
          }
      }
      • @Value 로 접속 설정 정보를 가져온다.

      • ProducerFactory : 키로 String, 값은 Object를 갖도록 팩토리를 지정합니다.

      • Broker의 엔드포인트, 전송될 key, 전송될 값을 직렬화하는 방법 설정. 보통 JsonSerializer가 많이 쓰입니다.

      • KafkaTemplate : Spring Kafka에서 Kafka 브로커로 메시지를 전송하는 데 사용되는 핵심 클래스입니다.

        • 주요 메서드로 send, sendDefault 등 토픽에 데이터를 전송하는 메서드들이 있습니다.

      메시지 객체

      Kafka에 produce하고 consume할 매시지 객체

      @Data
      @NoArgsConstructor
      @AllArgsConstructor
      public class testEntity {
          private String title;
          private String contents;
          private LocalDateTime time;
      }

      ProducerController

      메시지를 전송하기 위한 컨트롤러 입니다.

      @Slf4j
      @RestController
      @RequestMapping("/api")
      public class ProducerController {
      
          // kafka producer를 위한 KafkaTemplate를 지정
          private final KafkaTemplate<String, Object> kafkaProducerTemplate;
      
          public ProducerController(KafkaTemplate<String, Object> kafkaProducerTemplate) {
              this.kafkaProducerTemplate = kafkaProducerTemplate;
          }
      
          @PostMapping("produce")
          public ResponseEntity<?> produceMessage(@RequestBody TestEntity testEntity) {
              testEntity.setTime(LocalDateTime.now());
      
              // kafkaProducerTemplate.send를 이용하여 메시지를 전송
              // send(String topic, String data)은 지정된 토픽에 데이터를 전송
              // ListenableFuture로 전송 결과 확인
              ListenableFuture<SendResult<String, Object>> future = kafkaProducerTemplate.send(KafkaTopicConfig.DEFAULT_TOPIC, testEntity);
      
              // 메시지 처리는 비동기로 처리 
              future.addCallback(new ListenableFutureCallback<SendResult<String, Object>>() {
                  @Override
                  public void onFailure(Throwable ex) {
                      log.error("Fail to send message to broker: {}", ex.getMessage());
                  }
      
                  @Override
                  public void onSuccess(SendResult<String, Object> result) {
                      log.info("Send message with offset: {}, partition: {}", result.getRecordMetadata().offset(), result.getRecordMetadata().partition());
                  }
              });
      
              return ResponseEntity.ok(testEntity);
      
          }
      }

      Listenable Future

      • Google의 Guava 라이브러리

      • 비동기 처리의 결과를 나타내는 future의 확장판

      • 기본 future는 스레드가 대기하고 있다가 응답이 오면 변환합니다.

      • 콜백을 등록하여 비동기 작업이 완료되었을 때 특정 동작을 수행할 수 있습니다.

      future 는 ListenableFuture<SendResult<String, Object>> 타입의 객체로, Kafka에 메시지를 전송하는 작업의 결과를 저장합니다.

      addCallBack() 으로 작업 완료 결과를 처리하는 콜백을 등록하는데 ListenableFutureCallback 인터페이스가 제공하는

      메서드 onSuccess ,onFailure 를 사용해서 정상적인 케이스와 비정상 케이스로 나누어 처리되도록 합니다.

      • onSuccess -> 메시지의 오프셋과 파티션을 로그에 기록

      • onFailure -> 실패 메시지를 로그에 기록

      Kotlin에선 Coroutine 을 사용합니다.

      • 순차적인 코드 스타일 -> 좋은 가독성.

      • async와 await를 사용하여 비동기 작업을 쉽게 조합 가능

      • CoroutineScope -> 구조적 동시성

      import kotlinx.coroutines.*
      
      fun main() = runBlocking {
          val result = async {
              // 비동기 작업 수행
              "Result"
          }
          // 결과 처리
          println(result.await())
      }

      Consumer 설정

      @EnableKafka
      @Configuration
      public class KafakConsumerConfig {
      
          @Value("${kafka.bootstrap-servers}")
          private String bootstrapServer;
      
          private ConsumerFactory<String, Object> consumerFactory(String groupId) {
              Map<String, Object> props = new HashMap<>();
      
              props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer);
              props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
      
              JsonDeserializer<Object> jsonDeserializer = new JsonDeserializer<>();
              // Deserialize에 대해서 신뢰하는 패키지를 지정한다. "*"를 지정하면 모두 신뢰하게 된다.
              jsonDeserializer.addTrustedPackages("*");
      
              return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), jsonDeserializer);
          }
      
          @Bean
          public ConcurrentKafkaListenerContainerFactory<String, Object> defaultKafkaListenerContainerFactory() {
              ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
              factory.setConsumerFactory(consumerFactory("defaultGroup"));
              factory.setConcurrency(1);
              factory.setAutoStartup(true);
              return factory;
          }
      }
      • @EnableKafka가 있어야 @KafkaListener가 활성화 됩니다.

      • ConsumerFactory는 consumer를 생성합니다.

      • ConcurrentKafkaListenerContainerFactory : 동시에 Kafka Cluster로부터 메시지를 읽을 수 있도록 합니다.

        • consumerFactory를 파라미터로 전달 -> consumer 생성

        • setConcurrency : 동시에 읽을 consumer 개수 지정

        • setAutoStartup : true로 설정하면 서버가 부트업될 때 자동 실행됩니다.


      Partition Key가 있는 토픽 생성

      메시지 키가 있어야 hashing 방식으로 파티션에 메시지를 할당하고,키에 따라 들어온 순서대로 메시지가 적재됩니다.

      스크린샷 2024-07-21 오후 8.42.55.png

      Key 할당을 위한 생성

      # 토픽 키를 이용 설정
      kafka.topic-with-key=topic-key

      Topic 설정

      package com.schooldevops.kafkatutorials.configs;
      
      import org.apache.kafka.clients.admin.NewTopic;
      import org.apache.kafka.common.config.TopicConfig;
      import org.springframework.beans.factory.annotation.Autowired;
      import org.springframework.beans.factory.annotation.Value;
      import org.springframework.context.annotation.Configuration;
      import org.springframework.kafka.config.TopicBuilder;
      import org.springframework.kafka.core.KafkaAdmin;
      
      import javax.annotation.PostConstruct;
      
      @Configuration
      public class KafkaTopicConfig {
      
          public final static String DEFAULT_TOPIC = "DEF_TOPIC";
      
          @Value("${kafka.topic-with-key}")
          public String TOPIC_WITH_KEY;
      
          @Autowired
          private KafkaAdmin kafkaAdmin;
      
          private NewTopic defaultTopic() {
              return TopicBuilder.name(DEFAULT_TOPIC)
                      .partitions(2)
                      .replicas(2)
                      .build();
          }
      
          private NewTopic topicWithKey() {
              return TopicBuilder.name(TOPIC_WITH_KEY)
                      .partitions(2)
                      .replicas(2)
                      .build();
          }
      
          @PostConstruct
          public void init() {
              kafkaAdmin.createOrModifyTopics(defaultTopic());
              kafkaAdmin.createOrModifyTopics(topicWithKey());
          }
      }
      • createOrModifyTopics : 어플리케이션이 기동될때 토픽이 있다면 수정, 없다면 새로 생성

      Listener 생성

      ... 생략 
          @KafkaListener(topics = "${kafka.topic-with-key}", containerFactory = "defaultKafkaListenerContainerFactory")
          public void listenTopicWithKey(Object record) {
              log.info("Receive Message from {}, values {} with key", record);
      
          }
      ... 생략 
      • 팩토리는 꼭 쓰지 않아도 자동으로 생성됩니다.

      • 메시지를 수신하면 단순하게 로깅하는 코드.

      ProducerController

      default 토픽과 kafkaProducerTemplate의 send 메서드에 전달하는 값이 달라집니다. key 를 전달해서 키에 의한 파티셔닝이 이뤄지도록 합니다.

      ListenableFuture<SendResult<String, Object>> future = kafkaProducerTemplate.send(TOPIC_WITH_KEY, key, testEntity);


      스터디 회고

      대학 입학 후 약 2년을 진로 선택을 고민하며 보냈습니다.

      데보션영 2기 활동을 통해 개발자라는 직업에 진지하게 고민해볼 수 있게 되었고,

      많은 분들이 데보션 커뮤니티나 블로그, 기술 행사에서 지식을 공유하고 네트워킹 하는 모습들에 반해 올해 초부터 본격적으로 백엔드 공부를 시작했습니다.

      혼자서 하면서도 이게 맞는 건지 잘 모르겠는 것들이 많았을 때OpenLab 스터디 참여 기회를 얻게 되었습니다.

      스터디를 통해 막막하던 spring 공부에 어느정도 길이 보이기 시작했고, api 작성 경험, 테스트 코드 성공시키기 경험, 새롭게 알게 된 exposed, redis, kafka 등 많이 배우고 성장할 수 있었습니다.

      그리고 이 성장엔 아직 부족한 점이 너무 많은 학생과 진행하는 스터디임에도 열심히 알려주시고 도와주신 좋은 스터디장 배기도님과 다른 스터디원분들 덕분이었습니다!

      이번에 이렇게 좋은 기회로 배운 것들, 만난 분들 다 잊지 않고 열심히 공부해서 저도 언젠간 저 같은 학생에게 도움이 될 수 있는 멋진 소프트웨어 개발자가 되겠습니다. 😊

      댓글 0

      DEVOTEE를 활성화 시키면
      지금 작성한 댓글에 AI가 댓글을 달아줍니다.

      kchabin 님의 최신 블로그

      더보기

      DEVOTEE 추천 블로그

      동영상 기고하기