Kafka와 Beam으로 설계하는 실시간 이벤트 처리 아키텍처

Apache Kafka 이벤트 백본과 Apache Beam 통합 처리 모델을 결합해 실시간 스트리밍 파이프라인을 설계하고 운영하는 방법을 정리합니다.

2026-08-14 · 최초 발행 2024-04-29

Kafka와 Beam이 맡는 경계

대규모 데이터 처리에서 Kafka는 이벤트를 오래 보관하고 전달하는 로그 기반 백본이며, Beam은 그 이벤트를 배치와 스트리밍 구분 없이 처리 모델로 표현한다. 둘을 함께 두면 프로듀서와 소비자 사이의 전달 계층, 그리고 시간·상태·집계를 담당하는 처리 계층을 분리할 수 있다.

Kafka는 토픽과 파티션으로 확장성을 확보한다. Replication과 ISR을 통해 내결함성을 다루고, 적어도 한 번(at-least-once) 전달을 기본으로 하며 트랜잭션을 이용한 정확히 한 번(exactly-once) 처리도 지원한다. Producer·Consumer뿐 아니라 Kafka Connect, Schema Registry를 연결하면 이벤트 중심 아키텍처의 수집과 계약 관리 범위가 넓어진다.

Beam은 PCollection을 데이터 단위로, PTransform을 연산 단위로 사용한다. ParDo/UDF, Window, Watermark, Trigger, State와 Timer를 조합해 데이터 흐름을 정의하며, Flink·Spark·Dataflow 같은 Runner 위에서 같은 파이프라인을 실행할 수 있다. Runner의 체크포인트와 저장점은 처리 정확성과 복구 전략의 기반이 된다.

이 조합의 전형적인 흐름은 프로듀서가 Kafka 토픽에 이벤트를 기록하고, Beam 스트리밍 파이프라인이 이를 읽어 Serving·OLAP·데이터 레이크 싱크로 전달하는 구조다. 이때 스키마 진화, 역직렬화 검증, event_time 기준 윈도잉, 지연 데이터 처리, DLQ를 별도 운영 요소로 봐야 한다.

이벤트 백본과 처리 파이프라인의 연결

Kafka 클러스터는 파티션을 늘려 수평 확장하고, 리플리카와 ISR·컨트롤러 기반 리더 선출로 가용성을 유지한다. 로그 보존 정책은 Time, Size, Compaction을 기준으로 정하며 저장 비용과 검색 성능의 균형에 직접 영향을 준다.

프로듀서 측에서는 idempotence, acks=all, linger.ms, batch.size를 통해 처리량과 지연을 조정한다. transactional.id는 여러 파티션에 대한 쓰기를 원자적으로 묶고, 컨슈머 그룹의 오프셋 커밋과 결합하면 EOS를 구성할 수 있다.

메시지 계약은 Schema Registry에서 관리한다. Backward, Forward, Full 호환성 규칙을 적용할 수 있으며, Kafka Connect와 Debezium을 이용하면 DB CDC 수집을 자동화하고 커넥터 기반 운영 방식을 표준화할 수 있다.

Beam에서는 고정·슬라이딩·세션 윈도우와 워터마크, 처리 시각·이벤트 시각·커스텀 트리거를 사용해 지연 데이터를 다룬다. 키별 상태와 타이머는 애그리게이션, 패턴 감지, 조인에 쓰인다. 동일한 파이프라인을 FlinkRunner, SparkRunner, DataflowRunner에서 실행할 수 있지만, 정확히 한 번 의미론은 체크포인트·저장점과 싱크의 멱등 또는 트랜잭션 특성을 함께 맞춰야 한다.

Beam Pipeline (FlinkRunner)Kafka ClusterProducersAvro/JSONvalidinvalidExactly-onceschema evolutionApp/ServiceKafka ProducerCDC (Debezium)Kafka ConnectKafka Topic: rawConsumer Group: BeamSchema RegistryDecode & Schema ValidateWindow & AggregateSideOutput: DLQEnrich/Join (State/Timer)Sink: OLAP/Serving/StorageKafka Topic: dlqKafka Topic: enrichedWarehouse/Lake:BigQuery/Iceberg

입력 단계에서는 애플리케이션 로그와 이벤트, Debezium CDC가 Kafka의 raw 토픽으로 들어온다. Beam 컨슈머는 이를 배치 단위로 폴링해 스키마를 검증하고 정제한 뒤, 윈도우·워터마크·집계를 적용하고 상태 기반 조인이나 패턴 감지를 수행한다. 결과는 enriched 토픽과 데이터 웨어하우스 또는 데이터 레이크에 커밋된다.

