Kafka Streams 완벽 가이드: 실시간 데이터 처리부터 상태 저장, 윈도우 집계, 조인, ksqlDB까지

1장: 실시간 데이터 분석의 핵심, Kafka Streams 완벽 이해#

1.1 들어가며: 데이터가 도착하는 순간 분석한다#

전통적인 데이터 분석은 데이터를 먼저 저장한 다음 일정한 주기로 처리하는 방식이 일반적이었습니다. 하루 동안 발생한 주문을 모아 새벽에 집계하거나, 한 시간 동안 수집한 로그를 모아서 분석하는 방식입니다.

하지만 오늘날의 서비스에서는 이런 방식만으로 충분하지 않은 경우가 많습니다.

사용자가 결제를 완료하는 순간 이상 거래 여부를 판단해야 하고, IoT 센서에서 비정상적인 온도가 감지되는 순간 관리자에게 경고를 보내야 합니다. 사용자가 상품을 여러 번 조회하면 그 행동을 바탕으로 바로 추천 상품을 변경할 수도 있습니다.

이처럼 데이터가 발생한 이후를 기다리지 않고 데이터가 들어오는 순간 처리하는 것이 스트림 처리의 핵심입니다.

Apache Kafka는 이러한 이벤트 데이터를 안정적으로 전달하고 저장하는 데 강력하지만, Kafka에 저장된 데이터를 실제 애플리케이션의 비즈니스 로직에 따라 필터링하고 변환하고 집계하려면 별도의 처리 계층이 필요합니다.

Kafka Streams는 이 문제를 해결하기 위한 Kafka의 스트림 처리 라이브러리입니다.

Kafka Streams를 사용하면 별도의 대규모 스트림 처리 클러스터를 구축하지 않고도 일반적인 Java 애플리케이션 형태로 Kafka 데이터를 실시간 처리할 수 있습니다. Kafka의 파티션과 컨슈머 그룹을 기반으로 작업을 분산하고, 로컬 상태 저장소와 상태 복구 메커니즘을 이용해 상태 기반 연산도 수행할 수 있습니다.

이번 장에서는 Kafka Streams의 기본 개념부터 실무에서 자주 사용하는 필터링, 변환, 집계, 윈도우, 조인, 상태 저장 처리까지 단계적으로 살펴보겠습니다.

마지막에는 SQL 기반 스트림 처리 도구인 ksqlDB까지 살펴보면서 코드 기반 처리와 SQL 기반 처리의 차이도 이해해 보겠습니다.

1.2 스트림 처리란 무엇인가#

1.2.1 배치 처리와 스트림 처리의 차이#

배치 처리에서는 데이터를 일정한 단위로 모은 뒤 한꺼번에 처리합니다.

예를 들어 쇼핑몰에서 하루 동안 발생한 주문 100만 건을 밤 12시에 분석하여 다음 날 판매 통계를 생성할 수 있습니다.

반면 스트림 처리는 주문 이벤트가 발생할 때마다 데이터를 계속 처리합니다.

배치 처리

이벤트 → 이벤트 → 이벤트 → 이벤트
              ↓
         데이터 저장
              ↓
        일정 시간 대기
              ↓
         일괄 처리

스트림 처리는 다음과 같은 흐름에 가깝습니다.

이벤트 → 처리 → 결과
이벤트 → 처리 → 결과
이벤트 → 처리 → 결과
이벤트 → 처리 → 결과

따라서 스트림 처리는 데이터 발생과 처리 사이의 시간을 짧게 만들 수 있습니다.

1.2.2 스트림 처리의 핵심 특징#

