Streams DSL 애플리케이션에서 연산자 이름 짓기
Streams DSL 애플리케이션에서 연산자 이름 짓기
DSL로 토폴로지를 만들 때는 프로세서 이름이 자동으로 생성돼요. 그런데 이 이름들이 "몇 번째 프로세서냐"에 따라 달라져서, 토폴로지에서 연산 순서가 바뀌면 상태 저장소·체인지로그·리파티션 토픽의 이름까지 흔들릴 수 있어요. 이 페이지에서는 왜 상태 저장소 같은 것들에 직접 이름을 붙여야 하는지, 그리고 Named, Materialized 같은 클래스로 어떻게 붙이는지 설명해드릴게요.
출처: 문서
본문
Kafka Streams DSL 애플리케이션에서 연산자 이름 짓기
Kafka Streams DSL을 사용할 때 프로세서에 이름을 줄 수 있게 됐어요. PAPI(Processor API)에서는 Processors와 State Stores가 있고 각각 명시적으로 이름을 지정해야 해요.
DSL 레이어에는 연산자(operator)가 있어요. 단일 DSL 연산자는 여러 Processors와 State Stores, 그리고 필요하다면 repartition topics로 컴파일될 수 있어요. 하지만 Kafka Streams DSL에서는 이 모든 이름이 자동으로 생성돼요. 생성된 프로세서 이름, 상태 저장소 이름(따라서 체인지로그 토픽 이름), 리파티션 토픽 이름 사이에는 관계가 있어요. 상태 저장소와 체인지로그/리파티션 토픽의 이름은 "stateful"(상태 저장)인 반면 프로세서 이름은 "stateless"(무상태)라는 점을 유의해요.
이 stateful vs stateless 이름 구분은 토폴로지를 업데이트할 때 중요한 의미를 가져요. 내부 명명 방식이 DSL로 토폴로지를 만드는 것을 훨씬 더 직관적으로 만들지만, 몇 가지 트레이드오프가 있어요. 첫 번째 트레이드오프는 가독성 문제라고 볼 수 있어요. 더 심각한 트레이드오프는 DSL 연산자와 생성된 Processors, State Stores 체인지로그 토픽, 리파티션 토픽 사이의 관계 때문에 생기는 이름 이동(이름이 흔들리는 것)이에요.
가독성 문제
가독성 트레이드오프가 있다는 말은 토폴로지의 설명을 볼 때를 말해요. Topology#describe() 메서드로 토폴로지의 문자열 설명을 렌더링하면 프로세서가 무엇인지는 볼 수 있지만, 그 비즈니스 목적에 대한 맥락은 없어요. 예를 들어 다음의 간단한 토폴로지를 고려해보세요:
KStream<String,String> stream = builder.stream("input");
stream.filter((k,v) -> !v.equals("invalid_txn"))
.mapValues((v) -> v.substring(0,5))
.to("output");
Topology#describe()를 실행하면 다음 문자열이 나와요:
Topologies:
Sub-topology: 0
Source: KSTREAM-SOURCE-0000000000 (topics: [input])
--> KSTREAM-FILTER-0000000001
Processor: KSTREAM-FILTER-0000000001 (stores: [])
--> KSTREAM-MAPVALUES-0000000002
<-- KSTREAM-SOURCE-0000000000
Processor: KSTREAM-MAPVALUES-0000000002 (stores: [])
--> KSTREAM-SINK-0000000003
<-- KSTREAM-FILTER-0000000001
Sink: KSTREAM-SINK-0000000003 (topic: output)
<-- KSTREAM-MAPVALUES-0000000002
이 보고서에서 다양한 연산자가 무엇인지 볼 수 있지만, 여기서 더 넓은 맥락은 무엇일까요? 예를 들어 KSTREAM-FILTER-0000000001을 보면 filter 연산이라는 것을 알 수 있어요. 즉 주어진 프레디킷과 일치하지 않는 레코드는 버려진다는 뜻이에요. 그런데 그 프레디킷의 의미는 무엇일까요? 또한 소스와 싱크 노드의 토픽 이름은 볼 수 있지만, 토픽 이름이 의미 있게 지어지지 않았다면 어떨까요? 그러면 이 토픽들 뒤의 비즈니스 목적을 추측해야만 해요.
또한 여기서 번호 매기기에 주목해보세요: 소스 노드는 0000000000으로 끝나며 토폴로지의 첫 번째 프로세서임을 나타내요. filter는 0000000001로 끝나며 토폴로지의 두 번째 프로세서임을 나타내요. Kafka Streams에서 이제 KStream과 KTable 모두에 Named라는 새 파라미터를 받는 오버로드된 메서드가 있어요. Named 클래스를 사용하면 DSL 사용자가 토폴로지의 프로세서에 의미 있는 이름을 제공할 수 있어요.
이제 모든 프로세서에 이름을 붙인 토폴로지를 살펴볼게요:
KStream<String,String> stream =
builder.stream("input", Consumed.as("Customer_transactions_input_topic"));
stream.filter((k,v) -> !v.equals("invalid_txn"), Named.as("filter_out_invalid_txns"))
.mapValues((v) -> v.substring(0,5), Named.as("Map_values_to_first_6_characters"))
.to("output", Produced.as("Mapped_transactions_output_topic"));
Topologies:
Sub-topology: 0
Source: Customer_transactions_input_topic (topics: [input])
--> filter_out_invalid_txns
Processor: filter_out_invalid_txns (stores: [])
--> Map_values_to_first_6_characters
<-- Customer_transactions_input_topic
Processor: Map_values_to_first_6_characters (stores: [])
--> Mapped_transactions_output_topic
<-- filter_out_invalid_txns
Sink: Mapped_transactions_output_topic (topic: output)
<-- Map_values_to_first_6_characters
이제 토폴로지 설명을 보고 각 프로세서가 토폴로지에서 어떤 역할을 하는지 쉽게 이해할 수 있어요. 하지만 Kafka Streams 애플리케이션의 재시작 사이에 유지되는 stateful 연산자, 상태 저장소, 체인지로그 토픽, 리파티션 토픽이 있을 때 프로세서 노드에 이름을 붙여야 하는 또 다른 이유가 있어요.
이름 변경
생성된 이름은 토폴로지에서 만들어지는 위치에 번호가 매겨져요. 이름 생성 전략은 KSTREAM|KTABLE->연산자 이름<->숫자 접미사< 형태예요. 숫자는 토폴로지에서 연산자의 순서를 나타내는 전역 증가 숫자예요. 생성된 숫자에는 항상 10자리 문자열이 되도록 다양한 수의 "0"이 접두사로 붙어요. 이는 연산을 추가/제거하거나 순서를 바꾸면 프로세서의 위치가 이동하고, 프로세서의 이름도 이동한다는 뜻이에요. 대부분의 프로세서는 메모리에만 존재하므로, 이 이름 이동은 많은 토폴로지에 문제가 없어요. 하지만 이름 이동은 stateful 연산자나 리파티션 토픽이 있는 토폴로지에는 영향을 미쳐요. 다음은 상태가 있는 다른 토폴로지예요:
KStream<String,String> stream = builder.stream("input");
stream.groupByKey()
.count()
.toStream()
.to("output");
이 토폴로지 설명은 다음을 생성해요:
Topologies:
Sub-topology: 0
Source: KSTREAM-SOURCE-0000000000 (topics: [input])
--> KSTREAM-AGGREGATE-0000000002
Processor: KSTREAM-AGGREGATE-0000000002 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000001])
--> KTABLE-TOSTREAM-0000000003
<-- KSTREAM-SOURCE-0000000000
Processor: KTABLE-TOSTREAM-0000000003 (stores: [])
--> KSTREAM-SINK-0000000004
<-- KSTREAM-AGGREGATE-0000000002
Sink: KSTREAM-SINK-0000000004 (topic: output)
<-- KTABLE-TOSTREAM-0000000003
위 토폴로지 설명에서 상태 저장소 이름이 KSTREAM-AGGREGATE-STATE-STORE-0000000001임을 볼 수 있어요. 이제 집계에서 일부 레코드를 걸러내기 위해 filter를 추가하면 어떻게 되는지 봐요:
KStream<String,String> stream = builder.stream("input");
stream.filter((k,v)-> v !=null && v.length() >= 6 )
.groupByKey()
.count()
.toStream()
.to("output");
그리고 해당 토폴로지:
Topologies:
Sub-topology: 0
Source: KSTREAM-SOURCE-0000000000 (topics: [input])
--> KSTREAM-FILTER-0000000001
Processor: KSTREAM-FILTER-0000000001 (stores: [])
--> KSTREAM-AGGREGATE-0000000003
<-- KSTREAM-SOURCE-0000000000
Processor: KSTREAM-AGGREGATE-0000000003 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000002])
--> KTABLE-TOSTREAM-0000000004
<-- KSTREAM-FILTER-0000000001
Processor: KTABLE-TOSTREAM-0000000004 (stores: [])
--> KSTREAM-SINK-0000000005
<-- KSTREAM-AGGREGATE-0000000003
Sink: KSTREAM-SINK-0000000005 (topic: output)
<-- KTABLE-TOSTREAM-0000000004
count 연산 앞에 연산을 추가했기 때문에 상태 저장소(그리고 체인지로그 토픽) 이름이 변경된 것을 주목하세요. 이 이름 변경은 업데이트된 토폴로지를 롤링 재배포할 수 없다는 것을 의미해요. 또한 시작 시 체인지로그 토픽이 변경되었고 새 체인지로그 토픽에는 데이터가 없으므로, 집계를 다시 계산하려면 Streams Reset Tool을 사용해야 해요. 다행히 이 상황을 해결하는 쉬운 방법이 있어요. 생성된 이름에 의존하는 대신 상태 저장소에 사용자 정의 이름을 주면, 토폴로지 변경이 상태 저장소 이름을 흔들 걱정을 하지 않아도 돼요. Joined, StreamJoined, Grouped 클래스로 리파티션 토픽에 이름을 붙일 수 있었고, Materialized로 상태 저장소와 체인지로그 토픽에 이름을 붙일 수 있었어요. 하지만 이 DSL 토폴로지 연산에 이름을 붙이는 중요성을 다시 강조할 가치가 있어요. 상태 저장소에 특정 이름을 주는 DSL 코드는 이렇게 됩니다:
KStream<String,String> stream = builder.stream("input");
stream.filter((k, v) -> v != null && v.length() >= 6)
.groupByKey()
.count(Materialized.as("Purchase_count_store"))
.toStream()
.to("output");
그리고 토폴로지는:
Topologies:
Sub-topology: 0
Source: KSTREAM-SOURCE-0000000000 (topics: [input])
--> KSTREAM-FILTER-0000000001
Processor: KSTREAM-FILTER-0000000001 (stores: [])
--> KSTREAM-AGGREGATE-0000000002
<-- KSTREAM-SOURCE-0000000000
Processor: KSTREAM-AGGREGATE-0000000002 (stores: [Purchase_count_store])
--> KTABLE-TOSTREAM-0000000003
<-- KSTREAM-FILTER-0000000001
Processor: KTABLE-TOSTREAM-0000000003 (stores: [])
--> KSTREAM-SINK-0000000004
<-- KSTREAM-AGGREGATE-0000000002
Sink: KSTREAM-SINK-0000000004 (topic: output)
<-- KTABLE-TOSTREAM-0000000003
이제 상태 저장소 앞에 프로세서를 추가해도 저장소 이름과 그 체인지로그 토픽 이름은 변하지 않아요. 이것은 토폴로지를 프로세서 추가·제거로 인한 변경에 더 견고하고 탄력적으로 만들어요.
결론
DSL을 사용할 때 처리 노드에 이름을 붙이는 것은 좋은 습관이며, 애플리케이션에 리파티션 토픽과 상태 저장소(그리고 그에 따른 체인지로그 토픽) 같은 "stateful" 프로세서가 있을 때는 더욱 중요해요.
DSL 토폴로지에 이름을 붙일 때 기억할 몇 가지 포인트:
- 기존 토폴로지가 있고 상태 저장소(및 체인지로그 토픽)와 리파티션 토픽에 이름을 붙이지 않았다면, 이름을 붙이는 것을 권장해요. 하지만 이것은 토폴로지에 변화를 주는 변경(breaking change)이므로, 모든 애플리케이션 인스턴스를 종료하고 변경을 하고 Streams Reset Tool을 실행해야 해요. 처음에는 불편할 수 있지만, 토폴로지 변경으로 인한 예상치 못한 오류로부터 애플리케이션을 보호하는 노력은 가치가 있어요.
- 새 토폴로지라면, 토폴로지의 영속적인 부분인 상태 저장소(체인지로그 토픽)와 리파티션 토픽에 이름을 붙이세요. 이렇게 하면 배포 시 Kafka Streams 애플리케이션을 깨뜨릴 수 있는 토폴로지 변경으로부터 보호돼요. 처음에 stateless 프로세서에 이름을 추가하고 싶지 않다면 괜찮아요. 나중에 언제든지 돌아와서 이름을 추가할 수 있으니까요.
토폴로지 이름 변경이 애플리케이션을 깨뜨리는 것을 방지하기 위해 Kafka Streams 애플리케이션의 중요한 부분에 이름을 붙이는 빠른 참조:
| 연산 | 명명 클래스 |
|---|---|
| 집계 리파티션 토픽 | Grouped |
| KStream-KStream Join 리파티션 토픽 | StreamJoined |
| KStream-KTable Join 리파티션 토픽 | Joined |
| KStream-KStream Join 상태 저장소 | StreamJoined |
| 상태 저장소 (집계 및 KTable-KTable 조인용) | Materialized |
| Stream/Table 비상태 연산 | Named |
모범 사례를 더 강화하기 위해 Kafka Streams는 ensure.explicit.internal.resource.naming 구성 옵션을 제공해요:
Properties props = new Properties();
props.put(StreamsConfig.ENSURE_EXPLICIT_INTERNAL_RESOURCE_NAMING_CONFIG, true);
이 파라미터는 모든 내부 토픽, 상태 저장소, 체인지로그 토픽이 명시적으로 정의된 이름을 가지도록 보장해요. 이 구성이 활성화되면 이러한 구성 요소 중 하나라도 자동 생성 이름에 의존하면 Kafka Streams 애플리케이션이 시작되지 않아요. 새 프로세서나 변환이 추가되어도 수동으로 정의된 이름은 변하지 않으므로, 이는 토폴로지 업데이트 전반에 걸쳐 안정성을 보장해요. 명시적 명명을 강제하는 것은 특히 프로덕션 환경에서 중요해요. 신뢰할 수 있는 스트림 처리 애플리케이션을 유지하는 데 일관성과 하위 호환성이 필수적이기 때문이에요.
더 알아보기
- Streams DSL — DSL 연산자들이 실제로 어떻게 구성되는지 봐요.
- Processor API — PAPI에서 이름을 명시해야 하는 이유를 봐요.
- Application Reset Tool — 이름 변경 토폴로지에서 상태를 재계산하는 방법을 봐요.