Kafka Streams 소개

Kafka Streams 소개

Kafka Streams는 미션 크리티컬한 실시간 애플리케이션과 마이크로서비스를 만들 때 아주 편한 클라이언트 라이브러리예요. 핵심 포인트는 "별도의 처리 클러스터가 필요 없다"는 거예요. 표준 Java/Scala 애플리케이션을 만드는 간단함에 카프카 서버 측 클러스터 기술의 장점이 더해진다고 보시면 돼요.

출처: 문서

본문

Kafka Streams

미션 크리티컬한 실시간 애플리케이션과 마이크로서비스를 쓰는 가장 쉬운 방법

Kafka Streams는 입력과 출력 데이터가 카프카 클러스터에 저장되는 애플리케이션과 마이크로서비스를 구축하기 위한 클라이언트 라이브러리예요. 클라이언트 측에서 표준 Java와 Scala 애플리케이션을 작성·배포하는 단순함과, 카프카의 서버 측 클러스터 기술의 이점을 결합해요.

Streams API 둘러보기

이 문서에는 Streams API를 소개하는 흥미로운 투어가 있어요: 1) Intro to Streams, 2) Creating a Streams Application, 3) Transforming Data Pt. 1, 4) Transforming Data Pt. 2.

왜 Kafka Streams를 좋아하게 될까요!

  • 탄력적(elastic), 고도로 확장 가능, 장애 허용
  • 컨테이너, VM, 베어메탈, 클라우드에 배포 가능
  • 소규모, 중규모, 대규모 사용 사례에 모두 적합
  • 카프카 보안과 완전히 통합
  • 표준 Java와 Scala 애플리케이션 작성
  • 정확히 한 번(exactly-once) 처리 의미론
  • 별도의 처리 클러스터 불필요
  • Mac, Linux, Windows에서 개발

Kafka Streams 사용 사례

  • The New York Times는 실시간으로 게시된 콘텐츠를 독자들이 접근할 수 있게 만드는 다양한 애플리케이션과 시스템에 저장·배포하기 위해 Apache Kafka와 Kafka Streams API를 사용해요.
  • Pinterest는 광고 인프라의 실시간 예측 예산 시스템을 구동하기 위해 대규모로 Apache Kafka와 Kafka Streams API를 사용해요. Kafka Streams 덕분에 지출 예측이 그 어느 때보다 정확해졌어요.
  • 유럽 선도 온라인 패션 리테일러인 Zalando는 모놀리식에서 마이크로서비스 아키텍처로 전환하는 데 도움이 되는 ESB(Enterprise Service Bus)로 Kafka를 사용해요. 이벤트 스트림 처리를 위해 Kafka를 사용하면 기술 팀이 거의 실시간 비즈니스 인텔리전스를 할 수 있어요.
  • LINE은 서비스들이 서로 통신하는 중앙 데이터허브로 Apache Kafka를 사용해요. 매일 수천억 개의 메시지가 생성되어 다양한 비즈니스 로직, 위협 탐지, 검색 인덱싱, 데이터 분석을 실행하는 데 사용돼요. LINE은 Kafka Streams를 이용해 신뢰성 있게 토픽을 변환·필터링하여 서브 토픽을 소비하는 컨슈머가 효율적으로 소비하게 하면서도, 정교하지만 최소한의 코드 베이스 덕분에 유지보수 용이성을 유지해요.
  • 네덜란드 3대 은행 중 하나인 Rabobank의 디지털 신경계인 Business Event Bus는 Apache Kafka로 구동돼요. 점점 더 많은 금융 프로세스와 서비스가 사용하며, 그중 하나가 Rabo Alerts예요. 이 서비스는 금융 이벤트가 발생하면 고객에게 실시간으로 알리고 Kafka Streams로 구축되어 있어요.
  • ironSource는 업계를 선도하는 게임 성장 플랫폼을 통해 흐르는 초당 수백만 개의 이벤트의 비동기 메시징을 위한 백본 인프라로 Apache Kafka를 사용해 게임 성장을 지원해요. 또한 Kafka Streams API를 사용해 예산 관리, 모니터링, 알림 같은 여러 실시간 사용 사례를 처리해요.
  • Kpow는 Apache Kafka를 위한 풍부하고 데이터 중심 UI와 안전한 API를 제공하는 엔터프라이즈급 도구킷이에요. 엔지니어에게 카프카 클러스터, 스키마 레지스트리(Schema Registry), Kafka Connect, Kafka Streams 애플리케이션에 대한 깊은 가시성과 제어를 제공하도록 설계되었어요.
  • La Redoute는 비즈니스 이벤트를 통해 애플리케이션을 디커플링하는 중앙 신경계로 Kafka를 사용해요. Kafka Connect, Kafka Streams, KSQL을 결합한 AI 파이프라인과 함께 거의 실시간 데이터 보고, 분석을 가져오는 분산·이벤트 중심 아키텍처를 가능하게 해요.
  • Nuuly는 프론트엔드 고객 경험을 배포 센터의 실시간 재고 관리·운영과 통합하기 위해 Kafka를 중앙 신경계로 사용해요. Kafka Streams와 Kafka Connect를 데이터 사이언스·머신러닝과 결합하여 즉각적인 비즈니스 인텔리전스와 개인화된 렌탈 경험을 제공해요.
  • Recursion은 신약 발견 노력을 위한 데이터 파이프라인을 구동하기 위해 Kafka Streams를 사용해요.
  • Salesforce는 pub/sub 아키텍처 시스템을 구현하고 멀티테넌트 시스템에 엔터프라이즈급 이벤트 중심 레이어를 안전하게 추가하기 위해 Apache Kafka를 채택했어요. Kafka를 마이크로서비스 아키텍처의 중앙 신경계로 삼아, Kafka Streams 애플리케이션이 고객을 위한 유용한 실시간 인사이트를 생성하는 다양한 연산을 수행해요.
  • Schrödinger는 예측 모델링, 데이터 분석, 협업 서비스에 데이터를 공급함으로써 물리 기반 컴퓨팅 플랫폼을 구동해 화학 공간의 빠른 탐색을 가능하게 해요. Kafka는 분산 고속 이벤트 버스로 사용되고, Kafka Connect와 Kafka Streams는 엔터프라이즈 정보 솔루션 LiveDesign이 사용하는 스트리밍 Change Data Capture 프레임워크의 기본 구성 요소예요.