스트림 처리에는 다음과 같은 특징이 있습니다.

  • 실시간성: 이벤트가 발생한 직후 처리할 수 있습니다.
  • 저지연 처리: 데이터가 처리되기까지의 시간을 짧게 유지할 수 있습니다.
  • 연속 처리: 데이터가 계속 발생하는 동안 지속적으로 처리합니다.
  • 상태 저장: 이전 이벤트의 처리 결과를 저장하고 이후 이벤트와 결합할 수 있습니다.
  • 확장성: 여러 파티션과 처리 인스턴스를 이용해 작업을 분산할 수 있습니다.
  • 이벤트 기반 처리: 특정 이벤트가 발생했을 때 후속 작업을 실행할 수 있습니다.

실시간 로그 분석, 금융 거래 탐지, 추천 시스템, IoT 모니터링, 알림 시스템, 실시간 대시보드 등이 대표적인 활용 사례입니다.

1.3 Kafka Streams는 무엇인가#

Kafka Streams는 Apache Kafka에 포함된 스트림 처리 라이브러리입니다.

중요한 점은 Kafka Streams가 별도의 중앙 집중식 스트림 처리 서버를 반드시 필요로 하는 시스템이 아니라는 것입니다.

개발자는 Kafka Streams API를 사용하여 Java 애플리케이션을 만들고 이를 여러 인스턴스로 실행할 수 있습니다. Kafka Streams는 Kafka의 토픽과 파티션, 컨슈머 그룹 등의 기능을 활용하여 처리 작업을 분산합니다.

전체적인 구조는 다음과 같습니다.

Kafka Topic
    ↓
Kafka Streams Application
    ↓
Filter / Map / Join / Aggregate
    ↓
State Store
    ↓
Output Topic

예를 들어 다음과 같은 이벤트가 있다고 가정해 보겠습니다.

{
  "userId": "user-100",
  "productId": "product-200",
  "price": 120000
}

Kafka Streams는 이 데이터를 읽은 후 특정 가격 이상의 상품만 추출하거나, 사용자별 구매 횟수를 계산하거나, 다른 토픽의 사용자 정보와 결합할 수 있습니다.

1.4 Kafka Streams의 핵심 개념#

1.4.1 KStream#

KStream은 Kafka에 계속 추가되는 이벤트의 흐름을 표현합니다.

예를 들어 주문 이벤트가 발생할 때마다 새로운 레코드가 추가된다면 이를 KStream으로 표현할 수 있습니다.

주문 1
주문 2
주문 3
주문 4
...

각 이벤트는 독립적인 레코드이며 같은 키를 가진 이벤트가 여러 번 등장할 수 있습니다.

1.4.2 KTable#

KTable은 특정 키를 기준으로 현재 상태를 표현하는 데 적합합니다.

예를 들어 다음과 같은 이벤트가 순서대로 발생한다고 가정해 보겠습니다.

user-1 → 서울
user-1 → 부산
user-1 → 대전

KTable에서는 최종 상태가 다음과 같이 표현될 수 있습니다.

user-1 → 대전

즉 KStream이 사건의 연속이라면 KTable은 그 사건들을 반영한 현재 상태를 표현하는 데 적합합니다.

1.4.3 StreamsBuilder#

Kafka Streams 애플리케이션은 StreamsBuilder를 이용하여 데이터 처리 흐름을 정의할 수 있습니다.

StreamsBuilder builder = new StreamsBuilder();

KStream<String, String> stream =
    builder.stream("input-topic");

이후 filter(), mapValues(), groupByKey(), count() 등의 연산을 연결하여 처리 파이프라인을 구성합니다.

1.5 Kafka Streams를 사용하는 이유#

1.5.1 Kafka와 자연스럽게 통합된다#

Kafka Streams는 Kafka를 중심으로 설계되어 있기 때문에 Kafka 토픽과 파티션 구조를 그대로 활용할 수 있습니다.

별도의 데이터 수집 시스템과 메시지 브로커 사이에서 복잡한 연결 구조를 만들 필요가 줄어듭니다.

1.5.2 애플리케이션 형태로 배포할 수 있다#

Kafka Streams 애플리케이션은 일반적인 JVM 애플리케이션처럼 패키징하고 배포할 수 있습니다.

