kafka quickstart

Description

링크드인에서 처음 개발된 분산 메시징 시스템

 

구조는 Pub-Sub 모델을 가지며, Publish-Subscribe 모델은 데이터를 만들어 내는 프로듀서와 소비하는 컨슈머, 그리고 이 둘 사이 중재 역할을 하는 Broker 로 구성된 느슨한 결합(loosely Coupled)의 시스템

 

프로듀서는 브로커를 통해 메시지를 발행함(Publish)

메시지 전달 대상을 명시하지는 않고, 메시지를 구독할 컨슈머가 브로커에 요청해서 가져가는 방식.

즉, 블로그를 작성하여 발행하면 블로그 글을 구독하는 구독자들이 따로 읽어가는 형태

카프카 클러스터로 메시지를 전송하는 프로듀서 API 와 메시지를 읽어갈 수 있는 컨슈머 API 도 제공됨.

카프카에서 프로듀서는 특정 Topic 으로 데이터 전달이 가능하고, 컨슈머는 특정 Topic 에서 읽어갈 수 있다.

카프카는 Scale Out 형태로 클러스터를 구성하며, 유통되는 데이터가 늘어나면 브로커의 부담(Load) 를 줄여야 되는데 이를 Apache Zookeeper 를 통해 진행한다.

카프카에 전송되는 프로듀서에 의해 발행된 메시지는 중복 저장되고, 장애가 발생해도 HA(고가용성)를 보장함.

프로듀서가 카프카에 메시지를 보내면, 또 다른 브로커가 프로듀서의 메시지를 중복해서 저장함.

하나의 브로커에서 장애가 발생해도, 중복 저장본을 컨슈머에게 전달할 수 있다.

 

features

  • 다중 프로듀서, 다중 컨슈머 (N:M)
  • 파일 시스템에 저장
    • 전통 메시지큐는 메모리상에 데이터를 저장 (생성된 메시지가 빠르게 소비될 것으로 간주)
    • 카프카는 메시지를 파일 시스템에 저장함.
    • 만약, 컨슈머에서 장애가 생기거나 네트워크 트래픽이 폭주해도 브로커에 영향을 주지 않고 처리 속도를 따라 갈 수 있다.
    • 컨슈머들이 데이터를 모았다가 처리하는 배치도 진행 가능
    • 컨슈머 쪽에서 장애가 발생할 경우, 이전 데이터로 부터 읽어올 수 있음(전에 읽었던 데이터라도)
  • 확장성
    • Scale Out 가능
    • 카프카 토픽에 메시지를 전달하는 프로듀서 역시 운영 중 증대 가능.
    • 카프카 토픽은 내부 파티션을 통해 세분화된 단위로 저장하며, 토픽 파티션 개수도 지정 가능
  • 고성능
    • 논문에 의하면, 카프카는 대용량 실시간 로그 처리에 특화됨.
    • activemq, rabbitmq 와 비교하여 탁월한 throughput 을 보장 (producer 성능)
    • activemq, rabbitmq 와 비교하여 탁월한 throughput 을 보장 (consumer 성능)
    • 특히, producer 쪽에서 성능 차이가 더 심함 (producer 는 몇백배, consumer 는 5~10배)
  • 컨슈머 Pull 방식
    • 기존 메시지큐 방식은 브로커가 컨슈머에게 메시지를 push 하는 방식.
    • 카프카는 컨슈머가 브로커로부터 메시지를 스스로 가져오는 pull 방식
      • 컨슈머의 처리량을 broker 가 고민할 필요가 없다!

