핵심 개념

핵심 개념 (Core Concepts)

Kafka Streams를 제대로 쓰려면 몇 가지 핵심 개념을 먼저 잡아두는 게 좋아요. 스트림 처리 토폴로지가 뭔지, 시간(이벤트 시간·처리 시간·수집 시간)의 의미, 스트림-테이블 이중성(duality), 집계·윈도잉, 상태 저장소, 그리고 정확히 한 번(exactly-once) 처리 보장까지 이 페이지에서 처음부터 차근차근 설명해드릴게요.

출처: 문서

본문

Kafka Streams는 카프카에 저장된 데이터를 처리하고 분석하기 위한 클라이언트 라이브러리예요. 이벤트 시간과 처리 시간을 제대로 구분하기, 윈도잉 지원, 간단하면서도 효율적인 애플리케이션 상태 관리와 실시간 쿼리 같은 중요한 스트림 처리 개념 위에 구축돼요.

Kafka Streams는 진입 장벽이 낮아요. 단일 머신에서 소규모 개념 증명(proof-of-concept)을 빠르게 작성하고 실행할 수 있고, 대용량 프로덕션 워크로드로 확장하려면 여러 머신에서 애플리케이션의 추가 인스턴스를 실행하기만 하면 돼요. Kafka Streams는 카프카의 병렬성 모델을 활용해 같은 애플리케이션의 여러 인스턴스의 로드 밸런싱을 투명하게 처리해요.

Kafka Streams의 몇 가지 하이라이트:

  • 간단하고 가벼운 클라이언트 라이브러리로 설계되어, 어떤 Java 애플리케이션에도 쉽게 임베드되고 스트리밍 애플리케이션을 위해 사용자가 가진 기존 패키징·배포·운영 도구와 통합될 수 있어요.
  • 내부 메시징 레이어로 Apache Kafka 자체 외에는 다른 시스템에 대한 외부 의존성이 없어요. 특히 카프카의 파티셔닝 모델을 사용해 강력한 순서 보장을 유지하면서 처리를 수평 확장해요.
  • 장애 허용 로컬 상태를 지원해 윈도우 조인과 집계 같은 매우 빠르고 효율적인 상태 저장 연산을 가능하게 해요.
  • 정확히 한 번(exactly-once) 처리 의미론을 지원해, 처리 중간에 Streams 클라이언트나 카프카 브로커에서 장애가 발생해도 각 레코드가 정확히 한 번만 처리되도록 보장해요.
  • 한 번에 하나의 레코드(one-record-at-a-time) 처리로 밀리초 처리 지연을 달성하고, 레코드의 순서가 뒤집혀 도착하는 경우에도 이벤트 시간 기반 윈도잉 연산을 지원해요.
  • 필요한 스트림 처리 프리미티브와 함께, 고수준 Streams DSL과 저수준 Processor API를 제공해요.

먼저 Kafka Streams의 핵심 개념을 요약할게요.

스트림 처리 토폴로지

**스트림(Stream)**은 Kafka Streams가 제공하는 가장 중요한 추상화예요. 무한하고 연속적으로 업데이트되는 데이터 집합을 나타내요. 스트림은 불변 데이터 레코드의 순서화되고 재생 가능하며 장애 허용적인 시퀀스이며, 데이터 레코드는 키-값 쌍으로 정의돼요.

**스트림 처리 애플리케이션(Stream processing application)**은 Kafka Streams 라이브러리를 사용하는 모든 프로그램이에요. 하나 이상의 **프로세서 토폴로지(processor topology)**를 통해 계산 로직을 정의하는데, 프로세서 토폴로지는 스트림(엣지)으로 연결된 스트림 프로세서(노드)의 그래프예요.

**스트림 프로세서(Stream processor)**는 프로세서 토폴로지의 노드예요. 토폴로지에서 업스트림 프로세서로부터 한 번에 하나의 입력 레코드를 받아 연산을 적용하고, 이후 다운스트림 프로세서에 하나 이상의 출력 레코드를 생성하는 데이터 변환 처리 단계를 나타내요.

