Streams 애플리케이션 실행하기
Streams 애플리케이션 실행하기
카프카 스트림즈 애플리케이션을 만들었다면 이제 "어떻게 실행하고, 어떻게 늘리고 줄이고, 장애가 나면 어떻게 되지?"라는 궁금증이 생길 거예요. 이 페이지는 그런 운영 관점의 질문을 한 번에 정리해드려요. 별다른 설정 없이 그냥 Java 앱처럼 실행하면 되고, 인스턴스를 늘리면 자동으로 스케일 아웃된다는 점이 핵심이에요.
출처: 문서
본문
Kafka Streams 라이브러리를 사용하는 Java 애플리케이션은 추가 구성이나 요구사항 없이 실행할 수 있어요. Kafka Streams는 또한 애플리케이션의 다양한 상태에 대한 알림을 받는 기능도 제공해요. 런타임 상태 모니터링에 대해서는 모니터링 가이드에서 다뤄요.
Kafka Streams 애플리케이션 시작하기
Java 애플리케이션을 fat JAR 파일로 패키징한 다음 다음과 같이 시작할 수 있어요:
# Start the application in class `com.example.MyStreamsApp`
# from the fat JAR named `path-to-app-fatjar.jar`.
$ java -cp path-to-app-fatjar.jar com.example.MyStreamsApp
애플리케이션을 시작하면 애플리케이션의 Kafka Streams 인스턴스 하나를 실행하는 거예요. 여러 인스턴스를 실행할 수 있어요. 흔한 시나리오는 애플리케이션의 여러 인스턴스가 병렬로 실행되는 거예요. 자세한 내용은 병렬성 모델을 참고해요.
애플리케이션 인스턴스가 실행을 시작하면 정의된 프로세서 토폴로지가 하나 이상의 스트림 태스크로 초기화돼요. 프로세서 토폴로지가 상태 저장소(state store)를 정의한다면 이것들도 초기화 기간에 구성돼요. 자세한 내용은 워크로드 리밸런스 중 상태 복원 섹션을 참고해요.
애플리케이션의 탄력적 스케일링
Kafka Streams는 스트림 처리 애플리케이션을 탄력적이고 확장 가능하게 만들어요. 애플리케이션 런타임 중에 다운타임이나 데이터 손실 없이 처리 용량을 동적으로 추가·제거할 수 있어요. 이것은 애플리케이션을 장애에 강하게 만들고, 필요에 따라 유지보수(예: 롤링 업그레이드)를 할 수 있게 해줘요.
이 탄력성에 대한 자세한 내용은 병렬성 모델 섹션을 참고해요. Kafka Streams는 카프카 와이어 프로토콜에 내장된 카프카 그룹 관리 기능을 활용해요. 이것이 Kafka Streams 애플리케이션의 탄력성을 가능하게 하는 기반이에요: 그룹의 멤버들이 카프카의 데이터 소비와 처리를 공동으로 조정·협력해요. 또한 Kafka Streams는 상태 저장(stateful) 처리를 제공하고, 애플리케이션 인스턴스가 언제든지 나오고 사라질 수 있는 환경에서 장애 허용 상태를 허용해요.
애플리케이션에 용량 추가하기
스트림 처리 애플리케이션에 더 많은 처리 용량이 필요하면, 스케일 아웃을 위해 스트림 처리 애플리케이션의 다른 인스턴스(예: 다른 머신에서)를 시작하면 돼요. 애플리케이션 인스턴스들은 서로를 인식하고 자동으로 처리 작업을 공유하기 시작해요. 더 구체적으로, 기존 인스턴스에서 새 인스턴스로 넘겨지는 것은 (일부) 기존 인스턴스가 실행하던 스트림 태스크예요. 스트림 태스크를 한 인스턴스에서 다른 인스턴스로 이동하면 처리 작업과 이 스트림 태스크의 내부 상태를 이동하게 돼요(태스크의 상태는 해당 체인지로그 토픽에서 상태를 복원함으로써 대상 인스턴스에 재생성돼요).
애플리케이션의 각 인스턴스는 각자 자신의 JVM 프로세스에서 실행돼요. 즉 각 인스턴스는 해당 JVM 프로세스에서 사용 가능한 모든 처리 용량(애플리케이션의 비-Kafka-Streams 부분이 사용할 수 있는 용량 제외)을 활용할 수 있어요. 이것이 추가 인스턴스를 실행하면 애플리케이션에 추가 처리 용량이 부여되는 이유예요. 새 인스턴스를 실행함으로써 추가할 정확한 용량은 물론 새 인스턴스가 실행되는 환경에 달려 있어요: 사용 가능한 CPU 코어, 메인 메모리와 Java 힙 공간, 로컬 저장소, 네트워크 대역폭 등요. 마찬가지로, 실행 중인 인스턴스 중 하나를 중지하면 해당 처리 용량을 제거·해제해요.
- 용량 추가 전: Kafka Streams 애플리케이션의 단일 인스턴스만 실행 중이에요. 이 시점에 애플리케이션의 카프카 컨슈머 그룹에는 단일 멤버(이 인스턴스)만 포함돼요. 모든 데이터는 이 단일 인스턴스가 읽고 처리해요.
- 용량 추가 후: 이제 Kafka Streams 애플리케이션의 추가 인스턴스 두 개가 실행 중이고, 자동으로 애플리케이션의 카프카 컨슈머 그룹에 합류해 총 3명의 현재 멤버가 돼요. 이 세 인스턴스는 자동으로 처리 작업을 서로 나눠요. 분할은 데이터가 읽히는 카프카 토픽 파티션을 기준으로 해요.
애플리케이션에서 용량 제거하기
처리 용량을 제거하려면 실행 중인 스트림 처리 애플리케이션 인스턴스를 중지할 수 있어요(예: 네 인스턴스 중 두 개를 종료). 그러면 자동으로 애플리케이션의 컨슈머 그룹을 떠나고, 나머지 인스턴스들이 자동으로 처리 작업을 인계받아요. 나머지 인스턴스들은 중지된 인스턴스가 실행하던 스트림 태스크를 인계받아요. 스트림 태스크를 한 인스턴스에서 다른 인스턴스로 이동하면 처리 작업과 이 스트림 태스크의 내부 상태를 이동하게 돼요. 태스크의 상태는 체인지로그 토픽에서 대상 인스턴스에 재생성돼요.
워크로드 리밸런스 중 상태 복원
태스크가 마이그레이션될 때, 애플리케이션 인스턴스가 처리를 재개하기 전에 태스크 처리 상태가 완전히 복원돼요. 이것은 올바른 처리 결과를 보장해요. Kafka Streams에서 상태 복원은 보통 해당 체인지로그 토픽을 재생(replay)해 상태 저장소를 재구성함으로써 이루어져요. 복제된 로컬 상태 저장소를 사용해 체인지로그 기반 복원 지연을 최소화하려면 num.standby.replicas를 지정할 수 있어요. 스트림 태스크가 애플리케이션 인스턴스에서 초기화되거나 재초기화될 때, 그 상태 저장소는 다음과 같이 복원돼요:
- 로컬 상태 저장소가 없으면, 체인지로그가 가장 이른 오프셋부터 현재 오프셋까지 재생돼요. 이것은 로컬 상태 저장소를 가장 최근 스냅샷으로 재구성해요.
- 로컬 상태 저장소가 있으면, 체인지로그가 이전에 체크포인트된 오프셋부터 재생돼요. 변경 사항이 적용되고 상태가 가장 최근 스냅샷으로 복원돼요. 이 방법은 체인지로그의 더 작은 부분을 적용하므로 시간이 덜 걸려요.
자세한 내용은 스탠바이 복제본을 참고해요.
버전 2.6부터 Streams는 워밍업 복제본(warmup replica)을 통해 태스크 복원의 대부분을 백그라운드에서 수행해요. 이것들은 태스크에 대해 많은 상태를 복원해야 하는 인스턴스에 할당돼요. 상태 저장 활성 태스크는 상태가 구성된 acceptable.recovery.lag(존재하는 경우) 내에 있게 된 후에만 인스턴스에 할당돼요. 이는 대부분의 경우 태스크 마이그레이션이 해당 태스크의 다운타임을 초래하지 않는다는 뜻이에요. 이미 따라잡은(caught up) 인스턴스에서 활성 상태를 유지하고, 마이그레이션 대상 인스턴스는 상태 복원 작업을 하는 동안 말이죠. Streams는 복원을 마친 워밍업 태스크를 주기적으로 검사하고 준비되면 활성 태스크로 전환해요.
참고로, 태스크 가용성의 유일한 예외는 어떤 인스턴스도 해당 태스크의 따라잡은 버전을 가지고 있지 않은 경우예요. 그 경우, 따라잡지 않은 인스턴스에 활성 태스크를 할당할 수밖에 없고, 체인지로그에서 태스크 상태를 복원하는 동안 추가 처리를 차단해야 해요. 애플리케이션에 고가용성이 중요하다면 스탠바이를 활성화하는 것을 적극 권장해요.
실행할 애플리케이션 인스턴스 수 결정하기
Kafka Streams 애플리케이션의 병렬성은 주로 입력 토픽이 몇 개의 파티션을 가지는지에 의해 결정돼요. 예를 들어, 애플리케이션이 10개의 파티션을 가진 단일 토픽에서 읽는다면 애플리케이션 인스턴스를 최대 10개 실행할 수 있어요. 더 많은 인스턴스를 실행할 수는 있지만, 그것들은 유휴 상태가 될 거예요.
토픽 파티션 수는 Kafka Streams 애플리케이션의 병렬성과 실행 중인 인스턴스 수의 상한선이에요.
애플리케이션 인스턴스 전반에 걸친 균형 잡힌 워크로드 처리를 달성하고 처리 핫스팟을 방지하려면 데이터와 처리 워크로드를 분산해야 해요:
- 데이터는 토픽 파티션 전체에 균등하게 분산되어야 해요. 예를 들어, 두 파티션이 각각 100만 개의 메시지를 가지는 것이, 한 파티션에 200만 개가 있고 다른 파티션에는 없는 것보다 나아요.
- 처리 워크로드는 토픽 파티션 전체에 균등하게 분산되어야 해요. 예를 들어, 메시지 처리 시간이 크게 다르다면, 처리 집약적인 메시지를 같은 파티션에 두는 것보다 파티션 전체에 퍼뜨리는 것이 좋아요.
크래시와 장애 처리
Kafka Streams 애플리케이션의 크래시와 장애 가능성을 줄이기 위해 할 수 있는 몇 가지가 있어요.
- Kafka Streams는 브로커 장애에 대한 복원력을 돕는 몇 가지 구성이 있어요. 그것들은 구성 가이드에서 찾을 수 있어요.
- 애플리케이션이 오류와 장애를 처리할 수 있는지 확인하세요. 여기에는 인가·역직렬화 오류 같은 오류를 처리하기 위한 올바른 예외 핸들러 구성, 그리고 "poison pill" 레코드를 처리하기 위해 데드 레터 큐(dead letter queue) 같은 전략 사용이 포함돼요.
Kafka Streams 애플리케이션이 크래시되거나 실패하면, 먼저 PENDING_ERROR 상태로 들어가 기존 리소스를 모두 정상적으로(gracefully) 닫은 다음 ERROR 상태로 전환돼요. PENDING_ERROR 상태는 복구 가능하지 않으며, 재시작해야만 애플리케이션이 RUNNING 상태로 돌아온다는 점을 유의해야 해요. 따라서 ERROR 상태와 함께 이 상태를 모니터링하는 것이 애플리케이션이 복구할 수 있도록 보장하는 데 중요해요.
더 알아보기
- 관리 토픽 — 내부 토픽과 리텐션 구성을 봐요.
- Streams 보안 — 보안 환경에서 실행할 때의 설정을 봐요.
- 핵심 개념 — 병렬성 모델과 태스크 개념을 제대로 이해해요.