AI Briefing

AI 실시간 추천 시스템을 위한 Flink 기반 스트림 조인 서비스 구축기

·2025.06.11 09:00

Flink KeyedProcessFunction으로 실시간 이벤트 조합과 무중단 배포를 구현했다.

Azar의 실시간 추천 시스템은 서로 다른 시점에 생성된 유저 이벤트를 조합해 최신 정보를 반영해야 했다. 이벤트가 일부 유실되거나 조합 규칙이 복잡해지면서, 단순한 파이프라인이 아니라 스트림 조인 서비스가 필요해졌다.

플랫폼은 Spark Streaming, Kafka Streams, Apache Flink를 비교해 Flink로 결정했다. 낮은 지연, 정밀한 시간 제어, 상태 관리, Exactly Once 처리 지원이 핵심 이유였고, 그중에서도 복잡한 시간 제어가 가능한 KeyedProcessFunction을 선택했다.

KeyedProcessFunction에서는 키별로 상태를 관리하고 TimerService로 타이머를 등록·연장·취소하며 조합 결과를 판단했다. 이벤트가 모두 도착하면 즉시 발행하고, 일부만 도착한 경우에는 타이머 만료 시 부분 발행, 필수 이벤트가 없으면 누락 처리하는 방식으로 동작했다. 여기에 heartbeat 이벤트를 더해 필요한 이벤트를 조금 더 기다리면서 조합 성공률을 높였다.

배포 측면에서는 Savepoint로 상태와 Kafka 오프셋, 타이머를 안전하게 복구했고, 복구 직후 만료된 타이머 문제는 CheckpointedFunction과 복구 시점 타임스탬프를 이용해 보완했다. 무중단 배포는 Blue-Green 전략과 Spinnaker Pipeline으로 자동화해 해결했다.

중복 전달 문제는 처음에 Kafka Sink의 Exactly Once, 별도 dedup Flink 애플리케이션까지 검토했지만, 최종적으로 Redis 기반 중복 제거로 정리했다. SET NX 또는 Lua script를 활용해 두 개의 독립된 서비스에서도 동시성 문제 없이 중복을 막았고, 중복 제거 지연은 평균 300ms에서 3ms 미만으로 크게 줄었다.

이 요약은 원문 이해를 돕기 위한 큐레이션입니다. 저작권은 원저작자에게 있으며, 정확한 내용과 맥락은 원문을 확인하세요.

요약 오류, 출처 표기 문제, 삭제 요청은 문의 · 건의로 알려주세요.