[Kafka] Kafka 설치/실행 및 테스트

카프카 도큐먼트 번역

설치 및 실행

ref. https://kafka.apache.org/documentation/#quickstart

1️⃣ 설치와 kafka server 실행하기
아래 사이트에서 Binary downloads에 있는 파일을 다운받는다.
https://kafka.apache.org/downloads

다운로드 받은 파일은 적절한 위치에 압축을 풀어준다.

$ tar -xzf kafka_2.13-2.6.0.tgz

압축이 해제되면 동명의 폴더가 생길 것이다. kafka는 이 폴더 내 bin 하위 스크립트 파일을 읽음으로써 실행된다.

$ cd kafka_2.13-2.6.0

kafka는 zookeeper 위에서 돌아가므로 zookeeper를 먼저 실행한다.

$ bin/zookeeper-server-start.sh config/zookeeper.properties

다음은 kafka를 실행한다.

$ bin/kafka-server-start.sh config/server.properties

2️⃣ Topic 생성하기
localhost:9092 카프카 서버에 quickstart-events란 토픽을 생성한다.

$ bin/kafka-topics.sh --create --topic quickstart-events --bootstrap-server localhost:9092

현재 만들어져 있는 토픽 확인하기

$ bin/kafka-topics.sh --list --bootstrap-server localhost:9092

특정 토픽의 설정 확인하기

$ bin/kafka-topics.sh --describe --topic quickstart-events --bootstrap-server localhost:9092

3️⃣ Producer, Consumer 실행하기
콘솔에서 Producer와 Consumer를 실행하여 실시간으로 토픽에 event를 추가하고 받을 수 있다.
터미널을 분할로 띄워서 진행해본다.

Producer

bin/kafka-console-producer.sh --topic quickstart-events --bootstrap-server localhost:9092

Consumer

bin/kafka-console-consumer.sh --topic quickstart-events --from-beginning --bootstrap-server localhost:9092

경축! 아무것도 안하여 에스천사게임즈가 새로운 모습으로 재오픈 하였습니다.
어린이용이며, 설치가 필요없는 브라우저 게임입니다.
https://s1004games.com

테스트

(IntelliJ + Gradle + Spring boot 환경에서 진행합니다.)

kafka 종속 추가

implementation 'org.springframework.kafka:spring-kafka'

application.yml

server:
  port: 8082

spring:
  application:
    name: third-point
  kafka:
    consumer:
      group-id: my-test
      enable-auto-commit: true
      auto-offset-reset: latest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      max-poll-records: 1000
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
    template:
      default-topic: quickstart-events
    bootstrap-servers: localhost:9092
  zipkin:
    sender:
      type: kafka

kafkaConfig

@Configuration
@ConfigurationProperties("spring.kafka")
@Data
public class KafkaConfig {

    private String bootstrapServers;
    private Producer producer;
    private Consumer consumer;
    private Template template;

    @Data
    public static class Producer {
        private String keySerializer;
        private String valueSerializer;
    }

    @Data
    public static class Template {
        private String defaultTopic;
    }

    @Data
    public static class Consumer {
        private String groupId;
        private String enableAutoCommit;
        private String autoOffsetReset;
        private String keyDeserializer;
        private String valueDeserializer;
        private String maxPollRecords;
    }



}

CustomConsumer

@Service
@RequiredArgsConstructor
@Slf4j
public class CustomConsumer {

    private final KafkaConfig config;
    private KafkaConsumer<String, String> consumer = null;

    @PostConstruct
    public void build(){
        Properties properties = new Properties();
        properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, config.getBootstrapServers());
        properties.setProperty(ConsumerConfig.GROUP_ID_CONFIG, config.getConsumer().getGroupId());
        properties.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, config.getConsumer().getKeyDeserializer());
        properties.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, config.getConsumer().getValueDeserializer());
        properties.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, config.getConsumer().getAutoOffsetReset());
        properties.setProperty(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, config.getConsumer().getMaxPollRecords());
        properties.setProperty(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, config.getConsumer().getEnableAutoCommit());
        consumer = new KafkaConsumer<>(properties);
    }

    @KafkaListener(topics = "${spring.kafka.template.default-topic}")
    public void consume(@Header(KafkaHeaders.RECEIVED_TOPIC) String topic, @Payload String payload){
        log.info("CONSUME TOPIC : "+topic);
        log.info("CONSUME PAYLOAD : "+payload);
    }

}