토폴로지에는 두 가지 특수 프로세서가 있어요:

  • 소스 프로세서 (Source Processor): 업스트림 프로세서가 없는 특수한 유형의 스트림 프로세서예요. 하나 이상의 카프카 토픽에서 레코드를 소비하고 그것들을 다운스트림 프로세서로 전달함으로써 토폴로지에 입력 스트림을 만들어요.
  • 싱크 프로세서 (Sink Processor): 다운스트림 프로세서가 없는 특수한 유형의 스트림 프로세서예요. 업스트림 프로세서에서 받은 모든 레코드를 지정된 카프카 토픽으로 보내요.

일반 프로세서 노드에서는 현재 레코드를 처리하는 동안 다른 원격 시스템에도 접근할 수 있다는 점을 유의해요. 따라서 처리된 결과를 카프카로 다시 스트리밍하거나 외부 시스템에 쓸 수 있어요.

Kafka Streams는 스트림 처리 토폴로지를 정의하는 두 가지 방법을 제공해요. Kafka Streams DSL은 map, filter, join, aggregations 같은 가장 일반적인 데이터 변환 연산을 기본 제공해요. 저수준 Processor API는 개발자가 커스텀 프로세서를 정의·연결하고 상태 저장소와 상호작용할 수 있게 해줘요.

프로세서 토폴로지는 스트림 처리 코드를 위한 단지 논리적 추상화일 뿐이에요. 런타임에 논리적 토폴로지는 병렬 처리를 위해 애플리케이션 내부에서 인스턴스화되고 복제돼요 (자세한 내용은 스트림 파티션과 태스크 참조).

시간 (Time)

스트림 처리에서 중요한 측면은 **시간(time)**의 개념과, 그것이 어떻게 모델링되고 통합되는지예요. 예를 들어 윈도잉 같은 일부 연산은 시간 경계를 기반으로 정의돼요.

스트림에서 흔한 시간의 개념:

  • 이벤트 시간 (Event time) — 이벤트나 데이터 레코드가 발생한 시점, 즉 "소스에서" 원래 생성된 시점이에요. 예: 이벤트가 자동차의 GPS 센서가 보고한 지리 위치 변경이라면, 연관된 이벤트 시간은 GPS 센서가 위치 변경을 캡처한 시간이에요.
  • 처리 시간 (Processing time) — 이벤트나 데이터 레코드가 스트림 처리 애플리케이션에 의해 처리되는 시점, 즉 레코드가 소비되는 시점이에요. 처리 시간은 원래 이벤트 시간보다 밀리초, 시간, 또는 며칠 등 더 나중일 수 있어요. 예: 자동차 센서가 보고한 지리 위치 데이터를 읽고 처리해 차량 관리 대시보드에 표시하는 분석 애플리케이션을 상상해보세요. 여기서 분석 애플리케이션의 처리 시간은 이벤트 시간 후 밀리초나 초(Apache Kafka와 Kafka Streams 기반 실시간 파이프라인의 경우) 또는 시간(Apache Hadoop이나 Apache Spark 기반 배치 파이프라인의 경우)일 수 있어요.
  • 수집 시간 (Ingestion time) — 이벤트나 데이터 레코드가 카프카 브로커에 의해 토픽 파티션에 저장된 시점이에요. 이벤트 시간과의 차이는 이 수집 타임스탬프가 레코드가 "소스에서" 생성될 때가 아니라 카프카 브로커가 레코드를 대상 토픽에 추가할 때 생성된다는 점이에요. 처리 시간과의 차이는 처리 시간이 스트림 처리 애플리케이션이 레코드를 처리하는 시점이라는 점이에요. 예: 레코드가 절대 처리되지 않으면 그 레코드에 대한 처리 시간의 개념은 없지만, 여전히 수집 시간은 있어요.

