아키텍처

아키텍처

Kafka Streams가 내부적으로 어떻게 돌아가는지 궁금하다면 이 페이지가 정답이에요. 카프카의 프로듀서·컨슈머 라이브러리 위에 올라가서, 카프카의 데이터 병렬성·분산 조정·장애 허용·운영 단순함을 그대로 활용해요. 스트림 파티션과 태스크가 어떻게 병렬성을 만들어내는지, 스레딩 모델, 로컬 상태 저장소, 장애 허용까지 큰 그림을 잡아드릴게요.

출처: 문서

본문

Kafka Streams는 카프카 프로듀서·컨슈머 라이브러리 위에 구축되고 카프카의 네이티브 기능을 활용해 애플리케이션 개발을 단순화해요. 데이터 병렬성, 분산 조정, 장애 허용, 운영 단순함을 제공하죠. 이 섹션에서는 Kafka Streams가 내부적으로 어떻게 동작하는지 설명해요.

스트림 파티션과 태스크

카프카의 메시징 레이어는 저장·전송을 위해 데이터를 파티셔닝해요. Kafka Streams는 처리를 위해 데이터를 파티셔닝해요. 두 경우 모두 이 파티셔닝이 데이터 지역성, 탄력성, 확장성, 고성능, 장애 허용을 가능하게 해요. Kafka Streams는 카프카 토픽 파티션을 기반으로 파티션과 태스크를 병렬성 모델의 논리적 단위로 사용해요. 병렬성 측면에서 Kafka Streams와 카프카 사이에는 밀접한 연결이 있어요:

  • 각 스트림 파티션은 전적으로 순서화된 데이터 레코드의 시퀀스이며 카프카 토픽 파티션에 매핑돼요.
  • 스트림의 데이터 레코드는 해당 토픽의 카프카 메시지에 매핑돼요.
  • 데이터 레코드의 키는 카프카와 Kafka Streams 양쪽에서 데이터의 파티셔닝, 즉 데이터가 토픽 내 특정 파티션으로 라우팅되는 방식을 결정해요.

애플리케이션의 프로세서 토폴로지는 여러 태스크로 나누어 스케일됩니다. 더 구체적으로 Kafka Streams는 애플리케이션의 입력 스트림 파티션을 기준으로 고정된 수의 태스크를 만들고, 각 태스크에 입력 스트림(즉 카프카 토픽)의 파티션 목록을 할당해요. 파티션에서 태스크로의 할당은 절대 변하지 않아서 각 태스크는 애플리케이션의 고정된 병렬성 단위예요. 그런 다음 태스크는 할당된 파티션을 기반으로 자신의 프로세서 토폴로지를 인스턴스화할 수 있어요. 또한 각 할당된 파티션에 대한 버퍼를 유지하고 이 레코드 버퍼에서 메시지를 한 번에 하나씩 처리해요. 그 결과 스트림 태스크는 수동 개입 없이 독립적이고 병렬로 처리될 수 있어요.

약간 단순화하면, 애플리케이션이 실행될 수 있는 최대 병렬성은 최대 스트림 태스크 수로 제한되며, 이는 다시 애플리케이션이 읽는 입력 토픽의 최대 파티션 수로 결정돼요. 예를 들어 입력 토픽에 5개의 파티션이 있다면 애플리케이션 인스턴스를 최대 5개 실행할 수 있어요. 이 인스턴스들은 토픽 데이터를 협력적으로 처리해요. 입력 토픽의 파티션 수보다 많은 앱 인스턴스를 실행하면, "초과" 앱 인스턴스는 실행되지만 유휴 상태로 유지돼요. 하지만 바쁜 인스턴스 중 하나가 내려가면, 유휴 인스턴스 중 하나가 전자의 작업을 재개해요.