CustomProducer

@Service
@RequiredArgsConstructor
@Slf4j
public class CustomProducer {

    private final KafkaConfig config;
    private KafkaProducer<String, String> producer = null;

    @PostConstruct
    public void build(){
        Properties properties = new Properties();
        properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, config.getBootstrapServers());
        properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, config.getProducer().getKeySerializer());
        properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, config.getProducer().getValueSerializer());
        producer = new KafkaProducer<>(properties);
    }

    public void send(String message) {
        ProducerRecord<String, String> record = new ProducerRecord<>(config.getTemplate().getDefaultTopic(), message);

        producer.send(record, new Callback() {
            @Override
            public void onCompletion(RecordMetadata metadata, Exception exception) {
                log.info("publish message: {}", message);
                if (exception!=null){
                    log.info(exception.getMessage());
                }
            }
        });
    }

}

CustomProducerTest

CustomProducer를 사용해 메세지를 전송하는 테스트를 작성한다.

@SpringBootTest
class CustomProducerTest {

    @Autowired
    private CustomProducer sut;

    @Test
    void test1(){
        sut.send("this message sent from spring boot application!");
    }

}

kafka consumer를 실행한 뒤 테스트를 수행해보자.

테스트 코드의 Producer가 전송한 메세지를 Consumer가 받았음을 알 수 있다.

CustomConsumerTest

콘솔의 producer에서 메세지를 보내보자

IDE 내 콘솔에 다음과 같은 로그가 찍힐 것이다.

