Processor API

Processor API

DSL이 고수준의 편리함이라면, Processor API는 저수준의 유연함이에요. 커스텀 프로세서를 직접 정의·연결하고 상태 저장소와 자유롭게 상호작용할 수 있게 해주죠. 이 페이지에서는 Processor 인터페이스로 커스텀 프로세서를 만드는 법, 상태 저장소의 종류와 생성, 장애 허용, 그리고 프로세서와 저장소를 Topology로 연결하는 방법까지 정리해드릴게요.

출처: 문서

본문

Processor API를 사용하면 개발자가 커스텀 프로세서를 정의·연결하고 상태 저장소와 상호작용할 수 있어요. Processor API로 받은 레코드를 한 번에 하나씩 처리하는 임의의 스트림 프로세서를 정의하고, 이 프로세서들을 관련 상태 저장소와 연결해 커스터마이즈된 처리 로직을 나타내는 프로세서 토폴로지를 구성할 수 있어요.

개요

Processor API는 무상태(stateless)와 상태 저장(stateful) 연산을 모두 구현하는 데 사용될 수 있으며, 후자는 상태 저장소를 통해 달성돼요.

: DSL과 Processor API 결합하기 — Applying processors (Processor API integration) 섹션에 설명된 대로 DSL의 편리함과 Processor API의 힘·유연성을 결합할 수 있어요. 사용 가능한 API 기능의 전체 목록은 Streams API docs를 참고해요.

스트림 프로세서 정의

스트림 프로세서는 프로세서 토폴로지에서 단일 처리 단계를 나타내는 노드예요. Processor API로 받은 레코드를 한 번에 하나씩 처리하는 임의의 스트림 프로세서를 정의하고, 이 프로세서들을 관련 상태 저장소와 연결해 프로세서 토폴로지를 구성할 수 있어요.

process() API 메서드를 제공하는 Processor 인터페이스를 구현해 커스터마이즈된 스트림 프로세서를 정의할 수 있어요. process() 메서드는 받은 각 레코드에 대해 호출돼요.