중요한 것은 Kafka Streams가 리소스 매니저가 아니라, 스트림 처리 애플리케이션이 실행되는 어디든 "실행"되는 라이브러리라는 점이에요. 애플리케이션의 여러 인스턴스는 같은 머신에서 또는 여러 머신에 걸쳐 실행되고, 태스크는 라이브러리에 의해 실행 중인 애플리케이션 인스턴스에 자동으로 분산될 수 있어요. 파티션에서 태스크로의 할당은 절대 변하지 않아요. 애플리케이션 인스턴스가 실패하면, 할당된 모든 태스크가 다른 인스턴스에서 자동으로 재시작되고 같은 스트림 파티션에서 계속 소비해요.

참고: 토픽 파티션은 태스크에, 태스크는 모든 인스턴스에 걸친 모든 스레드에, 로드 밸런싱과 상태 저장 태스크의 스티키니스(stickiness)를 절충하는 최선의 노력(best-effort) 방식으로 할당돼요. 이 할당을 위해 Kafka Streams는 StreamsPartitionAssignor 클래스를 사용하고 다른 assignor로 변경할 수 없어요. 다른 assignor를 사용하려고 하면 Kafka Streams가 그것을 무시해요.

스레딩 모델

Kafka Streams는 라이브러리가 애플리케이션 인스턴스 내 처리를 병렬화하는 데 사용할 수 있는 스레드 수를 사용자가 구성할 수 있게 해요. 각 스레드는 자신의 프로세서 토폴로지를 가진 하나 이상의 태스크를 독립적으로 실행할 수 있어요.

더 많은 스트림 스레드나 애플리케이션 인스턴스를 시작하는 것은 단지 토폴로지를 복제하고 그것이 카프카 파티션의 다른 부분 집합을 처리하게 하는 것뿐이에요. 효과적으로 처리를 병렬화하죠. 스레드 사이에는 공유 상태가 없으므로 스레드 간 조정이 필요 없다는 점을 유의할 만해요. 이것은 애플리케이션 인스턴스와 스레드 전반에 걸쳐 토폴로지를 병렬로 실행하는 것을 매우 간단하게 만들어요. 다양한 스트림 스레드 사이의 카프카 토픽 파티션 할당은 카프카의 조정 기능을 활용하는 Kafka Streams가 투명하게 처리해요.

위에서 설명한 것처럼 Kafka Streams로 스트림 처리 애플리케이션을 스케일링하는 것은 쉬워요: 추가 애플리케이션 인스턴스를 시작하기만 하면, Kafka Streams가 애플리케이션 인스턴스에서 실행되는 태스크 간에 파티션을 분산해요. 입력 카프카 토픽 파티션 수만큼 애플리케이션 스레드를 시작할 수 있어서, 애플리케이션의 모든 실행 인스턴스에 걸쳐 모든 스레드(더 정확히는 그것이 실행하는 태스크)가 처리할 입력 파티션을 하나 이상 가지게 돼요.

Kafka 2.8부터 카프카 스트림즈 클라이언트를 스케일링하는 것과 같은 방식으로 스트림 스레드를 스케일링할 수 있어요. 스트림 스레드를 추가하거나 제거하기만 하면 카프카 스트림즈가 파티션 재분배를 처리해요. 죽은 스트림 스레드를 대체하기 위해 스레드를 추가할 수도 있고, 실행 중인 스레드 수를 복구하기 위해 클라이언트를 재시작할 필요가 없어요.

로컬 상태 저장소

Kafka Streams는 소위 상태 저장소(state store)를 제공해요. 스트림 처리 애플리케이션이 데이터를 저장하고 쿼리하는 데 사용할 수 있는데, 상태 저장(stateful) 연산을 구현할 때 중요한 기능이에요. 예를 들어 Kafka Streams DSL은 join()이나 aggregate() 같은 상태 저장 연산자를 호출하거나 스트림을 윈도윙할 때 이 상태 저장소를 자동으로 생성·관리해요.

Kafka Streams 애플리케이션의 모든 스트림 태스크는 처리를 위해 필요한 데이터를 저장·쿼리하는 API로 접근할 수 있는 하나 이상의 로컬 상태 저장소를 임베드할 수 있어요. Kafka Streams는 이러한 로컬 상태 저장소에 대한 장애 허용과 자동 복구를 제공해요.

장애 허용