Hello Kafka Streams

아래 코드 예시는 탄력적이고 고도로 확장 가능하며 장애 허용적이고 상태 저장이며 대규모 프로덕션에서 실행할 준비가 된 WordCount 애플리케이션을 구현해요.

Java:

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.common.utils.Bytes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.Produced;
import org.apache.kafka.streams.state.KeyValueStore;

import java.util.Arrays;
import java.util.Properties;

public class WordCountApplication {

   public static void main(final String[] args) throws Exception {
       Properties props = new Properties();
       props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-application");
       props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker1:9092");
       props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
       props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

       StreamsBuilder builder = new StreamsBuilder();
       KStream<String, String> textLines = builder.stream("TextLinesTopic");
       KTable<String, Long> wordCounts = textLines
           .flatMapValues(textLine -> Arrays.asList(textLine.toLowerCase().split("\\W+")))
           .groupBy((key, word) -> word)
           .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("counts-store"));
       wordCounts.toStream().to("WordsWithCountsTopic", Produced.with(Serdes.String(), Serdes.Long()));

       KafkaStreams streams = new KafkaStreams(builder.build(), props);
       streams.start();
   }

}

Scala:

import java.util.Properties
import java.util.concurrent.TimeUnit

import org.apache.kafka.streams.kstream.Materialized
import org.apache.kafka.streams.scala.ImplicitConversions._
import org.apache.kafka.streams.scala._
import org.apache.kafka.streams.scala.kstream._
import org.apache.kafka.streams.{KafkaStreams, StreamsConfig}

object WordCountApplication extends App {
  import Serdes._

  val props: Properties = {
    val p = new Properties()
    p.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-application")
    p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker1:9092")
    p
  }

  val builder: StreamsBuilder = new StreamsBuilder
  val textLines: KStream[String, String] = builder.stream[String, String]("TextLinesTopic")
  val wordCounts: KTable[String, Long] = textLines
    .flatMapValues(textLine => textLine.toLowerCase.split("\\W+"))
    .groupBy((_, word) => word)
    .count()(Materialized.as("counts-store"))
  wordCounts.toStream.to("WordsWithCountsTopic")

  val streams: KafkaStreams = new KafkaStreams(builder.build(), props)
  streams.start()

  sys.ShutdownHookThread {
     streams.close(10, TimeUnit.SECONDS)
  }
}

더 알아보기