2020-08-23 15:27:02.340  INFO [third-point,3c74a95e62f7f6bd,9b79595f0e9204ec,true] 16164 --- [ntainer#0-0-C-1] com.example.kafka.CustomConsumer         : CONSUME TOPIC : quickstart-events
2020-08-23 15:27:02.340  INFO [third-point,3c74a95e62f7f6bd,9b79595f0e9204ec,true] 16164 --- [ntainer#0-0-C-1] com.example.kafka.CustomConsumer         : CONSUME PAYLOAD : to spring boot

ThirdController

Spring Boot Application의 CustomProducer와 CustomConsumer를 동시에 사용하도록 컨트롤러를 작성해보자.

@RestController
@RequestMapping("/third")
@RequiredArgsConstructor
@Slf4j
public class ThirdController {

    @Autowired
    private CustomProducer producer;

    @PostMapping("/publish")
    public void producer(@RequestBody Map<String, String> message){
        log.info(">>> start event publish");
        producer.send(message.get("msg"));
    }

}

터미널에서 요청을 보내보자

$ curl -XPOST 'http://localhost:8082/third/publish' -H 'Content-Type:application/json' -d '{"msg":"hello world!"}'

로그는 아래와 같다.

2020-08-23 15:21:36.098  INFO [third-point,d49790ceafd1d102,d49790ceafd1d102,true] 16164 --- [nio-8082-exec-7] com.example.controller.ThirdController   : >>> start event publish
2020-08-23 15:21:36.100  INFO [third-point,,,] 16164 --- [ad | producer-1] com.example.kafka.CustomProducer         : publish message: hello world!
2020-08-23 15:21:36.101  INFO [third-point,02f37a5c64e4874c,abda7d1fd8b18eea,false] 16164 --- [ntainer#0-0-C-1] com.example.kafka.CustomConsumer         : CONSUME TOPIC : quickstart-events
2020-08-23 15:21:36.101  INFO [third-point,02f37a5c64e4874c,abda7d1fd8b18eea,false] 16164 --- [ntainer#0-0-C-1] com.example.kafka.CustomConsumer         : CONSUME PAYLOAD : hello world!

CustomProducer를 통해 hello world!라는 메세지를 발행하고, 이 이벤트를 CustomConsumer측에서 받아오는 것이다.

 

[출처] https://velog.io/@hanblueblue/Kafka-%EC%84%A4%EC%B9%98%EC%8B%A4%ED%96%89-%EB%B0%8F-%ED%85%8C%EC%8A%A4%ED%8A%B8

 

 

 

본 웹사이트는 광고를 포함하고 있습니다.
광고 클릭에서 발생하는 수익금은 모두 웹사이트 서버의 유지 및 관리, 그리고 기술 콘텐츠 향상을 위해 쓰여집니다.
번호 제목 글쓴이 날짜 조회 수
공지 오라클 기본 샘플 데이터베이스 졸리운_곰 2014.01.02 86959
공지 [SQL컨셉] 서적 "SQL컨셉"의 샘플 데이타 베이스 SAMPLE DATABASE of ORACLE 가을의 곰을... 2013.02.10 79228
공지 [G_SQL] Sample Database 가을의 곰을... 2012.05.20 95967
101 [NoSQL] 성공적인 NoSQL 도입을 위한 키포인트 : NoSQL 데이터 모델링 file 졸리운_곰 2024.08.10 1082
100 [NoSQL] [MongoDB] NOSQL 데이터 모델링 기법 살펴보기 file 졸리운_곰 2024.08.09 1004
99 [NoSQL][MongoDB] Truncate a collection 졸리운_곰 2023.06.04 1273
98 [NoSQL] MongoDB 인증 모드 (password) 설정 졸리운_곰 2023.03.26 1422
97 [NoSQL] [Redis] Redis Persistence(영속성) 졸리운_곰 2021.04.11 1647
96 [NoSQL] [Cloud] Redis 설치, 사용 방법, 데이터 백업을 위한 RDB & AOF 개념 및 간단한 Redis 사용 사례 연구 file 졸리운_곰 2021.04.11 1719
95 [mongodb , 몽고디비] How to Use MongoDB Comparison Query Operators in Java 졸리운_곰 2021.02.19 1511
94 [mongodb, 몽고디비] How do you query for “is not null” in Mongo? 졸리운_곰 2021.02.19 1273
93 [MongoDB] 확장 검색 쿼리 - Aggregation 파이프라인 스테이지(Pipline Stage)와 표현식 졸리운_곰 2021.02.16 1192
92 [MongoDB] 확장 검색 쿼리 - 범용 Aggregation 졸리운_곰 2021.02.16 1353
91 [MongoDB] 확장 검색 쿼리 - Aggregation의 목적 및 작동방식 file 졸리운_곰 2021.02.16 1071
90 [MongoDB] 확장 검색 쿼리 - 집계 파이프 라인연산자 종류 졸리운_곰 2021.02.16 1143
89 [MongoDB] 확장 검색 쿼리 - Aggregation 파이프라인 스테이지(Pipline Stage)와 표현식 졸리운_곰 2021.02.16 1303
88 [mongodb, 몽고디비 쿼리] Mongodb query on substring of a field 졸리운_곰 2021.02.16 1308
87 [mongodb]Spring boot와 mongoDB 연동하기 그리고 REST API 졸리운_곰 2021.01.21 1709
86 [WebApp / Express] 간단한 MongoDB Middleware 만들기 졸리운_곰 2020.12.29 1155
85 [mongodb] Error: network error while attempting to run command 'isMaster' on host '127.0.0.1:27017' file 졸리운_곰 2020.12.16 1450
84 [mongodb] How to Update a Document in MongoDB using Java 졸리운_곰 2020.12.14 2340
83 [Java, MongoDB] Mapping a BSON MongoDB document to a MyClass.class object? 졸리운_곰 2020.12.09 1823
82 [mongodb] MongoDB CRUD 동작의 이해 졸리운_곰 2020.12.07 1293
대표 김성준 주소 : 경기 용인 분당수지 U타워 등록번호 : 142-07-27414
통신판매업 신고 : 제2012-용인수지-0185호 출판업 신고 : 수지구청 제 123호 개인정보보호최고책임자 : 김성준 sjkim70@stechstar.com
대표전화 : 010-4589-2193 [fax] 02-6280-1294 COPYRIGHT(C) stechstar.com ALL RIGHTS RESERVED