23.06.15
DEVOTEE를 활성화 시키면
지금 작성한 커뮤니티 글에 대해 1개의 댓글을 달아줍니다.
버튼을 누르면 글 수정 시 ChatGPT가 작성한 댓글이 수정됩니다.
| 컨텐츠 유형 | 제목 | 저장일 | 삭제 |
|---|
본인인증 로그인에 실패하였습니다.
회원이 아니시거나 본인인증 등록이
완료되지 않은 사용자입니다.
안녕하세요, kchabin이란 이름으로 활동 중인 강다빈입니다.😎
저번엔 오픈랩 코틀린 스터디 5회차 후기글을 작성했었는데요, 이번엔 마지막 6회차 스터디에서 배운 Kafka 내용 정리와 회고를 해보고자 합니다.
Kafka는 로그 기반의 분산 스트리밍 플랫폼입니다.
둘 다 스트림 처리에 사용할 수 있는 메시지 대기열 시스템입니다. Kafka는 대용량의 실시간 데이터를 처리하고 저장하는데 중점을 두고, RabbitMQ는 메시지를 송신자에서 수신자로 전달하는데 중점을 둡니다.
(AWS:Kafka와 RabbitMQ의 차이점은 무엇인가요?)
이미지 출처 : Part 1: RabbitMQ for beginners - What is RabbitMQ?
RabbitMQ로 전달된 메시지는 queue로 전달되지 않고, 전화교환원 같은 exchange로 전달됩니다.
여기서 key에 따라서 어떤 큐로 넣을지 달라지고, 이 exchange가 SOPF가 될 수 있어 이를 안 두는 방식으로 쓰기 위해 Kafka를 사용한다는 걸 지난 스터디에서 간단히 배웠습니다.
Kafka는 RabbitMQ와 메시지 큐 전달 방식이 다릅니다.
Broker : producer가 consumer에게 데이터를 스트리밍하게 해주는 kafka의 서버.
Topic = 카테고리로, 메시지 송수신용 통로 역할을 합니다. broker 내부에 존재합니다.
zookeeper는 kafka 메타데이터와 클러스터 상태를 관리합니다.
Partition : Topic 내부엔 여러 개의 partition이 존재합니다. kafka producer가 write 동작을 할 때, producer가 생성한 레코드는 broker -> topic -> partition 순으로 이동하여 저장됩니다.
partition 개수는 얼마만큼 병렬 처리를 할 것인가에 따라서 달라집니다.
kafka의 주요 기능들을 좀 더 알아보겠습니다.
Consumer가 읽은 곳에 offset이 가 있게 됩니다. 읽어가는 위치를 찍어놓고, 메시지가 들어올 때마다 producer offeset은 최신 위치를 가리키게 됩니다.
kafka는 offset을 통해서 레코드를 읽게 됩니다.
READ : 레코드를 파티션으로부터 읽어오기
PROCESS : 특정 처리 수행
COMMIT : 처리 완료 후 커밋을 브로커로 보내서 처리 완료 알림
partition 내부는 큐같아서, 순서를 보장하지만 partition끼리는 순서가 보장되지 않습니다.
그림처럼, partitioner가 producer가 생성한 레코드를 지정된 토픽의 특정 파티션으로 전달할 때 key의 유무에 따라 방식이 달라집니다.
Key가 있다 | 해시 값으로 결정 |
|---|---|
Key가 없다 | 라운드로빈으로 전달 -> 순서 보장 x |
hash값에 따라 파티셔닝을 할 때, 동일한 해시 결과가 나오는 키를 가진 레코드는 동일한 파티션에만 들어갑니다.
이 때문에 해시값을 사용하는 방식은 파티션 내 레코드의 순서가 보장됩니다.
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 endpoint for endpoint
kafka.bootstrap-servers=localhost:29092,localhost:29093,localhost:29094bootstrap-servers : broker를 위한 endpoint 경로
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);
}
} @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() : 키가 있는 토픽 객체를 생성하는 메서드.
@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;
}메시지를 전송하기 위한 컨트롤러 입니다.
@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 FutureGoogle의 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())
}@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로 설정하면 서버가 부트업될 때 자동 실행됩니다.
메시지 키가 있어야 hashing 방식으로 파티션에 메시지를 할당하고,키에 따라 들어온 순서대로 메시지가 적재됩니다.
# 토픽 키를 이용 설정
kafka.topic-with-key=topic-keypackage 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 : 어플리케이션이 기동될때 토픽이 있다면 수정, 없다면 새로 생성
... 생략
@KafkaListener(topics = "${kafka.topic-with-key}", containerFactory = "defaultKafkaListenerContainerFactory")
public void listenTopicWithKey(Object record) {
log.info("Receive Message from {}, values {} with key", record);
}
... 생략 팩토리는 꼭 쓰지 않아도 자동으로 생성됩니다.
메시지를 수신하면 단순하게 로깅하는 코드.
default 토픽과 kafkaProducerTemplate의 send 메서드에 전달하는 값이 달라집니다. key 를 전달해서 키에 의한 파티셔닝이 이뤄지도록 합니다.
ListenableFuture<SendResult<String, Object>> future = kafkaProducerTemplate.send(TOPIC_WITH_KEY, key, testEntity);대학 입학 후 약 2년을 진로 선택을 고민하며 보냈습니다.
데보션영 2기 활동을 통해 개발자라는 직업에 진지하게 고민해볼 수 있게 되었고,
많은 분들이 데보션 커뮤니티나 블로그, 기술 행사에서 지식을 공유하고 네트워킹 하는 모습들에 반해 올해 초부터 본격적으로 백엔드 공부를 시작했습니다.
혼자서 하면서도 이게 맞는 건지 잘 모르겠는 것들이 많았을 때OpenLab 스터디 참여 기회를 얻게 되었습니다.
스터디를 통해 막막하던 spring 공부에 어느정도 길이 보이기 시작했고, api 작성 경험, 테스트 코드 성공시키기 경험, 새롭게 알게 된 exposed, redis, kafka 등 많이 배우고 성장할 수 있었습니다.
그리고 이 성장엔 아직 부족한 점이 너무 많은 학생과 진행하는 스터디임에도 열심히 알려주시고 도와주신 좋은 스터디장 배기도님과 다른 스터디원분들 덕분이었습니다!
이번에 이렇게 좋은 기회로 배운 것들, 만난 분들 다 잊지 않고 열심히 공부해서 저도 언젠간 저 같은 학생에게 도움이 될 수 있는 멋진 소프트웨어 개발자가 되겠습니다. 😊
DEVOTEE를 활성화 시키면
지금 작성한 댓글에 AI가 댓글을 달아줍니다.