따라서 기존 Java 기반 백엔드 개발 환경과 결합하기 쉽습니다.

1.5.3 자동으로 작업을 분산할 수 있다#

Kafka 토픽이 여러 파티션으로 구성되어 있다면 Kafka Streams 애플리케이션의 여러 인스턴스가 파티션을 나누어 처리할 수 있습니다.

              Kafka Topic
          ┌────┬────┬────┬────┐
          │ P0 │ P1 │ P2 │ P3 │
          └─┬──┴─┬──┴─┬──┴─┬──┘
            ↓    ↓    ↓    ↓
          ┌────────┬────────┐
          │Streams │Streams │
          │   A    │   B    │
          └────────┴────────┘

데이터와 파티션의 구조를 적절하게 설계하면 처리량을 높이면서 수평 확장을 할 수 있습니다.

1.5.4 상태 저장 처리를 지원한다#

단순히 이벤트 하나를 보고 판단하는 것뿐만 아니라 이전 이벤트의 결과를 저장해 이후 이벤트 처리에 사용할 수 있습니다.

이 기능은 집계, 조인, 윈도우 처리 등에서 매우 중요합니다.

1.6 실시간 데이터 필터링#

가장 간단한 스트림 처리 중 하나는 필터링입니다.

특정 조건을 만족하는 이벤트만 다음 단계로 전달하는 방식입니다.

StreamsBuilder builder = new StreamsBuilder();

KStream<String, String> stream =
    builder.stream("input-topic");

KStream<String, String> filtered =
    stream.filter(
        (key, value) -> value.contains("ERROR")
    );

filtered.to("error-topic");

이 코드는 input-topic에서 메시지를 읽은 다음 ERROR라는 문자열이 포함된 이벤트만 error-topic으로 전달합니다.

실무에서는 다음과 같은 방식으로 사용할 수 있습니다.

전체 로그
   ↓
ERROR 로그만 필터링
   ↓
error-topic
   ↓
알림 시스템

예를 들어 애플리케이션 로그를 Kafka로 수집하고 Kafka Streams에서 오류 이벤트만 추출한 뒤 모니터링 시스템으로 전달하는 구조를 만들 수 있습니다.

1.7 실시간 데이터 변환#

필터링이 데이터를 선택하는 작업이라면 변환은 데이터 자체를 변경하는 작업입니다.

가장 간단한 예로 문자열을 대문자로 변환할 수 있습니다.

KStream<String, String> stream =
    builder.stream("input-topic");

KStream<String, String> transformed =
    stream.mapValues(String::toUpperCase);

transformed.to("output-topic");

mapValues()는 키는 유지하면서 값만 변환할 때 유용합니다.

좀 더 복잡한 데이터 구조에서는 map() 등을 사용하여 키와 값을 함께 변경할 수도 있습니다.

KStream<String, Order> orders =
    builder.stream("orders-topic");

KStream<String, OrderSummary> summaries =
    orders.mapValues(order ->
        new OrderSummary(
            order.getOrderId(),
            order.getUserId(),
            order.getAmount()
        )
    );

이처럼 Kafka Streams에서는 작은 변환 작업을 연결하여 하나의 데이터 처리 파이프라인을 만들 수 있습니다.

1.8 실시간 데이터 집계#

스트림 처리의 진정한 힘은 여러 이벤트를 하나의 의미 있는 결과로 만드는 데 있습니다.

예를 들어 사용자별 주문 횟수를 계산할 수 있습니다.

KStream<String, String> orders =
    builder.stream("orders-topic");

KTable<String, Long> counts =
    orders
        .groupByKey()
        .count();

counts
    .toStream()
    .to("order-count-topic");

여기서 중요한 개념은 키입니다.

Kafka Streams에서 같은 키를 가진 이벤트를 하나의 그룹으로 묶으려면 데이터가 올바른 키를 가지고 있어야 합니다.

예를 들어 다음과 같은 데이터가 있다면:

user-1 → order-101
user-2 → order-102
user-1 → order-103
user-1 → order-104

사용자별 집계 결과는 다음과 같이 만들 수 있습니다.

user-1 → 3
user-2 → 1

따라서 스트림 처리에서 파티션과 키 설계는 단순한 구현 세부사항이 아니라 성능과 데이터 처리 정확성에 직접 영향을 미치는 중요한 설계 요소입니다.

1.9 상태 저장 스트림 처리#

1.9.1 Stateless와 Stateful의 차이#

필터링이나 단순한 값 변환은 현재 이벤트만 보고 처리할 수 있습니다.

이를 Stateless 처리라고 합니다.

반면 다음과 같은 요구사항은 과거 데이터를 알아야 합니다.

최근 1분 동안 로그인 실패 횟수는?
최근 5분 동안 같은 카드로 몇 번 결제했는가?
오늘 이 상품은 몇 개 팔렸는가?
이 사용자가 최근에 어떤 상품을 조회했는가?

이러한 처리를 Stateful 스트림 처리라고 합니다.

Kafka Streams는 이러한 상태를 로컬 상태 저장소에 유지하고 Kafka의 내부 토픽 등을 활용하여 장애 상황에서 상태를 복구할 수 있도록 설계되어 있습니다.

1.10 윈도우 기반 집계#

스트림은 이론적으로 끝없이 이어집니다.

따라서 단순히 전체 데이터를 계속 합산하는 것이 아니라 일정한 시간 범위로 데이터를 묶어 분석해야 하는 경우가 많습니다.

이때 사용하는 것이 Windowing입니다.

예를 들어 1분 동안 발생한 주문 수를 계산할 수 있습니다.

KStream<String, String> orders =
    builder.stream("orders-topic");

KTable<Windowed<String>, Long> counts =
    orders
        .groupByKey()
        .windowedBy(
            TimeWindows.ofSizeWithNoGrace(
                Duration.ofMinutes(1)
            )
        )
        .count();

개념적으로는 다음과 같습니다.

10:00:00 ───────── 10:01:00
        1분 윈도우
              ↓
        주문 127건

10:01:00 ───────── 10:02:00
        1분 윈도우
              ↓
        주문 154건

윈도우를 이용하면 실시간 트래픽 분석, 분당 주문량, 초당 이벤트 수, 실시간 이상 탐지 등을 구현할 수 있습니다.

1.11 이벤트 시간과 윈도우를 이해해야 하는 이유#

실시간 시스템에서 단순히 애플리케이션이 메시지를 받은 시각만 사용하는 것은 충분하지 않을 수 있습니다.

예를 들어 네트워크 지연 때문에 다음과 같은 상황이 발생할 수 있습니다.

이벤트 발생
10:00:59

        ↓ 네트워크 지연

Kafka 도착
10:01:02

실제로 이벤트는 10:00:59에 발생했지만 시스템에는 10:01:02에 도착했습니다.

따라서 실시간 분석에서는 이벤트 시간, 처리 시간, 윈도우 종료 시점, 지연 이벤트 등을 함께 고려해야 합니다.

Kafka Streams의 윈도우 처리에서는 이러한 시간 개념과 Grace Period 등을 적절히 설정해야 실제 운영 환경에서 예상하지 못한 집계 결과를 줄일 수 있습니다.

1.12 스트림 조인#

실무에서는 하나의 토픽만 사용하는 경우보다 여러 이벤트를 결합해야 하는 경우가 많습니다.

예를 들어 주문 이벤트와 사용자 이벤트를 생각해 보겠습니다.

orders-topic

order-1001
user-10
상품 A
120,000원
users-topic

user-10
홍길동
서울

두 데이터를 결합하면 다음과 같은 결과를 만들 수 있습니다.

주문번호: order-1001
사용자: 홍길동
지역: 서울
상품: 상품 A
금액: 120,000원

Kafka Streams에서는 스트림과 스트림, 스트림과 테이블, 테이블과 테이블 등 여러 형태의 조인을 지원합니다.