역직렬화 실패나 스키마 불일치는 사이드아웃을 통해 DLQ 토픽에 기록하고 재처리 파이프라인을 분리한다. 지연 데이터에는 워터마크 허용 지연과 on-time/late 재발행 트리거 정책을 적용한다. Kafka Producer의 transactional.id, enable.idempotence=true 설정과 FlinkRunner의 체크포인트, Kafka Write 연산의 2단계 커밋 정렬이 맞물리면 정확히 한 번 보장을 구성할 수 있다.

운영 중 확인할 설정과 트레이드오프

Kafka 파티션 설계에서는 키 해시 편중을 관찰해야 한다. Hot partition이 생기면 키 설계나 Custom Partitioner로 완화한다. 안정성 측면에서는 min.insync.replicas≥2, acks=all, unclean.leader.election=false 설정을 권장한다. 페이지 캐시 활용을 전제로 디스크 IOPS와 네트워크 대역폭을 우선 확보하고, 동시 스로틀링도 설정한다.

Flink 기반 Beam 파이프라인은 checkpointing interval을 5–30s 범위에서 두고, timeout은 최대 처리 지연을 고려해 정한다. 소스 지연 분포를 분석한 뒤 allowed lateness를 설정해야 하며, 허용 지연이 과도하면 메모리 압박으로 이어질 수 있다. 오토스케일링을 활용하고 키 스큐는 사전 샤딩이나 키 살팅으로 완화한다.

Kafka TLS/SASL과 RBAC/ACL을 적용하고 Schema Registry 권한은 분리한다. 스키마 호환성 규칙은 릴리즈 속도와 트레이드오프가 있다. 엄격한 규칙은 안정성은 높이고 유연성은 낮춘다. 장기 보존 데이터는 Object Storage 기반 레이크로 오프로드하고 Tiered Storage 도입도 검토할 수 있다.

지표 Kafka에서 보는 요소 Beam에서 보는 요소 운영 시점
처리량(throughput) 파티션 수, 배치/압축, acks 병렬도, 바운디드 연산 최적화 P99 지연 목표에 맞춰 배치와 병렬도를 동조
지연(latency) linger.ms, fetch.min.bytes 윈도우/트리거 정책 서브초 vs 분 단위의 지연-비용 균형 탐색
일관성(consistency) 트랜잭션, 오프셋 커밋 체크포인트, 멱등 싱크 Flink + Kafka EOS 조합 권장
안정성(reliability) 리플리케이션, ISR 상태 백엔드, savepoint 재시작·롤링업그레이드 무중단 전략
운영 편의성(operability) Connect, 모니터링 Runner 관리 메트릭·로그·트레이스 관측성 표준화

스트림 처리 패턴

클릭스트림 분석에서는 세션 윈도우로 전환 퍼널을 계산하고, 지연 이벤트를 허용한 뒤 재발행 트리거를 통해 지표 정확도를 확보한다. 결과는 Kafka에서 Druid 또는 ClickHouse로 보내 서브초 대시보드로 제공할 수 있다.

금융 이상 징후 탐지에서는 거래 스트림을 State와 Timer 기반 패턴 감지로 처리한다. 모델 추론은 사이드 입력이나 외부 서비스 호출로 연결하고, 경보는 Kafka 토픽으로 즉시 푸시해 인시던트 자동화에 연계한다.

CDC 기반 레이크 적재는 Debezium → Kafka → Beam → Iceberg/Hudi/Mergetree 경로를 사용한다. Upsert와 Compact 파이프라인을 분리하고, 스키마 진화와 타임트래블 분석을 지원한다.

IoT 처리에서는 장치 이벤트를 고정 윈도우로 집계하고 외란·결측을 보정한다. 지연·중복 데이터는 멱등적으로 처리하며, 경량 온라인 피처 스토어인 Redis 같은 저장소를 이용해 저지연으로 제공할 수 있다.

이 구성을 적용하면 하드웨어·파티션·병렬도에 따라 서브초수초 P99 지연과 수십만수백만 EPS 처리가 가능하다. Kafka 트랜잭션과 Beam 체크포인트의 조합은 EOS 의미론을 구현하고 재처리 비용을 줄인다. Connect와 표준 파이프라인은 신규 소스·싱크 도입 리드타임을 수일→수시간 수준으로 단축하며, 장기 보존을 오브젝트 스토리지로 전환하면 워크로드에 따라 스토리지 비용 30%+ 절감을 기대할 수 있다.

Python Beam 파이프라인 예시