이벤트 시간과 수집 시간 사이의 선택은 실제로 카프카(카프카 스트림즈가 아니라)의 구성을 통해 이루어져요. Kafka 0.10.x부터 타임스탬프가 카프카 메시지에 자동으로 임베드돼요. 카프카 구성에 따라 이 타임스탬프는 이벤트 시간 또는 수집 시간을 나타내요. 해당 카프카 구성 설정은 브로커 수준 또는 토픽별로 지정할 수 있어요. Kafka Streams의 기본 타임스탬프 추출기(default timestamp extractor)는 이 임베드된 타임스탬프를 그대로 가져와요. 따라서 애플리케이션의 유효 시간 의미론은 이러한 임베드된 타임스탬프에 대한 유효 카프카 구성에 따라 달라져요.

Kafka Streams는 TimestampExtractor 인터페이스를 통해 모든 데이터 레코드에 타임스탬프를 할당해요. 이러한 레코드별 타임스탬프는 시간 측면에서 스트림의 진행을 설명하고 윈도우 연산 같은 시간 종속 연산에 활용돼요. 그 결과, 이 시간은 새 레코드가 프로세서에 도착할 때만 진행돼요. 우리는 이 데이터 중심 시간을 애플리케이션의 **스트림 시간(stream time)**이라고 부르며, 애플리케이션이 실제로 실행되는 벽시계(wall-clock) 시간과 구분해요. TimestampExtractor 인터페이스의 구체적인 구현은 스트림 시간 정의에 다른 의미론을 제공해요. 예를 들어 임베드된 타임스탬프 필드 같은 데이터 레코드의 실제 내용을 기반으로 타임스탬프를 가져와 이벤트 시간 의미론을 제공하거나, 현재 벽시계 시간을 반환해 스트림 시간에 처리 시간 의미론을 부여할 수 있어요. 따라서 개발자는 비즈니스 요구에 따라 시간의 다른 개념을 적용할 수 있어요.

마지막으로 Kafka Streams 애플리케이션이 레코드를 카프카에 쓸 때마다 이 새 레코드에도 타임스탬프를 할당해요. 타임스탬프가 할당되는 방식은 맥락에 따라 달라져요:

  • 입력 레코드를 처리해 새 출력 레코드를 생성할 때, 예를 들어 process() 함수 호출에서 트리거된 context.forward(), 출력 레코드 타임스탬프는 입력 레코드 타임스탬프에서 직접 상속돼요.
  • Punctuator#punctuate() 같은 주기적 함수로 새 출력 레코드를 생성할 때, 출력 레코드 타임스탬프는 스트림 태스크의 현재 내부 시간(context.timestamp()로 획득)으로 정의돼요.
  • 집계의 경우, 결과 업데이트 레코드의 타임스탬프는 결과에 기여하는 모든 입력 레코드의 최대 타임스탬프가 돼요.

Processor API에서 #forward()를 호출할 때 출력 레코드에 타임스탬프를 명시적으로 할당해 기본 동작을 변경할 수 있어요.

집계와 조인의 경우 타임스탬프는 다음 규칙을 사용해 계산돼요:

  • 왼쪽·오른쪽 입력 레코드가 있는 조인(스트림-스트림, 테이블-테이블)에서 출력 레코드의 타임스탬프는 max(left.ts, right.ts)로 할당돼요.
  • 스트림-테이블 조인에서 출력 레코드는 스트림 레코드의 타임스탬프를 할당받아요.
  • 집계의 경우 Kafka Streams는 키별로 모든 레코드에 대한 max 타임스탬프를 (비윈도우의 경우) 전역적으로 또는 윈도우별로 계산해요.
  • 무상태(stateless) 연산의 경우 입력 레코드 타임스탬프가 그대로 전달돼요. flatMap 및 여러 레코드를 생성하는 자매 연산의 경우 모든 출력 레코드가 해당 입력 레코드의 타임스탬프를 상속해요.

스트림과 테이블의 이중성

실제로 스트림 처리 사용 사례를 구현할 때는 보통 스트림과 데이터베이스 둘 다 필요해요. 실제에서 매우 흔한 예는 들어오는 고객 거래 스트림을 데이터베이스 테이블의 최신 고객 정보로 강화(enrich)하는 전자상거래 애플리케이션이에요. 다시 말해, 스트림은 어디에나 있고 데이터베이스도 어디에나 있어요.