예를 들어 두 스트림의 이벤트를 시간 범위 내에서 조인할 수 있습니다.

KStream<String, String> orders =
    builder.stream("orders-topic");

KStream<String, String> payments =
    builder.stream("payments-topic");

KStream<String, String> joined =
    orders.join(
        payments,
        (order, payment) ->
            order + " / " + payment,
        JoinWindows.ofTimeDifferenceWithNoGrace(
            Duration.ofMinutes(5)
        ),
        StreamJoined.with(
            Serdes.String(),
            Serdes.String(),
            Serdes.String()
        )
    );

joined.to("order-payment-topic");

여기서 중요한 것은 조인 키와 파티션 구조입니다.

조인에 사용되는 키가 올바르게 구성되지 않았다면 Kafka Streams는 데이터를 다시 파티션하는 repartition 과정을 수행할 수 있습니다.

따라서 실무에서는 조인 자체뿐 아니라 어떤 키를 사용하고 어떤 파티션 구조로 데이터를 배치할 것인지를 함께 설계해야 합니다.

1.13 실무 스트리밍 파이프라인 만들기#

1.13.1 실시간 이상 거래 탐지#

금융 시스템에서는 거래 이벤트를 Kafka로 전달하고 Kafka Streams에서 이상 패턴을 분석할 수 있습니다.

결제 시스템
    ↓
Kafka
    ↓
Kafka Streams
    ├─ 금액 필터
    ├─ 사용자별 집계
    ├─ 시간 윈도우
    └─ 이상 패턴 탐지
          ↓
     의심 거래 토픽
          ↓
       알림 시스템

예를 들어 다음 조건을 조합할 수 있습니다.

  • 짧은 시간에 반복되는 결제
  • 평소보다 큰 금액의 결제
  • 비정상적인 지역에서 발생한 거래
  • 동일 계정에서 동시에 발생하는 여러 거래

실제 금융 시스템에서는 이러한 규칙을 여러 데이터와 결합하여 더욱 정교하게 구현할 수 있습니다.

1.13.2 실시간 사용자 행동 분석#

웹 서비스에서는 사용자의 행동 자체가 이벤트입니다.

상품 조회
   ↓
장바구니 추가
   ↓
상품 조회
   ↓
구매

이러한 이벤트를 Kafka로 전달하고 Kafka Streams에서 사용자별 행동을 집계하면 실시간 사용자 프로필이나 추천 시스템의 입력 데이터로 활용할 수 있습니다.

사용자 이벤트
     ↓
Kafka
     ↓
Kafka Streams
     ↓
사용자별 행동 집계
     ↓
추천 시스템
     ↓
추천 결과

1.13.3 IoT 센서 데이터 처리#

IoT 환경에서는 센서 데이터가 끊임없이 발생합니다.

{
  "deviceId": "sensor-100",
  "temperature": 82.4,
  "timestamp": 1780000000
}

Kafka Streams는 이러한 데이터를 실시간으로 분석하여 정상 범위를 벗어난 센서를 찾아낼 수 있습니다.

센서
 ↓
Kafka
 ↓
Kafka Streams
 ├─ 데이터 검증
 ├─ 필터링
 ├─ 시간 윈도우
 ├─ 평균값 계산
 └─ 이상치 탐지
       ↓
    경고 이벤트

스마트 팩토리, 물류, 에너지 관리, 차량 관제 등의 시스템에서 이러한 구조를 활용할 수 있습니다.

1.14 Kafka Streams의 처리 보장#

실시간 시스템에서는 메시지를 어떻게 처리할 것인가도 중요한 문제입니다.

Kafka Streams 계열 시스템에서는 대표적으로 다음과 같은 처리 의미론을 고려할 수 있습니다.

1.14.1 At-least-once#

메시지가 유실되지 않도록 처리하지만 장애 상황에서 일부 메시지가 다시 처리될 수 있습니다.

