빠른 시작
빠른 시작 (Quick Start)
시간을 아끼고 싶다면 이 페이지가 제일 좋아요. 실제로 돌아가는 Kafka Streams 데모 애플리케이션을 로컬에 띄워서 입력 토픽에 데이터를 넣고, 단어 빈도(word count)가 출력 토픽에 흘러나오는 걸 그대로 지켜볼 수 있어요. 카프카가 켜져 있다면 앞 두 단계를 건너뛰고 바로 진행해도 돼요.
출처: 문서
본문
Kafka Streams 데모 애플리케이션 실행하기
이 튜토리얼은 완전히 처음 시작해서 기존 카프카 데이터가 없다고 가정해요. 하지만 이미 카프카를 시작했다면 처음 두 단계를 건너뛰어도 돼요.
Kafka Streams는 입력 및/또는 출력 데이터가 카프카 클러스터에 저장되는, 미션 크리티컬한 실시간 애플리케이션과 마이크로서비스를 구축하기 위한 클라이언트 라이브러리예요. Kafka Streams는 클라이언트 측에서 표준 Java와 Scala 애플리케이션을 작성·배포하는 단순함과, 카프카의 서버 측 클러스터 기술의 이점을 결합해 애플리케이션을 확장성·탄력성·장애 허용·분산성 등을 가지게 해요.
이 퀵스타트 예제는 이 라이브러리로 코딩된 스트리밍 애플리케이션을 실행하는 방법을 보여줘요. [WordCountDemo](https://github.com/apache/kafka/blob/4.3/streams/examples/src/main/java/org/apache/kafka/streams/examples/wordcount/WordCountDemo.java) 예제 코드의 핵심은 다음과 같아요:
// Serializers/deserializers (serde) for String and Long types
final Serde<String> stringSerde = Serdes.String();
final Serde<Long> longSerde = Serdes.Long();
// Construct a `KStream` from the input topic "streams-plaintext-input", where message values
// represent lines of text (for the sake of this example, we ignore whatever may be stored
// in the message keys).
KStream<String, String> textLines = builder.stream(
"streams-plaintext-input",
Consumed.with(stringSerde, stringSerde)
);
KTable<String, Long> wordCounts = textLines
// Split each text line, by whitespace, into words.
.flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+")))
// Group the text words as message keys
.groupBy((key, value) -> value)
// Count the occurrences of each word (message key).
.count();
// Store the running counts as a changelog stream to the output topic.
wordCounts.toStream().to("streams-wordcount-output", Produced.with(Serdes.String(), Serdes.Long()));
이것은 입력 텍스트에서 단어 발생 히스토그램을 계산하는 WordCount 알고리즘을 구현해요. 하지만 이전에 봤을 수도 있는, 유한한 데이터(bounded data)에 동작하는 다른 WordCount 예제들과 달리, WordCount 데모 애플리케이션은 무한하고 경계 없는(unbounded) 데이터 스트림에 동작하도록 설계되었기 때문에 조금 다르게 동작해요. 유한한 변형과 비슷하게, 이것은 단어의 개수를 추적·업데이트하는 상태 저장(stateful) 알고리즘이에요. 그러나 잠재적으로 무한한 입력 데이터를 가정해야 하므로, "모든" 입력 데이터를 처리했는지 알 수 없기 때문에 더 많은 데이터를 계속 처리하면서 주기적으로 현재 상태와 결과를 출력해요.
첫 단계로, 카프카를 시작하고(아직 시작하지 않았다면) 카프카 토픽에 입력 데이터를 준비한 다음, 그것을 Kafka Streams 애플리케이션이 처리하도록 할 거예요.
1단계: 코드 다운로드
4.3.1 릴리스를 다운로드하고 압축을 풀어요. 여러 다운로드 가능한 Scala 버전이 있고, 여기서는 권장 버전(2.13)을 선택해요:
$ tar -xzf kafka_2.13-4.3.1.tgz
$ cd kafka_2.13-4.3.1
2단계: 카프카 서버 시작
클러스터 UUID 생성:
$ KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
로그 디렉터리 포맷:
$ bin/kafka-storage.sh format --standalone -t $KAFKA_CLUSTER_ID -c config/server.properties
카프카 서버 시작:
$ bin/kafka-server-start.sh config/server.properties
3단계: 입력 토픽 준비 및 카프카 프로듀서 시작
다음으로, streams-plaintext-input이라는 입력 토픽과 streams-wordcount-output이라는 출력 토픽을 만들어요:
$ bin/kafka-topics.sh --create \
--bootstrap-server localhost:9092 \
--replication-factor 1 \
--partitions 1 \
--topic streams-plaintext-input
Created topic "streams-plaintext-input".
참고: 출력 스트림이 체인지로그 스트림이므로(아래 애플리케이션 출력 설명 참조) 컴팩션이 활성화된 상태로 출력 토픽을 만들어요:
$ bin/kafka-topics.sh --create \
--bootstrap-server localhost:9092 \
--replication-factor 1 \
--partitions 1 \
--topic streams-wordcount-output \
--config cleanup.policy=compact
Created topic "streams-wordcount-output".
생성된 토픽은 같은 kafka-topics 도구로 describe할 수 있어요:
$ bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --exclude-internal
Topic:streams-wordcount-output PartitionCount:1 ReplicationFactor:1 Configs:cleanup.policy=compact,segment.bytes=1073741824
Topic: streams-wordcount-output Partition: 0 Leader: 0 Replicas: 0 Isr: 0
Topic:streams-plaintext-input PartitionCount:1 ReplicationFactor:1 Configs:segment.bytes=1073741824
Topic: streams-plaintext-input Partition: 0 Leader: 0 Replicas: 0 Isr: 0
4단계: Wordcount 애플리케이션 시작
다음 명령은 WordCount 데모 애플리케이션을 시작해요:
$ bin/kafka-run-class.sh org.apache.kafka.streams.examples.wordcount.WordCountDemo
데모 애플리케이션은 streams-plaintext-input 입력 토픽에서 읽고, 읽은 각 메시지에 WordCount 알고리즘 계산을 수행하고, 현재 결과를 계속해서 streams-wordcount-output 출력 토픽에 써요. 따라서 결과가 카프카에 다시 기록되므로 로그 항목 외에는 STDOUT 출력이 없어요.
선택적으로 group.protocol=streams를 사용해 KIP-1071에서 도입된 새 리밸런싱 프로토콜을 활성화할 수 있어요. 이것은 태스크 할당 로직을 클라이언트에서 브로커로 옮겨 리밸런스 시간을 줄이고 stop-the-world 리밸런스를 없애요.
$ echo "group.protocol=streams" > streams.properties
$ bin/kafka-run-class.sh org.apache.kafka.streams.examples.wordcount.WordCountDemo streams.properties
이제 별도 터미널에서 콘솔 프로듀서를 시작해 이 토픽에 입력 데이터를 쓸 수 있어요:
$ bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic streams-plaintext-input
그리고 별도 터미널에서 콘솔 컨슈머로 출력 토픽을 읽어 WordCount 데모 애플리케이션의 출력을 확인할 수 있어요:
$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic streams-wordcount-output \
--from-beginning \
--formatter-property print.key=true \
--formatter-property print.value=true \
--formatter-property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer \
--formatter-property value.deserializer=org.apache.kafka.common.serialization.LongDeserializer
5단계: 데이터 처리
이제 콘솔 프로듀서를 사용해 한 줄의 텍스트를 입력하고 Enter를 눌러 streams-plaintext-input 입력 토픽에 메시지를 써볼게요. 이것은 입력 토픽에 새 메시지를 보내는데, 메시지 키는 null이고 메시지 값은 방금 입력한 문자열로 인코딩된 텍스트 줄이에요(실제로는 애플리케이션의 입력 데이터는 이 퀵스타트처럼 수동으로 입력되는 대신 보통 카프카로 계속해서 스트리밍돼요):
$ bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic streams-plaintext-input
>all streams lead to kafka
이 메시지는 Wordcount 애플리케이션에 의해 처리되고, 다음 출력 데이터가 streams-wordcount-output 토픽에 기록되어 콘솔 컨슈머에 의해 출력돼요:
$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic streams-wordcount-output \
--from-beginning \
--formatter-property print.key=true \
--formatter-property print.value=true \
--formatter-property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer \
--formatter-property value.deserializer=org.apache.kafka.common.serialization.LongDeserializer
all 1
streams 1
lead 1
to 1
kafka 1
여기서 첫 번째 열은 java.lang.String 형식의 카프카 메시지 키이고 카운트되는 단어를 나타내며, 두 번째 열은 java.lang.Long 형식의 메시지 값으로 단어의 최신 카운트를 나타내요.
이제 콘솔 프로듀서로 streams-plaintext-input 입력 토픽에 메시지 하나를 더 써볼게요. "hello kafka streams" 텍스트 줄을 입력하고 Enter를 눌러요. 터미널은 다음과 같을 거예요:
$ bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic streams-plaintext-input
>all streams lead to kafka
>hello kafka streams
콘솔 컨슈머가 실행 중인 다른 터미널에서 WordCount 애플리케이션이 새 출력 데이터를 쓴 것을 관찰할 수 있어요:
$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic streams-wordcount-output \
--from-beginning \
--formatter-property print.key=true \
--formatter-property print.value=true \
--formatter-property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer \
--formatter-property value.deserializer=org.apache.kafka.common.serialization.LongDeserializer
all 1
streams 1
lead 1
to 1
kafka 1
hello 1
kafka 2
streams 2
여기서 마지막에 출력된 kafka 2와 streams 2는 카운트가 1에서 2로 증가된 kafka와 streams 키의 업데이트를 나타내요. 입력 토픽에 입력 메시지를 더 쓸 때마다 WordCount 애플리케이션이 계산한 가장 최근 단어 카운트를 나타내는 새 메시지가 streams-wordcount-output 토픽에 추가되는 것을 관찰할 거예요. 이 퀵스타트를 마무리하기 전에 콘솔 프로듀서의 streams-plaintext-input 입력 토픽에 마지막 입력 텍스트 줄 "join kafka summit"을 입력해볼게요:
$ bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic streams-plaintext-input
>all streams lead to kafka
>hello kafka streams
>join kafka summit
streams-wordcount-output 토픽은 이후에 해당하는 업데이트된 단어 카운트를 보여줄 거예요 (마지막 세 줄 참조):
$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic streams-wordcount-output \
--from-beginning \
--formatter-property print.key=true \
--formatter-property print.value=true \
--formatter-property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer \
--formatter-property value.deserializer=org.apache.kafka.common.serialization.LongDeserializer
all 1
streams 1
lead 1
to 1
kafka 1
hello 1
kafka 2
streams 2
join 1
kafka 3
summit 1
보시다시피 Wordcount 애플리케이션의 출력은 사실 연속적인 업데이트 스트림이에요. 각 출력 레코드(위 원본 출력의 각 줄)는 "kafka" 같은 레코드 키 하나의 업데이트된 카운트예요. 같은 키를 가진 여러 레코드에 대해, 각 이후 레코드는 이전 레코드의 업데이트예요.
아래 두 다이어그램은 배후에서 실제로 일어나는 일을 보여줘요. 첫 번째 열은 count에 대해 단어 발생을 카운트하는 KTable<String, Long>의 현재 상태의 진화를 보여줘요. 두 번째 열은 KTable에 대한 상태 업데이트에서 발생하여 출력 카프카 토픽 streams-wordcount-output으로 보내지는 변경 레코드(change record)를 보여줘요.
먼저 텍스트 줄 "all streams lead to kafka"가 처리돼요. 각 새 단어가 새 테이블 항목을 만들면서(초록 배경으로 강조) KTable이 구축되고, 해당 변경 레코드가 다운스트림 KStream으로 보내져요.
두 번째 텍스트 줄 "hello kafka streams"이 처리될 때, 처음으로 KTable의 기존 항목이 업데이트되는 것(여기서는 "kafka"와 "streams" 단어)을 관찰해요. 그리고 다시 변경 레코드가 출력 토픽으로 보내져요.
이런 식으로 계속됩니다 (세 번째 줄이 처리되는 방식의 그림은 생략해요). 이것이 출력 토픽의 내용이 위에서 보여준 것인 이유를 설명해요. 변경의 전체 기록을 포함하기 때문이에요.
이 구체적인 예의 범위를 넘어서, Kafka Streams가 여기서 하는 것은 테이블과 체인지로그 스트림 사이의 이중성(duality)을 활용하는 거예요 (여기서: 테이블 = KTable, 체인지로그 스트림 = 다운스트림 KStream): 테이블의 모든 변경을 스트림에 게시할 수 있고, 전체 체인지로그 스트림을 처음부터 끝까지 소비하면 테이블의 내용을 재구성할 수 있어요.
6단계: 애플리케이션 정리
이제 Ctrl-C를 사용해 콘솔 컨슈머, 콘솔 프로듀서, Wordcount 애플리케이션, 카프카 브로커를 순서대로 중지할 수 있어요.
더 알아보기
- Write a streams app — 더 자세한 단계별 애플리케이션 작성 튜토리얼을 봐요.
- 핵심 개념 — 여기서 만난 KStream/KTable, 스트림·테이블 이중성 개념을 제대로 이해해요.
- Introduction — Kafka Streams의 등장 배경과 쓰임을 봐요.