Streams 앱 작성하기

Streams 앱 작성하기 (Write a streams app)

이 튜토리얼은 카프카 스트림즈 애플리케이션을 처음부터 만들어보는 실습이에요. Maven 프로젝트를 세팅하고, 데이터를 그대로 옮기는 Pipe, 텍스트를 단어로 쪼개는 LineSplit, 그리고 단어 빈도를 세는 WordCount까지 세 개 프로그램을 단계별로 작성해요. 코드로 직접 익히고 싶다면 이 페이지가 제일 좋아요.

출처: 문서

본문

튜토리얼: Kafka Streams 애플리케이션 작성하기

이 가이드에서는 Kafka Streams를 사용한 스트림 처리 애플리케이션을 작성하기 위해 자신의 프로젝트를 처음부터 세팅할 거예요. 아직 안 해봤다면 먼저 퀵스타트를 읽어 Kafka Streams로 작성된 Streams 애플리케이션을 실행하는 방법을 먼저 보는 것을 적극 권장해요.

Maven 프로젝트 세팅하기

Kafka Streams Maven 아키타입(archetype)을 사용해 Streams 프로젝트 구조를 다음 명령으로 만들 거예요:

$ mvn archetype:generate \
-DarchetypeGroupId=org.apache.kafka \
-DarchetypeArtifactId=streams-quickstart-java \
-DarchetypeVersion=4.3.1 \
-DgroupId=streams.examples \
-DartifactId=streams-quickstart \
-Dversion=0.1 \
-Dpackage=myapps

groupId, artifactId, package 파라미터에 다른 값을 사용해도 돼요. 위 파라미터 값을 사용한다고 가정하면, 이 명령은 다음과 같은 프로젝트 구조를 만들어요:

$ tree streams-quickstart
streams-quickstart
|-- pom.xml
|-- src
    |-- main
        |-- java
        |   |-- myapps
        |       |-- LineSplit.java
        |       |-- Pipe.java
        |       |-- WordCount.java
        |-- resources
            |-- log4j.properties

프로젝트에 포함된 pom.xml 파일에는 이미 Streams 의존성이 정의되어 있어요. 생성된 pom.xml은 Java 11을 대상으로 한다는 점을 유의해요. src/main/java 아래에 Streams 라이브러리로 작성된 여러 예시 프로그램이 이미 있어요. 우리는 그러한 프로그램을 처음부터 작성할 것이므로, 이 예시들을 삭제할 수 있어요:

$ cd streams-quickstart
$ rm src/main/java/myapps/*.java

첫 번째 Streams 애플리케이션 작성: Pipe

이제 코딩할 시간이에요! 좋아하는 IDE를 열고 이 Maven 프로젝트를 가져오거나, 텍스트 편집기를 열고 src/main/java/myapps 아래에 java 파일을 만들어도 돼요. 이름을 Pipe.java로 지어볼게요:

package myapps;

public class Pipe {

    public static void main(String[] args) throws Exception {

    }
}

main 함수를 채워서 이 파이프 프로그램을 작성할 거예요. IDE가 보통 자동으로 추가할 수 있으므로 import 문은 단계별로 나열하지 않을 거예요. 하지만 텍스트 편집기를 사용한다면 import를 수동으로 추가해야 해요. 이 섹션 끝에 import 문이 있는 완전한 코드 스니펫을 보여드릴게요.

Streams 애플리케이션을 작성하는 첫 단계는 StreamsConfig에 정의된 다양한 Streams 실행 구성을 지정하는 java.util.Properties 맵을 만드는 것이에요. 설정해야 할 몇 가지 중요한 구성 값은 StreamsConfig.BOOTSTRAP_SERVERS_CONFIG(카프카 클러스터에 대한 초기 연결을 수립하는 데 사용할 호스트/포트 쌍 목록)와 StreamsConfig.APPLICATION_ID_CONFIG(같은 카프카 클러스터와 통신하는 다른 애플리케이션과 구분하기 위한 Streams 애플리케이션의 고유 식별자)예요:

Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "streams-pipe");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");    // assuming that the Kafka broker this application is talking to runs on local machine with port 9092

추가로 같은 맵에서 다른 구성을 커스터마이즈할 수 있어요. 예를 들어 레코드 키-값 쌍의 기본 직렬화·역직렬화 라이브러리:

props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

Kafka Streams 구성의 전체 목록은 이 테이블을 참고해요.

다음으로 Streams 애플리케이션의 계산 로직을 정의할 거예요. Kafka Streams에서 이 계산 로직은 연결된 프로세서 노드의 topology로 정의돼요. 토폴로지 빌더를 사용해 그러한 토폴로지를 구성할 수 있어요:

final StreamsBuilder builder = new StreamsBuilder();

그리고 이 토폴로지 빌더를 사용해 streams-plaintext-input이라는 카프카 토픽에서 소스 스트림을 만들 수 있어요:

KStream<String, String> source = builder.stream("streams-plaintext-input");

이제 소스 카프카 토픽 streams-plaintext-input에서 레코드를 지속적으로 생성하는 KStream을 얻었어요. 레코드는 String 타입의 키-값 쌍으로 구성돼요. 이 스트림으로 할 수 있는 가장 간단한 일은 streams-pipe-output이라는 다른 카프카 토픽에 쓰는 것이에요:

source.to("streams-pipe-output");

위 두 줄을 한 줄로 연결할 수도 있다는 점을 유의해요:

builder.stream("streams-plaintext-input").to("streams-pipe-output");

다음을 수행해 이 빌더로 어떤 topology가 만들어지는지 검사할 수 있어요:

final Topology topology = builder.build();

그리고 그 설명을 표준 출력으로 출력해요:

System.out.println(topology.describe());

여기서 멈추고 컴파일·실행하면 다음 정보가 출력돼요:

$ mvn clean package
$ mvn exec:java -Dexec.mainClass=myapps.Pipe
Sub-topologies:
  Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000(topics: streams-plaintext-input) --> KSTREAM-SINK-0000000001
    Sink: KSTREAM-SINK-0000000001(topic: streams-pipe-output) <-- KSTREAM-SOURCE-0000000000
Global Stores:
  none

위와 같이, 구성된 토폴로지가 두 개의 프로세서 노드, 소스 노드 KSTREAM-SOURCE-0000000000와 싱크 노드 KSTREAM-SINK-0000000001을 가진다는 것을 보여줘요. KSTREAM-SOURCE-0000000000은 카프카 토픽 streams-plaintext-input에서 레코드를 계속 읽어 다운스트림 노드 KSTREAM-SINK-0000000001로 파이프하고, KSTREAM-SINK-0000000001은 받은 각 레코드를 순서대로 streams-pipe-output이라는 다른 카프카 토픽에 써요 (--><-- 화살표는 이 노드의 다운스트림·업스트림 프로세서 노드, 즉 토폴로지 그래프의 "자식"과 "부모"를 나타내요). 또한 이 간단한 토폴로지에 연결된 전역 상태 저장소가 없다는 것을 보여줘요 (상태 저장소는 다음 섹션에서 더 이야기할게요).

코드로 토폴로지를 만드는 동안 언제든지 위처럼 토폴로지를 설명할 수 있다는 점을 유의해요. 따라서 사용자는 토폴로지에 정의된 계산 로직을 원하는 대로 될 때까지 인터랙티브하게 "시도하고 맛볼" 수 있어요. 한 카프카 토픽에서 다른 카프카 토픽으로 무한한 스트리밍 방식으로 데이터를 파이프하는 이 간단한 토폴로지가 완성되었다고 가정하면, 위에서 방금 구성한 두 구성 요소, 즉 java.util.Properties 인스턴스에 지정된 구성 맵과 Topology 객체로 Streams 클라이언트를 구성할 수 있어요:

final KafkaStreams streams = new KafkaStreams(topology, props);

start() 함수를 호출하면 이 클라이언트의 실행을 트리거할 수 있어요. 이 클라이언트에서 close()가 호출될 때까지 실행은 멈추지 않아요. 예를 들어, 카운트다운 래치(countdown latch)가 있는 셧다운 훅을 추가해 사용자 인터럽트를 잡고 프로그램이 종료될 때 클라이언트를 닫을 수 있어요:

final CountDownLatch latch = new CountDownLatch(1);

// attach shutdown handler to catch control-c
Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown-hook") {
    @Override
    public void run() {
        streams.close();
        latch.countDown();
    }
});

try {
    streams.start();
    latch.await();
} catch (Throwable e) {
    System.exit(1);
}
System.exit(0);

지금까지의 완전한 코드는 다음과 같아요:

package myapps;

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.Topology;

import java.util.Properties;
import java.util.concurrent.CountDownLatch;

public class Pipe {

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

        final StreamsBuilder builder = new StreamsBuilder();

        builder.stream("streams-plaintext-input").to("streams-pipe-output");

        final Topology topology = builder.build();

        final KafkaStreams streams = new KafkaStreams(topology, props);
        final CountDownLatch latch = new CountDownLatch(1);

        // attach shutdown handler to catch control-c
        Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown-hook") {
            @Override
            public void run() {
                streams.close();
                latch.countDown();
            }
        });

        try {
            streams.start();
            latch.await();
        } catch (Throwable e) {
            System.exit(1);
        }
        System.exit(0);
    }
}

localhost:9092에 카프카 브로커가 실행 중이고, 그 브로커에 streams-plaintext-inputstreams-pipe-output 토픽이 생성되어 있다면, Maven을 사용해 IDE나 명령줄에서 이 코드를 실행할 수 있어요:

$ mvn clean package
$ mvn exec:java -Dexec.mainClass=myapps.Pipe

Streams 애플리케이션을 실행하고 그 계산 결과를 관찰하는 상세한 방법은 Play with a Streams Application 섹션을 읽어보세요. 이 섹션의 나머지에서는 다루지 않을게요.

두 번째 Streams 애플리케이션 작성: Line Split

두 개의 핵심 구성 요소인 StreamsConfigTopology로 Streams 클라이언트를 구성하는 방법을 배웠어요. 이제 현재 토폴로지를 확장해 실제 처리 로직을 추가해볼게요. 먼저 기존 Pipe.java 클래스를 복사해 다른 프로그램을 만들 수 있어요:

$ cp src/main/java/myapps/Pipe.java src/main/java/myapps/LineSplit.java

그리고 원래 프로그램과 구분하기 위해 클래스 이름과 application id 구성을 변경해요:

public class LineSplit {

    public static void main(String[] args) throws Exception {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "streams-linesplit");
        // ...
    }
}

소스 스트림의 각 레코드가 String 타입의 키-값 쌍이므로, 값 문자열을 텍스트 줄로 취급하고 FlatMapValues 연산자로 단어로 쪼개볼게요:

KStream<String, String> source = builder.stream("streams-plaintext-input");
KStream<String, String> words = source.flatMapValues(new ValueMapper<String, Iterable<String>>() {
            @Override
            public Iterable<String> apply(String value) {
                return Arrays.asList(value.split("\\W+"));
            }
        });

이 연산자는 source 스트림을 입력으로 받고, 소스 스트림의 각 레코드를 순서대로 처리해 그 값 문자열을 단어 목록으로 쪼개고 각 단어를 새 레코드로 생성해 words라는 새 스트림을 생성해요. 이것은 이전에 받은 레코드나 처리된 결과를 추적할 필요가 없는 무상태(stateless) 연산자예요. JDK 8을 사용한다면 람다 표현식을 사용해 위 코드를 단순화할 수 있어요:

KStream<String, String> source = builder.stream("streams-plaintext-input");
KStream<String, String> words = source.flatMapValues(value -> Arrays.asList(value.split("\\W+")));

그리고 마지막으로 단어 스트림을 streams-linesplit-output이라는 다른 카프카 토픽에 쓸 수 있어요. 다시 말하지만, 이 두 단계는 다음과 같이 연결될 수 있어요 (람다 표현식 사용 가정):

KStream<String, String> source = builder.stream("streams-plaintext-input");
source.flatMapValues(value -> Arrays.asList(value.split("\\W+")))
      .to("streams-linesplit-output");

System.out.println(topology.describe())로 이 확장된 토폴로지를 설명하면 다음을 얻을 거예요:

$ mvn clean package
$ mvn exec:java -Dexec.mainClass=myapps.LineSplit
Sub-topologies:
  Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000(topics: streams-plaintext-input) --> KSTREAM-FLATMAPVALUES-0000000001
    Processor: KSTREAM-FLATMAPVALUES-0000000001(stores: []) --> KSTREAM-SINK-0000000002 <-- KSTREAM-SOURCE-0000000000
    Sink: KSTREAM-SINK-0000000002(topic: streams-linesplit-output) <-- KSTREAM-FLATMAPVALUES-0000000001
  Global Stores:
    none

위에서 볼 수 있듯이, 새 프로세서 노드 KSTREAM-FLATMAPVALUES-0000000001이 원래 소스와 싱크 노드 사이의 토폴로지에 주입됐어요. 그것은 소스 노드를 부모로, 싱크 노드를 자식으로 가져요. 다시 말해, 소스 노드가 가져온 각 레코드는 먼저 새로 추가된 KSTREAM-FLATMAPVALUES-0000000001 노드로 이동해 처리되고, 결과적으로 하나 이상의 새 레코드가 생성돼요. 그것들은 계속 싱크 노드로 이동해 카프카에 다시 쓰여요. 이 프로세서 노드는 어떤 저장소와도 연결되지 않으므로(즉 (stores: [])) "stateless(무상태)"라는 점을 유의해요.

완전한 코드는 다음과 같아요 (람다 표현식 사용 가정):

package myapps;

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.kstream.KStream;

import java.util.Arrays;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;

public class LineSplit {

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

        final StreamsBuilder builder = new StreamsBuilder();

        KStream<String, String> source = builder.stream("streams-plaintext-input");
        source.flatMapValues(value -> Arrays.asList(value.split("\\W+")))
              .to("streams-linesplit-output");

        final Topology topology = builder.build();
        final KafkaStreams streams = new KafkaStreams(topology, props);
        final CountDownLatch latch = new CountDownLatch(1);

        // ... same as Pipe.java above
    }
}

세 번째 Streams 애플리케이션 작성: Wordcount

이제 소스 텍스트 스트림에서 쪼개진 단어의 발생을 세어 토폴로지에 "상태 저장(stateful)" 계산을 추가하는 한 단계를 더 나아가볼게요. 유사한 단계에 따라 LineSplit.java 클래스를 기반으로 다른 프로그램을 만들어볼게요:

public class WordCount {

    public static void main(String[] args) throws Exception {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "streams-wordcount");
        // ...
    }
}

단어를 세기 위해 먼저 flatMapValues 연산자를 수정해 모든 단어를 소문자로 처리할 수 있어요 (람다 표현식 사용 가정):

source.flatMapValues(new ValueMapper<String, Iterable<String>>() {
    @Override
    public Iterable<String> apply(String value) {
        return Arrays.asList(value.toLowerCase(Locale.getDefault()).split("\\W+"));
    }
});

카운트 집계를 수행하려면 먼저 groupBy 연산자로 값 문자열, 즉 소문자 단어로 스트림에 키를 지정하고 싶다고 명시해야 해요. 이 연산자는 새 그룹화된 스트림을 생성하며, 이후 count 연산자로 집계되어 그룹화된 각 키에 대한 실행 카운트를 생성해요:

KTable<String, Long> counts =
source.flatMapValues(new ValueMapper<String, Iterable<String>>() {
            @Override
            public Iterable<String> apply(String value) {
                return Arrays.asList(value.toLowerCase(Locale.getDefault()).split("\\W+"));
            }
        })
      .groupBy(new KeyValueMapper<String, String, String>() {
           @Override
           public String apply(String key, String value) {
               return value;
           }
        })
      // Materialize the result into a KeyValueStore named "counts-store".
      // The Materialized store is always of type <Bytes, byte[]> as this is the format of the inner most store.
      .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>> as("counts-store"));

count 연산자에는 실행 카운트가 counts-store라는 상태 저장소에 저장되어야 한다고 지정하는 Materialized 파라미터가 있다는 점을 유의해요. 이 counts-store 저장소는 실시간으로 쿼리할 수 있으며, 자세한 내용은 개발자 매뉴얼에 설명되어 있어요.

또한 counts KTable의 체인지로그 스트림을 streams-wordcount-output이라는 다른 카프카 토픽에 쓸 수 있어요. 결과가 체인지로그 스트림이므로 출력 토픽 streams-wordcount-output은 로그 컴팩션이 활성화되도록 구성해야 해요. 이번에는 값 타입이 더 이상 String이 아니라 Long이므로, 기본 직렬화 클래스로는 Kafka에 쓰는 것이 더 이상 가능하지 않아요. Long 타입에 대한 오버라이드된 직렬화 메서드를 제공해야 해요. 그렇지 않으면 런타임 예외가 발생해요:

counts.toStream().to("streams-wordcount-output", Produced.with(Serdes.String(), Serdes.Long()));

streams-wordcount-output 토픽에서 체인지로그 스트림을 읽으려면 값 역직렬화를 org.apache.kafka.common.serialization.LongDeserializer로 설정해야 한다는 점을 유의해요. 자세한 내용은 Play with a Streams Application 섹션에서 찾을 수 있어요. JDK 8의 람다 표현식을 사용할 수 있다고 가정하면 위 코드는 다음과 같이 단순화할 수 있어요:

KStream<String, String> source = builder.stream("streams-plaintext-input");
source.flatMapValues(value -> Arrays.asList(value.toLowerCase(Locale.getDefault()).split("\\W+")))
      .groupBy((key, value) -> value)
      .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("counts-store"))
      .toStream()
      .to("streams-wordcount-output", Produced.with(Serdes.String(), Serdes.Long()));

System.out.println(topology.describe())로 이 확장된 토폴로지를 다시 설명하면 다음을 얻을 거예요:

$ mvn clean package
$ mvn exec:java -Dexec.mainClass=myapps.WordCount
Sub-topologies:
  Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000(topics: streams-plaintext-input) --> KSTREAM-FLATMAPVALUES-0000000001
    Processor: KSTREAM-FLATMAPVALUES-0000000001(stores: []) --> KSTREAM-KEY-SELECT-0000000002 <-- KSTREAM-SOURCE-0000000000
    Processor: KSTREAM-KEY-SELECT-0000000002(stores: []) --> KSTREAM-FILTER-0000000005 <-- KSTREAM-FLATMAPVALUES-0000000001
    Processor: KSTREAM-FILTER-0000000005(stores: []) --> KSTREAM-SINK-0000000004 <-- KSTREAM-KEY-SELECT-0000000002
    Sink: KSTREAM-SINK-0000000004(topic: counts-store-repartition) <-- KSTREAM-FILTER-0000000005
  Sub-topology: 1
    Source: KSTREAM-SOURCE-0000000006(topics: counts-store-repartition) --> KSTREAM-AGGREGATE-0000000003
    Processor: KSTREAM-AGGREGATE-0000000003(stores: [counts-store]) --> KTABLE-TOSTREAM-0000000007 <-- KSTREAM-SOURCE-0000000006
    Processor: KTABLE-TOSTREAM-0000000007(stores: []) --> KSTREAM-SINK-0000000008 <-- KSTREAM-AGGREGATE-0000000003
    Sink: KSTREAM-SINK-0000000008(topic: streams-wordcount-output) <-- KTABLE-TOSTREAM-0000000007
Global Stores:
  none

위에서 볼 수 있듯이, 토폴로지는 이제 두 개의 연결되지 않은 서브-토폴로지를 포함해요. 첫 번째 서브-토폴로지의 싱크 노드 KSTREAM-SINK-0000000004는 두 번째 서브-토폴로지의 소스 노드 KSTREAM-SOURCE-0000000006이 읽는 리파티션 토픽 counts-store-repartition에 쓰게 돼요. 리파티션 토픽은 집계 키(여기서는 값 문자열)로 소스 스트림을 "셔플"하는 데 사용돼요. 추가로, 첫 번째 서브-토폴로지 안에서 그룹화 KSTREAM-KEY-SELECT-0000000002 노드와 싱크 노드 사이에 무상태 KSTREAM-FILTER-0000000005 노드가 주입되어 집계 키가 비어 있는 중간 레코드를 걸러내요.

두 번째 서브-토폴로지에서 집계 노드 KSTREAM-AGGREGATE-0000000003counts-store라는 이름의 상태 저장소와 연결돼요 (그 이름은 count 연산자에서 사용자가 지정한 것입니다). 업스트림 스트림 소스 노드에서 각 레코드를 받으면 집계 프로세서는 먼저 연결된 counts-store 저장소를 쿼리해 해당 키의 현재 카운트를 얻고, 1을 증가시킨 다음, 새 카운트를 저장소에 다시 써요. 키에 대한 각 업데이트된 카운트는 또한 다운스트림 KTABLE-TOSTREAM-0000000007 노드로 파이프되는데, 이 노드는 이 업데이트 스트림을 레코드 스트림으로 해석한 다음 싱크 노드 KSTREAM-SINK-0000000008로 더 파이프해 카프카에 다시 써요.

완전한 코드는 다음과 같아요 (람다 표현식 사용 가정):

package myapps;

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.Topology;
import org.apache.kafka.streams.kstream.KStream;
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.Locale;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;

public class WordCount {

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

        final StreamsBuilder builder = new StreamsBuilder();

        KStream<String, String> source = builder.stream("streams-plaintext-input");
        source.flatMapValues(value -> Arrays.asList(value.toLowerCase(Locale.getDefault()).split("\\W+")))
              .groupBy((key, value) -> value)
              .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("counts-store"))
              .toStream()
              .to("streams-wordcount-output", Produced.with(Serdes.String(), Serdes.Long()));

        final Topology topology = builder.build();
        final KafkaStreams streams = new KafkaStreams(topology, props);
        final CountDownLatch latch = new CountDownLatch(1);

        // ... same as Pipe.java above
    }
}

더 알아보기