즉 중복 처리가 발생할 가능성이 있습니다.

1.14.2 Exactly-once#

읽기, 처리, 쓰기 과정에서 중복 결과가 발생하지 않도록 처리하는 의미론입니다.

다만 Exactly-once가 모든 Kafka Streams 애플리케이션에서 자동으로 보장되는 것은 아닙니다. 애플리케이션과 Kafka의 설정 및 사용 방식에 따라 적절한 처리 보장 수준을 구성해야 합니다.

따라서 실무에서는 단순히 "Kafka Streams는 Exactly-once를 지원한다"라고 이해하기보다 어떤 처리 보장 수준을 사용할 것인지 명시적으로 설계한다는 관점이 중요합니다.

1.15 Kafka Streams와 ksqlDB#

Kafka Streams가 Java 또는 Scala 코드로 스트림 처리 로직을 작성하는 방식이라면, ksqlDB는 SQL을 사용하여 Kafka 기반 스트림 처리를 수행하는 방식입니다.

현재 ksqlDB는 Kafka Streams 위에서 동작하며 SQL 문장을 Kafka Streams 애플리케이션으로 변환하여 실행하는 구조를 사용합니다.

개념적으로 다음과 같습니다.

                 Kafka
                   │
        ┌──────────┴──────────┐
        ↓                     ↓
 Kafka Streams              ksqlDB
 Java / Scala                SQL
        ↓                     ↓
        └──────────┬──────────┘
                   ↓
             처리 결과

SQL에 익숙한 개발자나 데이터 엔지니어라면 ksqlDB를 이용하여 비교적 빠르게 스트리밍 애플리케이션을 만들 수 있습니다.

1.16 ksqlDB의 Stream과 Table#

ksqlDB에서는 Kafka의 이벤트를 Stream과 Table이라는 개념으로 표현합니다.

Stream은 변경되지 않는 이벤트가 계속 추가되는 형태이고, Table은 키를 기준으로 현재 상태를 표현하는 데 적합합니다.

예를 들어 결제 이벤트는 Stream으로 표현할 수 있습니다.

결제 발생
결제 발생
결제 발생
결제 발생
...

반면 사용자별 현재 잔액이나 현재 상태 같은 정보는 Table 형태가 자연스럽습니다.

이 구분은 Kafka Streams의 KStream과 KTable을 이해할 때도 매우 중요한 개념입니다.

1.17 ksqlDB로 스트림 만들기#

다음과 같이 Kafka 토픽을 ksqlDB의 Stream으로 등록할 수 있습니다.

CREATE STREAM user_clicks (
    user_id VARCHAR,
    url VARCHAR
) WITH (
    KAFKA_TOPIC = 'user-clicks',
    VALUE_FORMAT = 'JSON'
);

이후 SQL을 사용하여 필요한 데이터를 계속 처리할 수 있습니다.

예를 들어 특정 URL을 방문한 이벤트만 추출할 수 있습니다.

CREATE STREAM product_clicks AS
SELECT user_id, url
FROM user_clicks
WHERE url LIKE '%/products/%'
EMIT CHANGES;

CREATE STREAM AS SELECT는 기존 스트림을 기반으로 새로운 스트림을 만들고 처리 결과를 지속적으로 출력하는 방식입니다.

1.18 ksqlDB로 실시간 집계하기#

사용자별 클릭 횟수를 1분 단위로 집계하는 예제는 다음과 같이 작성할 수 있습니다.

CREATE TABLE click_counts AS
SELECT
    user_id,
    COUNT(*) AS click_count
FROM user_clicks
WINDOW TUMBLING (
    SIZE 1 MINUTE
)
GROUP BY user_id
EMIT CHANGES;

이 쿼리는 이벤트가 들어올 때마다 지속적으로 결과를 계산합니다.

즉 일반적인 SQL처럼 한 번 실행하고 끝나는 쿼리가 아니라 계속 변화하는 데이터를 대상으로 지속적으로 실행되는 스트리밍 쿼리라는 점이 중요합니다.