Processor 인터페이스에는 또한 task 구성 단계 동안 Kafka Streams 라이브러리가 호출하는 init() 메서드가 있어요. 프로세서 인스턴스는 이 메서드에서 필요한 초기화를 수행해야 해요. init() 메서드는 ProcessorContext 인스턴스를 전달하는데, 이것은 소스 카프카 토픽과 파티션, 해당 메시지 오프셋 등을 포함한 현재 처리 중인 레코드의 메타데이터에 접근을 제공해요. 이 컨텍스트 인스턴스를 사용해 punctuation 함수를 스케줄하고(ProcessorContext#schedule()), 새 레코드를 다운스트림 프로세서로 전달하고(ProcessorContext#forward()), 현재 처리 진행 상황의 커밋을 요청할 수 있어요(ProcessorContext#commit()). init()에서 설정한 리소스는 close() 메서드에서 정리할 수 있어요. Kafka Streams는 close()init()을 다시 호출해 단일 Processor 객체를 재사용할 수 있다는 점을 유의해요.

Processor 인터페이스는 네 개의 제네릭 파라미터 KIn, VIn, KOut, VOut을 받아요. 이것들은 프로세서 구현이 처리할 수 있는 입력·출력 타입을 정의해요. KInVInprocess()에 전달될 Record의 키·값 타입을 정의해요. 마찬가지로 KOutVOutProcessorContext#forward()가 받아들일 결과 Record의 전달될 키·값 타입을 정의해요. 프로세서가 레코드를 전혀 전달하지 않으면(또는 null 키나 값만 전달하면) 출력 제네릭 타입 인자를 Void로 설정하는 것이 모범 사례예요. 공통 상위 클래스를 공유하지 않는 여러 타입을 전달해야 한다면 출력 제네릭 타입 인자를 Object로 설정해야 해요.

Processor#process()ProcessorContext#forward() 메서드 모두 Record<K, V> 데이터 클래스 형태의 레코드를 처리해요. 이 클래스는 카프카 레코드의 주요 구성 요소인 키, 값, 타임스탬프, 헤더에 접근하게 해줘요. 레코드를 전달할 때 생성자를 사용해 처음부터 새 Record를 만들거나, 편리한 빌더 메서드를 사용해 Record의 속성 중 하나를 교체하고 나머지를 복사할 수 있어요. 예를 들어 inputRecord.withValue(newValue)inputRecord에서 키, 타임스탬프, 헤더를 복사하면서 출력 레코드의 값을 newValue로 설정해요. 이것은 inputRecord를 변경하지 않고 얕은 복사본을 만든다는 점을 유의해요. 이것은 얕은 복사일 뿐이므로, 프로그램의 다른 곳에서 키·값·헤더를 변경할 계획이라면 그 필드들의 깊은 복사를 직접 만들어야 해요.

Processor#process()로 들어오는 레코드를 처리하는 것 외에, 프로세서의 init() 메서드에서 ProcessorContext#schedule()를 호출하고 Punctuator를 전달해 주기적 호출("punctuation"이라고 함)을 스케줄하는 옵션이 있어요. PunctuationType은 punctuation 스케줄링에 어떤 시간 개념을 사용할지 결정해요: 스트림-시간(stream-time) 또는 벽시계-시간(wall-clock-time) (기본적으로 스트림-시간은 TimestampExtractor를 통해 이벤트-시간을 나타내도록 구성됨). 스트림-시간이 사용되면, punctuate()는 순수하게 데이터에 의해 촉발돼요. 스트림-시간이 입력 데이터에서 파생된 타임스탬프에 의해 결정(그리고 진행)되기 때문이에요. 새 입력 데이터가 도착하지 않으면 스트림-시간은 진행되지 않으므로 punctuate()는 호출되지 않아요.

예를 들어 PunctuationType.STREAM_TIME 기반으로 10초마다 Punctuator 함수를 스케줄하고, 1초(첫 번째 레코드)부터 60초(마지막 레코드)까지의 연속 타임스탬프를 가진 60개 레코드의 스트림을 처리한다면, punctuate()는 6번 호출돼요. 이것은 이러한 레코드를 실제로 처리하는 데 걸리는 시간과 무관해요. 이 60개 레코드를 처리하는 데 1초가 걸리든, 1분이 걸리든, 1시간이 걸리든 punctuate()는 6번 호출돼요.

벽시계-시간(PunctuationType.WALL_CLOCK_TIME)이 사용되면 punctuate()는 순수하게 벽시계 시간에 의해 촉발돼요. 위 예시를 재사용하면, Punctuator 함수가 PunctuationType.WALL_CLOCK_TIME 기반으로 스케줄되고 이 60개 레코드가 20초 안에 처리되었다면, punctuate()는 2번 호출돼요 (10초마다 한 번). 이 60개 레코드가 5초 안에 처리되었다면 punctuate()는 전혀 호출되지 않아요. 같은 프로세서 안에서 init() 메서드에서 ProcessorContext#schedule()를 여러 번 호출해 다른 PunctuationType 유형으로 여러 Punctuator 콜백을 스케줄할 수 있다는 점을 유의해요.

주의: 스트림-시간은 Streams가 레코드를 처리할 때만 진행돼요. 처리할 레코드가 없거나, Streams가 Task Idling 구성 때문에 새 레코드를 기다리고 있다면 스트림 시간은 진행되지 않고, PunctuationType.STREAM_TIME을 지정했다면 punctuate()는 촉발되지 않아요. 이 동작은 구성된 타임스탬프 추출기와 무관해요. 즉 WallclockTimestampExtractor를 사용해도 punctuate()의 벽시계 촉발이 활성화되지는 않아요.

예시

다음 예시 Processor는 간단한 단어-카운트 알고리즘을 정의하고 다음 작업을 수행해요:

  • init() 메서드에서 1000 시간 단위마다 punctuation을 스케줄하고(시간 단위는 보통 밀리초이며, 이 예시에서는 1초마다 punctuation이 된다는 뜻) 이름 "Counts"로 로컬 상태 저장소를 검색해요.
  • process() 메서드에서 받은 각 레코드에 대해 값 문자열을 단어로 쪼개고 그 카운트를 상태 저장소에 업데이트해요 (이 섹션의 뒤에서 이야기할 거예요).
  • punctuate() 메서드에서 로컬 상태 저장소를 반복하고 집계된 카운트를 다운스트림 프로세서로 보내고(뒤에서 이야기할 거예요) 현재 스트림 상태를 커밋해요.
public class WordCountProcessor implements Processor<String, String, String, String> {
private KeyValueStore<String, Integer> kvStore;
@Override
public void init(final ProcessorContext<String, String> context) {
    context.schedule(Duration.ofSeconds(1), PunctuationType.STREAM_TIME, timestamp -> {
        try (final KeyValueIterator<String, Integer> iter = kvStore.all()) {
            while (iter.hasNext()) {
                final KeyValue<String, Integer> entry = iter.next();
                context.forward(new Record<>(entry.key, entry.value.toString(), timestamp));
            }
        }
    });
    kvStore = context.getStateStore("Counts");
}

@Override
public void process(final Record<String, String> record) {
    final String[] words = record.value().toLowerCase(Locale.getDefault()).split("\\W+");

    for (final String word : words) {
        final Integer oldValue = kvStore.get(word);

        if (oldValue == null) {
            kvStore.put(word, 1);
        } else {
            kvStore.put(word, oldValue + 1);
        }
    }
}

@Override
public void close() {
    // close any resources managed by this processor
    // Note: Do not close any StateStores as these are managed by the library
}
}

참고: 상태 저장소와의 상태 저장 처리 — 위에서 정의된 WordCountProcessorprocess() 메서드에서 현재 받은 레코드에 접근할 수 있고, 상태 저장소를 활용해 처리 상태를 유지할 수 있어요. 예를 들어 집계와 조인 같은 상태 저장 처리 요구를 위해 최근에 도착한 레코드를 기억할 수 있어요. 자세한 내용은 상태 저장소 문서를 참고해요.

프로세서 유닛 테스트

Kafka Streams는 프로세서 유닛 테스트를 작성하는 데 도움이 되는 test-utils 모듈을 함께 제공해요.

상태 저장소 (State Stores)

상태 저장 Processor를 구현하려면 프로세서에 하나 이상의 상태 저장소를 제공해야 해요 (무상태 프로세서는 상태 저장소가 필요 없어요). 상태 저장소는 최근에 받은 입력 레코드를 기억하고, 롤링 집계를 추적하고, 입력 레코드를 중복 제거하는 데 사용될 수 있어요. 상태 저장소의 또 다른 기능은 NodeJS 기반 대시보드나 Scala·Go로 구현된 마이크로서비스 같은 다른 애플리케이션에서 인터랙티브하게 쿼리될 수 있다는 것이에요.

Kafka Streams에서 사용 가능한 상태 저장소 유형은 기본적으로 장애 허용이 활성화되어 있어요.

상태 저장소 정의 및 생성

사용 가능한 저장소 유형 중 하나를 사용하거나 자신의 커스텀 저장소 유형을 구현할 수 있어요. Stores 팩토리를 사용해 기존 저장소 유형을 활용하는 것이 일반적인 관행이에요.

Kafka Streams를 사용할 때는 보통 코드에서 상태 저장소를 직접 만들거나 인스턴스화하지 않는다는 점을 유의해요. 오히려 소위 StoreBuilder를 만들어 상태 저장소를 간접적으로 정의해요. 이 빌더는 Kafka Streams가 필요할 때와 필요한 곳에 애플리케이션 인스턴스에 실제 상태 저장소를 인스턴스화하는 팩토리로 사용돼요.

다음 저장소 유형이 기본으로 제공돼요.

저장소 유형 저장 엔진 장애 허용? 설명
영속 KeyValueStore<K, V> RocksDB 예 (기본 활성화) 대부분의 사용 사례에 권장되는 저장소 유형. 데이터를 로컬 디스크에 저장. 저장 용량: 관리되는 로컬 상태는 애플리케이션 인스턴스의 메모리(힙 공간)보다 클 수 있지만, 사용 가능한 로컬 디스크 공간에는 맞아야 함. RocksDB 설정은 미세 조정 가능, RocksDB 구성 참고. 사용 가능한 영속 저장소 변형: plain 키-값 저장소(값만 — 상태에 임베드된 레코드 타임스탬프 없음), timestamped 키-값 저장소, versioned 키-값 저장소, 윈도우 저장소, 세션 저장소. 헤더 인지 변형도 아래에 설명됨. plain 키-값 저장소(임베드된 레코드 타임스탬프 없음)가 필요하면 persistentKeyValueStore. put/get/delete와 범위 쿼리를 지원하는 영속 키-(값/타임스탬프) 저장소가 필요하면 persistentTimestampedKeyValueStore. put/get/delete와 timestamped get 연산을 지원하는 영속 버전 키-(값/타임스탬프) 저장소가 필요하면 persistentVersionedKeyValueStore. 영속 plain 윈도우 저장소 또는 영속 timestamped 윈도우 저장소가 필요하면 각각 persistentWindowStore 또는 persistentTimestampedWindowStore. 영속 세션 저장소가 필요하면 persistentSessionStore. 헤더: 상태에 레코드 헤더를 영속하려면 해당 StoreBuilder 팩토리와 함께 WithHeaders 저장소 공급자를 사용(아래 State Stores의 Headers 참고). plain 영속 키-값 또는 plain 영속 윈도우 저장소에는 WithHeaders 공급자가 없음. WithHeaders 공급자는 영속 timestamped 키-값, 영속 timestamped 윈도우, 세션 저장소에만 존재. 헤더를 유지하고 put/get/delete와 범위 쿼리를 지원하는 영속 키-(값/타임스탬프) 저장소가 필요하면 timestampedKeyValueStoreWithHeadersBuilder와 함께 persistentTimestampedKeyValueStoreWithHeaders. 헤더를 유지하는 영속 timestamped 윈도우 저장소가 필요하면 timestampedWindowStoreWithHeadersBuilder와 함께 persistentTimestampedWindowStoreWithHeaders. 헤더를 유지하는 영속 세션 저장소가 필요하면 sessionStoreWithHeadersBuilder와 함께 persistentSessionStoreWithHeaders.
인메모리 KeyValueStore<K, V> 예 (기본 활성화) 데이터를 메모리에 저장. 저장 용량: 관리되는 로컬 상태는 애플리케이션 인스턴스의 메모리(힙 공간)에 맞아야 함. 로컬 디스크 공간을 사용할 수 없거나 앱 인스턴스 재시작 사이에 로컬 디스크 공간이 지워지는 환경에서 유용. 사용 가능한 인메모리 저장소 변형: plain 키-값 저장소, 윈도우 저장소, 세션 저장소. put/get/delete와 범위 쿼리를 지원하는 키-(값/타임스탬프) 저장소가 필요하면 TimestampedKeyValueStore. 윈도우된 키-(값/타임스탬프) 쌍을 저장해야 하면 TimestampedWindowStore. 현재 내장 인메모리 versioned 키-값 저장소는 없음.
// Creating a persistent key-value store:
// here, we create a `KeyValueStore<String, Long>` named "persistent-counts".
import org.apache.kafka.streams.state.StoreBuilder;
import org.apache.kafka.streams.state.Stores;
// Using a `KeyValueStoreBuilder` to build a `KeyValueStore`.
StoreBuilder<KeyValueStore<String, Long>> countStoreSupplier =
Stores.keyValueStoreBuilder(
Stores.persistentKeyValueStore("persistent-counts"),
Serdes.String(),
Serdes.Long());
KeyValueStore<String, Long> countStore = countStoreSupplier.build();
// Creating an in-memory key-value store:
// here, we create a `KeyValueStore<String, Long>` named "inmemory-counts".
import org.apache.kafka.streams.state.StoreBuilder;
import org.apache.kafka.streams.state.Stores;
// Using a `KeyValueStoreBuilder` to build a `KeyValueStore`.
StoreBuilder<KeyValueStore<String, Long>> countStoreSupplier =
Stores.keyValueStoreBuilder(
Stores.inMemoryKeyValueStore("inmemory-counts"),
Serdes.String(),
Serdes.Long());
KeyValueStore<String, Long> countStore = countStoreSupplier.build();

장애 허용 상태 저장소

상태 저장소를 장애 허용적으로 만들고 데이터 손실 없이 상태 저장소 마이그레이션을 허용하기 위해, 상태 저장소는 배후에서 카프카 토픽에 지속적으로 백업될 수 있어요. 예를 들어 애플리케이션에 용량을 탄력적으로 추가·제거할 때 상태 저장 스트림 태스크를 한 머신에서 다른 머신으로 마이그레이션하는 경우요. 이 토픽은 때때로 상태 저장소의 연관 체인지로그 토픽 또는 그 체인지로그라고 불려요. 예를 들어 머신 장애를 겪으면 상태 저장소와 애플리케이션의 상태가 그 체인지로그에서 완전히 복원될 수 있어요. 이 백업 기능은 상태 저장소에 대해 활성화하거나 비활성화할 수 있어요.

장애 허용 상태 저장소는 컴팩션된 체인지로그 토픽으로 백업돼요. 이 토픽을 컴팩션하는 목적은 토픽이 무한정 커지는 것을 방지하고, 연관 카프카 클러스터에서 소비되는 저장소를 줄이고, 상태 저장소를 체인지로그 토픽에서 복원해야 할 때 복구 시간을 최소화하기 위해서예요.

장애 허용 윈도우 상태 저장소는 컴팩션과 삭제를 모두 사용하는 토픽으로 백업돼요. 체인지로그 토픽으로 보내지는 메시지 키의 구조 때문에, 윈도우 저장소의 체인지로그 토픽에는 이 삭제·컴팩션 조합이 필요해요. 윈도우 저장소의 메시지 키는 "일반" 키와 윈도우 타임스탬프를 포함하는 복합 키예요. 이러한 복합 키 유형의 경우 체인지로그 토픽이 범위를 벗어나 커지는 것을 방지하기 위해 컴팩션만 활성화하는 것은 충분하지 않아요. 삭제가 활성화되면, 만료된 이전 윈도우는 로그 세그먼트가 만료됨에 따라 카프카의 로그 클리너에 의해 정리돼요. 기본 리텐션 설정은 Windows#maintainMs() + 1일이에요. StreamsConfig에서 StreamsConfig.WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG를 지정해 이 설정을 오버라이드할 수 있어요.

상태 저장소에서 Iterator를 열면 작업이 끝났을 때 리소스를 회수하기 위해 반복자에 close()를 호출해야 해요. 또는 try-with-resources 문 안에서 반복자를 사용할 수도 있어요. 반복자를 닫지 않으면 OOM 오류가 발생할 수 있어요.

상태 저장소의 장애 허용 활성화·비활성화 (저장소 체인지로그)

enableLogging()disableLogging()을 통해 저장소의 변경 로깅을 활성화·비활성화함으로써 상태 저장소의 장애 허용을 활성화·비활성화할 수 있어요. 필요하면 연관 토픽의 구성을 미세 조정할 수도 있어요.

장애 허용 비활성화 예시:

import org.apache.kafka.streams.state.StoreBuilder;
import org.apache.kafka.streams.state.Stores;

StoreBuilder<KeyValueStore<String, Long>> countStoreSupplier = Stores.keyValueStoreBuilder(
  Stores.persistentKeyValueStore("Counts"),
    Serdes.String(),
    Serdes.Long())
  .withLoggingDisabled(); // disable backing up the store to a changelog topic

주의: 체인지로그가 비활성화되면 연결된 상태 저장소는 더 이상 장애 허용이 아니고 스탠바이 복제본을 가질 수 없어요.

추가 체인지로그-토픽 구성을 포함한 장애 허용 활성화 예시: kafka.log.LogConfig의 어떤 로그 구성이든 추가할 수 있어요. 인식되지 않는 구성은 무시돼요.

import org.apache.kafka.streams.state.StoreBuilder;
import org.apache.kafka.streams.state.Stores;

Map<String, String> changelogConfig = new HashMap();
// override min.insync.replicas
changelogConfig.put(TopicConfig.MIN_IN_SYNC_REPLICAS_CONFIG, "1")

StoreBuilder<KeyValueStore<String, Long>> countStoreSupplier = Stores.keyValueStoreBuilder(
  Stores.persistentKeyValueStore("Counts"),
    Serdes.String(),
    Serdes.Long())
  .withLoggingEnabled(changelogConfig); // enable changelogging, with custom changelog settings

Timestamped 상태 저장소

KTables는 기본적으로 항상 타임스탬프를 저장해요. timestamped 상태 저장소는 스트림 처리 의미론을 개선하고 소스 KTable에서 순서가 뒤집힌 데이터 처리를 가능하게 하며, 순서가 뒤집힌 조인·집계 감지, 인터랙티브 쿼리에서 최신 업데이트의 타임스탬프 가져오기를 가능하게 해요.

timestamped 상태 저장소는 타임스탬프 유무와 관계없이 쿼리할 수 있어요.

업그레이드 참고: 모든 사용자는 인스턴스당 단일 롤링 바운스로 업그레이드해요.

  • Processor API 사용자의 경우 기존 애플리케이션에서 아무것도 변하지 않으며 timestamped 저장소를 사용할 옵션이 있어요.
  • DSL 연산자의 경우 저장소 데이터는 백그라운드에서 지연(lazily) 업그레이드돼요.
  • 커스텀 XxxBytesStoreSupplier를 제공하면 업그레이드가 발생하지 않지만 TimestampedBytesStore 인터페이스를 구현해 옵트인할 수 있어요. 이 경우 이전 형식이 유지되고 Streams는 읽기/쓰기 시 타임스탬프를 제거/추가하는 프록시 저장소를 사용해요.

상태 저장소의 헤더

Kafka 레코드 헤더를 키·값과 함께 RocksDB 기반 상태로 구체화(materialize)할 수 있어요. plain 영속 키-값 저장소는 임베드된 레코드 타임스탬프 없이 값을 유지하며, timestamped 키-값·윈도우 또는 세션 의미론에 대한 공급자는 저장소 유형에 따라 타임스탬프를 노출해요. 다운스트림 처리가 이전 입력의 레코드 헤더에 접근해야 할 때 사용해요 — 예를 들어 Processor API로 구현한 집계나 조인이 출력에 헤더를 전파해야 할 때요.

헤더 인지 저장소에는 영속 RocksDB 기반 공급자만 존재해요 (Stores 팩토리 이름이 persistent로 시작하고 WithHeaders로 끝남).

이름이 WithHeaders로 끝나는 Stores 메서드를 사용하며, 각각 해당 StoreBuilder 팩토리와 쌍을 이뤄요. 예를 들어 persistentTimestampedKeyValueStoreWithHeaderstimestampedKeyValueStoreWithHeadersBuilder와, persistentTimestampedWindowStoreWithHeaderstimestampedWindowStoreWithHeadersBuilder와, persistentSessionStoreWithHeaderssessionStoreWithHeadersBuilder와 쌍으로 사용해요.

키-값 및 윈도우 읽기는 값, 타임스탬프, 관련 헤더를 결합한 ValueTimestampHeaders를 반환해요. SessionStoreWithHeaders는 세션 집계를 AggregationWithHeaders, 즉 집계된 값과 해당 세션에 연결된 헤더로 저장해요.

업그레이드 참고: 롤링 바운스, 체인지로그 호환성, 지연 온디스크 마이그레이션, 성능 트레이드오프, 다운그레이드 제약은 Kafka Streams 업그레이드 가이드에서 다뤄요.

버전 키-값 상태 저장소

버전 키-값 상태 저장소는 Kafka Streams 3.5부터 사용 가능해요. 키당 단일 레코드 버전(값과 타임스탬프)을 저장하는 대신, 버전 상태 저장소는 키당 여러 레코드 버전을 저장할 수 있어요. 이것은 버전 상태 저장소가 timestamped 검색 연산을 지원해 지정된 타임스탬프 시점의 (키당) 최신 레코드를 반환할 수 있게 해줘요.

VersionedBytesStoreSupplierversionedKeyValueStoreBuilder에 전달하거나 자신의 VersionedKeyValueStore를 구현해 영속 버전 상태 저장소를 만들 수 있어요.

각 버전 저장소에는 이전 레코드 버전을 얼마나 오래 유지할지 지정하는 연관 고정 길이 히스토리 리텐션 파라미터가 있어요. 특히 버전 저장소는 쿼리되는 타임스탬프가 현재 관찰된 스트림 시간의 히스토리 리텐션 내에 있는 timestamped 검색 연산에 대해 정확한 결과를 반환하도록 보장해요.

히스토리 리텐션은 또한 그것의 유예 기간(grace period) 역할을 해요. 이것은 저장소에 대한 순서가 뒤집힌 쓰기가 얼마나 과거까지 받아들여질지 결정해요. 버전 저장소는 쓰기와 연관된 타임스탬프가 현재 관찰된 스트림 시간보다 유예 기간보다 더 오래된 경우 쓰기(삽입, 업데이트, 삭제)를 받아들이지 않아요. 이 맥락의 스트림 시간은 키별이 아니라 파티션별로 추적되므로, 한 키의 레코드가 다른 키의 레코드에 대해 순서가 뒤집혀 도착하는 것을 수용하도록 유예 기간(즉 히스토리 리텐션)을 충분히 높게 설정하는 것이 중요해요.

버전 키-값 상태 저장소의 메모리 사용량이 비버전 키-값 저장소보다 높으므로 RocksDB 메모리 설정을 그에 따라 조정하고 싶을 수 있어요. 성능이 비버전 저장소보다 나쁠 것으로 예상되므로 버전 저장소로 애플리케이션을 벤치마킹하는 것도 권장돼요.

버전 저장소는 현재 캐싱이나 인터랙티브 쿼리를 지원하지 않아요. 또한 윈도우 저장소와 글로벌 테이블은 버전화될 수 없어요.

업그레이드 참고: 버전 상태 저장소는 옵트인만 가능해요. 비버전에서 버전 저장소로의 자동 업그레이드는 발생하지 않아요. 원래 저장소가 업그레이드 대상 버전 저장소와 같은 체인지로그 토픽 형식을 가지는 한, 영속 비버전 키-값 저장소에서 영속 버전 키-값 저장소로의 업그레이드가 지원돼요. 영속 키-값 저장소와 timestamped 키-값 저장소 모두 영속 버전 키-값 저장소와 같은 체인지로그 토픽 형식을 공유하므로 둘 다 업그레이드 자격이 있어요.

영속 비버전 키-값 저장소를 사용하는 애플리케이션을 영속 버전 키-값 저장소를 사용하도록 업그레이드하려면 다음 절차를 수행할 수 있어요:

  1. 모든 애플리케이션 인스턴스를 중지하고, 업그레이드할 저장소의 로컬 상태 디렉터리를 모두 지운다.
  2. 원하는 곳에서 버전 저장소를 사용하도록 애플리케이션 코드를 업데이트한다.
  3. 관련 상태 저장소의 체인지로그 토픽 구성에서 min.compaction.lag.ms 값을 원하는 히스토리 리텐션 이상으로 설정한다. 컴팩션 중 브로커 벽시계 시간 사용을 위한 버퍼로 히스토리 리텐션에 하루를 더하는 것을 권장한다.
  4. 애플리케이션 인스턴스를 재시작하고 버전 저장소가 체인지로그에서 상태를 재구축할 시간을 허용한다.

읽기 전용 상태 저장소

읽기 전용 상태 저장소는 입력 토픽의 데이터를 구체화해요. 또한 장애 허용을 위해 입력 토픽을 사용하므로 추가 체인지로그 토픽이 없어요 (입력 토픽이 체인지로그로 재사용됨). 따라서 입력 토픽은 로그 컴팩션으로 구성되어야 해요. 다른 프로세서는 상태 저장소의 내용을 수정하지 말아야 하며, 유일한 작성자는 연관 "상태 업데이트 프로세서"여야 해요. 다른 프로세서는 읽기 전용 저장소의 내용을 읽을 수 있어요.

참고: 처리 중 조회를 위해 읽기 전용 상태 저장소를 사용할 때 파티셔닝 요구사항에 주의하세요. 원래 체인지로그 토픽이 읽기 전용 상태 저장소를 읽는 프로세서와 공동 파티션(copartition)되도록 확인하고 싶을 수 있어요.

커스텀 상태 저장소 구현

내장 상태 저장소 유형을 사용하거나 자신의 것을 구현할 수 있어요. 저장소를 위해 구현할 주요 인터페이스는 org.apache.kafka.streams.processor.StateStore예요. Kafka Streams에는 KeyValueStoreVersionedKeyValueStore 같은 몇 가지 확장 인터페이스도 있어요.

커스터마이즈된 org.apache.kafka.streams.processor.StateStore 구현은 org.apache.kafka.streams.processor.StateRestoreCallback 또는 org.apache.kafka.streams.processor.BatchingStateRestoreCallback 인터페이스를 통해 상태를 복원하는 로직도 제공해야 한다는 점을 유의해요. 이러한 인터페이스를 인스턴스화하는 방법에 대한 세부사항은 javadocs에서 찾을 수 있어요.

또한 org.apache.kafka.streams.state.StoreBuilder 인터페이스를 구현해 저장소의 "빌더"를 제공해야 하며, Kafka Streams가 이 빌더를 사용해 저장소 인스턴스를 만들어요.

프로세서 컨텍스트 접근

Defining a Stream Processor 섹션에서 언급했듯이 ProcessorContext는 punctuation 함수 스케줄링, 현재 처리된 상태 커밋 같은 처리 워크플로를 제어해요. 이 객체는 applicationId, taskId, stateDir 같은 애플리케이션과 관련된 메타데이터, 그리고 topic, partition, offset 같은 RecordMetadata에 접근하는 데도 사용될 수 있어요.

프로세서와 상태 저장소 연결

이제 프로세서(WordCountProcessor)와 상태 저장소가 정의되었으므로, Topology 인스턴스를 사용해 이 프로세서와 상태 저장소를 연결함으로써 프로세서 토폴로지를 구성할 수 있어요. 추가로, 지정된 카프카 토픽으로 소스 프로세서를 추가해 토폴로지에 입력 데이터 스트림을 생성하고, 지정된 카프카 토픽으로 싱크 프로세서를 추가해 토폴로지에서 출력 데이터 스트림을 생성할 수 있어요.

구현 예시:

Topology builder = new Topology();
// add the source processor node that takes Kafka topic "source-topic" as input
builder.addSource("Source", "source-topic")
    // add the WordCountProcessor node which takes the source processor as its upstream processor
    .addProcessor("Process", () -> new WordCountProcessor(), "Source")
    // add the count store associated with the WordCountProcessor processor
    .addStateStore(countStoreBuilder, "Process")
    // add the sink processor node that takes Kafka topic "sink-topic" as output
    // and the WordCountProcessor node as its upstream processor
    .addSink("Sink", "sink-topic", "Process");

이 예시의 빠른 설명:

  • addSource 메서드를 사용해 "Source"라는 이름의 소스 프로세서 노드가 토폴로지에 추가되고, "source-topic"이라는 카프카 토픽 하나가 공급돼요.
  • 사전 정의된 WordCountProcessor 로직을 가진 "Process"라는 이름의 프로세서 노드가 addProcessor 메서드를 사용해 "Source" 노드의 다운스트림 프로세서로 추가돼요.
  • 사전 정의된 영속 키-값 상태 저장소가 countStoreBuilder를 사용해 만들어져 "Process" 노드와 연결돼요.
  • addSink 메서드를 사용해 싱크 프로세서 노드가 추가되어 토폴로지를 완성하며, "Process" 노드를 업스트림 프로세서로 받고 별도의 "sink-topic" 카프카 토픽에 써요 (사용자는 addSink의 다른 오버로드 변형을 사용해 업스트림 프로세서에서 받은 각 레코드에 대해 쓸 카프카 토픽을 동적으로 결정할 수도 있어요).

어떤 경우에는 프로세서를 토폴로지에 추가할 때 상태 저장소를 동시에 추가·연결하는 것이 더 편리할 수 있어요. Topology#addStateStore()를 호출하는 대신 ProcessorSupplierConnectedStoreProvider#stores()를 구현하면 됩니다:

Topology builder = new Topology();
// add the source processor node that takes Kafka "source-topic" as input
builder.addSource("Source", "source-topic")
    // add the WordCountProcessor node which takes the source processor as its upstream processor.
    // the ProcessorSupplier provides the count store associated with the WordCountProcessor
    .addProcessor("Process", new ProcessorSupplier<String, String, String, String>() {
        public Processor<String, String, String, String> get() {
            return new WordCountProcessor();
        }

        public Set<StoreBuilder<?>> stores() {
            final StoreBuilder<KeyValueStore<String, Long>> countsStoreBuilder =
                Stores
                    .keyValueStoreBuilder(
                        Stores.persistentKeyValueStore("Counts"),
                        Serdes.String(),
                        Serdes.Long()
                    );
            return Collections.singleton(countsStoreBuilder);
        }
    }, "Source")
    // add the sink processor node that takes Kafka topic "sink-topic" as output
    // and the WordCountProcessor node as its upstream processor
    .addSink("Sink", "sink-topic", "Process");

이것은 프로세서가 상태 저장소를 "소유"하게 해, 토폴로지를 구성하는 사용자로부터 그 사용을 효과적으로 캡슐화해요. 상태 저장소를 공유하는 여러 프로세서는 StoreBuilder가 같은 instance인 한 이 기술로 같은 저장소를 제공할 수 있어요.

이러한 토폴로지에서 "Process" 스트림 프로세서 노드는 "Source" 노드의 다운스트림 프로세서이자 "Sink" 노드의 업스트림 프로세서로 간주돼요. 따라서 "Source" 노드가 카프카에서 새로 가져온 레코드를 다운스트림 "Process" 노드로 전달할 때마다, WordCountProcessor#process() 메서드가 촉발되어 레코드를 처리하고 연관 상태 저장소를 업데이트해요. WordCountProcessor#punctuate() 메서드에서 context#forward()가 호출될 때마다 집계 레코드가 "Sink" 프로세서 노드를 통해 카프카 토픽 "sink-topic"으로 전송돼요. WordCountProcessor 구현에서 키-값 저장소에 접근할 때 같은 저장소 이름 "Counts"를 참조해야 한다는 점을 유의해요. 그렇지 않으면 상태 저장소를 찾을 수 없다는 런타임 예외가 발생해요. 상태 저장소가 Topology 코드의 프로세서와 연결되지 않았다면, 프로세서의 init() 메서드에서 접근해도 런타임 예외가 발생해 이 프로세서에서 상태 저장소에 접근할 수 없다고 표시돼요.

Topology#addProcessor 함수는 ProcessorSupplier를 인자로 받으며, 공급자 패턴은 ProcessorSupplier#get()이 호출될 때마다 새 Processor 인스턴스가 반환되어야 한다는 것을 요구해요. 단일 Processor 객체를 만들고 ProcessorSupplier#get()에서 같은 객체 참조를 반환하는 것은 공급자 패턴을 위반하고 런타임 예외를 초래해요. 따라서 Topology에 싱글턴 Processor 인스턴스를 제공하지 말아야 한다는 것을 기억하세요. ProcessorSupplierProcessorSupplier#get()이 호출될 때마다 항상 새 인스턴스를 생성해야 해요.

이제 애플리케이션에서 프로세서 토폴로지를 완전히 정의했으니, Kafka Streams 애플리케이션 실행을 진행할 수 있어요.

더 알아보기