
백엔드
Kafka에서 S3로 실시간 데이터 수집 파이프라인 설계와 구축기
두줄요약
Kafka 소비 결과를 Parquet으로 변환해 S3에 적재하는 실시간 수집 파이프라인을 설계하고 구축했습니다. 또한 Flush, 커밋, 모니터링 체계를 통해 누락 없이 안정적으로 운영하는 방법을 정리했습니다.
문제 상황
- 배치 중심 데이터 수집 구조에서 실시간 이벤트 수집 필요성 대두
- MariaDB Trigger 기반 로그 적재의 스키마 변경 부담과 RDB 저장 한계
- 광고·주문 등 이벤트 데이터를 장기 보관하고 분석할 수 있는 수집 채널 필요
원인 분석
- 컬럼 추가마다 Trigger와 로그 테이블을 함께 수정해야 하는 운영 복잡도
- 데이터 증가에 따른 스토리지·인덱스 관리 비용 상승
- S3 Sink Connector의 JSON 중심 제약과 지연, 비즈니스 로직 반영 한계
해결 방법
- Kafka Consumer를 직접 개발해 Parquet 변환 후 S3 적재
- at-least-once 전략과 수동 커밋으로 메시지 누락 방지
- 시간·버퍼 크기 기준 Flush, 종료·리밸런스 시 Flush & Commit 처리
- 레코드/배치 단위 실패 처리와 Slack·OpenSearch·Airflow 기반 모니터링 구성
성능/운영 포인트
- 고빈도·저빈도 토픽을 모두 지연 없이 반영하는 Flush 정책
- 리밸런스 시 커밋 누락으로 인한 중복 처리 방지
- 메타데이터 저장과 일일 DAG 점검으로 누락 구간 추적
적용해볼 점
- CDC 기반 로그 수집으로 DB 의존도 완화
- 원형 데이터 S3 저장 후 Athena로 조회하는 레이크하우스 패턴
- 향후 Iceberg 도입으로 트랜잭션과 스키마 진화 보완