ksqlDB는 필터링, 변환, 집계뿐 아니라 조인, 윈도우, Materialized View, Push Query 등 다양한 스트림 처리 기능을 제공합니다.

1.19 Kafka Streams와 ksqlDB 중 무엇을 사용할까#

두 기술은 경쟁 관계라기보다 서로 다른 방식으로 Kafka 스트림을 처리하는 도구라고 이해하는 편이 좋습니다.

구분 Kafka Streams ksqlDB
기본 방식 Java / Scala SQL
개발 형태 애플리케이션 코드 스트리밍 SQL
복잡한 비즈니스 로직 적합 상대적으로 제한적
빠른 데이터 처리 구현 가능 매우 편리
사용자 정의 코드 자유도가 높음 UDF 등을 통해 확장
Kafka 통합 매우 높음 매우 높음
데이터 변환·필터링 적합 매우 편리
집계·조인 적합 SQL로 간결하게 구현
외부 라이브러리 활용 자유로움 상대적으로 제한적

Kafka Streams는 복잡한 애플리케이션 로직이나 사용자 정의 처리, 외부 서비스 연동 등이 중요한 경우 유용합니다.

반대로 데이터 변환이나 집계처럼 SQL로 자연스럽게 표현할 수 있는 문제라면 ksqlDB가 빠른 개발에 적합할 수 있습니다. 공식 문서에서도 ksqlDB는 SQL로 자연스럽게 표현되는 처리에, Kafka Streams는 복잡한 사용자 정의 로직이나 외부 서비스 연동 등에 적합한 선택지로 설명하고 있습니다.

1.20 실무에서 반드시 고려해야 할 설계 포인트#

Kafka Streams를 실제 시스템에 적용할 때는 API 사용법만 알아서는 충분하지 않습니다.

1.20.1 키 설계#

Kafka Streams의 많은 연산은 키와 파티션 구조의 영향을 받습니다.

따라서 어떤 데이터를 어떤 키로 묶을 것인지 먼저 결정해야 합니다.

1.20.2 파티션 설계#

처리량을 높이려면 적절한 파티션 수와 데이터 분산 구조가 필요합니다.

특정 키에 데이터가 지나치게 집중되면 특정 파티션에 부하가 몰리는 문제가 발생할 수 있습니다.

1.20.3 상태 저장소 관리#

집계나 조인처럼 상태가 필요한 작업에서는 상태 저장소의 크기와 복구 방식, 디스크 사용량 등을 고려해야 합니다.

1.20.4 윈도우와 지연 이벤트#

실시간 데이터에서는 이벤트가 항상 순서대로 도착한다고 가정해서는 안 됩니다.

네트워크 지연과 시스템 장애 등으로 늦게 도착하는 이벤트를 어떻게 처리할 것인지 설계해야 합니다.

1.20.5 처리 보장 수준#

모든 시스템이 Exactly-once를 필요로 하는 것은 아닙니다.

비즈니스 요구사항과 성능, 운영 복잡도를 고려하여 적절한 처리 보장 수준을 선택해야 합니다.

1.21 Kafka Streams로 만드는 실시간 데이터 파이프라인#

지금까지 배운 내용을 하나의 구조로 합쳐보겠습니다.

                  이벤트 발생
                      ↓
              ┌─────────────┐
              │    Kafka    │
              └──────┬──────┘
                     ↓
             ┌───────────────┐
             │ Kafka Streams │
             └───────┬───────┘
                     ↓
       ┌─────────────┼─────────────┐
       ↓             ↓             ↓
    Filtering    Transformation  Aggregation
       │             │             │
       └─────────────┼─────────────┘
                     ↓
               Window / Join
                     ↓
               State Store
                     ↓
              ┌──────┴──────┐
              ↓             ↓
          Kafka Topic    외부 시스템

