배치는 스트리밍의 특별한 경우입니다. 비즈니스가 몇 시간 대신 몇 초 내에 반응해야 할 때, 지속적인 데이터 흐름을 위한 아키텍처가 필요합니다.

대시보드는 누군가가 볼 때쯤이면 이미 오래되었습니다. 사기 탐지는 야간 배치 작업으로 실행되어 다음 날 아침에 사기를 잡습니다. 재고 수량은 시간 단위로 업데이트되어 과매출을 초래합니다. 센서 데이터는 수집되지만 야간 ETL에서 분석될 때까지 조치가 취해지지 않습니다. 데이터가 소스에서 처리 및 소비자에게 초단위 지연으로 지속적으로 흐르는 시스템이 필요합니다 — 실시간 분석, 라이브 알림, 스트리밍 AI 추론, 시스템 간 즉각적인 동기화.
실시간 스트리밍 아키텍처는 데이터를 개별 배치가 아닌 연속적이고 무한한 흐름으로 처리합니다. 이벤트 생성자는 스트리밍 플랫폼(Kafka, Kinesis, Pulsar)에 게시합니다. 스트림 프로세서는(Flink, Kafka Streams, 커스텀 소비자) 이벤트를 실시간으로 변환, 풍부화, 필터링, 집계합니다. 처리된 결과는 소비자에게 푸시됩니다: 실시간 대시보드(WebSocket), 검색 인덱스(Elasticsearch), 분석 데이터베이스(ClickHouse), 다운스트림 서비스. Change Data Capture (CDC)는 기존 데이터베이스가 애플리케이션 변경 없이 이벤트 소스로 참여할 수 있게 합니다.
Explore more design patterns and system architectures
MicrocosmWorks는 multi-consumer replay, 장기 retention periods, cross-cloud portability가 필요한 팀에게 Kafka를 추천합니다. Kafka의 log-based architecture는 무제한의 consumer groups가 동일한 데이터 스트림을 독립적으로 다시 읽는 것을 지원하기 때문입니다. Kinesis는 AWS ecosystem에 긴밀하게 통합된 fully managed service를 원하고 데이터 retention 요구 사항이 7일 미만이며 consumer applications가 10개 미만일 때 더 나은 선택입니다. 저희는 올바른 권장 사항을 제공하기 위해 architecture assessment 중에 고객의 특정 요구 사항—throughput, retention, consumer patterns, operational maturity—을 평가합니다.
MicrocosmWorks는 멱등성(idempotent) 프로듀서, 트랜잭션(transactional) 소비자, 그리고 Redis와 같은 빠른 조회 캐시에 저장된 이벤트 지문(event fingerprints)을 사용하는 중복 제거(deduplication) 계층의 조합을 통해 exactly-once 시맨틱을 구현합니다. Kafka 기반 시스템의 경우, 우리는 소비자 오프셋(consumer offsets)과 프로듀서 쓰기(producer writes)를 원자적으로 커밋하는 Kafka의 내장 transactional API를 활용하며, 반면 사용자 정의 스트리밍 파이프라인의 경우 소비자 측에서 중복 제거(deduplication)를 통해 outbox pattern을 구현합니다. 우리는 안전망(safety net)으로서 소비자를 항상 멱등성(idempotent)을 가지도록 설계하여, exactly-once 메커니즘에 엣지 케이스(edge-case) 실패가 발생하더라도 이벤트 재처리(reprocessing)는 동일한 결과를 생성하도록 합니다.
MicrocosmWorks는 일반적으로 ingestion, processing, 그리고 sink writing을 포함하는 스트리밍 파이프라인에 대해 50-200ms의 종단 간(end-to-end) 지연 시간(레이턴시)을 제공하며, Apache Flink 또는 Kafka Streams와 같은 인메모리 스트림 프로세서를 사용하는 더 간단한 passthrough 또는 filtering 워크로드의 경우 10ms 미만도 달성 가능합니다. 가장 큰 지연 시간(레이턴시) 발생 요인은 일반적으로 네트워크 홉, 직렬화 오버헤드, 그리고 싱크 쓰기 배치(batching)이며, 이는 고객의 지연 시간(latency)과 처리량(throughput) 간의 트레이드오프 선호도에 따라 조정합니다. 아키텍처 설계 시, 저희는 파이프라인 단계별로 명시적인 지연 시간(latency) SLO를 설정하고, 프로덕션 환경에서 p50, p95, p99 지연 시간(레이턴시)을 추적하는 모니터링 대시보드를 구축합니다.
MicrocosmWorks는 역방향 및 전방향 호환성 규칙을 적용하는 스키마 레지스트리(일반적으로 Confluent Schema Registry 또는 AWS Glue Schema Registry)를 구현하여, 생산자가 기존 소비자를 손상시키지 않고 데이터 형식을 발전시킬 수 있도록 보장합니다. 저희는 명시적인 스키마 버전 관리를 사용하는 Avro 또는 Protobuf 직렬화를 이용하여, 모든 메시지가 자체 설명적이며 메시지가 생성된 이후 스키마가 변경되었더라도 역직렬화될 수 있도록 합니다. 저희 CI/CD 파이프라인에는 제안된 스키마 변경이 다운스트림 소비자를 손상시킬 경우 배포를 차단하는 자동화된 스키마 호환성 검사가 포함됩니다.
MicrocosmWorks는 프로덕션 스트리밍 플랫폼을 안정적으로 유지 관리하기 위해 분산 시스템, 스트림 처리 프레임워크, 인프라 자동화 경험을 가진 최소 2~3명의 엔지니어를 권장합니다. 이러한 전문 지식을 사내에서 구축하기를 원치 않는 기업을 위해, 저희는 시간당 $15~$40의 비용으로 관리형 스트리밍 플랫폼 지원을 제공합니다. 이를 통해 저희 팀이 클러스터 운영, 성능 튜닝 및 사고 대응을 처리하는 동안, 고객사의 개발자는 스트림 처리 애플리케이션 구축에 집중할 수 있습니다. 저희는 또한 4-8주간의 참여를 통해 기존 엔지니어링 팀의 Kafka, Flink 또는 Kinesis 운영 역량을 향상시키는 교육 프로그램을 제공합니다.
아키텍처는 네 개의 계층으로 구성됩니다. 이벤트 소스는 데이터를 생성합니다 — 애플리케이션 이벤트, 데이터베이스 CDC 스트림, IoT 원격 측정, 사용자 클릭스트림, 외부 API 웹훅. 스트리밍 플랫폼(Kafka)은 내구성 있고, 순서가 있으며, 재생 가능한 이벤트 저장소를 제공합니다. 스트림 프로세서는 토픽에서 소비하고, 변환(필터링, 풍부화, 윈도우 집계, 조인)을 적용하며, 출력 토픽 또는 싱크로 생성합니다. 소비자는 처리된 스트림에 구독합니다 — WebSocket 서버는 브라우저로 푸시하고, 커넥터는 데이터베이스로 싱크하며, 경고 엔진은 규칙을 평가하고 알림을 발송합니다.
| 계층 | 기술 |
|---|---|
| 스트리밍 | Apache Kafka (MSK, Confluent), Kinesis, Apache Pulsar, Redpanda |
| CDC | Debezium, AWS DMS, Maxwell |
| 처리 | Apache Flink, Kafka Streams, Benthos, 커스텀 소비자 |
| 실시간 전달 | WebSocket (Socket.io), SSE, GraphQL Subscriptions |
| 분석 | ClickHouse, Apache Druid, Elasticsearch, TimescaleDB |
| 관측성 | Kafka 지연 모니터링 (Burrow), Flink 메트릭, 커스텀 지연 추적 |
| 사용 시기 | 피해야 할 시기 |
|---|---|
| 비즈니스 결정이 초단위 데이터 신선도를 필요로 할 때(사기, 모니터링, 거래) | 시간별/일별 신선도가 비즈니스 요구를 충족하는 배치 처리 |
| 여러 소비자가 동일한 이벤트 스트림을 필요로 할 때(팬아웃, 분리된 시스템) | 단일 프로듀서와 단일 소비자가 있는 경우 — 간단한 큐로 충분합니다 |
| 디버깅, 재처리 또는 새로운 소비자 구축을 위한 이벤트 재생이 필요할 때 | 데이터 볼륨이 낮고(< 1K 이벤트/분) 스트리밍 인프라를 정당화하지 않을 때 |
| 기존 데이터베이스를 코드 변경 없이 다운스트림 시스템과 동기화하기 위해 CDC가 필요할 때 | 팀이 분산 시스템에 대한 경험이 부족할 때 — 스트리밍은 상당한 운영 복잡성을 추가합니다 |
MW는 "재생 원칙"으로 스트리밍 시스템을 설계합니다 — 모든 스트림은 특정 시점부터 재생 가능해야 하며, 새로운 소비자가 과거 데이터를 백필할 수 있고 기존 소비자가 버그 수정 후 재처리할 수 있게 합니다. 우리의 Kafka 배포에는 스키마 진화 정책(기본적으로 하위 호환), 소비자 지연 경고(비즈니스에 가시적인 지연이 되기 전), 자동 재시도를 위한 데드레터 토픽이 포함됩니다. 우리는 비디오 분석, IoT 원격 측정, 실시간 대시보드를 위해 초당 500K+ 이벤트를 처리하는 스트리밍 파이프라인을 구축했습니다.
하나의 코드베이스, 수백 개의 테넌트, 데이터 유출 제로 — 모든 확장 가능한 SaaS 비즈니스의 기반입니다.