따라서 모든 스트림 처리 기술은 스트림과 테이블에 대한 일급 지원을 제공해야 해요. 카프카의 Streams API는 스트림과 테이블에 대한 핵심 추상화를 통해 그러한 기능을 제공해요. 흥미로운 관찰은 스트림과 테이블 사이에 실제로 밀접한 관계가 있다는 것이에요, 이른바 **스트림-테이블 이중성(stream-table duality)**이에요. 카프카는 이 이중성을 여러 가지 방식으로 활용해요: 예를 들어 애플리케이션을 탄력적으로 만들고, 장애 허용 상태 저장 처리를 지원하고, 애플리케이션의 최신 처리 결과에 대해 인터랙티브 쿼리를 실행하기 위해서요. 그리고 내부 사용을 넘어, Kafka Streams API는 개발자도 자신의 애플리케이션에서 이 이중성을 활용할 수 있게 해줘요.

집계 같은 개념을 논의하기 전에 먼저 테이블을 더 자세히 소개하고 앞서 언급한 스트림-테이블 이중성에 대해 이야기해야 해요. 본질적으로 이 이중성은 스트림을 테이블로 볼 수 있고 테이블을 스트림으로 볼 수 있다는 뜻이에요. 예를 들어 카프카의 로그 컴팩션 기능이 이 이중성을 활용해요.

테이블의 간단한 형태는 키-값 쌍의 컬렉션으로, 맵(map) 또는 연관 배열(associative array)이라고도 불러요.

스트림-테이블 이중성은 스트림과 테이블 사이의 밀접한 관계를 설명해요.

  • 테이블로서의 스트림 (Stream as Table): 스트림은 테이블의 체인지로그(changelog)로 간주될 수 있어요. 스트림의 각 데이터 레코드는 테이블의 상태 변경을 캡처해요. 따라서 스트림은 위장된 테이블이고, 체인지로그를 처음부터 끝까지 재생해 테이블을 재구성함으로써 "실제" 테이블로 쉽게 바꿀 수 있어요. 마찬가지로 더 일반적인 비유로, 스트림의 데이터 레코드를 집계하는 것 — 페이지뷰 이벤트 스트림에서 사용자별 총 페이지뷰 수를 계산하는 것처럼 — 은 테이블을 반환해요 (여기서 키와 값은 각각 사용자와 그에 해당하는 페이지뷰 수).
  • 스트림으로서의 테이블 (Table as Stream): 테이블은 어느 시점에서 스트림의 각 키에 대한 최신 값의 스냅샷으로 간주될 수 있어요 (스트림의 데이터 레코드는 키-값 쌍). 따라서 테이블은 위장된 스트림이고, 테이블의 각 키-값 항목을 반복함으로써 "실제" 스트림으로 쉽게 바꿀 수 있어요.

이것을 예로 설명해볼게요. 사용자별 총 페이지뷰 수를 추적하는 테이블을 상상해보세요 (아래 다이어그램의 첫 번째 열). 시간이 지나면서 새 페이지뷰 이벤트가 처리될 때마다 테이블의 상태가 그에 따라 업데이트돼요. 여기서 서로 다른 시점 사이의 상태 변경 — 그리고 테이블의 서로 다른 개정(revision) — 은 체인지로그 스트림(두 번째 열)으로 표현될 수 있어요.

흥미롭게도 스트림-테이블 이중성 때문에, 같은 스트림이 원래 테이블을 재구성하는 데 사용될 수 있어요 (세 번째 열).

같은 메커니즘은 예를 들어 CDC(change data capture)를 통한 데이터베이스 복제와, Kafka Streams 내부에서 장애 허용을 위해 이른바 상태 저장소(state store)를 머신 간에 복제하는 데 사용돼요. 스트림-테이블 이중성은 너무 중요한 개념이라 Kafka Streams가 KStream, KTable, GlobalKTable 인터페이스를 통해 명시적으로 모델링해요.

집계 (Aggregations)