terminology

  • topic (메시지 스트림 추상화)
    • producer 와 consumer 가 만나는 지점.
  • partition
    • producer 가 topic 에 메시지를 전달하면, 카프카 클러스터는 토픽을 좀 더 세분화한 단위로 관리.
    • producer 는 메시지가 어떤 파티션에 저장되는지 알 필요 없음.
    • 각 파티션은 카프카 클러스터를 구성하는 브로커들이 나눠가짐.
    • 특정 파티션으로 전달된 메시지는 offset 을 가짐.
      • offset : 해당 파티션에서 몇번째 메시지인지를 의미하는 ID 개념
    • 복제 개념
      • 고가용성 보장을 목적.
      • topic 생성 시 replication factor 설정을 통해, 파티션의 복제본을 저장할 수 있다.
      • replication factor - n 개로 설정하면, n 개의 파티션 데이터 복사본이 생성되고, 카프카 브로커가 겹치지 않게 나눠가짐 (n개 복사본은 replica 라고 하고, 1개가 리더가 되어 클라이언트 요청 담당)
      • 리더 변경사항을 따르며, 복제본을 담당하는 팔로워는 ISR(In-Sync Replica) 라고 함.
      • 만약 리더가 장애가 난 경우, 팔로워들 중 하나가 클라이언트 요청을 담다.ㅇ
    • 파티션의 리더와 팔로워는 다른 브로커에 할당해야 함.
      • 고가용성 보장을 위해서.
  • producer, consumer, consumer group
    • producer : 메시지 생성해서 카프카에 전달하는 client
      • 유의할 점은 서로 다른 파티션에 전송된 메시지는 순서 보장이 안됨.
      • 첫번째 메시지가 0번 파티션, 두번째 메시지가 1번 파티션인 경우, 두번째 메시지가 먼저 소비될 수 있음.
      • 카프카로 전송된 메시지는 같은 파티션인 경우에만 순서 보장됨.
    • consumer : 메시지를 카프카로부터 읽어가는 client
      • consumer 는 자체 그룹을 형성한다. (consumer group)
      • topic 의 파티션은 group 당 하나의 컨슈머의 소비만 소비될 수 있음.
      • 파티션와 컨슈머의 연결을 소유권(ownership) 이라고 함.
      • group 에서 컨슈머 추가/제거가 발생한 경우, 리밸런싱이 이뤄짐.
      • 컨슈머가 파티션 수보다 많은 경우, 파티션 개수만큼의 컨슈머만 동작하고 나머지는 논다.
      • group 은 각 파티션에 대한 offset 을 가짐.
      • group 에 할당된 offset 변경하는 작업을 offset commit 이라고 함.
  • segment
    • 카프카로 전송된 메시지는 클러스터 내부에서 세그먼트 파일 형태로 저장됨.
    • 파일 형태로 저장하여 메시지의 영속성을 얻게됨.
    • 나중에 메시지를 재소비할 경우, 파일 시스템에 저장된 메시지를 읽어서 컨슈머에 전달.
    • 파일 시스템은 영속성이 보장되지만 속도가 느린 문제가 존재
      • 시간이 흐르며, OS와 디스크의 최적화가 지행되고 순차읽기와 미리읽기 를 통해 제법 빠르게 사용 가능.
    • 파일 시스템에 세그먼트를 쓸 때 OS의 페이지 캐시를 사용.
      • 사용하지 않은 메모리를 파일 시스템의 캐시로 사용.
      • 사용자가 요청 안해도 미리 읽기 (Read Ahead) 를 통해 앞으로 읽을 가능성 있는 뒤쪽 내용을 미리 메모리로 읽는 최적화를 진행.
    • 카프카 내부에서는 버퍼 캐시를 하지 않아 JVM 에서 발생하는 GC(Garbage Collection) 오버헤드도 줄인다.

 

Microservice 로 오는 각 서비스 로그 활용을 위해서 kafka 를 사용할 예정인데,

kafka tutorial 을 아래와 같이 진행한다.

Tutorial

Go to link

chapter 3 - create a topic to store events

  • 이벤트는 예를 들어, 아래와 같은 것들이 있음
    • 지불 내역
    • 위치 변경
    • 주문 목록
    • 센서 측정
  • topic 이 폴더, events 가 파일

아래 코드는 topic 하나 생성한 것임.

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

response message

Created topic quickstart-events.

토픽 생성 확인

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

response message (partition count 확인 가능)

Topic: quickstart-events        TopicId: iCSsJU1yTCSKOuAPqdzBOw PartitionCount: 1       ReplicationFactor: 1    Configs: segment.bytes=1073741824
        Topic: quickstart-events        Partition: 0    Leader: 0       Replicas: 0     Isr: 0

chapter 4 - write some events to topic

console-producer.sh 파일을 통해 quickstart-events 토픽 내 아래 문장을 전달

  • this is my first event
  • this is my second event
  • this is my third event
  • tutorial is finished
$ bin/kafka-console-producer.sh --topic quickstart-events --bootstrap-server localhost:9092

this is my first event
>this is my second event
>this is my third event
>tutorial is finished

chapter 5 - read the events

read events from topic (quickstart-events)

  • --from-beginning 이 시작부터 끝까지 다 읽는다는 의미인듯
$ bin/kafka-console-consumer.sh --topic quickstart-events --from-beginning --bootstrap-server localhost:9092

response

this is my first eent
this is my second event
this is my third event
tutorial is finished
^CProcessed a total of 4 messages

chapter 8 - terminate the kafka environment

  • producer, consumer 중지
  • kafka broker 중지
  • zookeeper server 중지

만약, 로컬 환경에서 데이터까지 삭제한다고 하면

$ rm -rf /tmp/kafka-logs /tmp/zookeeper

 

 

참고

https://soft.plusblog.co.kr/3

 

[Kafka] #1 - 아파치 카프카(Apache Kafka)란 무엇인가?

데이터 파이프라인(Data Pipeline)을 구축할 때 가장 많이 고려되는 시스템 중 하나가 '카프카(Kafka)' 일 것이다. 아파치 카프카(Apache Kafka)는 링크드인(LinkedIn)에서 처음 개발된 분산 메시징 시스템이

soft.plusblog.co.kr

 

+ Recent posts