Kafka 3.x 클러스터에는 raw, enriched, dlq 토픽이 있어야 하며 Schema Registry는 선택 사항이다. Flink 1.16+ 클러스터에서는 Checkpointing을 활성화한다. Apache Beam Python 2.5x 이상과 confluent-kafka가 필요하며, 실행 기준 Runner는 FlinkRunner다.

pip install "apache-beam[gcp,aws]==2.56.0" confluent-kafka fastavro
# 버전은 환경에 맞게 조정, 최신 정보 확인 필요
# python >= 3.9
import json
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.io.kafka import ReadFromKafka, WriteToKafka
from apache_beam.transforms.window import FixedWindows
from apache_beam.transforms.trigger import AfterWatermark, AccumulationMode

bootstrap_servers = "kafka-broker:9092"
raw_topic = "raw"
enriched_topic = "enriched"
dlq_topic = "dlq"

class ParseAndValidate(beam.DoFn):
    def process(self, kv):
        key, value = kv
        try:
            record = json.loads(value.decode("utf-8"))
            # 필수 필드 검증
            if "event_time" not in record or "user_id" not in record:
                raise ValueError("missing field")
            yield beam.pvalue.TaggedOutput("valid", (key, record))
        except Exception as e:
            err = {"error": str(e), "raw": value.decode("utf-8")}
            yield beam.pvalue.TaggedOutput("invalid", (key, json.dumps(err).encode("utf-8")))

class Enrich(beam.DoFn):
    def process(self, kv):
        key, rec = kv
        rec["enriched"] = True
        yield (key, json.dumps(rec).encode("utf-8"))

def to_kv_for_kafka(element):
    key, val = element
    if isinstance(val, dict):
        val = json.dumps(val).encode("utf-8")
    return (key, val if isinstance(val, (bytes, bytearray)) else str(val).encode("utf-8"))

options = PipelineOptions([
    "--runner=FlinkRunner",
    "--flink_master=localhost:8081",
    "--checkpointing_interval=10000",  # 10s
    "--max_bundle_size=5000",
    "--max_bundle_time_mills=1000",
])

consumer_conf = {"bootstrap.servers": bootstrap_servers, "group.id": "beam-cg", "auto.offset.reset": "earliest"}
producer_conf = {
    "bootstrap.servers": bootstrap_servers,
    "enable.idempotence": "true",
    "acks": "all",
    "linger.ms": "10",
    "batch.num.messages": "10000",
    # EOS는 Flink 체크포인트와 결합 시 활성화
    "transaction.timeout.ms": "600000"
}

with beam.Pipeline(options=options) as p:
    msgs = (p
        | "ReadKafka" >> ReadFromKafka(
            consumer_config=consumer_conf,
            topics=[raw_topic],
            with_metadata=False)
    )

    parsed = (msgs
        | "ParseValidate" >> beam.ParDo(ParseAndValidate()).with_outputs("valid", "invalid")
    )

    valid = parsed.valid | "KeyByUser" >> beam.Map(lambda kv: (kv[1]["user_id"], kv[1]))

    windowed = (valid
        | "FixedWindow1m" >> beam.WindowInto(
            FixedWindows(60),
            trigger=AfterWatermark(late=AfterWatermark.late),
            accumulation_mode=AccumulationMode.DISCARDING)
    )

    enriched = (windowed
        | "Enrich" >> beam.ParDo(Enrich())
    )

    # 결과 쓰기: Exactly-once는 FlinkRunner + 체크포인트 + WriteToKafka 조합에서 달성
    _ = (enriched
        | "WriteEnriched" >> WriteToKafka(
            producer_config=producer_conf,
            topic=enriched_topic)
    )

    dlq = parsed.invalid | "EnsureKV" >> beam.Map(lambda kv: kv)
    _ = (dlq
        | "WriteDLQ" >> WriteToKafka(
            producer_config=producer_conf,
            topic=dlq_topic)
    )
python pipeline.py \
  --runner=FlinkRunner \
  --flink_master=localhost:8081 \
  --checkpointing_interval=10000

Exactly-once 보장은 Flink 체크포인트의 일관성과 Kafka 트랜잭션 또는 멱등 쓰기를 함께 사용해야 한다. BigQuery, S3, Iceberg 같은 싱크는 멱등 업서트와 2단계 커밋 지원 여부를 확인한다. 핵심 유스케이스부터 표준 파이프라인으로 정착시킨 뒤 커넥터와 스키마 거버넌스 체계를 확립하고, 비용과 성능을 반복적으로 튜닝하는 방식이 적합하다.

Apache KafkaApache Beam이벤트 스트리밍실시간 처리빅데이터