집계 연산은 하나의 입력 스트림이나 테이블을 받아, 여러 입력 레코드를 단일 출력 레코드로 결합함으로써 새 테이블을 생성해요. 집계의 예로는 카운트나 합(sum) 계산이 있어요.

Kafka Streams DSL에서 aggregation의 입력 스트림은 KStream이나 KTable일 수 있지만, 출력 스트림은 항상 KTable이에요. 이것은 Kafka Streams가 값이 생성·발행된 후 추가 레코드가 순서가 뒤집혀 도착할 때 집계 값을 업데이트할 수 있게 해줘요. 그러한 순서가 뒤집힌 도착이 발생하면, 집계하는 KStream 또는 KTable이 새 집계 값을 발행해요. 출력이 KTable이므로, 새 값은 후속 처리 단계에서 같은 키를 가진 이전 값을 덮어쓰는 것으로 간주돼요.

윈도잉 (Windowing)

윈도잉을 사용하면 aggregationsjoins 같은 상태 저장 연산을 위해 같은 키를 가진 레코드를 소위 윈도우로 그룹화하는 방법을 제어할 수 있어요. 윈도우는 레코드 키별로 추적돼요.

Windowing operationsKafka Streams DSL에서 사용 가능해요. 윈도우로 작업할 때 윈도우에 대한 **유예 기간(grace period)**을 지정할 수 있어요. 이 유예 기간은 주어진 윈도우에 대해 Kafka Streams가 순서가 뒤집힌 데이터 레코드를 얼마나 오래 기다릴지 제어해요. 윈도우의 유예 기간이 지난 후 레코드가 도착하면, 그 레코드는 폐기되어 해당 윈도우에서 처리되지 않아요. 구체적으로, 레코드의 타임스탬프가 윈도우에 속한다고 지시하지만 현재 스트림 시간이 윈도우 끝 + 유예 기간보다 크면 레코드는 폐기돼요.

순서가 뒤집힌 레코드는 실제 세계에서 항상 가능하며 애플리케이션에서 제대로 고려되어야 해요. 순서가 뒤집힌 레코드가 어떻게 처리되는지는 유효 time semantics에 달려 있어요. 처리 시간의 경우 의미론은 "레코드가 처리되고 있는 시점"이며, 정의상 어떤 레코드도 순서가 뒤집힐 수 없으므로 순서가 뒤집힌 레코드의 개념은 적용되지 않아요. 따라서 순서가 뒤집힌 레코드는 이벤트 시간에서만 그렇게 간주될 수 있어요. 두 경우 모두 Kafka Streams는 순서가 뒤집힌 레코드를 제대로 처리할 수 있어요.

상태 (States)

일부 스트림 처리 애플리케이션은 상태가 필요 없어요. 즉 메시지의 처리가 다른 모든 메시지의 처리와 독립적이에요. 그러나 상태를 유지할 수 있으면 정교한 스트림 처리 애플리케이션에 많은 가능성이 열려요: 입력 스트림을 조인하거나 데이터 레코드를 그룹화·집계할 수 있어요. Kafka Streams DSL은 많은 그러한 상태 저장(stateful) 연산자를 제공해요.

Kafka Streams는 소위 **상태 저장소(state store)**를 제공하며, 스트림 처리 애플리케이션이 데이터를 저장하고 쿼리하는 데 사용할 수 있어요. 이것은 상태 저장 연산을 구현할 때 중요한 기능이에요. Kafka Streams의 모든 태스크는 처리를 위해 필요한 데이터를 저장·쿼리하기 위해 API로 접근할 수 있는 하나 이상의 상태 저장소를 임베드해요. 이 상태 저장소는 영속 키-값 저장소, 인메모리 해시맵, 또는 다른 편리한 데이터 구조일 수 있어요. Kafka Streams는 로컬 상태 저장소에 대한 장애 허용과 자동 복구를 제공해요.

