이벤트 저장소를 Athena에서 Snowflake로 이전합니다
에서 Singular, 우리는 전 세계 수백만 모바일 기기에서 광고 노출, 클릭 및 앱 설치 데이터를 수집하는 파이프라인을 보유하고 있습니다. 이 방대한 데이터는 시간별·일별로 집계됩니다. 다양한 마케팅 지표로 강화해 고객이 캠페인’ 성과를 분석하고 ROI를 확인할 수 있게 제공합니다.
결론적으로 우리는 초당 수만 건의 이벤트를 수신하고 매일 수십 테라바이트의 데이터를 처리하며, 수 페타바이트에 달하는 데이터 세트를 관리합니다.
이 방대하고 복잡한 데이터 웨어하우스를 Athena에서 Snowflake로 마이그레이션하는 것은 매우 복잡한 과정이었습니다. 이 글에서는 우리가 왜 이러한 어려운 단계를 밟기로 결정했는지, 어떻게 진행했는지, 그리고 그 과정에서 얻은 교훈들을 살펴보겠습니다.
눈송이 vs. 아테나
저희는 2018년부터 사용자 수준 이벤트 저장소로 Athena를 사용하기 시작했습니다. 컴퓨팅과 스토리지를 분리하고, S3와 쉽게 통합되며, 설치에 노력이 거의 들지 않아 훌륭한 솔루션처럼 보였습니다. 마치 마법 같았죠.
우리는 Kinesis 스트림을 통해 사용자 이벤트를 스트리밍하고 데이터를 고객 및 날짜별로 파티션하여 S3에 업로드했습니다. 파일은 S3에 저장할 때 Athena’s 모범 사례 가능한 한 많이.
하지만 출시 후 1년이 지나면서 문제가 발생하기 시작했고, 결국 다른 해결책을 모색하게 되었습니다.
수요에 따른 컴퓨팅 자원 부족
첫 번째 고통스러운 문제는 필요했던 규모 와 즉각적인 고성능 컴퓨팅. Athena는 멀티테넌트 서비스라 쿼리 컴퓨팅 파워를 제어할 수 없습니다. 프로덕션 파이프라인이 시간당 약 4,000개 쿼리와 고객당 하루 한 번의 대용량 쿼리를 실행하면서, 컴퓨팅 자원 부족으로 피크 시간에 장애가 발생했습니다.
부실한 자원 관리
컴퓨팅 리소스를 사용량에 따라 분할하고 제어하는 것은 불가능합니다. 예를 들어, 일부 데이터 세트는 다른 데이터 세트보다 더 빠르게 쿼리해야 하는 경우(시간별 쿼리 vs. 대규모 일일 쿼리)가 있습니다. 고객에게 성능을 보장하고 싶었지만, 컴퓨팅 리소스를 제어할 방법이 없었습니다.
기본 제공 클러스터링 기능 부족
고객별 데이터를 S3에 업로드하기 전에 일정 시간 동안(또는 데이터 크기에 따라) 버퍼링해야 했습니다. 버퍼링은 스트리밍 파이프라인에 적합한 방식이 아니며, 데이터 처리량이 많을 때는 성능 저하를 초래했습니다. 또한 버퍼링은 집계 보고서의 데이터 업데이트 주기에도 영향을 미쳤습니다.
특정 기록을 삭제하고 업데이트하는 번거로운 작업
EU 개인정보보호법(GDPR)을 준수하기 위해 최종 사용자의 요청이 있을 경우 특정 기록 삭제를 지원해야 할 의무가 있습니다. Athena에서 파일 내 특정 기록을 삭제하는 과정은 길고 복잡합니다.
스노우플레이크를 선택하는 이유는 무엇인가요?
우리가 마이그레이션을 결정한 이유는 Snowflake가 설계상 이러한 핵심 기능을 지원하기 때문입니다
- 쿼리 컴퓨팅 파워를 제어할 수 있습니다 warehouse’s 사용을 올바르게 계획하세요. 다양한 인스턴스로 서로 다른 쿼리를 별도 창고에 분산시켜, 필요한 곳에 더 많은 리소스를 제공할 수 있습니다.
- Snowflake’s 마이크로 파티션에 데이터를 클러스터링하세요 클러스터 키를 정의하고 데이터를 인제스트할 때 해당 키에 따라 정렬하세요. 또한 사용할 수 있습니다 자동 클러스터링 클러스터 키에 따라 마이크로 파티션을 최적화하세요. 쿼리 사용량에 맞게 데이터를 현명하게 클러스터링하면 쿼리를 크게 최적화할 수 있습니다.
- Snowflake에서 삭제는 간단합니다, 주로 클러스터 키를 기반으로 할 때 그렇습니다’ 특정 레코드 업데이트는 여전히 비용이 많이 들지만 Athena보다 훨씬 효율적이며, 컴퓨팅 파워를 분리하고 필요에 따라 확장할 수 있습니다.
마이그레이션 전 단계: 개념 증명
Snowflake를 사용한 POC를 계획할 때, 우리 시스템에 필요한 모든 유형의 사용 사례를 테스트하고 싶었습니다. 그래서 POC를 읽기, 쓰기, 삭제의 세 가지 주요 범주로 나누었습니다.
설정 요구 사항에는 모든 고객의 하루치 데이터(약 100TB), 멀티테넌트 테이블, 인덱싱(클러스터 키 없음, 클러스터 키 1개, 클러스터 키 2개)이 포함되었습니다.
우리는 네 가지 성공 기준을 세웠습니다
- 데이터 로딩 시간 준수: 데이터 수집 파이프라인의 개념 증명(POC)을 통해 빠르고 저렴한 처리 속도를 입증해야 했습니다.
- 기존 쿼리 지원: 현재 사용 중인 쿼리를 Snowflake에서 실행할 수 있음을 입증해야 했습니다.
- 회의 쿼리 SLA: 쿼리 소요 시간을 테스트한 결과, Snowflake로 마이그레이션하는 것만이 개선 방법임을 확인했습니다.
- 비용 요구사항: Snowflake에서 프로덕션 로드를 실행하면 Athena보다 비용이 절감되는지 확인하고 싶었습니다(이는 주요 프로젝트 지표는 아니지만’ 비용 상승을 방지하는 것이 중요했습니다).
저희의 Snowflake POC는 모든 면에서 성공적이었습니다. 또한 Snowflake 팀은 매우 협조적이고 도움이 많이 되었습니다.
견고한 데이터 수집 파이프라인 구축
요구 사항
우리는 다음과 같은 모든 요구 사항을 충족하는 Snowflake 데이터 수집 메커니즘을 구축해야 한다는 것을 알고 있었습니다
- 견고하고 내구성이 뛰어나야 합니다. Snowflake로 수집되는 데이터 양은 단기간에 급격히 증가할 수 있으므로, 데이터 수집 파이프라인도 그에 맞춰 확장되어야 합니다. 고객사의 트래픽이 갑자기 급증하는 경우는 계획된 이벤트와 계획되지 않은 이벤트 모두에서 발생할 수 있습니다. 예를 들어, 중요한 뉴스가 보도되면 소셜 미디어 활동이 급증할 수 있습니다. 이와 같은 갑작스러운 트래픽 과부하는 과거에도 발생한 적이 있습니다. 따라서 전체 데이터 양이 갑자기 두 배로 늘어나는 상황을 예상하여 데이터 수집 파이프라인을 계획하고 구축해야 합니다. 손쉽게 확장할 수 있어야 합니다.
- 빠르게. 로드 프로세스는 SLA를 만족해야 합니다. 고객은 몇 분’ 지연으로 사용자 데이터를 받아야 합니다. Snowflake가 데이터 가용성을 그 임계값 이상 지연하지 않도록 해야 합니다.
- 가용성이 매우 높아야 합니다. 시스템 다운은 절대 용납할 수 없습니다. 새로운 데이터 전송이 한 시간 이상 지연되면 고객들이 이를 알아차리기 시작합니다.
- 가장 흔한 쿼리 사용에 최적화됨. 우리는 Snowflake를 시간당 4,000회 이상 쿼리해 집계 보고서를 최신으로 유지합니다. 최적화된 쿼리 계획을 위해 데이터 정렬·클러스터링을 최대한 준비해야 합니다.
- 데이터 보존·휴식 시 분리 등 모든 개인정보 요구사항을 준수합니다. 파트너마다 사용자 데이터 보존 규칙이 다르고, 마케팅 데이터를 완전 분리해야 할 경우도 있습니다. 우리는 이를 지원합니다. Athena에서는 S3 파일을 나눠 처리했지만 파일 크기 감소로 Athena’s 성능이 저하됐습니다. Snowflake에서는 파트너별 테이블로 분리·보존 규칙을 정의할 수 있습니다.
우리는 Snowflake 데이터 수집 파이프라인을 최대한 간소하고 자체 복구 기능이 있도록 구축하기로 결정했습니다
설계
저희 디자인을 좀 더 자세히 살펴보겠습니다
우리는 Singular스택(1)에서 Kinesis 스트림으로 사용자 수준 이벤트를 스트리밍합니다.
이벤트(예: 광고 클릭, 광고 조회, 앱 설치 또는 인앱 이벤트)는 스트림에서 일괄 처리 방식으로 읽어 .csv 파일(zstd로 압축됨)에 저장하고 Kubernetes에서 관리하는 Python 워커(2)를 사용하여 S3 파일에 업로드합니다.
모든 파일 생성은 Snowpipe(4)가 구독하는 SNS 주제(3)에 알림을 생성합니다.
Snowpipe는 파일이 버킷에 추가될 때마다 파일의 내용을 버퍼 테이블(5)로 삽입하고, COPY INTO 명령을 실행합니다.
버퍼 테이블에 저장된 데이터는 클러스터링되거나 정렬되어 있지 않으므로 집계 쿼리에 사용하면 효율적이지 않습니다.
데이터가 쿼리에 맞게 구성되고 최적화되도록 하기 위해 버퍼 테이블 위에 스트림(6)을 정의합니다.
그런 다음 Snowflake 주기적 작업(7)을 사용하여 스트림을 쿼리하고 데이터를 정렬하고 최종 테이블(8)에 삽입합니다.
스노우플레이크 서비스
다음은 저희가 사용한 주요 Snowflake 서비스 중 일부입니다
- Snowpipe: Snowflake에서 관리하는 큐 서비스로, 외부 소스(S3 등)에서 데이터를 Snowflake 테이블로 복사합니다. SNS 토픽으로 Snowpipe를 트리거하면 S3 버킷에 저장된 모든 파일이 Snowflake 테이블에 복사됩니다. Snowpipe의 컴퓨팅 비용은 Snowflake가 제공하고, 사용자 정의 웨어하우스를 사용하지 않습니다.
- 스트림: Snowflake 서비스로, 데이터베이스 테이블에 정의할 수 있습니다. 스트림을 읽을 때마다, 그것은’ 마지막 읽기 이후 테이블 데이터를 읽어 변경된 데이터를 활용해 조치를 취할 수 있습니다. 그것은’ Kafka/Kinesis와 같은 스트리밍 서비스이지만 Snowflake 테이블에 구현되었습니다.
- 작업: Snowflake 주기적 작업은 스트림 데이터를 필요 시 다른 테이블로 옮기는 훌륭한 도구입니다. We’re Snowflake 작업을 배치 메커니즘으로 활용해 마이크로 파티션 깊이를 최적화합니다.
모니터링, 자가 치유 및 보존 처리
Snowflake 파이프라인이 갑작스러운 데이터 폭주에도 견딜 수 있도록 하기 위해, 몇 가지 유용한 Snowflake 기능을 활용하는 것 외에도 자체적으로 몇 가지 도구를 개발했습니다.
- 수집에 최적화된 웨어하우스: 우리는 “heavy-duty”와 “lightweight” 수집 소스를 사용량에 따라 다른 웨어하우스로 분리합니다.
- 수평 자동 확장: 당사의 모든 데이터 웨어하우스는 멀티 클러스터로 있으며 필요에 따라 수평적으로 확장됩니다.
- 작업 실패 모니터: Snowflake의 작업 메커니즘은 간헐적 실패를 손쉽게 처리합니다. 작업이 실패하면 동일한 체크포인트에서 바로 중단됩니다. 우리에게 문제는 작업이 실행될 때마다 전체 스트림 데이터를 가져온다는 점이었습니다. 스트림에 행이 너무 많으면 작업이 시간 초과하거나 오래 걸릴 수 있습니다. 데이터 급증으로 인한 수집 지연은 관리 비용이 크게 증가합니다. 우리는 다음과 같이 문제를 해결했습니다: 내부 모니터 구축 Snowflake의 작업 이력 을 읽고 작업이 오래 걸리는지 확인합니다. 오래 걸리면 모니터가 작업을 중단하고 웨어하우스를 확장한 뒤 재실행합니다. 이후 원래 규모로 축소합니다. 기존 서비스로는 빠르게 구현하기 어려웠던 이 기능은 결국 생명줄이 되었습니다.
- 데이터 보존: 우리는 데이터 보존을 매일 주기적인 작업을 사용해 효율적으로 처리하고 x일보다 오래된 데이터(마케팅 소스/고객별 변경)를 삭제했습니다. Snowflake’s DELETE 명령은 효율적이고 비차단적이며, 특히 클러스터 키에서 실행할 때 (예: 전체 일 데이터 삭제).
매시간, 데이터 수집 파이프라인의 Python Celery 작업이 수집된 데이터에 대해 Snowflake에 쿼리를 보냅니다.
취급 규모
Singular 에서 Snowflake를 사용하는 다양한 방식은 고유한 쿼리 패턴으로 이어집니다
- Singular’s 집계 데이터 수집 파이프라인: 고객’s 데이터에 대해 매시간 집계 쿼리를 실행합니다.
- 고객 대상 ETL 프로세스: Singular’s ETL 솔루션은 원시 데이터를 매시간 쿼리하고 고객’ 데이터베이스로 푸시합니다.
- 고객 대상 수동 내보내기: 원시 데이터의 온디맨드 내보내기.
- Singular의 내부 BI.
- GDPR “사용자 삭제” 쿼리: 특정 행 업데이트를 매일 배치로 실행하여 처리 GDPR 사용자 삭제 요청.
우리는 매시간 수천 건의 쿼리를 실행하며, 다양한 데이터 볼륨에 대해 신속하게 처리해야 합니다. 또한, 데이터 클러스터링에 미치는 영향을 최소화하면서 Snowflake가 특정 행만 업데이트하도록 해야 합니다. 이러한 쿼리 프로세스는 Python Celery 워커에서 분산 방식으로 주기적으로 실행됩니다.
Snowflake로 마이그레이션을 시작했을 때, 사용량과 규모에 따라 가상 웨어하우스 크기를 관리해야 할 것이라고 예상했습니다. 분산 쿼리 패턴에서는 각 워커가 Snowflake에 직접 쿼리를 보내고, 각 쿼리마다 새로운 세션이 생성되는 것이 자연스러웠습니다. 따라서 곧 Snowflake의 동시 세션 제한(Snowflake가 실질적으로 제한 없이 처리하는 동시 쿼리와는 다름)에 부딪히게 되었습니다.
다행히 Snowflake에서는 세션 제한에 얽매일 필요가 없습니다. 웨어하우스 큐 크기에 의해서만 제한되는 비동기 쿼리를 실행할 수 있습니다. 저희는 이 점을 활용하여 태스크 실행에 사용할 수 있는 활성 세션 풀을 생성했습니다. 설계 단계부터 완전히 비동기적이고 상태 비저장 방식으로 구현하여 웨어하우스 사용량에 의해서만 제한되도록 했습니다. 연결을 관리하는 상태 비저장 Python ASGI 서버를 구축하는 것은 저희에게 매우 간단하고 효율적인 해결책이었습니다.
저희는 SqlAlchemy Snowflake 통합 , Uvicorn 및 FastAPI를 사용하여 Python Snowflake 연결 풀 서비스를 구현했습니다 . 조만간 이 서비스를 오픈 소스로 공개하여 커뮤니티와 공유할 계획이며, 이를 통해 Snowflake 기반 애플리케이션을 개발하는 다른 사람들이 세션 풀을 생성할 때 시간과 노력을 절약할 수 있기를 바랍니다.
현재 저희는 Snowflake 연결 풀을 사용하여 여러 Snowflake 가상 웨어하우스에서 시간당 4,000건 이상의 쿼리를 실행하고 있습니다. 가상 웨어하우스는 필요에 따라 자동으로 수평 확장됩니다. GDPR 업데이트는 대규모 전용 웨어하우스에서 실행되므로 복잡한 UPDATE 쿼리도 효율적으로 처리됩니다. 가상 웨어하우스는 처리되는 요청이 없을 때 자동으로 일시 중단됩니다.
데이터 유효성 검사
Singular 고객들은 사업 운영을 위해 저희 데이터에 크게 의존하고 있습니다. 아테나(Athena)에서 스노우플레이크(Snowflake)로 마이그레이션하는 과정에서 데이터 불일치가 발생하는 것은 절대 용납할 수 없는 일이었습니다.
두 개의 데이터 수집 파이프라인을 병렬로 실행하면서, S3에 업로드된 각 파일에 대해 Athena와 Snowflake에서 실행되는 쿼리 결과를 비교하는 모니터링 시스템을 지속적으로 구축했습니다. 이를 통해 데이터 수집에 문제가 발생하면 즉시 파악할 수 있었습니다. 또한, 프로덕션 환경에서 실행한 모든 Athena 쿼리에 대해 Snowflake 쿼리를 실행하고 그 결과를 비교했습니다.
이 모니터링 시스템을 구축함으로써 데이터 마이그레이션이 올바르게 진행되고 있는지, 그리고 불일치가 없는지 확인할 수 있었습니다.
그 외 얻은 교훈
마이그레이션이 완료된 후, 우리는 신용 사용 내역과 설정을 더 자세히 검토하여 최적화할 수 있는 몇 가지 간단한 방법을 발견했습니다
- 사용량에 맞게 창고 사용을 계획합니다. 예를 들어, “heavy” 쿼리는 하루에 한 번(동시에 실행) 대규모 창고에서 자동 일시정지 설정으로 실행되고, 시간당 “light” 쿼리는 더 작은 창고에서 수평 클러스터 확장을 높게 설정해 실행됩니다. 이를 통해 일일 사용량을 비용 최적화하고(heavy 쿼리를 빠르게 실행하고 대부분 일시정지), 시간당 쿼리는 동시성 규모에 맞는 탄력성을 확보합니다.
- Snowflake와 협업해 클라우드 서비스 비용을 파악하세요. 우리의 무거운 쿼리 패턴이 테이블 NDV 계산 시 컴파일 시간을 크게 늘린 것을 발견했습니다. Snowflake가 기능 설정을 변경해 주어 불필요한 컴파일 시간을 크게 절감했습니다(새 클러스터를 최신 테이블에만 쿼리하는 경우).
- 쿼리 최적화: Snowflake 프로파일링 도구를 사용하여 여러 쿼리의 실행 시간을 최적화했습니다 . 이 도구는 GDPR 업데이트와 같은 복잡한 쿼리 처리에 도움이 되었습니다.
다음 단계는 무엇인가요?
- 데이터 공유: Snowflake와 안전한 데이터 공유에 멋진 기회가 있습니다. Snowflake 팀과 함께 이 아이디어를 활용하고 있습니다. 분산 클린룸 을(를) 안전한 데이터 공유에 활용합니다.
- IP 난독화: IP 관련 필드에 30일 보관 정책을 적용해 보안 인증서를 유지했습니다. Snowflake 마스킹 정책과 세션 변수를 활용해 IP 필드를 암호화 저장하고 필요 시 복호화하는 효율적인 흐름을 구축했습니다. We’re 솔루션을 심층 분석한 블로그 포스트를 곧 공개할 예정입니다.
- 연결 풀 오픈소스화: 우리는 향후 몇 달 안에 파이썬 Snowflake 연결 풀 서비스를 오픈소스할 계획입니다. 이것이 도움이 될 것이라면—예를 들어, 매시간 대량 쿼리를 실행한다면—연락 주세요!
이 블로그는 원래 Snowflake’s 블로그.