Kafka Streams는 카프카에 네이티브로 통합된 장애 허용 기능 위에 구축돼요. 카프카 파티션은 고가용성이고 복제되므로, 스트림 데이터가 카프카에 영속되면 애플리케이션이 실패하고 재처리해야 해도 사용 가능해요. Kafka Streams의 태스크는 장애를 처리하기 위해 카프카 컨슈머 클라이언트가 제공하는 장애 허용 기능을 활용해요. 태스크가 실패한 머신에서 실행되고 있다면, Kafka Streams는 애플리케이션의 남아 있는 실행 인스턴스 중 하나에서 태스크를 자동으로 재시작해요.

또한 Kafka Streams는 로컬 상태 저장소도 장애에 강하도록 보장해요. 각 상태 저장소에 대해 복제된 체인지로그 카프카 토픽을 유지하며 상태 업데이트를 추적해요. 이 체인지로그 토픽도 파티셔닝되어 각 로컬 상태 저장소 인스턴스(그리고 저장소에 접근하는 태스크)가 고유한 전용 체인지로그 토픽 파티션을 가지게 해요. 체인지로그 토픽에서 로그 컴팩션을 활성화해 오래된 데이터를 안전하게 제거하여 토픽이 무한정 커지는 것을 방지해요. 태스크가 실패한 머신에서 실행되고 다른 머신에서 재시작되면, Kafka Streams는 새로 시작된 태스크에서 처리를 재개하기 전에 해당 체인지로그 토픽을 재생함으로써 연결된 상태 저장소를 장애 전의 내용으로 복원하는 것을 보장해요. 그 결과 장애 처리는 최종 사용자에게 완전히 투명해요.

태스크 (재)초기화 비용은 보통 상태 저장소의 체인지로그 토픽 재생에 의한 상태 복원 시간에 주로 의존한다는 점을 유의해요. 이 복원 시간을 최소화하기 위해 사용자는 로컬 상태의 스탠바이 복제본(즉 상태의 완전히 복제된 사본)을 가지도록 애플리케이션을 구성할 수 있어요. 태스크 마이그레이션이 발생하면 Kafka Streams는 그러한 스탠바이 복제본이 이미 존재하는 애플리케이션 인스턴스에 태스크를 할당해 태스크 (재)초기화 비용을 최소화해요. Kafka Streams Configs 섹션의 num.standby.replicas를 참고해요. 2.6부터 Kafka Streams는 완전히 따라잡은(caught-up) 로컬 상태 사본을 가진 인스턴스가 존재한다면 태스크가 항상 그러한 인스턴스에만 할당되도록 보장해요. 스탠바이 태스크는 장애 시 따라잡은 인스턴스가 존재할 가능성을 높여줘요.

스탠바이 복제본을 랙 인지(rack awareness)로 구성할 수도 있어요. 구성되면 Kafka Streams는 활성 태스크와 다른 "랙"에 스탠바이 태스크를 분산시키려고 시도해, 활성 태스크의 랙이 실패할 때 더 빠른 복구 시간을 가져요. Kafka Streams Developer Guide 섹션의 rack.aware.assignment.tags를 참고해요.

카프카 컨슈머의 랙을 설정하는 client.rack 클라이언트 구성도 있어요. 브로커도 broker.rack으로 랙을 설정했다면, rack.aware.assignment.strategy(Kafka Streams Developer Guide 참조)로 랙 인지 태스크 할당을 활성화해, 같은 랙의 클라이언트에 태스크를 할당하여 교차 랙 트래픽을 줄이는 태스크 할당을 계산할 수 있어요. client.rack은 활성 랙과 다른 랙에 스탠바이 태스크를 분산하는 데도 사용할 수 있으며, rack.aware.assignment.tags와 비슷한 기능이에요. 현재 rack.aware.assignment.tag가 스탠바이 태스크 분산에서 우선하므로, 두 구성이 모두 있으면 rack.aware.assignment.tag가 활성 랙과 다른 랙에 스탠바이 태스크를 분산하는 데 사용돼요. 더 많은 태그 키를 구성할 수 있기 때문이에요.

더 알아보기