이 구조를 이해하면 Kafka Streams의 대부분의 기능이 하나의 흐름으로 연결됩니다.

Kafka는 이벤트를 안정적으로 전달하고 저장합니다.

Kafka Streams는 그 이벤트를 실시간으로 처리합니다.

그리고 처리 결과를 다시 Kafka 토픽에 저장하거나 다른 시스템으로 전달합니다.

이 구조가 바로 Kafka 기반 이벤트 스트리밍 아키텍처의 핵심입니다.

1.22 마무리: Kafka를 메시지 브로커에서 실시간 처리 플랫폼으로 확장하기#

Kafka를 처음 배울 때는 보통 Producer와 Consumer를 이용하여 메시지를 보내고 받는 것부터 시작합니다.

하지만 Kafka의 활용 범위는 단순한 메시지 전달에 그치지 않습니다.

Kafka Streams를 이용하면 Kafka에 쌓이는 이벤트를 실시간으로 필터링하고 변환하고 집계할 수 있습니다.

여기에 상태 저장과 윈도우를 결합하면 최근 몇 분 동안의 이벤트를 분석할 수 있고, 조인을 이용하면 서로 다른 이벤트 소스의 데이터를 하나로 결합할 수 있습니다.

그리고 ksqlDB를 이용하면 이러한 스트림 처리 로직을 SQL로 표현할 수 있습니다.

결국 다음과 같은 흐름으로 발전하게 됩니다.

Kafka
 ↓
이벤트 저장
 ↓
Kafka Streams
 ↓
실시간 처리
 ↓
상태 저장 / 윈도우 / 조인
 ↓
실시간 분석
 ↓
서비스 / 알림 / 추천 / 데이터 파이프라인

Kafka Streams를 이해한다는 것은 단순히 특정 API를 암기하는 것이 아닙니다.

이벤트가 발생하고, 이동하고, 변환되고, 집계되고, 다시 새로운 이벤트로 만들어지는 전체 데이터 흐름을 이해하는 것입니다.

이 관점을 익히면 Kafka를 단순한 메시지 큐가 아니라 실시간 이벤트 기반 시스템을 구축하기 위한 핵심 플랫폼으로 바라볼 수 있게 됩니다.

1.23 자체 점검 문제#

  1. 스트림 처리와 배치 처리의 차이를 설명하십시오.
  2. Kafka Streams가 별도의 중앙 스트림 처리 클러스터 없이 애플리케이션 형태로 동작할 수 있는 이유를 설명하십시오.
  3. KStream과 KTable의 차이를 설명하고 각각 어떤 데이터에 적합한지 예를 들어보십시오.
  4. Kafka Streams에서 filter(), mapValues(), groupByKey(), count()가 각각 어떤 역할을 하는지 설명하십시오.
  5. Stateless 처리와 Stateful 처리의 차이를 설명하십시오.
  6. 윈도우 기반 집계가 필요한 이유와 실제 활용 사례를 설명하십시오.
  7. 이벤트 시간과 처리 시간의 차이가 실시간 스트림 처리에서 중요한 이유를 설명하십시오.
  8. Kafka Streams에서 조인을 수행할 때 키와 파티션 설계가 중요한 이유를 설명하십시오.
  9. At-least-once와 Exactly-once 처리 의미론의 차이를 설명하십시오.
  10. Kafka Streams와 ksqlDB의 차이를 설명하고 각각 어떤 상황에서 사용할 수 있는지 설명하십시오.
  11. 사용자 행동 데이터를 Kafka로 수집하여 실시간 추천 시스템으로 전달하는 스트리밍 파이프라인을 설계하십시오.
  12. IoT 센서 데이터를 Kafka Streams로 처리하여 이상 징후를 탐지하는 시스템의 구조를 설계하십시오.
  13. 자신이 개발하는 서비스에서 Kafka Streams를 적용할 수 있는 기능을 하나 선정하고 데이터 흐름을 설계하십시오.

이 페이지의 목차