Kafka Streams는 상태 저장소를 만든 스트림 처리 애플리케이션 외부의 메서드, 스레드, 프로세스 또는 애플리케이션에 의한 상태 저장소의 직접 읽기 전용 쿼리를 허용해요. 이것은 **인터랙티브 쿼리(Interactive Queries)**라는 기능으로 제공돼요. 모든 저장소에는 이름이 있으며 Interactive Queries는 기본 구현의 읽기 연산만 노출해요.

처리 보장 (Processing Guarantees)

스트림 처리에서 가장 자주 묻는 질문 중 하나는 "내 스트림 처리 시스템이 처리 중간에 장애가 발생해도 각 레코드가 정확히 한 번만 처리되는 것을 보장하나요?"예요. 정확히 한 번 스트림 처리를 보장하지 못하면 데이터 손실이나 중복을 용납할 수 없는 많은 애플리케이션에게는 치명적이에요. 그 경우 보통 스트림 처리 파이프라인에 추가로 배치 중심 프레임워크를 사용하는데, 이를 람다 아키텍처(Lambda Architecture)라고 해요. 0.11.0.0 이전에 카프카는 적어도 한 번(at-least-once) 전달 보장만 제공했고, 따라서 그것을 백엔드 저장소로 활용하는 모든 스트림 처리 시스템은 엔드투엔드 정확히 한 번 의미론을 보장할 수 없었어요. 사실 정확히 한 번 처리를 지원한다고 주장하는 스트림 처리 시스템조차, 카프카를 소스/싱크로 읽고 쓰는 한, 그 애플리케이션은 파이프라인 전체에서 중복이 생성되지 않을 것을 실제로 보장할 수 없었어요.

0.11.0.0 릴리스부터 카프카는 프로듀서가 트랜잭션적이고 멱등적으로(idempotent) 서로 다른 토픽 파티션에 메시지를 보낼 수 있게 지원을 추가했고, Kafka Streams는 이 기능들을 활용해 엔드투엔드 정확히 한 번 처리 의미론을 추가했어요. 더 구체적으로, 소스 카프카 토픽에서 읽은 어떤 레코드에 대해서도 그 처리 결과가 출력 카프카 토픽과 상태 저장 연산의 상태 저장소에 정확히 한 번 반영되는 것을 보장해요. Kafka Streams의 엔드투엔드 정확히 한 번 보장이 다른 스트림 처리 프레임워크가 주장하는 보장과의 핵심 차이는, Kafka Streams가 기본 카프카 저장 시스템과 긴밀하게 통합되어 입력 토픽 오프셋의 커밋, 상태 저장소의 업데이트, 출력 토픽으로의 쓰기가 부작용이 있을 수 있는 외부 시스템으로 카프카를 취급하는 대신 원자적으로 완료되도록 보장한다는 것뿐이에요. Kafka Streams 내부에서 이것이 어떻게 이루어지는지에 대한 자세한 내용은 KIP-129를 참고해요.

2.6.0 릴리스부터 Kafka Streams는 "exactly-once v2"라는 개선된 정확히 한 번 처리 구현을 지원하며, 브로커 버전 2.5.0 이상이 필요해요. 이 구현은 클라이언트 스레드와 사용된 네트워크 연결 같은 클라이언트와 브로커 리소스 활용을 줄여 더 효율적이고, 더 높은 처리량과 개선된 확장성을 가능하게 해요. 3.0.0 릴리스부터 정확히 한 번의 첫 번째 버전은 폐기되었어요. 지금부터 사용자는 정확히 한 번 처리를 위해 exactly-once v2를 사용하고, 필요하다면 브로커를 업그레이드해 준비하시길 권장해요. 브로커와 Kafka Streams 내부에서 이것이 어떻게 이루어지는지에 대한 자세한 내용은 KIP-447을 참고해요.

Kafka Streams 애플리케이션을 실행할 때 정확히 한 번 의미론을 활성화하려면 processing.guarantee 구성 값(기본값은 at_least_once)을 StreamsConfig.EXACTLY_ONCE_V2로 설정해요 (브로커 버전 2.5 이상 필요). 자세한 내용은 Kafka Streams Configs 섹션을 참고해요.

순서 뒤집힘 처리 (Out-of-Order Handling)

