Flink SQL 도입과 운영 경험
CPU 96개짜리 레거시 Flink 앱을 Flink SQL과 GitOps 배포로 대체했다.
Azar Matching Dev Team은 CPU 96개를 쓰던 레거시 Flink 앱의 복잡도와 운영 부담을 줄이기 위해 Flink SQL 도입을 검토했다. 핵심 매치 이벤트 조인은 이미 별도 Flink 앱으로 분리해 둔 상태였고, 그 뒤의 조건부 이벤트 발행과 Redis 플래그 저장을 SQL 기반으로 옮기는 방향이 가장 현실적이었다.
대안으로는 하나의 Flink App으로 다시 합치기, 여러 앱으로 쪼개기, Flink SQL을 쓰기 등이 있었고, 팀은 생산성과 운영 효율을 이유로 Flink SQL을 선택했다. 선택 근거는 크게 세 가지였다.
- High Availability: Checkpoint와 Savepoint, JobManager HA, TaskManager 재분배로 장애 대응이 가능했다.
- 고급 스트리밍 기능:
JOIN,UNION, window, event time, watermark를 SQL로 다룰 수 있었다. - 확장성: UDF와 Custom Connector로 Redis Connector를 직접 만들고, 빌트인으로 부족한 ARRAY 교집합도 보완했다.
비교 대상이었던 ksqlDB는 stateful 처리에서 failover 시 changelog replay 부담이 크고, 복제본에서도 동일 연산을 수행해 리소스가 두 배로 들 수 있다는 점이 아쉬웠다. Spark Structured Streaming은 에코시스템과 확장성은 좋았지만 마이크로 배치 특성상 레코드 단위 지연이 있고, 팀 내 Spark 경험이 거의 없어 Custom Sink가 필요한 상황에서 선택하기 어려웠다.
운영 환경은 Kubernetes 위에 Session mode Flink Cluster를 올리는 방식으로 구성했다. 기존 Application mode와 달리 이미 떠 있는 클러스터에 job을 제출하는 구조라, high-availability.type: kubernetes, high-availability.storageDir: s3://..., kubernetes.service-account 같은 설정을 넣고, JobManager는 jobmanager.rpc.address=$(POD_IP)로 서로 다른 주소를 쓰게 해 HA를 맞췄다.
쿼리 배포는 전용 UI 대신 GitHub Actions와 Flink SQL Gateway API를 조합한 GitOps 방식으로 처리했다. 레포지토리에서 job별 폴더와 SQL 파일을 관리하고, Python으로 작성한 액션이 해당 SQL을 읽어 Gateway REST API를 호출하는 구조라 구현과 테스트가 단순했다.
운영하면서는 몇 가지 대표적인 장애 패턴도 정리했다.
- JobManager / TaskManager fail: TaskManager는 Kubernetes QoS 정책으로 재시작되는 경우가 있었지만, 다른 TaskManager로 재분배되어 작업이 이어졌다.
- 잘못된 데이터 인입:
json.ignore-parse-errors로 JSON 파싱 오류를 무시하고,JSON_VALUE는DEFAULT ... ON ERROR로 방어했다. - 자원 부족: TaskManager CPU 100% 초과나 메모리 부족 시 리소스와 parallelism을 조정해 대응했다.
- 클러스터 재시작 후 일부 job fail: timeout과 retry 설정이 짧아 재시도가 너무 빨리 끝나는 문제를 수정했다.
- 쿼리 조건 변경: savepoint 복원은 단순 조건 변경엔 유효하지만, window 조건이 바뀌면 state 호환성이 깨져 직접 앱이 더 나을 수 있다.
모니터링은 Flink metric으로 구성했다. numRunningJobs로 비정상 종료를 빠르게 감지하고, taskmanager.cpu.load, taskmanager.memory.used, busyTimeMsPerSecond, Kafka의 records-lag-max 같은 지표로 부하와 지연을 추적했다.
부록에는 Kafka에서 로그인 이벤트를 받아 10초마다 지난 1분간의 로그인 수를 집계해 다시 Kafka로 내보내는 HOP window 예제가 포함돼 있다. 이 예시는 WATERMARK, HOP_START, GROUP BY HOP(...)를 사용해 Flink SQL만으로 상태를 가진 스트리밍 앱을 작성할 수 있음을 보여준다.
이 요약은 원문 이해를 돕기 위한 큐레이션입니다. 저작권은 원저작자에게 있으며, 정확한 내용과 맥락은 원문을 확인하세요.
요약 오류, 출처 표기 문제, 삭제 요청은 문의 · 건의로 알려주세요.