Streams 애플리케이션 작성하기

Streams 애플리케이션 작성하기

Kafka Streams로 애플리케이션을 만드는 것 자체는 어렵지 않아요. 어떤 Java(Scala) 애플리케이션이든 Kafka Streams 라이브러리를 쓰기만 하면 "스트림즈 애플리케이션"이 돼요. 이 페이지에서는 애플리케이션의 뼈대인 프로세서 토폴로지(processor topology)를 정의하고, 필요한 의존성을 추가하고, KafkaStreams 인스턴스를 만들어 실행·종료하는 기본 흐름을 잡아드릴게요.

출처: 문서

본문

Kafka Streams 라이브러리를 사용하는 모든 Java 또는 Scala 애플리케이션은 Kafka Streams 애플리케이션으로 간주돼요. Kafka Streams 애플리케이션의 계산 로직은 프로세서 토폴로지(processor topology)로 정의되는데, 이는 스트림 프로세서(노드)와 스트림(엣지)으로 이루어진 그래프예요.

Kafka Streams API로 프로세서 토폴로지를 정의할 수 있어요:

  • Kafka Streams DSLmap, filter, join, aggregations 같은 가장 일반적인 데이터 변환 연산을 기본 제공하는 고수준 API예요. Kafka Streams를 처음 접하는 개발자에게 권장되는 시작점이며, 많은 사용 사례와 스트림 처리 요구를 충족해요. Scala 애플리케이션을 작성한다면 Java DSL을 직접 쓰는 대신 Java/Scala 상호 운용 보일러플레이트를 많이 없애주는 Kafka Streams DSL for Scala 라이브러리를 사용할 수 있어요.
  • Processor API — 프로세서를 추가·연결하고 상태 저장소(state store)와 직접 상호작용할 수 있게 해주는 저수준 API예요. DSL보다 훨씬 더 많은 유연성을 제공하지만, 애플리케이션 개발자가 더 많은 수작업(예: 더 많은 코드 줄)을 해야 하는 비용이 있어요.

라이브러리 및 Maven 아티팩트

이 섹션은 Kafka Streams 애플리케이션을 작성할 때 사용할 수 있는 Kafka Streams 관련 라이브러리를 나열해요.

Kafka Streams 애플리케이션에 대해 다음 라이브러리에 대한 의존성을 정의할 수 있어요:

Group ID Artifact ID Version 설명
org.apache.kafka kafka-streams 4.3.1 (필수) Kafka Streams 기본 라이브러리
org.apache.kafka kafka-clients 4.3.1 (필수) Kafka 클라이언트 라이브러리. 내장 시리얼라이저/디시리얼라이저 포함
org.apache.kafka kafka-streams-scala 4.3.1 (선택) Scala Kafka Streams 애플리케이션을 작성하기 위한 Kafka Streams DSL for Scala 라이브러리. SBT를 사용하지 않을 때는 아티팩트 ID에 애플리케이션이 사용하는 Scala 버전(_2.12, _2.13)을 접미사로 붙여야 해요

: 시리얼라이저/디시리얼라이저에 대한 자세한 내용은 데이터 타입과 직렬화 섹션을 참고해요.

Maven을 사용할 때의 예시 pom.xml 스니펫:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams</artifactId>
    <version>4.3.1</version>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>4.3.1</version>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams-scala_2.13</artifactId>
    <version>4.3.1</version>
</dependency>

애플리케이션 코드 안에서 Kafka Streams 사용하기

애플리케이션 코드의 어디에서든 Kafka Streams를 호출할 수 있지만, 보통은 애플리케이션의 main() 메서드(또는 그 변형) 안에서 호출해요. 애플리케이션 안에서 처리 토폴로지를 정의하는 기본 요소는 아래에 설명돼요.

먼저 KafkaStreams 인스턴스를 만들어야 해요.

  • KafkaStreams 생성자의 첫 번째 인자는 토폴로지(DSL은 StreamsBuilder#build(), Processor API는 Topology)를 받아 토폴로지를 정의해요.
  • 두 번째 인자는 java.util.Properties 인스턴스로, 이 특정 토폴로지의 구성을 정의해요.

코드 예시:

import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.kstream.StreamsBuilder;
import org.apache.kafka.streams.processor.Topology;

// 빌더를 사용해 실제 처리 토폴로지를 정의한다. 예를 들어 어떤 입력 토픽에서 읽을지,
// 어떤 스트림 연산(filter, map 등)을 호출할지 등을 지정한다.
// 자세한 내용은 이 개발자 가이드의 후속 섹션에서 다룬다.

StreamsBuilder builder = ...;  // DSL 사용 시
Topology topology = builder.build();
//
// 또는
//
Topology topology = ...; // Processor API 사용 시

// 구성을 사용해 애플리케이션에 카프카 클러스터 위치, 기본으로 사용할
// 시리얼라이저/디시리얼라이저, 보안 설정 등을 알려준다.
Properties props = ...;

KafkaStreams streams = new KafkaStreams(topology, props);

이 시점에 내부 구조가 초기화되지만 처리는 아직 시작되지 않아요. KafkaStreams#start() 메서드를 호출해 Kafka Streams 스레드를 명시적으로 시작해야 해요:

// Kafka Streams 스레드 시작
streams.start();

이 스트림 처리 애플리케이션의 다른 인스턴스가 다른 곳(예: 다른 머신)에서 실행되고 있다면, Kafka Streams는 기존 인스턴스의 태스크를 방금 시작한 새 인스턴스로 투명하게 재할당해요. 자세한 내용은 스트림 파티션과 태스크스레딩 모델을 참고해요.

예상치 못한 예외를 잡으려면, 애플리케이션을 시작하기 전에 java.lang.Thread.UncaughtExceptionHandler를 설정할 수 있어요. 이 핸들러는 스트림 스레드가 예상치 못한 예외로 종료될 때마다 호출돼요:

streams.setUncaughtExceptionHandler((Thread thread, Throwable throwable) -> {
  // 여기서 throwable/exception을 검사하고 적절한 조치를 취해야 한다!
});

애플리케이션 인스턴스를 중지하려면 KafkaStreams#close() 메서드를 호출해요:

// Kafka Streams 스레드 중지
streams.close();

SIGTERM에 응답해 애플리케이션이 정상 종료(graceful shutdown)되도록 하려면, 셧다운 훅(shutdown hook)을 추가하고 KafkaStreams#close를 호출하는 것이 권장돼요.

Java에서의 셧다운 훅 예시:

// Kafka Streams 스레드를 중지하는 셧다운 훅 추가.
// 선택적으로 `close`에 타임아웃을 제공할 수 있다.
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));

애플리케이션이 중지된 후에는 Kafka Streams가 이 인스턴스에서 실행 중이던 태스크를 사용 가능한 나머지 인스턴스로 마이그레이션해요.

Streams 애플리케이션 테스트하기

Kafka Streams는 애플리케이션을 테스트하는 데 도움이 되는 test-utils 모듈을 함께 제공해요. 자세한 내용은 테스트 문서를 참고해요.

더 알아보기