각 레코드가 정확히 한 번 처리될 것이라는 보장 외에도, 많은 스트림 처리 애플리케이션이 직면하는 또 다른 문제는 비즈니스 로직에 영향을 줄 수 있는 순서가 뒤집힌 데이터를 어떻게 처리할까예요. Kafka Streams에서 타임스탬프와 관련해 순서가 뒤집힌 데이터 도착을 초래할 수 있는 두 가지 원인이 있어요:

  • 토픽-파티션 내에서 레코드의 타임스탬프가 그 오프셋과 함께 단조 증가하지 않을 수 있어요. Kafka Streams는 항상 토픽-파티션 내 레코드를 오프셋 순서대로 처리하려고 하므로, 같은 토픽-파티션에서 더 작은 타임스탬프(더 큰 오프셋)를 가진 레코드보다 더 큰 타임스탬프(더 작은 오프셋)를 가진 레코드가 더 일찍 처리될 수 있어요.
  • 여러 토픽-파티션을 처리할 수 있는 스트림 태스크 내에서, 사용자가 모든 파티션에 어떤 버퍼된 데이터가 포함되기를 기다리지 않고 가장 작은 타임스탬프를 가진 파티션에서 다음 레코드를 선택하도록 애플리케이션을 구성하면, 나중에 다른 토픽-파티션에 대한 일부 레코드를 가져올 때 그 타임스탬프가 다른 토픽-파티션에서 가져온 이미 처리된 레코드보다 작을 수 있어요.

무상태(stateless) 연산의 경우 순서가 뒤집힌 데이터는 과거 처리된 레코드의 이력을 보지 않고 한 번에 하나의 레코드만 고려하므로 처리 로직에 영향을 주지 않아요. 그러나 집계와 조인 같은 상태 저장 연산의 경우 순서가 뒤집힌 데이터가 처리 로직을 잘못되게 만들 수 있어요. 사용자가 그러한 순서가 뒤집힌 데이터를 처리하려면 일반적으로 애플리케이션이 기다리는 동안 상태를 유지하며 더 오래 기다리도록 허용해야 해요, 즉 지연·비용·정확성 사이에서 트레이드오프 결정을 하는 거예요. Kafka Streams에서 구체적으로 사용자는 윈도우 집계를 위해 윈도우 연산자를 구성해 그러한 트레이드오프를 달성할 수 있어요 (자세한 내용은 개발자 가이드에서 찾을 수 있어요). 조인에 관해서는 사용자가 버전 상태 저장소(versioned state store)를 사용해 순서가 뒤집힌 데이터에 대한 우려를 해결할 수 있지만, 순서가 뒤집힌 데이터는 기본적으로 처리되지 않아요:

  • 스트림-스트림 조인의 경우 세 가지 유형(inner, outer, left) 모두 순서가 뒤집힌 레코드를 올바르게 처리해요.
  • 스트림-테이블 조인의 경우 버전 저장소를 사용하지 않으면 순서가 뒤집힌 레코드는 처리되지 않아요 (즉 Streams 애플리케이션은 순서가 뒤집힌 레코드를 확인하지 않고 모든 레코드를 오프셋 순서로 처리해요), 따라서 예측할 수 없는 결과를 만들 수 있어요. 버전 저장소를 사용하면 테이블에서 타임스탬프 기반 조회를 수행해 스트림 측 순서가 뒤집힌 데이터를 제대로 처리해요. 테이블 측 순서가 뒤집힌 데이터는 여전히 처리되지 않아요.
  • 테이블-테이블 조인의 경우 버전 저장소를 사용하지 않으면 순서가 뒤집힌 레코드는 처리되지 않아요. 그러나 조인 결과는 체인지로그 스트림이므로 결국 일관적(eventually consistent)이 돼요. 버전 저장소를 사용하면 테이블-테이블 조인 의미론이 오프셋 기반 의미론에서 타임스탬프 기반 의미론으로 바뀌고 순서가 뒤집힌 레코드가 그에 따라 처리돼요.

더 알아보기