Streams 애플리케이션 구성하기
Streams 애플리케이션 구성하기
카프카 스트림즈 애플리케이션을 만들었다면 이제 구성을 잡아야 해요. 구성 방법은 간단해요 — java.util.Properties 인스턴스에 파라미터를 넣으면 돼요. 이 페이지에서는 필수 구성(application.id, bootstrap.servers)부터 복원성(resiliency)을 위한 권장 구성, 그리고 각종 선택 구성 파라미터의 의미를 하나씩 풀어서 설명해드릴게요.
출처: 문서
본문
Streams를 사용하기 전에 카프카와 Kafka Streams 구성 옵션을 설정해야 해요. java.util.Properties 인스턴스에 파라미터를 지정해 Kafka Streams를 구성할 수 있어요.
java.util.Properties인스턴스를 만든다.- 파라미터를 설정한다. 예를 들어:
import java.util.Properties;
import org.apache.kafka.streams.StreamsConfig;
Properties settings = new Properties();
// Set a few key parameters
settings.put(StreamsConfig.APPLICATION_ID_CONFIG, "my-first-streams-application");
settings.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker1:9092");
// Any further settings
settings.put(... , ...);
구성 파라미터 참조
이 섹션은 가장 일반적인 Streams 구성 파라미터를 포함해요. 전체 참조는 Streams Javadocs를 봐요.
필수 구성 파라미터
application.idbootstrap.servers
복원성을 위한 권장 구성 파라미터
acksreplication.factormin.insync.replicasnum.standby.replicas
선택 구성 파라미터
acceptable.recovery.lagdefault.deserialization.exception.handler(deprecated)default.key.serdedefault.production.exception.handler(deprecated)default.timestamp.extractordefault.value.serdedeserialization.exception.handlerdsl.store.formatenable.metrics.pushensure.explicit.internal.resource.naminggroup.protocollog.summary.interval.msmax.task.idle.msmax.warmup.replicasnum.standby.replicasnum.stream.threadsprobing.rebalance.interval.msprocessing.exception.handlerprocessing.exception.handler.global.enabled(deprecated)processing.guaranteeprocessor.wrapper.classproduction.exception.handlerrack.aware.assignment.non_overlap_costrack.aware.assignment.strategyrack.aware.assignment.tagsrack.aware.assignment.traffic_costreplication.factorrocksdb.config.setterstate.dirtask.assignor.classtopology.optimization
Kafka 컨슈머·프로듀서 구성 파라미터: 명명, 기본값, Kafka Streams가 제어하는 파라미터(enable.auto.commit 등)
필수 구성 파라미터
다음은 필수 Streams 구성 파라미터예요.
| 파라미터 이름 | 중요도 | 설명 | 기본값 |
|---|---|---|---|
application.id |
필수 | 스트림 처리 애플리케이션의 식별자. 카프카 클러스터 내에서 고유해야 함 | 없음 |
bootstrap.servers |
필수 | 카프카 클러스터에 대한 초기 연결을 수립하는 데 사용할 호스트/포트 쌍 목록 | 없음 |
application.id
(필수) 애플리케이션 ID. 각 스트림 처리 애플리케이션은 고유한 ID를 가져야 해요. 애플리케이션의 모든 인스턴스에 같은 ID를 주어야 해요. 영숫자 문자, .(점), -(하이픈), _(밑줄)만 사용하는 것을 권장해요. 예: "hello_world", "hello_world-v1.0.0"
이 ID는 애플리케이션이 사용하는 리소스를 다른 것과 격리하기 위해 다음 위치에서 사용돼요:
- 기본 카프카 컨슈머·프로듀서
client.id접두사 - 조정을 위한 카프카 컨슈머
group.id - 상태 디렉터리에서 서브디렉터리의 이름 (
state.dir참고) - 내부 카프카 토픽 이름의 접두사
팁: 애플리케이션이 업데이트될 때, 내부 토픽과 상태 저장소의 기존 데이터를 재사용하고 싶지 않다면
application.id를 변경해야 해요. 예를 들어application.id에 버전 정보를my-app-v1.0.0,my-app-v1.0.2처럼 임베드할 수 있어요.
bootstrap.servers
(필수) 카프카 부트스트랩 서버. 기본 프로듀서·컨슈머 클라이언트가 카프카 클러스터에 연결하는 데 사용하는 것과 같은 설정이에요. 예: "kafka-broker1:9092,kafka-broker2:9092".
복원성을 위한 권장 구성 파라미터
브로커 장애에 대한 복원성을 위해 명시적으로 구성해야 하는 여러 카프카 및 Kafka Streams 구성 옵션이 있어요:
| 파라미터 이름 | 해당 클라이언트 | 기본값 | 다음과 같이 설정 고려 |
|---|---|---|---|
acks |
Producer (version <=2.8) | acks="1") |
acks="all" |
replication.factor (broker version 2.3 또는 이전) |
Streams | -1 |
3 (broker 2.4+: 브로커 구성 default.replication.factor=3 확인) |
min.insync.replicas |
Broker | 1 |
2 |
num.standby.replicas |
Streams | 0 |
1 |
복제 팩터를 3으로 높이면 내부 Kafka Streams 토픽이 최대 2개의 브로커 장애를 견딜 수 있도록 보장돼요. 기본값에서 권장값으로 옮겨가는 트레이드오프는 일부 성능과 더 많은 저장 공간(복제 팩터 3에서 3배)이 더 큰 복원성을 위해 희생된다는 것이에요.
acks
리더가 요청을 완료한 것으로 간주하기 전에 받아야 하는 승인(acknowledgment) 수예요. 전송된 레코드의 내구성을 제어해요. 가능한 값:
acks="0"— 프로듀서는 서버의 승인을 기다리지 않고 레코드는 즉시 소켓 버퍼에 추가되어 전송된 것으로 간주돼요. 이 경우 서버가 레코드를 받았다는 보장이 없고, 프로듀서는 일반적으로 어떤 장애도 알지 못해요. 각 레코드에 대해 반환된 오프셋은 항상-1로 설정돼요.acks="1"— 리더는 레코드를 자신의 로컬 로그에 쓰고 모든 팔로워의 전체 승인을 기다리지 않고 응답해요. 리더가 레코드를 승인한 직후, 팔로워가 복제하기 전에 즉시 실패하면 레코드는 손실돼요.acks="all"(3.0 릴리스부터 기본값) — 리더가 in-sync 복제본의 전체 집합이 레코드를 승인할 때까지 기다려요. 이것은 in-sync 복제본이 하나라도 살아 있으면 레코드가 손실되지 않을 것을 보장해요. 사용 가능한 가장 강한 보장이에요.
자세한 내용은 카프카 Producer 문서를 참고해요.
replication.factor
여기의 설명을 참고해요.
min.insync.replicas
프로듀서가 acks="all"로 구성된 경우 복제에 사용 가능한 최소 in-sync 복제본 수 (토픽 구성 참고).
num.standby.replicas
여기의 설명을 참고해요.
Properties streamsSettings = new Properties();
// for broker version 2.3 or older
//streamsSettings.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);
// for version 2.8 or older
//streamsSettings.put(StreamsConfig.producerPrefix(ProducerConfig.ACKS_CONFIG), "all");
streamsSettings.put(StreamsConfig.topicPrefix(TopicConfig.MIN_IN_SYNC_REPLICAS_CONFIG), 2);
streamsSettings.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1);
선택 구성 파라미터
다음은 중요도 순으로 정렬된 선택 Streams javadocs예요:
- High: 기본값이 프로덕션 사용에 적합하지 않을 가능성이 높은 파라미터. 프로덕션 사용을 위해 이 파라미터를 다시 검토하는 것을 적극 권장해요.
- Medium: 이 파라미터의 기본값은 많은 경우 프로덕션에서 동작해야 하지만, 예를 들어 성능을 튜닝하기 위해 변경되는 일이 드물지 않아요.
- Low: 이 파라미터의 값을 변경할 필요는 거의 없어야 해요. 아주 특정한 문제를 해결하고 싶을 때만 변경하는 것이 좋아요.
| 파라미터 이름 | 중요도 | 설명 | 기본값 |
|---|---|---|---|
acceptable.recovery.lag |
Medium | 인스턴스가 따라잡힌(caught-up) 것으로 간주되어 활성 태스크에 준비되기 위한 최대 허용 lag(따라잡아야 할 오프셋 수) | 10000 |
allow.os.group.write.access |
Low | Kafka Streams가 만든 상태 저장소 디렉터리에 OS 그룹의 쓰기 접근을 허용 | false |
application.server |
Low | 단일 Kafka Streams 애플리케이션 내에서 상태 저장소 위치를 발견하는 데 사용할 수 있는 임베드된 사용자 정의 엔드포인트를 가리키는 host:port 쌍. 애플리케이션의 각 인스턴스마다 값이 달라야 함 | 빈 문자열 |
buffered.records.per.partition |
Low | 파티션당 버퍼링할 최대 레코드 수 | 1000 |
statestore.cache.max.bytes |
Medium | 모든 스레드에 걸쳐 레코드 캐시에 사용할 최대 메모리 바이트 수 | 10485760 |
cache.max.bytes.buffering (Deprecated. statestore.cache.max.bytes 사용) |
Medium | 모든 스레드에 걸쳐 레코드 캐시에 사용할 최대 메모리 바이트 수 | 10485760 |
client.id |
Medium | 요청 시 서버에 전달할 ID 문자열. (Kafka Streams가 내부적으로 사용하는 컨슈머/프로듀서 클라이언트에 전달됨) | 빈 문자열 |
commit.interval.ms |
Low | 태스크의 위치(소스 토픽의 오프셋)를 저장하는 빈도(밀리초) | 30000 (30초) (at-least-once) / 100 (exactly-once) |
default.deserialization.exception.handler (Deprecated. deserialization.exception.handler 사용) |
Medium | DeserializationExceptionHandler 인터페이스를 구현하는 예외 처리 클래스 |
LogAndFailExceptionHandler |
default.key.serde |
Medium | 레코드 키의 기본 시리얼라이저/디시리얼라이저 클래스, Serde 인터페이스 구현. 사용자가 설정하거나 모든 serde를 명시적으로 전달해야 함 (default.value.serde 참고) |
null |
default.production.exception.handler (Deprecated. production.exception.handler 사용) |
Medium | ProductionExceptionHandler 인터페이스를 구현하는 예외 처리 클래스 |
DefaultProductionExceptionHandler |
default.timestamp.extractor |
Medium | TimestampExtractor 인터페이스를 구현하는 타임스탬프 추출기 클래스. Timestamp Extractor 참고 |
FailOnInvalidTimestamp |
default.value.serde |
Medium | 레코드 값의 기본 시리얼라이저/디시리얼라이저 클래스, Serde 인터페이스 구현. 사용자가 설정하거나 모든 serde를 명시적으로 전달해야 함 (default.key.serde 참고) |
null |
default.dsl.store (Deprecated. dsl.store.suppliers.class 사용) |
Low | DSL 연산자가 사용하는 기본 상태 저장소 유형 | "ROCKS_DB" |
deserialization.exception.handler |
Medium | DeserializationExceptionHandler 인터페이스를 구현하는 예외 처리 클래스 |
LogAndContinueExceptionHandler |
dsl.store.suppliers.class |
Low | 저장소 구현 유형을 명시적으로 구성하지 않은 모든 상태 저장 DSL 연산자가 사용할 기본 상태 저장소 구현을 정의. org.apache.kafka.streams.state.DslStoreSuppliers 인터페이스를 구현해야 함 |
BuiltInDslStoreSuppliers.RocksDBDslStoreSuppliers |
dsl.store.format |
Low | DSL 연산자가 헤더 인지 상태 저장소를 구체화(materialize)할지 제어. 대소문자 구분 없음. 허용 값: default(연산자별 기존 timestamped 또는 plain 저장소 변형 사용), headers(값·타임스탬프와 함께 레코드 헤더를 영속할 수 있는 헤더 인지 저장소 선택) |
default |
ensure.explicit.internal.resource.naming |
High | 내부 토픽(예: 체인지로그·리파티션 토픽)과 관련 상태 저장소를 포함한 토폴로지의 모든 내부 리소스에 명시적 명명을 강제할지 여부. 활성화되면 내부 리소스에 자동 생성 이름이 있으면 애플리케이션이 시작을 거부 | false |
log.summary.interval.ms |
Low | 요약 정보를 기록하기 위한 출력 간격(밀리초) (음수면 비활성화) | 120000 (2분) |
enable.metrics.push |
Low | 이 클라이언트와 일치하는 클라이언트 메트릭 구독이 클러스터에 있으면 클라이언트 메트릭을 클러스터로 푸시할지 여부 | true |
max.task.idle.ms |
Medium | 조인·머지가 순서가 뒤집힌 결과를 생성할 수 있는지 제어. 구성 값은 스트림 태스크가 일부(전부는 아닌) 입력 파티션에 완전히 따라잡았을 때 유휴 상태로 머무를 최대 시간(밀리초)으로, 프로듀서가 추가 레코드를 보낼 때까지 기다려 여러 입력 스트림에 걸친 순서가 뒤집힌 레코드 처리를 피하기 위함. 기본(0)은 프로듀서가 더 많은 레코드를 보내기를 기다리지 않지만, 브로커에 이미 있는 데이터는 가져오기를 기다려요. 이 기본값은 브로커에 이미 있는 레코드에 대해 Streams가 타임스탬프 순서로 처리한다는 뜻이에요. -1로 설정하면 유휴 상태를 완전히 비활성화하고 로컬에서 사용 가능한 데이터를 처리하며, 순서가 뒤집힌 처리를 만들 수 있음 |
0 |
max.warmup.replicas |
Medium | 한 번에 할당할 수 있는 최대 워밍업 복제본 수 (구성된 num.standbys를 초과하는 추가 스탠바이) | 2 |
metric.reporters |
Low | 메트릭 리포터로 사용할 클래스 목록 | 빈 목록 |
metrics.num.samples |
Low | 메트릭 계산을 위해 유지되는 샘플 수 | 2 |
metrics.recording.level |
Low | 메트릭의 최고 기록 레벨 | INFO |
metrics.sample.window.ms |
Low | 메트릭 샘플이 계산되는 시간 창(밀리초) | 30000 (30초) |
num.standby.replicas |
High | 각 태스크의 스탠바이 복제본 수 | 0 |
num.stream.threads |
Medium | 스트림 처리를 실행하는 스레드 수 | 1 |
probing.rebalance.interval.ms |
Low | 충분히 따라잡힌 워밍업 복제본을 검색하기 위한 리밸런스 촉발 전 최대 대기 시간(밀리초) | 600000 (10분) |
processing.exception.handler |
Medium | ProcessingExceptionHandler 인터페이스를 구현하는 예외 처리 클래스 |
LogAndFailProcessingExceptionHandler |
processing.guarantee |
Medium | 처리 모드. "at_least_once" 또는 "exactly_once_v2" (EOS version 2, 브로커 버전 2.5+ 필요). Processing Guarantee 참고 |
"at_least_once" |
processor.wrapper.class |
Medium | ProcessorWrapper 인터페이스를 구현하는 클래스 또는 클래스 이름. 토폴로지를 만들 때 전달해야 하며, 적절한 생성자에 TopologyConfig로 전달하지 않으면 적용되지 않음. DSL 애플리케이션은 StreamsBuilder#new(TopologyConfig) 생성자, PAPI 애플리케이션은 Topology#new(TopologyConfig) 생성자를 사용해야 함 |
— |
production.exception.handler |
Medium | ProductionExceptionHandler 인터페이스를 구현하는 예외 처리 클래스 |
DefaultProductionExceptionHandler |
poll.ms |
Low | 입력을 기다리며 차단할 시간(밀리초) | 100 |
rack.aware.assignment.strategy |
Low | 랙 인지 할당에 사용되는 전략. 허용 값: "none"(기본), "min_traffic", "balance_suttopology". Rack Aware Assignment Strategy 참고 |
"none" |
rack.aware.assignment.tags |
Low | 스탠바이 복제본을 Kafka Streams 클라이언트 전체에 분산하는 데 사용되는 태그 키 목록. 구성되면 Kafka Streams는 다른 태그 값을 가진 클라이언트에 스탠바이 태스크 분산을 최선으로 시도. Rack Aware Assignment Tags 참고 | 빈 목록 |
rack.aware.assignment.non_overlap_cost |
Low | 기존 할당에서 태스크를 이동하는 비용. Rack Aware Assignment Non-Overlap-Cost 참고 | null |
rack.aware.assignment.non_overlap_cost |
Low | 교차 랙 트래픽과 관련된 비용. Rack Aware Assignment Traffic-Cost 참고 | null |
replication.factor |
Medium | 애플리케이션이 만들 체인지로그 토픽과 리파티션 토픽의 복제 팩터. 기본값 -1(브로커 기본 복제 팩터 사용)은 브로커 버전 2.4 이상 필요 |
-1 |
repartition.purge.interval.ms |
Low | 리파티션 토픽에서 완전히 소비된 레코드를 삭제하는 빈도(밀리초). 마지막 퍼지 이후 이 값 이상 시간이 지나면 퍼지가 발생하지만 나중에 지연될 수 있음 | 30000 (30초) |
retry.backoff.ms |
Low | 요청을 재시도하기 전의 시간(밀리초) | 100 |
rocksdb.config.setter |
Medium | RocksDB 구성 | null |
state.cleanup.delay.ms |
Low | 파티션이 마이그레이션되었을 때 상태를 삭제하기 전에 기다릴 시간(밀리초) | 600000 (10분) |
state.cleanup.dir.max.age.ms |
Low | 애플리케이션 시작 시 로컬 상태 디렉터리와 체크포인트 파일을 정리하기 위한 시간 기반 임계값. 적어도 state.cleanup.dir.max.age.ms 동안 수정되지 않은 상태 디렉터리는 제거됨 |
-1 (비활성화) |
state.dir |
High | 상태 저장소의 디렉터리 위치 | /${java.io.tmpdir}/kafka-streams |
task.assignor.class |
Medium | TaskAssignor 인터페이스를 구현하는 태스크 assignor 클래스 또는 클래스 이름 |
고가용성 태스크 assignor |
task.timeout.ms |
Medium | 태스크가 내부 오류와 재시도로 인해 오류가 발생할 때까지 멈출 수 있는 최대 시간(밀리초). 0 ms 타임아웃의 경우 태스크는 첫 번째 내부 오류에서 오류를 발생. 0 ms보다 큰 타임아웃의 경우 태스크는 오류가 발생하기 전에 적어도 한 번 재시도 |
300000 (5분) |
topology.optimization |
Medium | Kafka Streams가 토폴로지를 최적화해야 하는지와 어떤 최적화를 적용할지 알려주는 구성. 허용 값: StreamsConfig.NO_OPTIMIZATION (none), StreamsConfig.OPTIMIZE (all), 또는 특정 최적화의 쉼표 구분 목록: StreamsConfig.REUSE_KTABLE_SOURCE_TOPICS (reuse.ktable.source.topics), StreamsConfig.MERGE_REPARTITION_TOPICS (merge.repartition.topics), StreamsConfig.SINGLE_STORE_SELF_JOIN (single.store.self.join) |
"NO_OPTIMIZATION" |
upgrade.from |
Medium | 롤링 업그레이드 중 업그레이드하는 버전. Upgrade From 참고 | null |
windowstore.changelog.additional.retention.ms |
Low | 로그에서 데이터가 조기 삭제되지 않도록 windows maintainMs에 더해지는 값. 클록 드리프트 허용 | 86400000 (1일) |
window.size.ms (Deprecated. Window Serdes 참고) |
Low | 윈도우 종료 시간 계산을 위해 디시리얼라이저에 윈도우 크기 설정 | null |
windowed.inner.class.serde (Deprecated. Window Serdes 참고) |
Low | 윈도우 레코드의 내부 클래스 serde. Serde 인터페이스를 구현해야 함 |
null |
acceptable.recovery.lag
인스턴스가 따라잡히고 활성 태스크를 받을 수 있는 것으로 간주되기 위한 최대 허용 lag(체인지로그에서 따라잡아야 할 총 오프셋 수)예요. Streams는 상태 저장소가 acceptable recovery lag 내에 있는 인스턴스(존재하는 경우)에만 상태 저장 활성 태스크를 할당하고, 아직 따라잡지 못한 인스턴스에 대해서는 워밍업 복제본을 할당해 백그라운드에서 상태를 복원해요. 주어진 워크로드에 대해 잘 1분 미만의 복구 시간에 해당해야 해요. 최소 0이어야 해요.
참고: 이 값을
Long.MAX_VALUE로 설정하면 워밍업 복제본과 태스크 고가용성을 효과적으로 비활성화해, Streams가 워밍업 없이 즉시 균형 잡힌 할당을 만들고 태스크를 새 인스턴스로 마이그레이션할 수 있게 해요.
deserialization.exception.handler (deprecated: default.deserialization.exception.handler)
역직렬화 예외 핸들러를 사용하면 역직렬화에 실패하는 레코드 예외를 관리할 수 있어요. 이것은 손상된 데이터, 잘못된 직렬화 로직, 또는 처리되지 않은 레코드 유형으로 인해 발생할 수 있어요. 구현된 예외 핸들러는 레코드와 발생한 예외에 따라 FAIL 또는 CONTINUE를 반환해야 해요. FAIL을 반환하면 Streams가 종료되어야 한다는 신호이고, CONTINUE는 Streams가 문제를 무시하고 계속 처리해야 한다는 신호예요. 다음 라이브러리 내장 예외 핸들러가 사용 가능해요:
- LogAndContinueExceptionHandler: 이 핸들러는 역직렬화 예외를 기록한 다음 처리 파이프라인이 더 많은 레코드를 계속 처리하도록 신호를 보내요. 이 로그-및-건너뛰기 전략은 역직렬화에 실패하는 레코드가 있으면 실패하는 대신 Kafka Streams가 진행할 수 있게 해줘요.
- LogAndFailExceptionHandler: 이 핸들러는 역직렬화 예외를 기록한 다음 처리 파이프라인이 더 많은 레코드 처리를 중지하도록 신호를 보내요.
라이브러리 제공 핸들러 외에 요구에 맞는 자신의 커스터마이즈된 예외 핸들러를 제공할 수도 있어요. 예를 들어 손상된 레코드를 격리 토픽("데드 레터 큐"라고 생각하면 됩니다)으로 전달해 추가 처리를 하도록 선택할 수 있어요. 이렇게 하려면 Producer API를 사용해 손상된 레코드를 격리 토픽에 직접 써요. 더 구체적으로, Streams 클라이언트 밖에 별도의 KafkaProducer 객체를 만들고 이 객체와 데드 레터 큐 토픽 이름을 Properties 맵에 전달하면, configure 함수 호출에서 검색할 수 있어요. 이 접근 방식의 단점은 "수동" 쓰기가 Kafka Streams 런타임 라이브러리에서 보이지 않는 부작용이라, Streams API의 엔드투엔드 처리 보장 혜택을 받지 못한다는 점이에요:
public class SendToDeadLetterQueueExceptionHandler implements DeserializationExceptionHandler {
KafkaProducer<byte[], byte[]> dlqProducer;
String dlqTopic;
@Override
public DeserializationHandlerResponse handle(final ErrorHandlerContext context,
final ConsumerRecord<byte[], byte[]> record,
final Exception exception) {
log.warn("Exception caught during Deserialization, sending to the dead queue topic; " +
"taskId: {}, topic: {}, partition: {}, offset: {}",
context.taskId(), record.topic(), record.partition(), record.offset(),
exception);
dlqProducer.send(new ProducerRecord<>(dlqTopic, record.timestamp(), record.key(), record.value(), record.headers())).get();
return DeserializationHandlerResponse.CONTINUE;
}
@Override
public void configure(final Map<String, ?> configs) {
dlqProducer = .. // get a producer from the configs map
dlqTopic = .. // get the topic name from the configs map
}
}
production.exception.handler (deprecated: default.production.exception.handler)
프로덕션 예외 핸들러를 사용하면 너무 큰 레코드를 생성하려는 것처럼 브로커와 상호작용하려고 할 때 발생하는 예외를 관리할 수 있어요. 기본적으로 카프카는 이러한 예외가 발생할 때 항상 실패하는 DefaultProductionExceptionHandler를 제공하고 사용해요.
예외 핸들러는 레코드와 발생한 예외에 따라 FAIL, CONTINUE, 또는 RETRY를 반환할 수 있어요. FAIL을 반환하면 Streams가 종료되어야 한다는 신호예요. CONTINUE는 Streams가 문제를 무시하고 계속 처리해야 한다는 신호예요. RetriableException의 경우 핸들러가 RETRY를 반환해 런타임에 실패한 레코드 전송을 재시도하라고 알릴 수 있어요 (참고: RetriableException이 아닌 것에 RETRY를 반환하면 FAIL로 취급돼요). 항상 너무 큰 레코드를 무시하는 예외 핸들러를 제공하고 싶다면 다음과 같이 구현할 수 있어요:
import java.util.Properties;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.common.errors.RecordTooLargeException;
import org.apache.kafka.streams.errors.ProductionExceptionHandler;
import org.apache.kafka.streams.errors.ProductionExceptionHandler.ProductionExceptionHandlerResponse;
public class IgnoreRecordTooLargeHandler implements ProductionExceptionHandler {
public void configure(Map<String, Object> config) {}
public ProductionExceptionHandlerResponse handle(final ErrorHandlerContext context,
final ProducerRecord<byte[], byte[]> record,
final Exception exception) {
if (exception instanceof RecordTooLargeException) {
return ProductionExceptionHandlerResponse.CONTINUE;
} else {
return ProductionExceptionHandlerResponse.FAIL;
}
}
}
Properties settings = new Properties();
// other various kafka streams settings, e.g. bootstrap servers, application id, etc
settings.put(StreamsConfig.PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG,
IgnoreRecordTooLargeHandler.class);
default.timestamp.extractor
타임스탬프 추출기는 ConsumerRecord 인스턴스에서 타임스탬프를 뽑아내요. 타임스탬프는 스트림의 진행을 제어하는 데 사용돼요.
기본 추출기는 FailOnInvalidTimestamp예요. 이 추출기는 카프카 버전 0.10 이후 카프카 프로듀서 클라이언트가 카프카 메시지에 자동으로 임베드하는 내장 타임스탬프를 가져와요. 카프카의 서버 측 log.message.timestamp.type 브로커와 message.timestamp.type 토픽 파라미터의 설정에 따라 이 추출기는 다음을 제공해요:
log.message.timestamp.type이CreateTime("producer time", 기본값)으로 설정되면 이벤트 시간 처리 의미론. 이것은 카프카 프로듀서가 원래 메시지를 보낸 시간을 나타내요. 카프카 공식 프로듀서 클라이언트를 사용하면 타임스탬프는 epoch 이후 밀리초를 나타내요.log.message.timestamp.type이LogAppendTime("broker time")으로 설정되면 수집 시간 처리 의미론. 이것은 카프카 브로커가 원래 메시지를 받은 시간을 epoch 이후 밀리초로 나타내요.
FailOnInvalidTimestamp 추출기는 레코드에 유효하지 않은(즉 음수) 내장 타임스탬프가 포함되어 있으면 예외를 던져요. Kafka Streams가 이 레코드를 처리하지 않고 조용히 버릴 것이기 때문이에요. 유효하지 않은 내장 타임스탬프는 여러 이유로 발생할 수 있어요: 예를 들어 새 Kafka 0.10 메시지 형식을 아직 지원하지 않는 0.10 이전 카프카 프로듀서 클라이언트나 타사 프로듀서 클라이언트가 쓴 토픽을 소비하는 경우, 그리고 카프카 클러스터를 0.9에서 0.10으로 업그레이드한 후 모든 0.9로 생성된 데이터가 0.10 메시지 타임스탬프를 포함하지 않는 상황이에요.
유효하지 않은 타임스탬프가 있는 데이터를 가지고 있고 처리하고 싶다면 두 가지 대안 추출기가 있어요. 둘 다 내장 타임스탬프로 동작하지만 유효하지 않은 타임스탬프를 다르게 처리해요:
- LogAndSkipOnInvalidTimestamp: 이 추출기는 경고 메시지를 기록하고 유효하지 않은 타임스탬프를 Kafka Streams에 반환하며, Kafka Streams는 처리하지 않고 레코드를 조용히 버려요. 이 로그-및-건너뛰기 전략은 입력 데이터에 유효하지 않은 내장 타임스탬프가 있는 레코드가 있으면 실패하는 대신 Kafka Streams가 진행할 수 있게 해줘요.
- UsePartitionTimeOnInvalidTimestamp: 이 추출기는 유효하면(즉 음수가 아니면) 레코드의 내장 타임스탬프를 반환해요. 유효한 내장 타임스탬프가 없으면 현재 레코드와 같은 토픽 파티션의 레코드에서 이전에 추출한 유효한 타임스탬프를 타임스탬프 추정치로 반환해요. 타임스탬프를 추정할 수 없으면 예외를 던져요.
또 다른 내장 추출기는 WallclockTimestampExtractor예요. 이 추출기는 실제로 소비된 레코드에서 타임스탬프를 "추출"하지 않고 시스템 시계에서 현재 시간(밀리초)을 반환해요(System.currentTimeMillis()라고 생각하면 됩니다). 이것은 Streams가 소위 이벤트의 처리 시간(processing-time)을 기반으로 동작한다는 뜻이에요.
메시지 페이로드에 임베드된 타임스탬프를 가져오는 것처럼 자신의 타임스탬프 추출기를 제공할 수도 있어요. 유효한 타임스탬프를 추출할 수 없으면 예외를 던지거나, 음수 타임스탬프를 반환하거나, 타임스탬프를 추정할 수 있어요. 음수 타임스탬프를 반환하면 데이터 손실이 발생해요 — 해당 레코드는 처리되지 않고 조용히 버려져요. 새 타임스탬프를 추정하려면 previousTimestamp로 제공된 값(즉 Kafka Streams 타임스탬프 추정치)을 사용할 수 있어요. 다음은 커스텀 TimestampExtractor 구현 예시예요:
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.streams.processor.TimestampExtractor;
// Extracts the embedded timestamp of a record (giving you "event-time" semantics).
public class MyEventTimeExtractor implements TimestampExtractor {
@Override
public long extract(final ConsumerRecord<Object, Object> record, final long previousTimestamp) {
// `Foo` is your own custom class, which we assume has a method that returns
// the embedded timestamp (milliseconds since midnight, January 1, 1970 UTC).
long timestamp = -1;
final Foo myPojo = (Foo) record.value();
if (myPojo != null) {
timestamp = myPojo.getTimestampInMillis();
}
if (timestamp < 0) {
// Invalid timestamp! Attempt to estimate a new timestamp,
// otherwise fall back to wall-clock time (processing-time).
if (previousTimestamp >= 0) {
return previousTimestamp;
} else {
return System.currentTimeMillis();
}
}
}
}
그런 다음 Streams 구성에서 커스텀 타임스탬프 추출기를 다음과 같이 정의해요:
import java.util.Properties;
import org.apache.kafka.streams.StreamsConfig;
Properties streamsConfiguration = new Properties();
streamsConfiguration.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, MyEventTimeExtractor.class);
default.key.serde
레코드 키의 기본 Serializer/Deserializer 클래스, 사용자가 설정하지 않으면 null이에요. Kafka Streams에서 직렬화·역직렬화는 데이터를 구체화(materialize)해야 할 때마다 발생해요. 예를 들어:
- 데이터가 카프카 토픽에서 읽히거나 쓰일 때마다 (예:
StreamsBuilder#stream()과KStream#to()메서드를 통해). - 데이터가 상태 저장소에서 읽히거나 쓰일 때마다.
이것은 데이터 타입과 직렬화에서 더 자세히 다뤄요.
default.value.serde
레코드 값의 기본 Serializer/Deserializer 클래스, 사용자가 설정하지 않으면 null이에요. Kafka Streams에서 직렬화·역직렬화는 데이터를 구체화해야 할 때마다 발생해요. 예를 들어:
- 데이터가 카프카 토픽에서 읽히거나 쓰일 때마다.
- 데이터가 상태 저장소에서 읽히거나 쓰일 때마다.
이것은 데이터 타입과 직렬화에서 더 자세히 다뤄요.
dsl.store.format
상태 저장소를 구체화하는 모든 DSL 연산자가 사용하는 상태 저장소 형식을 선택해요. 허용 값은 DEFAULT와 HEADERS(대소문자 구분 없음)이며, 기본값은 DEFAULT예요.
DEFAULT: 연산자별 기존 timestamped 또는 plain 저장소 변형을 사용해요. 기존 애플리케이션은 영향받지 않아요.HEADERS: 값과 타임스탬프와 함께 레코드 헤더를 영속할 수 있는 헤더 인지 저장소(KIP-1271 도입)를 사용해요.
이 구성은 전역적이에요. 연산자별 커스터마이즈는 Materialized.withStoreType(...)로 커스텀 DslStoreSuppliers를 제공하거나 명시적인 헤더 인지 저장소 공급자를 공급함으로써 가능해요. dsl.store.format은 저장소 구현(RocksDB vs 인메모리)을 선택하는 dsl.store.suppliers.class와 직교하며, 둘은 독립적으로 설정할 수 있어요. 허용 문자열 값은 DEFAULT와 HEADERS(대소문자 구분 없음)예요. 이것들은 PLAIN, TIMESTAMPED, HEADERS 상수가 있는 DslStoreFormat Java enum과 다르며, DslStoreFormat.DEFAULT는 enum 상수로 존재하지 않아요.
마이그레이션 절차, 체인지로그 호환성, 복원 동작, 레코드별 오버헤드에 대해서는 KIP-1271을 참고해요.
현재 제한사항: dsl.store.format=HEADERS는 상태 저장소 형식을 변경해요. DSL 연산자가 출력 레코드용 헤더를 어떻게 만드는지는 정의하지 않아요. 일부 연산자는 구체화된 저장소에 빈 헤더를 쓰고, suppress()와 left/outer 스트림-스트림 조인에 사용되는 버퍼 저장소는 헤더 인지가 아니에요. 자세한 내용은 Stateful transformations와 Streams 업그레이드 가이드를 참고해요.
ensure.explicit.internal.resource.naming
내부 토픽(예: 체인지로그·리파티션 토픽)과 관련 상태 저장소를 포함한 토폴로지의 모든 내부 리소스에 명시적 명명을 강제할지 여부. 활성화되면 내부 리소스에 자동 생성 이름이 있으면 애플리케이션이 시작을 거부해요.
group.protocol
조정에 사용되는 Kafka Streams 클라이언트가 사용하는 그룹 프로토콜이에요. 클라이언트가 카프카 브로커와 같은 그룹의 다른 클라이언트와 어떻게 통신할지 결정해요. 기본값은 클래식 컨슈머 그룹 프로토콜인 "classic"이에요. 새 Kafka Streams 그룹 프로토콜을 활성화하려면 "streams"(브로커 측 활성화 필요)로 설정할 수 있어요.
rack.aware.assignment.non_overlap_cost
이 구성은 StickyTaskAssignor 또는 HighAvailabilityTaskAssignor가 계산한 원래 할당에서 태스크를 이동하는 비용을 설정해요. rack.aware.assignment.traffic_cost와 함께 교차 랙 트래픽 최소화와 기존 할당에서 태스크 이동 최소화 중 무엇을 선호할지 제어해요. 이 구성이 rack.aware.assignment.traffic_cost보다 큰 값으로 설정되면 최적화기는 태스크 assignor가 계산한 기존 할당을 유지하려고 해요. 최적화기는 기존 할당 유지 선호와 트래픽 비용 최소화 선호에 이 두 구성의 비율을 고려해요. 예를 들어 rack.aware.assignment.non_overlap_cost를 10, rack.aware.assignment.traffic_cost를 1로 설정하는 것이 각각 100과 50으로 설정하는 것보다 기존 할당을 유지할 가능성이 높아요.
기본값은 null이며, 이는 서로 다른 assignor의 기본 non_overlap_cost가 사용된다는 뜻이에요. StickyTaskAssignor에서는 기본값이 10이고 rack.aware.assignment.traffic_cost 기본값이 1이므로, StickyTaskAssignor에서 스티키니스 유지가 선호돼요. HighAvailabilityTaskAssignor에서는 기본값이 1이고 rack.aware.assignment.traffic_cost 기본값이 10이므로, HighAvailabilityTaskAssignor에서 교차 랙 트래픽 최소화가 선호돼요.
rack.aware.assignment.strategy
이 구성은 브로커에서 클라이언트로의 교차 트래픽을 줄일 수 있도록 Kafka Streams가 랙 인지 태스크 할당에 사용하는 전략을 설정해요. 이 구성은 브로커에 broker.rack이 설정되고 Kafka Streams 측에 client.rack이 설정된 경우에만 효과가 있어요. 이 구성의 두 가지 설정:
none. 기본값이며 랙 인지 태스크 할당이 비활성화된다는 뜻이에요.min_traffic. 랙 인지 태스크 할당자가 교차 랙 트래픽을 최소화하려는 할당을 계산한다는 뜻이에요.balance_subtopology. 랙 인지 태스크 할당자가 같은 서브토폴로지의 태스크를 다른 클라이언트에 균형을 맞추고 그 위에 교차 랙 트래픽을 최소화하려는 할당을 계산한다는 뜻이에요.
이 구성은 rack.aware.assignment.non_overlap_cost와 rack.aware.assignment.traffic_cost와 함께 사용해 교차 랙 트래픽 감소와 기존 할당 유지의 균형을 잡을 수 있어요.
rack.aware.assignment.tags
이 구성은 스탠바이 복제본을 Kafka Streams 클라이언트 전체에 분산하는 데 사용되는 태그 키 목록을 설정해요. 구성되면 Kafka Streams는 다른 태그 값을 가진 클라이언트에 스탠바이 태스크 분산을 최선으로 시도해요.
Kafka Streams 클라이언트의 태그는 client.tag. 접두사로 설정할 수 있어요. 예:
Client-1 | Client-2
_______________________________________________________________________
client.tag.zone: eu-central-1a | client.tag.zone: eu-central-1b
client.tag.cluster: k8s-cluster1 | client.tag.cluster: k8s-cluster1
rack.aware.assignment.tags: zone,cluster | rack.aware.assignment.tags: zone,cluster
Client-3 | Client-4
_______________________________________________________________________
client.tag.zone: eu-central-1a | client.tag.zone: eu-central-1b
client.tag.cluster: k8s-cluster2 | client.tag.cluster: k8s-cluster2
rack.aware.assignment.tags: zone,cluster | rack.aware.assignment.tags: zone,cluster
위 예시에는 두 개의 존(eu-central-1a, eu-central-1b)과 두 개의 클러스터(k8s-cluster1, k8s-cluster2)에 걸친 네 개의 Kafka Streams 클라이언트가 있어요. Client-1에 있는 활성 태스크에 대해 Kafka Streams는 Client-4에 스탠바이 태스크를 할당할 거예요. Client-4는 Client-1과 다른 zone과 다른 cluster를 가지기 때문이에요.
rack.aware.assignment.traffic_cost
이 구성은 교차 랙 트래픽의 비용을 설정해요. rack.aware.assignment.non_overlap_cost와 함께 교차 랙 트래픽 최소화와 기존 할당에서 태스크 이동 최소화 중 무엇을 선호할지 제어해요. 이 구성이 rack.aware.assignment.non_overlap_cost보다 큰 값으로 설정되면 최적화기는 교차 랙 트래픽을 최소화하는 할당을 계산하려고 해요. 최적화기는 기존 할당 유지 선호와 트래픽 비용 최소화에 이 두 구성의 비율을 고려해요. 예를 들어 rack.aware.assignment.traffic_cost를 10, rack.aware.assignment.non_overlap_cost를 1로 설정하는 것이 각각 100과 50으로 설정하는 것보다 교차 랙 트래픽을 최소화할 가능성이 높아요.
기본값은 null이며, 이는 서로 다른 assignor의 기본 트래픽 비용이 사용된다는 뜻이에요. StickyTaskAssignor에서는 기본값이 1이고 rack.aware.assignment.non_overlap_cost 기본값이 10이에요. HighAvailabilityTaskAssignor에서는 기본값이 10이고 rack.aware.assignment.non_overlap_cost 기본값이 1이에요.
log.summary.interval.ms
이 구성은 요약 정보의 출력 간격을 제어해요. 0 이상이면 설정된 시간 간격에 따라 요약 로그가 출력되고, 0 미만이면 요약 출력이 비활성화돼요.
enable.metrics.push
Kafka Streams 메트릭은 클라이언트 메트릭과 유사하게 브로커로 푸시할 수 있어요. 추가로 Kafka Streams는 임베드된 각 클라이언트에 대해 메트릭 푸시를 개별적으로 활성화/비활성화할 수 있어요. 그러나 Kafka Streams 메트릭 푸시는 main-consumer와 admin 클라이언트에서 enable.metric.push가 활성화되어 있어야 해요.
max.task.idle.ms
이 구성은 Streams가 순서 있는(in-order) 처리 의미론을 제공하기 위해 데이터를 가져오기 위해 얼마나 기다릴지 제어해요.
여러 입력 파티션(조인이나 머지에서처럼)이 있는 태스크를 처리할 때 Streams는 다음 레코드를 어느 파티션에서 처리할지 선택해야 해요. 모든 입력 파티션에 로컬로 버퍼된 데이터가 있으면, Streams는 다음 레코드의 타임스탬프가 가장 낮은 파티션을 선택해요. 이것은 입력 파티션을 타임스탬프 순서로 결합하는 바람직한 효과가 있으며, 이는 일반적으로 스트리밍 조인이나 머지에서 원하는 것입니다. 그러나 파티션 중 하나에 대해 로컬로 버퍼된 데이터가 없으면, Streams는 해당 파티션의 다음 레코드가 나머지 파티션 레코드보다 낮거나 높은 타임스탬프를 가질지 알지 못해요.
고려할 두 가지 경우가 있어요: 해당 파티션에 Streams가 아직 가져오지 않은 브로커 데이터가 있거나, Streams가 해당 파티션에 대해 브로커에 완전히 따라잡았고 프로듀서가 마지막 배치 이후 새 레코드를 생성하지 않은 경우예요.
기본값 0은 Streams가 파티션에 대해 로컬로 버퍼된 데이터가 없지만 브로커에서 사용 가능한 데이터가 있다고 감지할 때 태스크 처리를 지연시켜요. 구체적으로, 로컬 버퍼에 빈 파티션이 있지만 Streams가 해당 파티션에 대해 0이 아닌 lag를 가질 때요. 그러나 Streams가 브로커에 따라잡는 즉시, 파티션 중 하나에 데이터가 없어도 처리를 계속해요. 즉, 새 데이터가 생성될 때까지 기다리지 않아요. 이 기본값은 직관적으로 올바른 조인 의미론을 위해 일부 처리량을 희생하도록 설계되었어요.
0보다 큰 구성 값은 Streams가 따라잡았지만 빈 파티션이 있으면 추가로 기다릴 밀리초 수를 나타내요. 다시 말해, 느린 프로듀서의 경우 입력 파티션에 새 데이터가 생성될 때까지 기다려 데이터의 순서 있는 처리를 보장하는 시간이에요.
-1 구성 값은 Streams가 타임스탬프로 다음 레코드를 선택하기 전에 빈 파티션을 버퍼링하기 위해 결코 기다리지 않는다는 뜻이며, 순서가 뒤집힌 처리를 도입하는 대가로 최대 처리량을 달성해요.
max.warmup.replicas
한 번에 할당할 수 있는 최대 워밍업 복제본 수(구성된 num.standbys를 초과하는 추가 스탠바이)로, 태스크가 재할당된 다른 인스턴스에서 워밍업되는 동안 한 인스턴스에서 태스크를 사용 가능하게 유지하는 목적이에요. 고가용성에 사용할 수 있는 추가 브로커 트래픽과 클러스터 상태를 제한하는 데 사용돼요. 이것을 늘리면 Streams가 한 번에 더 많은 태스크를 워밍업해, 활성 태스크로 전환하기에 충분한 상태를 복원하는 재할당 워밍업 시간을 빨라지게 해요. 최소 1이어야 해요.
하나의 워밍업 복제본은 하나의 Stream Task에 대응한다는 점을 유의해요. 또한 각 워밍업 태스크는 리밸런스 동안에만(보통 probing.rebalance.interval.ms 구성이 지정하는 빈도로 발생하는 소위 프로빙 리밸런스 동안) 활성 태스크로 승격될 수 있다는 점을 유의해요. 즉, 활성 태스크가 한 Kafka Streams 인스턴스에서 다른 인스턴스로 마이그레이션될 수 있는 최대 비율은 (max.warmup.replicas / probing.rebalance.interval.ms)로 결정될 수 있어요.
num.standby.replicas
스탠바이 복제본 수예요. 스탠바이 복제본은 로컬 상태 저장소의 섀도 복사본이에요. Kafka Streams는 저장소당 지정된 수의 복제본을 만들고 충분한 인스턴스가 실행 중인 한 최신 상태로 유지하려고 해요. 스탠바이 복제본은 태스크 장애 조치(failover)의 지연을 최소화하는 데 사용돼요. 실패한 인스턴스에서 이전에 실행된 태스크는 스탠바이 복제본이 있는 인스턴스에서 재시작되는 것이 선호되어, 체인지로그에서 로컬 상태 저장소 복원 과정을 최소화할 수 있어요. Kafka Streams가 스탠바이 복제본을 사용해 장애 조치 시 태스크 재개 비용을 최소화하는 방법에 대한 자세한 내용은 State 섹션에서 찾을 수 있어요.
권장 사항: 즉각적인 장애 조치, 즉 고가용성을 얻으려면 스탠바이 수를 1로 늘리세요. 스탠바이 수를 늘리면 클라이언트 측 저장 공간이 더 필요해요. 예를 들어 스탠바이 1개면 2배 공간이 필요해요.
참고: n개의 스탠바이 태스크를 활성화하면 n+1개의 KafkaStreams 인스턴스를 준비해야 해요.
num.stream.threads
Kafka Streams 애플리케이션 인스턴스의 스트림 스레드 수를 지정해요. 스트림 처리 코드는 이 스레드에서 실행돼요. Kafka Streams 스레딩 모델에 대한 자세한 내용은 Threading Model을 참고해요.
probing.rebalance.interval.ms
충분히 복원되어 따라잡은 것으로 간주되는 워밍업 복제본을 검색하기 위한 리밸런스 촉발 전 최대 대기 시간이에요. Streams는 따라잡히고 acceptable.recovery.lag 내에 있는 인스턴스(존재하는 경우)에만 상태 저장 활성 태스크를 할당해요. 프로빙 리밸런스는 워밍업 복제본의 최신 총 lag를 쿼리하고 준비되면 활성 태스크로 전환하는 데 사용돼요. 워밍업 태스크가 있는 한, 그리고 할당이 균형 잡힐 때까지 계속 촉발돼요. 최소 1분이어야 해요.
processing.exception.handler
처리 예외 핸들러를 사용하면 레코드 처리 중에 발생하는 예외를 관리할 수 있어요. 구현된 예외 핸들러는 레코드와 발생한 예외에 따라 FAIL 또는 CONTINUE를 반환해야 해요. FAIL을 반환하면 Streams가 종료되어야 한다는 신호이고, CONTINUE는 Streams가 문제를 무시하고 계속 처리해야 한다는 신호예요.
참고: 기본적으로 이 핸들러는 일반 스트림 처리 태스크에만 적용돼요. 글로벌 저장소/KTable 처리에 대한 예외 처리를 활성화하려면(권장됨) 아래
processing.exception.handler.global.enabled를 보세요. 글로벌 예외 처리가 비활성화되면(기본), 글로벌 저장소/KTable 처리 중 발생하는 예외는 구성된 uncaught exception handler로 전파돼요.
다음 라이브러리 내장 예외 핸들러가 사용 가능해요:
- LogAndContinueProcessingExceptionHandler: 이 핸들러는 처리 예외를 기록한 다음 처리 파이프라인이 더 많은 레코드를 계속 처리하도록 신호를 보내요. 이 로그-및-건너뛰기 전략은 처리에 실패하는 레코드가 있으면 실패하는 대신 Kafka Streams가 진행할 수 있게 해줘요.
- LogAndFailProcessingExceptionHandler: 이 핸들러는 처리 예외를 기록한 다음 처리 파이프라인이 더 많은 레코드 처리를 중지하도록 신호를 보내요.
라이브러리 제공 핸들러 외에 요구에 맞는 자신의 커스터마이즈된 예외 핸들러를 제공할 수도 있어요. 예를 들어 손상된 레코드를 격리 토픽("데드 레터 큐")으로 전달해 추가 처리를 하도록 선택할 수 있어요. 이렇게 하려면 Producer API를 사용해 손상된 레코드를 격리 토픽에 직접 써요. 더 구체적으로, Streams 클라이언트 밖에 별도의 KafkaProducer 객체를 만들고 이 객체와 데드 레터 큐 토픽 이름을 Properties 맵에 전달하면 configure 함수 호출에서 검색할 수 있어요. 이 접근 방식의 단점은 "수동" 쓰기가 Kafka Streams 런타임 라이브러리에서 보이지 않는 부작용이라, Streams API의 엔드투엔드 처리 보장 혜택을 받지 못한다는 점이에요:
public class SendToDeadLetterQueueExceptionHandler implements ProcessingExceptionHandler {
KafkaProducer<byte[], byte[]> dlqProducer;
String dlqTopic;
@Override
public ProcessingHandlerResponse handle(final ErrorHandlerContext context,
final Record record,
final Exception exception) {
log.warn("Exception caught during message processing, sending to the dead queue topic; " +
"processor node: {}, taskId: {}, source topic: {}, source partition: {}, source offset: {}",
context.processorNodeId(), context.taskId(), context.topic(), context.partition(), context.offset(),
exception);
dlqProducer.send(new ProducerRecord<>(dlqTopic, null, record.timestamp(), (byte[]) record.key(), (byte[]) record.value(), record.headers()));
return ProcessingHandlerResponse.CONTINUE;
}
@Override
public void configure(final Map<String, ?> configs) {
dlqProducer = .. // get a producer from the configs map
dlqTopic = .. // get the topic name from the configs map
}
}
참고: 위 예시는 DLQ 토픽으로의 수동 프로덕션을 보여줘요. 다음 예시는 내장 DLQ 지원을 사용하는 권장 접근 방식을 보여줘요.
커스텀 처리 예외 핸들러는 사용자 로직이 예외를 던질 때 처리를 계속할지 실패할지 결정할 수 있어요. DLQ 동작이 필요하면 핸들러 응답에서 DLQ 레코드를 반환해요.
커스텀 예외 핸들러 구현
다음 예시는 실패한 레코드를 구성된 DLQ 토픽으로 전달해요:
public class DlqProcessingExceptionHandler implements ProcessingExceptionHandler {
private String deadLetterQueueTopic;
@Override
public Response handleError(final ErrorHandlerContext context,
final Record<?, ?> record,
final Exception exception) {
// Example: forward the raw record to a DLQ topic
ProducerRecord<byte[], byte[]> dlqRecord =
new ProducerRecord<>(deadLetterQueueTopic,
null,
context.timestamp(),
context.sourceRawKey(),
context.sourceRawValue());
// Applications may choose how to construct DLQ records. For example,
// they may forward the raw key/value bytes, transform the payload,
// or add headers with error metadata.
return Response.resume(List.of(dlqRecord));
}
@Override
public void configure(final Map<String, ?> configs) {
// Retrieve the DLQ topic name from the configs map, or any other source
deadLetterQueueTopic = (String) configs.get("my.dlq.topic.config.key");
}
}
커스텀 예외 핸들러를 활성화하고 DLQ 토픽을 구성하려면:
Properties props = new Properties();
props.put(
StreamsConfig.PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG,
DlqProcessingExceptionHandler.class
);
// Optional: if your custom handler reads the DLQ topic from StreamsConfig,
// set it here. Otherwise, configure the topic name via your own properties.
// props.put(
// StreamsConfig.ERRORS_DEAD_LETTER_QUEUE_TOPIC_NAME_CONFIG,
// "dlq-topic"
// );
processing.exception.handler.global.enabled (deprecated)
글로벌 저장소/KTable 처리 중 발생하는 예외에 대해 구성된 ProcessingExceptionHandler가 호출되는지 제어해요. true로 설정하면(권장) processing.exception.handler가 지정한 핸들러가 글로벌 저장소/KTable 처리 중 발생하는 예외에 대해 호출돼요. false(기본값)로 설정하면 글로벌 저장소/KTable의 예외가 처리 예외 핸들러를 호출하지 않고 대신 구성된 uncaught exception handler로 전파돼요.
기본값: false. 폐기됨: 이 구성은 5.0 릴리스에서 제거를 위해 폐기되었어요. 구성이 제거되면 처리 예외 핸들러가 글로벌 상태/KTable 처리 중 적용되고 더 이상 비활성화할 수 없어요. 따라서 나중에 하위 호환 문제를 피하려면 지금 이 구성을 활성화하는 것이 좋아요.
중요한 참고 사항:
- DLQ(Dead Letter Queue) 기능은 글로벌 저장소/KTable에 대해 지원되지 않아요. 글로벌 저장소/KTable 예외의 경우 레코드 메타데이터가 기록되고 레코드는 DLQ로 전송되지 않아요.
- 이 기능이 활성화되면 내장 핸들러(
LogAndContinueProcessingExceptionHandler또는LogAndFailProcessingExceptionHandler)를 사용하거나ProcessingExceptionHandler의 커스텀 구현을 제공할 수 있어요. - 자세한 내용은 KIP-1270을 참고해요.
예시 구성:
Properties streamsSettings = new Properties();
// Configure the processing exception handler
streamsSettings.put(StreamsConfig.PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG,
LogAndContinueProcessingExceptionHandler.class);
// Enable exception handling for Global KTables
streamsSettings.put(StreamsConfig.PROCESSING_EXCEPTION_HANDLER_GLOBAL_ENABLED_CONFIG, true);
processing.guarantee
사용해야 하는 처리 보장이에요. 가능한 값은 "at_least_once"(기본)와 "exactly_once_v2"(EOS 버전 2)예요. 폐기된 구성 옵션은 "exactly_once"(EOS alpha)와 "exactly_once_beta"(EOS 버전 2)예요. "exactly_once_v2"(또는 폐기된 "exactly_once_beta")를 사용하려면 브로커 버전 2.5 이상이 필요하고, 폐기된 "exactly_once"를 사용하려면 브로커 버전 0.11.0 이상이 필요해요. 정확히 한 번 처리가 활성화되면 파라미터 commit.interval.ms의 기본값이 100ms로 변경된다는 점을 유의해요. 또한 컨슈머는 isolation.level="read_committed"로, 프로듀서는 enable.idempotence=true가 기본으로 구성돼요. 기본적으로 정확히 한 번 처리는 프로덕션에 권장되는 적어도 세 개의 브로커 클러스터를 요구한다는 점을 유의해요. 개발을 위해 브로커 설정 transaction.state.log.replication.factor와 transaction.state.log.min.isr을 원하는 브로커 수로 조정해 이 구성을 변경할 수 있어요. 자세한 내용은 Processing Guarantees를 참고해요.
권장 사항: EOS를 어떤 복제 팩터로도 기술적으로 사용할 수 있지만, 3보다 낮은 복제 팩터를 사용하면 EOS를 효과적으로 무효화해요. 따라서 복제 팩터 3(min.in.sync.replicas=2와 함께)을 사용하는 것을 적극 권장해요. 이 권장 사항은 모든 토픽(즉 __transaction_state, __consumer_offsets, Kafka Streams 내부 토픽, 사용자 토픽)에 적용돼요.
processor.wrapper.class
ProcessorWrapper 인터페이스를 구현하는 클래스 또는 클래스 이름이에요. 이 기능을 사용하면 DSL 연산자를 위해 Streams가 만든 것과 커스텀 프로세서 구현을 모두 포함해 컴파일된 토폴로지의 어떤 프로세서도 래핑할 수 있어요. 이것은 그렇지 않으면 숨겨진 DSL 연산자 프로세서 컨텍스트에 접근할 수 있게 해주고, 단일 구성으로 전체 애플리케이션 토폴로지에 추가 디버깅 정보를 주입할 수 있게 해주므로 로깅이나 트레이싱 구현에 유용해요.
중요: 이것은 토폴로지를 만들 때 전달해야 하며, 적절한 토폴로지 구축 생성자에 전달하지 않으면 적용되지 않아요. DSL 애플리케이션은 StreamsBuilder#new(TopologyConfig) 생성자, PAPI 애플리케이션은 Topology#new(TopologyConfig) 생성자를 사용해야 해요.
replication.factor
이것은 로컬 상태가 사용되거나 스트림이 집계를 위해 리파티션될 때 Kafka Streams가 만드는 내부 토픽의 복제 팩터를 지정해요. 복제는 장애 허용에 중요해요. 복제가 없으면 단일 브로커 장애조차도 스트림 처리 애플리케이션의 진행을 막을 수 있어요. 소스 토픽과 유사한 복제 팩터를 사용하는 것이 좋아요.
권장 사항: 내부 Kafka Streams 토픽이 최대 2개의 브로커 장애를 견딜 수 있도록 복제 팩터를 3으로 높이세요. 더 많은 저장 공간(복제 팩터 3에서 3배)도 필요하다는 점을 유의하세요.
rocksdb.config.setter
RocksDB 구성이에요. Kafka Streams는 영속 저장소의 기본 저장 엔진으로 RocksDB를 사용해요. RocksDB의 기본 구성을 변경하려면 RocksDBConfigSetter를 구현하고 rocksdb.config.setter로 커스텀 클래스를 제공할 수 있어요.
다음은 RocksDB가 소비하는 메모리 크기를 조정하는 예시예요.
public static class CustomRocksDBConfig implements RocksDBConfigSetter {
// This object should be a member variable so it can be closed in RocksDBConfigSetter#close.
private org.rocksdb.Cache cache = new org.rocksdb.LRUCache(16 * 1024L * 1024L);
@Override
public void setConfig(final String storeName, final Options options, final Map<String, Object> configs) {
// See #1 below.
BlockBasedTableConfig tableConfig = (BlockBasedTableConfig) options.tableFormatConfig();
tableConfig.setBlockCache(cache);
// See #2 below.
tableConfig.setBlockSize(16 * 1024L);
// See #3 below.
tableConfig.setCacheIndexAndFilterBlocks(true);
options.setTableFormatConfig(tableConfig);
// See #4 below.
options.setMaxWriteBufferNumber(2);
}
@Override
public void close(final String storeName, final Options options) {
// See #5 below.
cache.close();
}
}
Properties streamsSettings = new Properties();
streamsConfig.put(StreamsConfig.ROCKSDB_CONFIG_SETTER_CLASS_CONFIG, CustomRocksDBConfig.class);
예시에 대한 참고:
BlockBasedTableConfig tableConfig = (BlockBasedTableConfig) options.tableFormatConfig();— 새 것을 만드는 대신 기존 테이블 구성에 대한 참조를 얻어, 중요한 최적화인BloomFilter같은 기본값을 실수로 덮어쓰지 않게 해요.tableConfig.setBlockSize(16 * 1024L);— RocksDB GitHub의 지침에 따라 기본 블록 크기를 수정해요.tableConfig.setCacheIndexAndFilterBlocks(true);— 인덱스와 필터 블록이 무한정 커지지 않게 해요. 자세한 내용은 RocksDB GitHub를 참고해요.options.setMaxWriteBufferNumber(2);— RocksDB GitHub의 고급 옵션을 참고해요.cache.close();— 메모리 누수를 피하려면org.rocksdb.RocksObject를 확장해서 구성한 객체는 닫아야 해요. 자세한 내용은 RocksJava 문서를 참고해요.
state.dir
상태 디렉터리예요. Kafka Streams는 상태 디렉터리 아래에 로컬 상태를 영속해요. 각 애플리케이션은 호스팅 머신에 상태 디렉터리 아래에 위치하는 서브디렉터리를 가지며, 서브디렉터리의 이름은 애플리케이션 ID예요. 애플리케이션과 연결된 상태 저장소는 이 서브디렉터리 아래에 만들어져요. 단일 머신에서 같은 애플리케이션의 여러 인스턴스를 실행할 때 이 경로는 각 인스턴스마다 고유해야 해요.
task.assignor.class
org.apache.kafka.streams.processor.assignment.TaskAssignor 인터페이스를 구현하는 태스크 assignor 클래스 또는 클래스 이름이에요. 기본값은 고가용성 태스크 assignor예요. Apache Kafka에서 제공되는 한 가지 대안 구현은 KIP-441 이전의 기본 태스크 assignor였고 상태 저장 태스크 가용성을 희생하면서 태스크 이동을 최소화하는 org.apache.kafka.streams.processor.assignment.assignors.StickyTaskAssignor예요. 태스크 할당 알고리즘의 대안 구현은 커스텀 TaskAssignor를 구현하고 이 구성에 커스텀 태스크 assignor 클래스 이름을 설정해 애플리케이션에 플러그인할 수 있어요.
topology.optimization
Kafka Streams가 토폴로지를 최적화해야 하는지와 어떤 최적화를 적용할지 알려주는 구성이에요. 허용 값: StreamsConfig.NO_OPTIMIZATION (none), StreamsConfig.OPTIMIZE (all), 또는 특정 최적화의 쉼표 구분 목록: StreamsConfig.REUSE_KTABLE_SOURCE_TOPICS (reuse.ktable.source.topics), StreamsConfig.MERGE_REPARTITION_TOPICS (merge.repartition.topics), StreamsConfig.SINGLE_STORE_SELF_JOIN (single.store.self.join).
Streams 라이브러리의 업그레이드 중에 토폴로지 구조가 예상치 않게 변하지 않도록 프로덕션 코드에서 구성에 특정 최적화를 나열하는 것을 권장해요.
이 최적화는 리파티션 토픽 이동/감소와 소스 KTable의 체인지로그로 소스 토픽 재사용을 포함해요. 이 최적화는 애플리케이션의 의미론을 변경하지 않고 카프카의 네트워크 트래픽과 저장소를 절약해요. 활성화하는 것을 권장해요.
중요: 최적화를 활성화하려면 두 단계가 필요해요. 둘 다 필요해요 — 구성만 설정하는 것으로는 충분하지 않아요:
Properties객체에서topology.optimization을StreamsConfig.OPTIMIZE(또는 특정 최적화의 쉼표 구분 목록)로 설정.- 토폴로지를 만들 때 오버로드된
StreamsBuilder.build(Properties)메서드에 같은Properties객체를 전달.
예를 들어:
Properties properties = new Properties();
properties.put(StreamsConfig.TOPOLOGY_OPTIMIZATION_CONFIG, StreamsConfig.OPTIMIZE);
// Step 2: pass properties to build() — this is required for optimizations to take effect
Topology topology = streamsBuilder.build(properties);
KafkaStreams myStream = new KafkaStreams(topology, properties);
Properties 객체를 전달하지 않고 streamsBuilder.build()를 호출하면 구성이 설정되어 있어도 최적화가 적용되지 않아요.
upgrade.from
업그레이드하는 버전이에요. 일부 버전으로 롤링 업그레이드를 수행할 때 이 구성을 설정하는 것이 중요하며, 업그레이드 가이드에 설명되어 있어요. 인스턴스를 바운스하고 더 새 버전으로 업그레이드하기 전에 이 구성을 적절한 버전으로 설정해야 해요. 모두가 새 버전을 사용하게 되면 이 구성을 제거하고 두 번째 롤링 바운스를 해야 해요. 3.4 미만의 어떤 버전에서도 3.4 이상으로 업그레이드할 때만 이 구성을 설정하고 두 번 바운스 업그레이드 경로를 따르면 돼요.
Kafka 컨슈머, 프로듀서 및 admin 클라이언트 구성 파라미터
내부적으로 사용되는 카프카 컨슈머·프로듀서·admin 클라이언트에 대한 파라미터를 지정할 수 있어요. 컨슈머·프로듀서·admin 클라이언트 설정은 StreamsConfig 인스턴스에 파라미터를 지정함으로써 정의돼요.
이 예시에서 카프카 컨슈머 세션 타임아웃이 Streams 설정에서 60000밀리초로 구성돼요:
Properties streamsSettings = new Properties();
// Example of a "normal" setting for Kafka Streams
streamsSettings.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker-01:9092");
// Customize the Kafka consumer settings of your Streams application
streamsSettings.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 60000);
명명
일부 컨슈머·프로듀서·admin 클라이언트 구성 파라미터는 같은 파라미터 이름을 사용하며, Kafka Streams 라이브러리 자체도 임베드된 클라이언트와 같은 이름을 공유하는 일부 파라미터를 사용해요. 예를 들어 send.buffer.bytes와 receive.buffer.bytes는 TCP 버퍼를 구성하는 데 사용되고, request.timeout.ms와 retry.backoff.ms는 클라이언트 요청의 재시도를 제어해요. 파라미터 이름에 consumer., producer., admin. 접두사를 붙여 중복 이름을 피할 수 있어요 (예: consumer.send.buffer.bytes, producer.send.buffer.bytes).
Properties streamsSettings = new Properties();
// same value for consumer, producer, and admin client
streamsSettings.put("PARAMETER_NAME", "value");
// different values for consumer and producer
streamsSettings.put("consumer.PARAMETER_NAME", "consumer-value");
streamsSettings.put("producer.PARAMETER_NAME", "producer-value");
streamsSettings.put("admin.PARAMETER_NAME", "admin-value");
// alternatively, you can use
streamsSettings.put(StreamsConfig.consumerPrefix("PARAMETER_NAME"), "consumer-value");
streamsSettings.put(StreamsConfig.producerPrefix("PARAMETER_NAME"), "producer-value");
streamsSettings.put(StreamsConfig.adminClientPrefix("PARAMETER_NAME"), "admin-value");
다른 접두사를 추가해 컨슈머 구성을 더 분리할 수 있어요:
main.consumer.— 스트림 소스의 기본 컨슈머인 main consumer용.restore.consumer.— 상태 저장소 복구를 담당하는 restore consumer용.global.consumer.— 글로벌 KTable 구축에 사용되는 global consumer용.
예를 들어 다른 컨슈머 설정을 건드리지 않고 restore consumer 구성만 설정하고 싶다면 간단히 restore.consumer.를 사용할 수 있어요.
Properties streamsSettings = new Properties();
// same config value for all consumer types
streamsSettings.put("consumer.PARAMETER_NAME", "general-consumer-value");
// set a different restore consumer config. This would make restore consumer take restore-consumer-value,
// while main consumer and global consumer stay with general-consumer-value
streamsSettings.put("restore.consumer.PARAMETER_NAME", "restore-consumer-value");
// alternatively, you can use
streamsSettings.put(StreamsConfig.restoreConsumerPrefix("PARAMETER_NAME"), "restore-consumer-value");
main.consumer.에도 같은 것이 적용돼요. 하나의 컨슈머 유형 구성만 지정하고 싶으면요. 추가로 내부 리파티션/체인지로그 토픽을 구성하려면 topic. 접두사와 함께 표준 토픽 구성 중 어떤 것이든 사용할 수 있어요.
Properties streamsSettings = new Properties();
// Override default for both changelog and repartition topics
streamsSettings.put("topic.PARAMETER_NAME", "topic-value");
// alternatively, you can use
streamsSettings.put(StreamsConfig.topicPrefix("PARAMETER_NAME"), "topic-value");
기본값
Kafka Streams는 기본 클라이언트 구성 중 일부에 대해 다른 기본값을 사용하며, 아래에 요약되어 있어요. 이 구성에 대한 자세한 설명은 Producer Configs와 Consumer Configs를 참고해요.
| 파라미터 이름 | 해당 클라이언트 | Streams 기본값 |
|---|---|---|
auto.offset.reset |
Consumer | earliest |
linger.ms |
Producer | 100 |
max.poll.records |
Consumer | 1000 |
client.id |
<application.id>-<random-UUID> |
EOS가 활성화되면 다른 파라미터의 기본값이 다음과 같아요.
| 파라미터 이름 | 해당 클라이언트 | Streams 기본값 |
|---|---|---|
transaction.timeout.ms |
Producer | 10000 |
delivery.timeout.ms |
Producer | Integer.MAX_VALUE |
Kafka Streams가 제어하는 파라미터
일부 파라미터는 사용자가 구성할 수 없어요. 기본값과 다른 값을 제공하면 그 값은 무시돼요. 다음은 이러한 파라미터 중 일부의 목록이에요.
| 파라미터 이름 | 해당 클라이언트 | Streams 기본값 |
|---|---|---|
allow.auto.create.topics |
Consumer | false |
group.id |
Consumer | application.id |
enable.auto.commit |
Consumer | false |
partition.assignment.strategy |
Consumer | StreamsPartitionAssignor |
EOS가 활성화되면 다른 파라미터가 다음 값으로 설정돼요.
| 파라미터 이름 | 해당 클라이언트 | Streams 기본값 |
|---|---|---|
isolation.level |
Consumer | READ_COMMITTED |
enable.idempotence |
Producer | true |
client.id
Kafka Streams는 client.id 파라미터를 사용해 내부 클라이언트의 파생된 client ID를 계산해요. client.id를 설정하지 않으면 Kafka Streams는 그것을 <application.id>-<random-UUID>로 설정해요. 이 값은 다음 내부 클라이언트의 client ID를 파생하는 데 사용될 거예요.
| 클라이언트 | client.id |
|---|---|
| Consumer | <client.id>-StreamThread-<threadIdx>-consumer |
| Restore consumer | <client.id>-StreamThread-<threadIdx>-restore-consumer |
| Global consumer | <client.id>-global-consumer |
| Producer | Non-EOS 및 EOS v2: <client.id>-StreamThread-<threadIdx>-producer |
EOS v1: <client.id>-StreamThread-<threadIdx>-<taskId>-producer |
|
| Admin | <client.id>-admin |
enable.auto.commit
컨슈머 자동 커밋이에요. 적어도 한 번(at-least-once) 처리 의미론을 보장하고 자동 커밋을 끄기 위해 Kafka Streams는 이 컨슈머 구성 값을 false로 오버라이드해요. 컨슈머는 Kafka Streams 라이브러리나 사용자가 현재 처리 상태를 커밋하기로 결정할 때만 commitSync 호출을 통해 명시적으로 커밋할 거예요.
더 알아보기
- Streams 애플리케이션 작성하기 — 구성 제공 방법을 봐요.
- Streams DSL — DSL 연산자에서 구성이 어떻게 쓰이는지 봐요.
- 개발자 가이드 전체 — 구성 관련 다른 문서를 봐요.