apache-spark

15 개의 포스트

kakao4분 읽기큐레이션 요약

개인화된 Airflow 테스트 환경 구축 및 운영 경험

수천 개의 Airflow DAG를 운영하는 카카오 데이터서비스 조직은 테스트 과정의 반복 작업과 환경 간 차이, 리소스 충돌 문제를 해결하기 위해 PR 단위의 개인 Airflow 환경인 AirZone을 구축했습니다. AirZone은 GitHub PR 코멘트에서 생성·삭제를 요청하면 Kubernetes Job과 전용 Helm 차트로 격리된 Airflow를 배포합니다. 사용자는 인프라를 직접 구성하지 않고 실제 운영 환경과 유사한 하둡·인증 환경에서 DAG를 테스트할 수 있습니다. ## 기존 테스트 방식의 한계 - **로컬 Airflow** - Airflow뿐 아니라 하둡 인증, 연결 설정, Docker 환경까지 직접 구성해야 합니다. - 초기 구축 비용이 크고, 로컬 환경과 운영 환경의 차이로 인해 실제 배포 후 실패할 수 있습니다. - **개발용 Airflow** - 코드를 커밋하고 푸시한 뒤 GitHub webhook, submodule 업데이트, DAG 파일 처리 과정을 기다려야 합니다. - DAG를 수정할 때마다 동기화 지연이 반복되어 개발 속도를 떨어뜨렸습니다. - **테스트용 Airflow** - SSH 컨테이너에 로컬 파일을 복사해야 하므로 코드 수정 때마다 추가 작업이 필요합니다. - 실제 데이터와 하둡에 접근하려면 prod VPN을 연결해야 하는 불편도 있었습니다. - **Production Airflow에서의 테스트** - 테스트 DAG가 스케줄러와 워커 자원을 점유해 다른 프로젝트의 실행을 지연시킬 수 있습니다. - 과도한 리소스 사용으로 노드 장애가 발생하면 같은 노드의 다른 태스크까지 중단될 위험이 있습니다. ## AirZone의 핵심 요구사항 - Kubernetes나 Helm을 몰라도 브라우저에서 Airflow 환경을 생성하고 삭제할 수 있어야 합니다. - 사용자는 DAG 검증에만 집중하고, 인프라 구성은 AirZone이 담당해야 합니다. - Jupyter Notebook을 제공해 별도 로컬 환경 없이 코드를 수정할 수 있어야 합니다. - 운영 환경과 유사한 DAG 실행 환경과 하둡 인증 방식을 제공해야 합니다. - PR마다 독립된 Airflow를 생성해 사용자와 테스트 작업을 서로 격리해야 합니다. ## PR 단위의 격리된 환경 - 레포지터리명과 PR 번호를 조합해 Kubernetes 네임스페이스를 생성합니다. - Airflow 웹 서버, 스케줄러, PostgreSQL, Jupyter, DAG PVC, 로그가 PR별로 분리됩니다. - 한 PR의 테스트가 다른 프로젝트의 스케줄러·워커 자원을 침범하지 않습니다. - 리뷰어는 PR에 남은 링크로 특정 코드 상태의 실행 결과를 직접 확인할 수 있습니다. - PR이 종료되면 네임스페이스를 기준으로 관련 리소스를 쉽게 정리할 수 있습니다. ## 요청 처리와 배포 작업의 분리 - `airzone-api`는 PR 존재 여부, PR이 열려 있는지, 네임스페이스 중복 여부 등 요청의 유효성만 검증합니다. - 실제 Helm 설치와 헬스체크는 별도의 Kubernetes Job이 수행합니다. - API가 수 분이 걸리는 배포 작업을 직접 기다리지 않으므로 빠르게 응답할 수 있습니다. - Job별로 상태와 로그가 독립적으로 남아 실패 단계와 원인을 추적하기 쉽습니다. - 실패한 Job을 삭제한 뒤 새 Job을 생성하는 방식으로 배포를 재시도할 수 있습니다. - 생성 요청은 `create-airzone-{namespace}`, 삭제 요청은 `delete-airzone-{namespace}` 형식의 Job 이름을 사용합니다. ## GitHub PR 코멘트를 사용자 인터페이스로 활용 - PR 생성 이벤트를 webhook으로 받아 저장소, 브랜치, PR 번호, 요청자 정보를 확인합니다. - 사용자가 선택할 수 있도록 하둡 환경별 AirZone 생성 링크를 PR 코멘트에 남깁니다. - 생성 완료 결과와 접속 정보도 PR에 표시해 별도 플랫폼 없이 테스트를 시작할 수 있습니다. - 다만 Jupyter 토큰과 Kubernetes 네임스페이스 토큰처럼 민감한 정보는 공개 범위가 넓은 PR 대신 카카오워크로 전달합니다. - 요청 접수와 배포 완료 시점에 카카오워크 알림을 보내 진행 상태를 알립니다. ## 전용 Helm 차트로 구성한 Airflow 기존 운영용 Airflow 차트가 아닌 AirZone 전용 Helm 차트를 만들어 테스트 환경에 필요한 구성만 묶었습니다. - **Airflow 구성** - 웹 서버와 스케줄러를 배포합니다. - `KubernetesExecutor`, DAG 스캔 주기, 로그 설정, 하둡 관련 변수를 테스트 환경에 맞게 설정합니다. - **데이터베이스와 저장소** - 개인 환경용 PostgreSQL을 함께 배포합니다. - scheduler와 Jupyter가 같은 DAG 작업 디렉터리를 사용하도록 DAG PVC를 공유합니다. - **DAG 동기화** - PR의 head repository와 branch 정보를 받아 해당 코드만 동기화합니다. - Git 초기화 컨테이너 등을 통해 배포 환경에 테스트 대상 DAG를 준비합니다. - **인증과 보안** - 사용자·공용 principal, 키탭, Jupyter 토큰을 환경에 주입합니다. - 하둡 접근에 필요한 Kerberos 인증을 운영 환경과 유사하게 구성합니다. - dkos에서 제공하는 TLS 인증서도 테스트 환경에 반영합니다. - **운영 연동** - 여유 있는 노드 그룹과 같은 리전의 `storageClass`를 선택합니다. - Airflow 로그를 Elasticsearch와 Kibana에서 확인할 수 있도록 연결 정보를 주입합니다. - Jupyter를 함께 제공해 브라우저에서 DAG와 관련 코드를 수정할 수 있게 합니다. ## 자동 정리와 운영 구조 - 생성·삭제 요청은 Kubernetes Job으로 처리합니다. - PR 종료 후 남아 있는 AirZone은 매일 실행되는 CronJob이 자동으로 회수합니다. - 사용자에게는 “PR 코멘트의 링크를 누르는 기능”으로 단순하게 보이지만, 운영자는 요청·설치·헬스체크·알림·정리 단계를 각각 추적할 수 있습니다. - 운영 환경에서 불필요한 PGBouncer나 외부 DB 연결 등은 제외해 개인 테스트 환경의 복잡도와 비용을 줄였습니다. AirZone과 같은 구조를 도입할 때는 테스트 환경을 운영 환경과 최대한 유사하게 유지하되, PR 또는 브랜치 단위로 리소스를 격리하는 것이 중요합니다. 또한 긴 배포 작업은 API 요청과 분리하고, Kubernetes Job의 상태·로그·재시도 기능을 활용하면 장애 대응과 운영 추적이 훨씬 쉬워집니다.

원문 읽기(새 탭에서 열림)
line5분 읽기큐레이션 요약

임베딩 안정화로 검색 리랭킹의 콜드 스타트 문제를 해결하다: LINE Part Time Jobs 적용 사례

LINE Part Time Jobs는 기존 2타워 임베딩 기반 리랭킹의 콜드 스타트와 검색 쿼리 미반영 문제를 해결하기 위해, 임베딩 공간을 날짜별로 안정화하는 후처리 방식을 도입했습니다. 저차원 SVD와 직교 Procrustes 정렬을 통해 매일 재학습되는 임베딩의 연속성을 유지했고, 그 결과 오프라인 전환 nDCG가 약 9%, 클릭 nDCG가 약 4.5% 향상되었습니다. A/B 테스트에서도 서비스 전체 KPI 4.7%, 매출 6.5% 증가를 달성했습니다. ## LINE Part Time Jobs 검색 리랭킹 구조 - 검색 시스템은 다음 두 단계로 구성됩니다. - **검색(retrieval):** 사용자의 쿼리에 맞는 구인 공고 후보를 수집 - **리랭킹(reranking):** 수집된 후보를 사용자별로 재정렬 - 기존에는 별도 배치 파이프라인에서 생성한 사용자·아이템 2타워 임베딩을 활용했습니다. - 사용자 임베딩과 아이템 임베딩의 코사인 유사도를 계산해 검색 결과 순위를 정했습니다. - 이 방식은 실시간 연산 부담을 줄일 수 있지만 다음 한계가 있었습니다. - 검색 쿼리와 역, 거리 같은 화면별 정보가 임베딩에 충분히 반영되지 않음 - 검색 외 추천 모듈이나 LINE 공식 계정에서 발생한 행동까지 함께 포함됨 - 검색 화면에 특화된 사용자 의도를 정밀하게 반영하기 어려움 ## 전용 리랭킹 모델 도입 과정의 문제 ### 공고 교체로 인한 콜드 스타트 - LINE Part Time Jobs의 공고는 월초에 대부분 교체됩니다. - 새로운 공고에 대한 클릭·지원 데이터가 충분히 쌓이기 전에는 전용 리랭킹 모델이 학습할 데이터가 부족합니다. - 그 결과 공고 교체 직후 모델 성능이 크게 저하되는 콜드 스타트 문제가 발생했습니다. ### 매일 변하는 임베딩 공간 - 2타워 모델은 성능 유지를 위해 주기적으로 랜덤 가중치에서 처음부터 재학습됩니다. - 재학습할 때마다 임베딩 공간의 방향과 좌표계가 달라질 수 있습니다. - 따라서 오늘 생성한 임베딩과 어제 생성한 임베딩은 실제 의미가 비슷해도 벡터 좌표상 직접 비교하기 어렵습니다. - 학습 시점과 추론 시점에 서로 다른 버전의 임베딩을 사용하면 다운스트림 리랭킹 모델의 입력 분포가 달라져 성능이 떨어질 수 있습니다. ## 임베딩 안정화 방식 - 각 날짜의 임베딩을 전날 안정화된 임베딩 공간에 맞춰 정렬합니다. - 첫날 임베딩은 별도 변환 없이 기준으로 사용합니다. - 이후에는 전날 결과를 다음 날의 기준으로 삼아 임베딩 공간을 순차적으로 연결합니다. - 이 방식은 특정 기준일에 모든 임베딩을 맞추는 대신, 시간에 따른 공간의 연속성을 유지합니다. - 결과적으로 임베딩 피처와 다운스트림 모델의 업데이트 시점을 엄격히 일치시키지 않아도 됩니다. ## 저차원 SVD와 직교 Procrustes ### 저차원 SVD - 아이템 임베딩과 사용자 임베딩을 각각 행렬 \(T\), \(W\)로 표현합니다. - 2타워 모델의 점수는 \(TW^\top\)로 계산되지만, 이 대규모 행렬을 직접 분해하지는 않습니다. - 대신 저차원 SVD를 사용해 변환 행렬 \(M_T\), \(M_W\)를 구합니다. - 변환 결과는 다음과 같습니다. - 아이템 임베딩: \(T' = TM_T\) - 사용자 임베딩: \(W' = WM_W\) - 이를 통해 각 학습에서 생성된 임베딩을 보다 표준화된 저차원 표현으로 변환합니다. ### 직교 Procrustes 정렬 - 당일 임베딩과 전날 안정화된 임베딩이 최대한 일치하도록 직교 변환을 계산합니다. - 직교 변환은 회전과 반전만 수행하므로 벡터 간 거리와 내적 구조를 보존합니다. - 따라서 임베딩의 유사도 기반 점수 계산 특성을 유지하면서 일별 공간 차이를 보정할 수 있습니다. ## 대규모 데이터 처리를 위한 구현 - 데이터 규모가 크기 때문에 알고리즘을 Apache Spark 기반 분산 처리로 구현했습니다. - 저차원 SVD에서는 원 논문의 QR 분해 대신 숄레스키 분해를 사용했습니다. - Gramian 행렬 \(G = A^\top A\)를 계산 - \(G = R^\top R\) 형태로 숄레스키 분해 - QR 분해에서 필요한 상삼각 행렬 \(R\)을 효율적으로 획득 - 직교 Procrustes에서는 다음과 같이 처리했습니다. - 대규모 행렬곱 \(M = B^\top A\)는 Spark로 분산 계산 - \(M\)은 임베딩 차원 \(e \times e\)의 작은 행렬이므로 SVD는 단일 노드에서 NumPy로 계산 - 대규모 벡터 데이터와 소규모 변환 행렬을 구분해 계산 자원을 효율적으로 배분했습니다. ## 안정화 효과와 평가 결과 - 안정화 전에는 서로 다른 날짜의 임베딩 상관관계가 거의 0에 가까웠습니다. - 안정화 후에는 다음 수준의 유사도를 유지했습니다. - 일주일 후: 약 0.88 - 한 달 후: 약 0.87 - 안정화하지 않은 임베딩을 다운스트림 모델에 추가하면 공간 불일치로 nDCG가 약 1~5% 하락했습니다. - 안정화된 임베딩을 사용한 경우: - 전환 nDCG 약 9.0% 향상 - 클릭 nDCG 약 4.5% 향상 ## A/B 테스트 결과와 해석 - 안정화된 임베딩과 콜드 스타트 대응책을 결합한 모델을 온라인 실험했습니다. - 검색 화면 단독 KPI에서는 통계적으로 유의미한 개선이 뚜렷하지 않았습니다. - 그러나 서비스 전체 기준으로는 다음 성과를 얻었습니다. - KPI 4.7% 향상 - 매출 6.5% 향상 - 이는 임베딩이 검색 화면뿐 아니라 서비스 전반의 사용자 행동과 장기적인 선호를 반영했기 때문으로 분석됩니다. - 검색 이후 다른 페이지로 이동하거나 다른 추천 모듈에서 지원하는 행동까지 긍정적인 영향을 받은 것으로 보입니다. - 기존 2타워 모델 구조나 학습 파이프라인을 변경하지 않고 임베딩 후처리만 추가했다는 점도 운영상 중요한 장점입니다. ## 실용적인 결론 재학습마다 좌표계가 달라지는 임베딩을 다운스트림 모델의 피처로 사용할 때는 날짜별 공간 정렬이 효과적인 해결책이 될 수 있습니다. 특히 저차원 SVD와 직교 Procrustes를 결합하면 임베딩의 유사도 구조를 유지하면서 버전 불일치와 드리프트를 줄일 수 있으므로, 기존 모델을 크게 변경하기 어려운 대규모 추천·검색 시스템에 적용하기 적합합니다.

원문 읽기(새 탭에서 열림)
line5분 읽기큐레이션 요약

총 용량 1EB 초과! 서로 역사가 다른 두 HDFS를 어떻게 연결할까? 데이터 플랫폼 연계 중 직면한 과제와 설계 결정

LY Corporation은 구 LINE과 구 Yahoo Japan의 HDFS 클러스터를 통합해 1EB를 넘는 데이터 플랫폼을 운영하고 있습니다. 두 플랫폼은 같은 HDFS를 사용했지만 Namespace 구성, 권한 관리, 사용자 접근 방식이 달라 운영·연계 설계에서 서로 다른 과제가 발생했습니다. 대규모 환경에서는 단순한 용량 증설보다 NameNode 메타데이터, 소규모 파일, Balancer 트래픽, 클라이언트 접근 경로를 함께 관리해야 한다는 것이 글의 결론입니다. ## 구 LINE과 구 Yahoo Japan의 운영 모델 차이 - **구 LINE** - 여러 분석 환경을 통합한 Hadoop 3.x 기반 플랫폼을 구축했습니다. - 사용자가 HDFS 권한이나 Apache Ranger를 직접 다루지 않도록 웹 포털을 제공했습니다. - DB·테이블 중심의 데이터 카탈로그와 역할 기반 권한 신청·승인 구조를 사용했습니다. - BI, 리포팅, ETL 등 다양한 사용 방식을 지원했지만, 기존 클러스터를 Namespace 단위로 통합하면서 관리 복잡성이 증가했습니다. - **구 Yahoo Japan** - Hadoop 0.2.x 시절부터 제한된 목적의 대규모 분석 플랫폼을 운영했습니다. - 접근 인터페이스를 최소화하고 사용 방식에 일정한 제약을 둬 안정적인 운영과 사용자 지원을 우선했습니다. - 단일 Namespace가 커지면서 NameNode 확장에 한계가 생기자 RBF(Router-Based Federation)를 도입했습니다. - 권한 관리는 여전히 HDFS Permission(POSIX 유사 모델)을 사용해 현대적인 세밀한 권한 관리에는 제약이 있습니다. - 같은 HDFS라도 데이터 위치, 사용자 접근 경로, 권한 변경의 영향 범위, 운영팀의 개입 방식이 크게 달랐습니다. ## Namespace 구성과 접근 방식의 차이 - 두 플랫폼 모두 NameNode 병목을 피하기 위해 여러 Namespace로 HDFS를 분할했습니다. - 각 Namespace는 2~4대의 NameNode로 이중화했지만, Namespace와 DataNode를 묶는 방식은 달랐습니다. ### 구 LINE: ViewFS와 DataNode 공유 - ViewFS를 사용해 클라이언트의 마운트 테이블에서 실제 Namespace와 NameNode를 결정했습니다. - 유연하고 단순하지만 모든 사용자와 서비스에 올바른 마운트 설정을 배포·유지해야 합니다. - 여러 Namespace가 DataNode를 공유해 자원 효율은 높지만 클러스터 구성과 운영이 복잡해졌습니다. - 통합 과정에서 NameNode와 DataNode 버전이 혼재한 점도 운영 부담을 키웠습니다. ### 구 Yahoo Japan: RBF 기반 서버 측 라우팅 - 사용자는 Router에 접속하고, Router가 적절한 Namespace와 NameNode로 요청을 분배합니다. - 클라이언트 설정을 단순화하고 여러 HDFS를 투명하게 연결하기 쉽습니다. - 대신 Router 계층의 가용성, 확장성, 네트워크 도달성, 진입점 제어가 중요합니다. - Observer NameNode를 활용해 읽기 부하도 분산했습니다. - 이 차이는 플랫폼 간 데이터 복사에서도 중요합니다. - ViewFS 환경은 클라이언트 마운트 테이블과 설정 배포가 핵심입니다. - RBF 환경은 Router를 통한 접근 경로와 Router 계층의 안정성이 핵심입니다. ## 구 LINE에서 발생한 용량과 네트워크 문제 - 전사적 활용이 확대되면서 예상보다 빠르게 HDFS 용량이 부족해졌습니다. - 신규 서버 납품 전까지 기존의 오래된 서버를 임시 재사용하고, 이후 서버를 교체하는 작업이 반복됐습니다. - DataNode를 대규모로 추가하거나 제거하면 HDFS Balancer와 블록 재배치가 대량의 네트워크 트래픽을 발생시켰습니다. - 따라서 서버 변경 작업은 HDFS뿐 아니라 네트워크 구성과 혼잡 가능성까지 고려해 네트워크팀과 협력해야 했습니다. ## 블록 증가와 NameNode 부하 - NameNode는 파일·디렉터리·블록 메타데이터를 메모리에 보관하므로 파일과 블록 수가 증가할수록 힙 사용량과 처리 부하가 커집니다. - 힙이 지나치게 커지면 GC 시간이 길어지고 NameNode 응답 지연과 불안정성이 발생합니다. - 특히 소규모 파일이 많은 Hive 테이블이 주요 원인이었습니다. - 해결을 위해: - FSImage를 정기적으로 덤프해 Hive 테이블로 저장했습니다. - 사용자별 경로, 파일 수, 블록 수, 데이터량을 분석했습니다. - 삭제나 스키마 변경이 필요 없는 테이블을 우선 선정했습니다. - 파일 병합으로 파일 수와 블록 수를 줄였습니다. - 파일 병합은 메타데이터 규모뿐 아니라 HDFS 요청 횟수도 줄여 잡 실행 속도와 NameNode 응답성을 개선했습니다. ## Namespace별 부하 특성과 Balancer 조정 - 부하는 Namespace마다 다르게 나타났습니다. - 임시 파일이 많은 Namespace에서는 평상시 부하는 낮지만 DataNode 추가 시 Balancer가 급격한 부하를 유발했습니다. - Apache Spark의 스테이징 파일은 짧은 시간에 생성·삭제되며 NameNode의 write lock을 빈번하게 발생시켰습니다. - Balancer의 블록 이동은 블록 정보 조회를 늘리고 read lock을 오래 유지해 파일 생성·삭제를 지연시킬 수 있었습니다. - 초기에는 Balancer 병렬도를 높였지만, DataNode 디스크 여유와 서비스 영향을 고려해 병렬도를 낮추고 처리 시간과 안정성 사이의 균형을 맞췄습니다. ## 조직 통합과 플랫폼 연계 - 조직 통합 후에는 서로 다른 설계 철학의 데이터 플랫폼을 연결해야 했습니다. - 주요 설계 과제는 다음과 같습니다. - 어느 플랫폼과 Namespace를 연결할 것인가 - 권한 관리 단위를 어떻게 맞출 것인가 - 어떤 접속 경로와 진입점을 사용할 것인가 - DistCP 등으로 데이터를 어떤 경로로 전송할 것인가 - 네트워크 도달성과 운영 책임을 어떻게 나눌 것인가 - 글의 후반부에서는 플랫폼 간 권한 모델 통합과 DistCP 기반 데이터 연계 방식을 다룰 예정이지만, 제공된 본문은 해당 설명이 시작되기 전에 끝나 있습니다. 대규모 HDFS를 운영할 때는 용량 증설만으로 문제를 해결하기 어렵습니다. 파일·블록 수를 지속적으로 관찰하고 소규모 파일을 줄이며, DataNode 변경과 Balancer의 네트워크 영향을 사전에 통제해야 합니다. 또한 플랫폼 통합 시에는 ViewFS와 RBF의 접근 모델 차이, 권한 체계, Router 및 클라이언트 설정까지 포함한 운영 경계를 함께 설계하는 것이 권장됩니다.

원문 읽기(새 탭에서 열림)
netflix4분 읽기큐레이션 요약

데이터 프로젝트: 넷플릭스 규모로 데이터 자산 관리

Netflix는 수백만 개의 테이블과 수만 개의 워크로드를 개별 자산·사용자 단위로 관리하면서 권한 변경과 워크플로 장애가 반복되는 문제를 겪었다. Data Projects는 관련 자산을 프로젝트로 묶고, 사람과 무관하게 유지되는 프로젝트 전용 신원(identity)을 부여해 권한과 실행 주체의 관리 단위를 상향한다. 이를 통해 조직 개편이나 담당자 변경에도 권한과 데이터 워크플로가 안정적으로 유지되도록 한다. ## 개별 자산 중심 권한 관리의 한계 - 기존에는 모든 테이블에 개별 ACL을 설정해야 했다. - 조직 개편이나 팀 통합 때 수백 개 테이블의 권한을 하나씩 수정해야 했다. - 대규모 권한 변경 요청이 지원팀에 집중됐다. - 관리 부담을 피하기 위해 테이블을 전사 공개하는 사례가 생겨 ACL의 의미가 약화됐다. - Data Projects는 여러 테이블과 워크로드를 하나의 프로젝트로 묶어 프로젝트 단위로 권한을 관리한다. ## 사람에게 종속된 워크로드 신원의 문제 - Maestro 워크플로, Spark 파이프라인, 데이터 이동 작업은 실행 시 신원이 필요했다. - 기존에는 워크플로 작성자의 사용자 계정으로 실행하는 경우가 많았다. - 담당자가 팀을 옮기거나 퇴사하면 계정 권한이 바뀌어 워크플로가 실패했다. - 다른 사람의 계정으로 교체해도 권한이 완전히 같지 않아 새로운 권한 오류가 연쇄적으로 발생했다. - 수만 개의 예약 워크로드를 운영하는 Netflix에서는 이러한 방식이 지속 가능하지 않았다. ## Data Projects의 기본 구조 - Data Project는 관련 데이터 자산을 묶어 관리·조회하는 논리적 컨테이너다. - 포함할 수 있는 자산에는 테이블, 워크플로, 시크릿 등이 있다. - 동시에 사람과 독립적으로 유지되는 합성(synthetic)·지속적(durable) 신원 역할을 한다. - 500개 테이블의 ACL을 각각 관리하는 대신, 하나의 프로젝트에 대한 권한을 관리한다. - 초기 목적은 접근 제어와 실행 신원 통합이지만, 향후 다른 데이터 플랫폼 관리 기능으로 확장될 수 있다. ## 프로젝트 기반 역할과 권한 - 프로젝트 소유 팀이 프로젝트의 grant를 관리한다. - 사용자, 그룹, 애플리케이션, CI 작업 등 다양한 주체를 grant로 추가할 수 있다. - 각 grant에는 프로젝트 내 작업 범위를 결정하는 역할이 부여된다. - 예를 들어: - `Contributor`: 프로젝트 자산에 대한 읽기·쓰기 권한 - `Viewer`: 읽기 전용 권한 - 팀원이 합류하거나 떠날 때 개별 자산 ACL 수백 개를 수정하지 않고 프로젝트 grant 하나만 변경하면 된다. ## Netflix 애플리케이션 ID와 AWS IAM 역할 - 모든 Data Project에는 Netflix 애플리케이션 신원이 provision된다. - 필요하면 AWS IAM 역할도 함께 제공된다. - Netflix 신원은 Maestro 같은 비동기 워크로드의 실행 주체가 된다. - AWS IAM 역할은 Amazon EMR의 Spark 작업 등 AWS 특화 작업에 사용된다. - IAM 역할은 암호학적으로 안전한 방식으로 프로젝트의 Netflix 신원으로 교환될 수 있다. - 권한이 충분한 프로젝트 구성원은 로컬 노트북이나 노트북 환경에서 프로젝트 신원을 가정해 실제 예약 작업과 동일한 권한으로 테스트·디버깅할 수 있다. ## ‘Gravity’를 통한 자산 자동 귀속 - 프로젝트 신원으로 실행된 워크로드가 새 자산을 만들면 해당 자산이 자동으로 프로젝트에 포함된다. - 예를 들어 Maestro 워크플로가 테이블 세 개를 생성하면 이 테이블들이 자동으로 프로젝트의 자산이 된다. - 생성 자산을 나중에 찾아 프로젝트에 수동 등록할 필요가 없다. - 프로젝트가 해당 워크로드가 만든 자산의 중심이 되어 관리 범위와 권한 적용이 자연스럽게 확장된다. ## Maestro와 신뢰된 워크로드 실행 - Maestro는 ETL, 데이터 이동, 머신러닝 학습 등 배치 분석 작업을 담당하는 Netflix의 핵심 오케스트레이터다. - 예약 작업은 원래 사용자가 실행 시점에 উপস্থিত하지 않아도 되므로, Maestro는 Trusted Workload Manager(TWM)로 지정됐다. - TWM은 관리하는 워크로드를 대신해 새로운 신원 토큰을 발급할 권한을 가진다. - 하나의 워크플로 실행은 데이터 웨어하우스 테이블 ACL, Netflix 리소스 정책, AWS IAM 정책을 모두 통과해야 할 수 있다. - 따라서 실행 신원이 불안정하면 전체 데이터 파이프라인이 실패한다. ## 프로젝트 기반 지속 가능한 신원 - 기존의 `maestro OBO alice@netflix.com` 방식은 Maestro와 개인 사용자의 권한을 결합했지만, 사용자 생명주기에 종속됐다. - Data Projects는 이를 팀이 소유하는 Netflix 애플리케이션 신원으로 대체한다. - 프로젝트 신원은 담당자의 휴가, 부서 이동, 퇴사와 무관하게 유지된다. - Maestro는 워크플로 실행 전 호출자가 해당 프로젝트를 사용할 권한이 있는지 검증한다. - 실행 중 생성된 테이블은 gravity를 통해 프로젝트에 자동 귀속되고 프로젝트 권한을 물려받는다. - 시크릿도 프로젝트 정책 범위에서 관리되므로 담당자 변경으로 자격 증명이 고립되지 않는다. - 결과적으로 권한 관리가 중앙화되고, 워크플로 실행이 안정적이며, 감사 가능성도 높아진다. 대규모 데이터 플랫폼에서는 개별 테이블과 사용자에 권한을 계속 부여하기보다, 팀·서비스·워크로드를 대표하는 프로젝트 단위의 소유권과 지속 가능한 실행 신원을 도입하는 것이 효과적이다. 특히 예약 작업과 조직 변화가 많은 환경에서는 프로젝트 단위 권한, 자동 자산 귀속, 사람과 분리된 서비스 신원을 함께 설계하는 것이 권장된다.

원문 읽기(새 탭에서 열림)
netflix4분 읽기큐레이션 요약

넷플릭스에서 카산드라 데이터 이동의 진화

넷플릭스는 하루 약 1,200건, 약 3PB의 Cassandra 데이터를 Iceberg로 옮기던 기존 Casspactor를 확장 가능한 계층형 엔진으로 교체했다. 기존 시스템은 여러 메타데이터 서비스에 의존하고 대규모 파티션과 다양한 Cassandra 데이터 모델을 제대로 처리하지 못해 안정성·비용·확장성 문제가 발생했다. 새 구조는 S3 백업을 단일 진실 공급원으로 사용하고, Spark DataFrame과 커넥터 팩토리를 기반으로 각 데이터 모델에 최적화된 커넥터를 구축한다. ## Casspactor의 역할과 한계 - Casspactor는 Cassandra의 SSTable과 메타데이터를 S3 백업에서 읽어 Iceberg 테이블로 변환했다. - 하루 약 1,200건의 데이터 이동과 약 3PB 규모의 전송을 처리하며 핵심 업무를 지원했다. - Cassandra 노드의 사이드카 프로세스가 SSTable과 메타데이터를 S3에 업로드하고, 작업 실행 시 엔진이 필요한 백업 구조를 구성했다. - 이후 SSTable을 다운로드하고 mutation compaction과 변환을 수행한 뒤 Iceberg에 기록했다. - 그러나 단일 커넥터로 설계되어 Key Value, Time Series, Graph 등 여러 데이터 추상화에 공통 기반을 제공하기 어려웠다. ## 분산된 메타데이터 의존성 문제 - Casspactor는 백업의 존재 여부, 완전성, 포함 데이터를 여러 독립 시스템의 메타데이터를 조합해 판단했다. - 각 시스템의 갱신 주기와 정확도, 장애 방식이 달라 실제 백업 상태와 Casspactor의 인식이 불일치할 수 있었다. - 메타데이터가 실제 백업과 어긋나면 오래되거나 잘못된 데이터를 조용히 읽는 문제가 발생했다. - Cassandra 클러스터 유지보수 중 비동기 스냅샷이 만들어졌고, 한 리전에 있는 모든 노드가 같은 시각에 스냅샷을 생성해야 한다는 제약도 있었다. - 노드 하나가 교체되면 리전 전체의 데이터 이동이 실패할 수 있었다. - 새 구조에서는 백업 파일 자체의 메타데이터를 직접 읽어 S3를 백업 존재 여부와 완전성을 판단하는 단일 진실 공급원으로 사용한다. ## 모든 커넥터가 물려받은 제약 - **대규모 파티션 처리 실패** - Key Value와 Time Series에서 흔한 넓은 파티션을 처리하지 못했다. - 일부 작업은 메모리 부족으로 종료됐다. - **데이터 모델 인식 부족** - 원시 Cassandra 테이블만 이동했기 때문에 Key Value 등의 커넥터가 별도의 후처리로 데이터 모델을 복원해야 했다. - 이로 인해 처리 비용과 복잡성, 장애 가능성이 증가했다. - **중간 테이블 증가** - Casspactor는 최종 결과 전에 중간 Iceberg 테이블을 생성했다. - Key Value 커넥터는 추가 중간 테이블과 스냅샷 테이블까지 필요했다. - 상위 데이터 추상화가 추가될수록 중간 저장 공간과 비용이 누적됐다. - **Time Travel 불가** - 여러 서비스를 조합해 백업 단위를 구성했기 때문에 클러스터 토폴로지나 Keyspace 스키마가 변경된 뒤 과거 백업을 복원하기 어려웠다. - **모놀리식 구조** - Casspactor는 재사용 가능한 엔진이 아니라 하나의 커넥터였다. - 데이터 모델별 목적형 커넥터를 공통 기반 위에 구축할 수 없었다. ## 새로운 계층형 아키텍처 - 새 구조는 Apache Cassandra Analytics, 넷플릭스의 Move Data 프레임워크, 내부 백업 표현 방식과 S3 클라이언트를 결합한다. - 가장 아래 계층인 Cassandra Analytics Wrapper가 S3 백업에서 원시 데이터를 읽는다. - 읽은 결과는 표준 Spark DataFrame으로 변환된다. - 상위 계층의 Connector Factory는 Java UDF와 변환 로직을 통해 데이터 추상화별 커넥터를 생성한다. - Key Value나 Time Series 커넥터는 일반적인 DataFrame을 입력으로 받아 자체 데이터 모델에 맞게 처리한다. - 핵심 읽기 엔진의 개선 사항은 모든 커넥터에 적용되고, 각 커넥터는 변환 로직에 집중할 수 있다. ## Spark 기반 처리와 운영 개선 - mutation compaction과 데이터 처리를 Spark Executor 수준으로 이동했다. - 대규모 또는 편향된 파티션을 과도한 셔플 없이 처리해 메모리 부족 문제를 줄였다. - Cassandra 백업에서 바로 Spark DataFrame을 생성하므로 중간 Iceberg 테이블이 필요하지 않다. - 중간 저장 비용과 다단계 파이프라인의 운영 복잡성이 감소한다. - 소스 테이블의 특성에 따라 작업 리소스를 자동 조정하는 auto-sizing 기능을 제공한다. - 엔지니어가 작업별 리소스를 수동으로 조정하지 않아도 성능과 비용을 최적화할 수 있다. - S3의 백업 메타데이터를 직접 사용해 여러 외부 서비스에 대한 의존성을 제거하고 안정성을 높였다. ## 실용적인 결론 대규모 Cassandra 데이터 이동 시스템은 백업 메타데이터를 여러 서비스에서 조합하기보다 실제 저장소를 단일 진실 공급원으로 삼는 편이 안정적이다. 또한 원시 데이터 추출 엔진과 데이터 모델별 변환 커넥터를 분리하면, 공통 성능 개선을 재사용하면서도 Key Value·Time Series 같은 각 도메인에 최적화된 처리를 구현할 수 있다.

원문 읽기(새 탭에서 열림)
toss6분 읽기큐레이션 요약

Spark Connect on Kubernetes #1: 견고한 Spark Connect 만들기

토스증권은 여러 사용자가 안정적으로 사용할 수 있는 Production급 Spark Connect를 Kubernetes에서 운영하고 있습니다. Spark Connect는 Driver를 애플리케이션마다 실행하는 대신 장기 실행 서버로 분리해 가벼운 클라이언트와 빠른 세션 생성을 제공하지만, 여러 세션이 하나의 SparkContext를 공유하면서 장애 전파와 리소스 경합 문제가 발생합니다. 이를 해결하기 위해 글로벌 장애 카운터를 사실상 비활성화하고, 결과 크기를 제한하며, 여러 Replica로 Driver와 SparkContext를 분리하는 전략을 사용합니다. ## Classic Spark의 구조와 Spark Connect의 등장 - Spark는 작업을 계획·지휘하는 **Driver**와 실제 연산을 수행하는 **Executor**로 구성됩니다. - Classic Spark의 배포 방식은 다음과 같습니다. - **Client mode**: 클라이언트 프로세스가 Driver 역할을 수행합니다. - **Cluster mode**: 작업 제출 시 클러스터에 Driver가 생성되고 작업 종료 후 사라집니다. - 두 방식 모두 애플리케이션마다 Driver가 하나씩 생성되고, 애플리케이션의 수명과 함께 종료됩니다. - Spark Connect는 Spark 3.4부터 도입됐으며, 4.0에서는 기존 Dataset/DataFrame API와 거의 동등한 수준에 도달했습니다. - 글의 구현과 설정은 Spark 4.1을 기준으로 합니다. ## Spark Connect의 동작 방식 - Spark Connect에서는 Driver를 애플리케이션별 프로세스가 아니라 **미리 실행해 둔 서버**로 운영합니다. - 클라이언트는 Spark 라이브러리와 JVM을 직접 포함하지 않는 Thin Client입니다. - 클라이언트의 DataFrame·SQL 연산은 다음 과정으로 처리됩니다. - 연산을 Unresolved Logical Plan으로 변환 - Protocol Buffer로 인코딩 - gRPC를 통해 서버로 전송 - 서버가 분석, 최적화, 스케줄링, 실행 수행 - 결과를 Arrow 기반으로 클라이언트에 스트리밍 - 구조적으로는 JDBC 클라이언트가 데이터베이스 서버에 질의하는 모델과 유사합니다. ## Spark Connect의 장점 - 클라이언트에 무거운 Spark 의존성이나 JVM이 없어도 됩니다. - Python, SQL, 노트북, BI 도구 등 다양한 클라이언트가 같은 서버에 접속할 수 있습니다. - Driver가 이미 실행 중이므로 매번 프로세스를 생성하고 리소스를 협상할 필요가 없습니다. - 클라이언트가 종료되거나 네트워크가 끊겨도 서버에서 실행 중인 작업은 보호할 수 있습니다. - 반면 하나의 장기 실행 서버에 여러 사용자가 접속하면서, Spark의 기존 “애플리케이션 하나에 워크로드 하나”라는 전제가 깨집니다. ## 공유 Driver가 만드는 단일 장애점 - 여러 세션이 하나의 SparkContext와 Driver JVM을 공유합니다. - Driver가 장애를 일으키면 해당 서버의 모든 세션, 실행 중인 Job, 캐시가 함께 사라집니다. - `spark.executor.maxNumFailures`는 Executor 실패를 애플리케이션 전체 단위로 누적합니다. - 기본 임계값은 `max(3, 2 × executor 수)`입니다. - 임계값을 초과하면 `stopApplication()`이 호출되고, 결과적으로 `sys.exit(11)`로 서버 전체가 종료됩니다. - 이 카운터는 다음 이유로 멀티세션 환경에서 위험합니다. - 개별 쿼리의 Task 실패가 아니라 Executor 실패를 전역적으로 집계합니다. - 시간이 지나도 실패 기록이 계속 누적됩니다. - 서로 다른 사용자의 실패가 합산됩니다. - 문제가 없는 세션도 장애를 함께 겪게 됩니다. ## 세션 격리와 리소스 경합의 한계 - `newSession()`은 SQL 네임스페이스 등 세션 상태만 분리합니다. - CPU, 메모리, Executor, Task 슬롯은 모든 세션이 공유합니다. - 한 사용자가 대규모 Job을 제출하면 다른 사용자의 쿼리 응답도 느려질 수 있습니다. - 기본 FIFO 스케줄링에서는 먼저 제출된 작업이 우선하며, 선점이 없어 이미 실행 중인 Task를 중단할 수 없습니다. - Fair Scheduler를 사용해도 Task 슬롯을 배분하는 순서만 조정할 뿐, 사용자별 CPU·메모리 격리는 제공하지 않습니다. - Spark Connect에서는 `spark.scheduler.pool`이 기본적으로 제대로 전파되지 않아 모든 쿼리가 Default Pool에 들어갑니다. - Classic Spark에서는 `setLocalProperty()`가 Driver 스레드에 직접 적용됩니다. - Spark Connect에서는 클라이언트와 Driver가 분리되어 서버의 요청 처리 스레드에 값을 별도로 설정해야 합니다. - 토스증권은 서버 스레드에 사용자별 Pool을 직접 설정하는 방식으로 이 문제를 보완했습니다. - 사용자별 Pool을 적용하려면 먼저 요청의 사용자를 식별해야 하며, 인증·인가와 연결됩니다. - 궁극적인 CPU·메모리 격리는 Spark 스케줄러가 아니라 Spark 외부의 리소스 관리 계층에서 해결해야 합니다. ## 고정된 서버 스케일 문제 - Spark Connect 서버는 이미지, Driver·Executor 리소스, Spark 설정이 고정된 상태로 실행됩니다. - Dynamic Resource Allocation으로 Executor 수는 조절할 수 있지만, 서버 자체의 기본 스펙은 실행 중 바뀌지 않습니다. - 서버를 필요에 따라 생성·교체하거나 팀 단위로 격리하는 문제는 후속 글에서 다룹니다. ## 글로벌 장애 카운터 비활성화 - 서버 전체를 종료시키는 Executor 실패 경로를 차단하기 위해 다음과 같이 설정합니다. - `spark.executor.maxNumFailures`: 사실상 무한대로 설정해 글로벌 종료 조건을 비활성화 - `spark.executor.failuresValidityInterval`: 오래된 실패 기록을 주기적으로 제거 - `spark.task.maxFailures`: 동일 Task의 반복 실패를 제한 - `spark.stage.maxConsecutiveAttempts`: Shuffle Fetch 실패로 Stage가 반복 실행되는 상황을 제한 - `task.maxFailures`는 OOM이나 예외처럼 동일 Task가 반복 실패하는 경우를 담당합니다. - `stage.maxConsecutiveAttempts`는 Shuffle Fetch 실패로 Stage 전체가 반복되는 경우를 담당합니다. - 이 방식으로 문제가 있는 쿼리만 실패시키고 서버와 다른 사용자의 세션은 유지할 수 있습니다. - 다만 실패 허용 횟수를 지나치게 낮추면 일시적인 장애에도 정상 쿼리가 실패할 수 있으므로 워크로드에 맞춰 여유를 둬야 합니다. ## Driver 메모리 보호와 결과 크기 제한 - Spark Connect에서는 쿼리 결과가 Driver를 거쳐 클라이언트로 스트리밍됩니다. - 사용자가 대규모 테이블을 `collect`하면 Driver 메모리가 고갈될 수 있습니다. - `spark.driver.maxResultSize`는 한 액션에서 반환되는 Task 결과의 누적 크기를 제한합니다. - 제한을 초과하면 Driver가 결과를 모두 가져오기 전에 Job을 중단하므로, 대규모 결과가 Driver 메모리에 유입되는 것을 막을 수 있습니다. - 기본값인 1GB는 애플리케이션 하나만 실행하는 환경의 값이므로, 여러 세션이 동시에 결과를 가져가는 멀티세션 서버에서는 더 보수적으로 설정해야 합니다. - 이 설정만으로 Driver OOM이나 노드 장애까지 막을 수는 없습니다. ## 여러 Replica를 통한 장애 영향 축소 - Driver 자체의 OOM이나 노드 소실처럼 설정으로 막을 수 없는 장애에 대비해 Spark Connect 서버를 여러 Replica로 구성합니다. - 각 Replica는 독립적인 다음 요소를 갖습니다. - SparkContext - Driver - Executor - 한 Replica가 장애로 종료되어도 장애 범위가 해당 Replica에 한정되고, 다른 Replica가 새로운 세션 요청을 처리할 수 있습니다. - 단일 서버의 장애가 Spark Connect 전체로 확산되는 구조를 여러 독립 실행 단위로 나누는 것이 핵심입니다. ## 실용적인 운영 방향 - 멀티세션 Spark Connect에서는 전역 장애 카운터를 그대로 두지 말고, Task·Stage·Job 단위의 실패 제한으로 문제 쿼리를 격리하는 것이 안전합니다. - `spark.driver.maxResultSize`를 동시 세션 수와 쿼리 특성에 맞게 보수적으로 설정해야 합니다. - 스케줄러 Pool만으로는 CPU·메모리 격리가 불가능하므로, 강한 격리가 필요하면 Replica나 Kubernetes 리소스 정책을 활용해야 합니다. - 단일 Driver를 그대로 공유하기보다 여러 Replica를 운영해 장애의 영향 범위를 줄이는 것이 Production 환경에 적합합니다.

원문 읽기(새 탭에서 열림)
slack4분 읽기큐레이션 요약

SSH에서 REST로: Slack EMR 데이터 파이프라인의 보안 주도 현대화

Slack은 700개가 넘는 EMR 데이터 파이프라인을 직접 SSH로 실행하던 구조에서 REST 기반 작업 제출 방식으로 전환했다. SSH는 보안 공격 표면, 키 관리, 장애 복구, 작업 관측성 측면에서 한계가 있었고 Spark on Kubernetes와 AWS 계정 분리 같은 현대화도 가로막았다. Slack은 YARN REST API와 YARN Distributed Shell을 활용해 8개 데이터 리전의 작업을 중단 없이 마이그레이션하고 SSH를 완전히 제거했다. ## SSH 기반 파이프라인의 확산 - 2017년경 Airflow가 EMR 마스터 노드에 SSH로 접속해 명령을 실행하는 방식으로 데이터 파이프라인을 구축했다. - 단순한 `SSHOperator` 패턴이 확산되면서 다음과 같은 작업까지 SSH로 실행됐다. - Spark 및 MapReduce 작업 - AWS CLI 명령 - 사용자 정의 Python 스크립트 - 2024년에는 700개 이상의 운영 작업이 SSH 기반으로 실행되고 있었다. - 검색 인덱싱, 분석, 비즈니스 인텔리전스 등 핵심 데이터 처리도 이 구조에 의존했다. ## SSH의 보안 및 운영상 문제 - **보안 위험** - 오케스트레이션 워커가 EMR 클러스터에 직접 SSH 접속해야 해 공격 표면이 커졌다. - SSH 키를 여러 워커에 배포하고 주기적으로 교체해야 했다. - 세밀한 감사 추적을 위해 여러 시스템의 로그를 상호 연관해야 했다. - 보안 그룹과 사용자 권한 설정이 복잡해졌다. - **운영 장애** - 작업이 EMR 마스터 노드에서 직접 실행되어 리소스 경쟁이 발생했다. - Kubernetes Pod가 재시작되면 SSH 연결이 끊겨 작업이 실패했다. - 연결이 끊긴 뒤에도 작업이 계속 실행되는 ‘좀비 작업’이 남을 수 있었다. - 연결 단절 후 작업의 성공·실패 상태를 안정적으로 확인하기 어려웠다. - **인프라 현대화 차단** - Spark on Kubernetes와 EMR on EKS 도입을 시작할 수 없었다. - 메인 AWS 계정의 EMR 클러스터를 자식 계정으로 이전하는 Whitecastle 프로젝트가 지연됐다. - 신뢰할 수 있는 작업 모니터링과 관측성을 구현하기 어려웠다. ## REST 기반 작업 제출의 장점 - SSH는 클라이언트와 서버 사이의 상태ful 연결을 유지해야 한다. - REST 방식에서는 작업의 생명주기를 서버가 관리한다. - `POST`: 작업을 제출하고 작업 ID를 받음 - `GET`: 작업 ID로 실행·완료·실패 상태를 조회 - `DELETE`: 필요할 때 작업을 취소 - Airflow나 Kubernetes Pod가 재시작되어도 작업 자체는 서버에서 계속 실행될 수 있다. - 클라이언트가 작업 상태를 다시 조회할 수 있어 연결 단절에 강하다. - 작업 취소, 리소스 관리, 로그 확인 등도 실행 엔진의 표준 기능으로 처리할 수 있다. ## YARN Distributed Shell을 활용한 해결책 - Spark는 Livy REST API, Hive는 HiveServer2를 사용할 수 있어 상대적으로 이전이 쉬웠다. - 반면 MapReduce와 `aws s3 sync`, `hadoop distcp` 같은 300개 이상의 임의 CLI 작업은 바로 사용할 REST API가 없었다. - 검토한 대안은 다음과 같았다. - 원격 명령 실행용 커스텀 래퍼 서비스 - Ansible이나 Salt 같은 원격 실행 프레임워크 - YARN에 새로운 작업 유형을 직접 개발 - 이러한 방법은 별도 보안 계층과 운영 인프라를 구축·유지해야 해 복잡도가 높았다. - YARN의 **Distributed Shell**은 임의의 셸 스크립트를 YARN 컨테이너에서 실행할 수 있도록 했다. - 기존 YARN REST API를 그대로 사용 - YARN의 인증·인가 체계 활용 - 별도 보안 서비스 불필요 - 오픈소스 표준 기반 - 컨테이너 리소스와 작업 생명주기 관리 지원 ## Distributed Shell의 실행 흐름 - 실행할 셸 스크립트를 S3에 업로드한다. - 예: `s3://bucket/command.sh` - 스크립트는 `aws s3 sync` 같은 임의 명령을 포함할 수 있다. - YARN REST 요청에 Distributed Shell의 `ApplicationMaster`와 스크립트 위치를 지정한다. - YARN이 컨테이너를 할당하고 S3에서 스크립트를 내려받아 실행한다. - 실행 과정에서 YARN이 다음을 담당한다. - 메모리와 vCore 등 리소스 제한 - 컨테이너 격리 - 재시도와 장애 복구 - 정상적인 작업 취소 - YARN UI를 통한 로그 및 상태 확인 ## 마이그레이션의 의미 - YARN Distributed Shell을 통해 REST API가 없던 CLI·MapReduce 작업까지 동일한 실행 모델로 통합할 수 있었다. - 작업 제출과 실행을 SSH 연결에서 분리해 클라이언트 재시작과 네트워크 단절에 대한 안정성을 높였다. - 700개 이상의 작업을 8개 데이터 리전에 걸쳐 중단 없이 이전하면서 SSH 의존성을 제거했다. - 결과적으로 보안 강화뿐 아니라 Kubernetes 기반 실행 환경, AWS 계정 분리, 표준화된 모니터링으로 나아갈 기반을 마련했다. 실용적으로는 원격 서버에 직접 접속해 명령을 실행하기보다, 작업 ID·상태 조회·취소를 제공하는 서버 측 실행 모델을 사용하는 것이 바람직하다. 특히 기존 작업이 단순 CLI 스크립트라면 복잡한 신규 실행 서비스를 만들기 전에 YARN Distributed Shell처럼 기존 플랫폼의 표준 기능을 우선 검토할 수 있다.

원문 읽기(새 탭에서 열림)
line원문

Hive에서 Iceberg로: 데이터 반영 속도 12배 향상의 비밀 (새 탭에서 열림)

LINE Plus는 수억 건에 달하는 상품 데이터를 처리하기 위해 기존에 사용하던 전체 데이터 복제(Full Dump) 방식의 ETL 구조를 탈피하고, Apache Iceberg와 Apache Flink를 결합한 증분(Incremental) 처리 구조를 도입했습니다. 이를 통해 데이터 규모가 커질수록 기하급수적으로 늘어나던 업데이트 비용과 시간을 대폭 절감하였으며, 결과적으로 데이터 반영 주기를 60분에서 5분으로 단축하여 약 12배의 성능 향상을 이루어냈습니다. 이 글은 대규모 데이터 환경에서 실시간성에 가까운 데이터 최신성을 확보하기 위한 기술적 여정과 엔진 선택의 근거를 상세히 다룹니다. **기존 전체 데이터 복제 방식의 한계** * **리소스 낭비와 지연:** 매번 수억 건의 전체 데이터를 다시 써야 하는 구조로 인해 데이터 규모가 커질수록 처리 비용이 증가하고, 사내 Hadoop 리소스 부족 시 업데이트 주기가 지연되는 문제가 발생했습니다. * **데이터 최신성 결여:** 스냅숏 기반의 추출 방식은 정합성은 보장하지만, 추출 작업에 걸리는 시간만큼 데이터가 과거 시점에 머물게 되어 라이브 서비스에서의 실시간 대응이 어려웠습니다. * **운영 DB 부하:** 대용량 데이터를 한꺼번에 추출할 때 발생하는 막대한 디스크 I/O와 Undo 세그먼트 팽창은 운영 환경의 성능 저하를 유발하는 고질적인 원인이 되었습니다. **Apache Iceberg를 통한 증분 처리 기반 마련** * **테이블 형식의 변화:** 기존 Hive의 디렉터리 기반 관리 방식에서 벗어나, 메타데이터를 이용해 스냅숏 단위로 파일을 추적하는 Iceberg 형식을 도입했습니다. * **행 단위 업데이트 지원:** 전체 데이터를 다시 쓸 필요 없이 변경된 행(row)만 선택적으로 업데이트(upsert)하거나 삭제(delete)할 수 있어, 데이터 규모와 상관없이 일정한 업데이트 비용을 유지할 수 있게 되었습니다. **Apache Flink 선택의 결정적 이유** * **스테이트풀(Stateful) 처리를 통한 최신성 보장:** Flink의 DataStream API를 활용해 `updatedate`를 상태값으로 관리함으로써, 컨슈머 랙 등으로 인해 뒤늦게 도착한 과거 데이터가 최신 데이터를 덮어쓰는 문제를 원천 차단했습니다. * **2단계 커밋(2PC) 기반의 정확히 한 번 처리:** Iceberg 테이블 쓰기와 Kafka 상태 메시지 발행을 하나의 트랜잭션으로 묶어, 데이터 누락이나 중복 없이 '전부 아니면 전무(All-or-Nothing)'의 정합성을 보장했습니다. * **강력한 장애 허용(Fault Tolerance):** 체크포인트 메커니즘을 통해 시스템 장애 발생 시에도 마지막 성공 지점부터 즉시 복구가 가능하며, 관리하던 상태값을 유실 없이 유지할 수 있습니다. **효율적인 운영을 위한 쿠버네티스 오퍼레이터 도입** * **운영 자동화:** 설정 작업을 수동으로 진행해야 하는 네이티브 쿠버네티스 방식 대신, Flink 쿠버네티스 오퍼레이터를 도입하여 라우팅, 웹 UI 구성 등 운영 요소를 커스텀 리소스로 추상화하고 관리를 자동화했습니다. * **격리 및 확장성:** 애플리케이션 모드를 통해 잡(job)별 클러스터 격리 수준을 높이고, 헬름(Helm) 차트를 이용해 손쉽게 배포 및 확장할 수 있는 환경을 구축했습니다. 대규모 데이터셋에서 실시간에 가까운 데이터 동기화와 엄격한 정합성이 모두 필요하다면, 단순한 배치 처리보다는 Flink와 Iceberg의 조합을 통한 증분 파이프라인 구축을 권장합니다. 특히 Flink의 2단계 커밋과 체크포인트 기능을 활용하면 분산 환경에서도 데이터 무결성을 보장하면서 시스템의 업데이트 주기를 획기적으로 단축할 수 있습니다.

pinterest원문

Pinterest의 Apache Spark에서 (새 탭에서 열림)

Pinterest는 대규모 Spark 환경에서 빈번하게 발생하는 OOM(Out-of-Memory) 오류를 해결하기 위해 'Auto Memory Retries' 기능을 도입했습니다. 이 시스템은 태스크 수준에서 리소스 요구량을 동적으로 판단하고, 실패 시 더 큰 메모리 프로필을 가진 실행기(Executor)에서 태스크를 재시도하도록 자동화합니다. 이를 통해 수동 튜닝의 번거로움을 줄이고 자원 효율성을 높여 전체적인 작업 실패율과 운영 비용을 획기적으로 낮추는 성과를 거두었습니다. ### 기존 Spark 리소스 관리의 한계와 문제점 * Pinterest의 Spark 클러스터는 하드웨어 대비 높은 메모리 요구량으로 인해 OOM 오류가 잦았으며, 전체 작업 실패 원인의 약 4.6%가 메모리 부족에서 기인했습니다. * 사용자가 모든 스테이지와 태스크의 메모리 요구량을 정확히 예측하여 수동으로 설정하는 것은 매우 어렵고 시간이 많이 소요되는 작업입니다. * 데이터 스큐(Skew) 현상으로 인해 같은 스테이지 내에서도 특정 태스크만 과도한 메모리를 사용하는 경우가 많아, 모든 태스크를 최대치에 맞춰 설정하면 심각한 자원 낭비가 발생합니다. * 제품 팀의 우선순위 문제로 인해 비용 절감을 위한 수동 최적화가 지속적으로 이루어지기 어려운 구조적 한계가 있었습니다. ### Auto Memory Retries의 단계별 대응 전략 * **CPU 할당량 증설을 통한 메모리 확보 (1단계):** 실행기에 1개 이상의 코어가 있는 경우, OOM 발생 시 첫 번째 재시도에서 태스크당 CPU 할당량(`spark.task.cpus`)을 두 배로 늘립니다. 이를 통해 실행기 내 동시 실행 태스크 수를 줄여 개별 태스크가 사용할 수 있는 공유 메모리 공간을 즉각적으로 확보합니다. * **물리적으로 큰 실행기 투입 (2단계):** CPU 조절만으로 해결되지 않거나 단일 태스크가 이미 실행기 전체 메모리를 사용 중인 경우, 물리적으로 더 큰 메모리를 가진 새로운 실행기를 동적으로 런칭합니다. * **하이브리드 확장 프로필 적용:** 기본 설정의 2배, 3배, 4배 크기의 리소스 프로필을 미리 등록하고 단계별로 순차 적용합니다. Apache Gluten을 사용하는 워크로드의 경우 Off-heap 메모리도 함께 증설하여 가속화된 연산을 지원합니다. ### 시스템 구현 및 Spark 엔진 확장 * **태스크 수준의 리소스 프로필:** 기존 Spark의 고정된 리소스 할당 방식에서 벗어나, `Task` 객체에 개별 리소스 프로필 ID(`taskRpId`)를 저장할 수 있도록 확장하여 동일한 TaskSet 내에서도 태스크마다 사양을 다르게 가질 수 있게 구현했습니다. * **스케줄링 로직 최적화:** `TaskSetManager`는 OOM 감지 시 즉시 상위 프로필을 할당하며, `TaskSchedulerImpl`은 증설된 CPU 속성을 가진 태스크를 기존 실행기에서 우선 실행할 수 있게 하여 리소스 재사용 속도를 높였습니다. * **동적 리소스 할당:** `ExecutorAllocationManager`가 상위 프로필을 필요로 하는 대기 태스크를 실시간으로 추적하고, 물리적으로 큰 실행기가 필요한 시점에 맞춰 Kubernetes 등에 자원을 요청합니다. * **사용자 경험 개선:** 사용자가 어떤 태스크가 더 많은 자원을 사용했는지 쉽게 파악할 수 있도록 Spark UI의 태스크 목록에 리소스 프로필 ID를 표시하는 기능을 추가했습니다. 효율적인 Spark 운영을 위해서는 모든 작업을 최대 메모리 요구량에 맞추기보다, 상위 90%(P90) 수준의 일반적인 설정으로 실행하고 예외적인 태스크만 'Auto Memory Retries'로 구제하는 탄력적 전략이 권장됩니다. 이는 데이터 스큐가 심한 대규모 파이프라인에서 운영 안정성을 확보함과 동시에 인프라 비용을 최적화할 수 있는 강력한 해법이 될 것입니다.

pinterest원문

Pinterest의 차세대 DB 인 (새 탭에서 열림)

Pinterest는 기존의 파편화된 배치 기반 DB 적재 시스템을 개선하기 위해 Iceberg와 CDC(Change Data Capture) 기술을 결합한 통합 프레임워크를 구축했습니다. 이 시스템은 데이터 지연 시간을 24시간 이상에서 수 분 단위로 단축하고, 변경된 데이터만 처리하는 방식으로 인프라 비용을 획기적으로 절감했습니다. 이를 통해 분석, 머신러닝, 규정 준수 등 현대적인 데이터 요구사항에 기민하게 대응할 수 있는 고성능 데이터 생태계를 마련했습니다. ### 통합 CDC 프레임워크의 계층 구조 * **CDC 레이어**: Debezium 및 TiCDC를 활용해 MySQL, TiDB, KVStore의 변경 사항을 1초 미만의 지연 시간으로 포착하여 Kafka에 기록합니다. * **스트리밍 레이어**: Flink 작업이 Kafka의 이벤트를 실시간으로 처리하여 S3에 위치한 'CDC Iceberg 테이블'에 추가 전용(Append-only) 방식으로 저장합니다. * **배치 레이어**: Spark 작업이 주기적으로(15~60분) CDC 테이블의 최신 변경 사항을 읽어 `Merge Into` 구문을 통해 최종 'Base Iceberg 테이블'에 업서트(Upsert)를 수행합니다. * **부트스트랩 및 유지보수**: 초기 데이터 로드를 위한 전용 파이프라인과 소형 파일 압축(Compaction) 및 스냅샷 만료 관리를 위한 유지보수 작업을 포함합니다. ### CDC 테이블과 베이스 테이블의 이원화 관리 * **CDC 테이블**: 모든 변경 이력을 담은 시계열 원장으로, 5분 미만의 지연 시간을 유지하며 원천 데이터의 변경 로그를 보존합니다. * **베이스 테이블**: 온라인 DB의 현재 상태를 그대로 반영하는 스냅샷 테이블입니다. CDC 테이블로부터 최신 레코드를 추출하여 정합성을 맞춥니다. * **동기화 로직**: `ROW_NUMBER()` 함수를 활용해 기본 키(PK)별로 가장 최신 업데이트(최근 타임스탬프 및 GTID 기준)를 식별한 후, 삭제 유형은 제거하고 나머지는 업데이트 또는 삽입합니다. ### 성능 및 비용 최적화 전략 * **Merge-on-Read (MOR) 방식 채택**: Copy-on-Write(COW) 방식은 업데이트 시 대규모 파일을 다시 작성해야 하므로 스토리지와 계산 비용이 높습니다. Pinterest는 비용 효율성을 극대화하기 위해 MOR 방식을 표준 전략으로 선택했습니다. * **기본 키 해시 버킷팅(Bucketing)**: 베이스 테이블을 PK의 해시값(예: `bucket(100, id)`)으로 파티셔닝하여 Spark가 업서트 작업을 병렬로 효율적으로 처리할 수 있도록 설계했습니다. * **증분 처리 효율성**: 매일 전체 테이블을 덤프하던 방식에서 변경된 데이터(통상 5% 미만)만 처리하는 방식으로 전환하여 연산 리소스 낭비를 차단했습니다. 방대한 양의 데이터베이스를 데이터 레이크로 통합할 때는 Iceberg의 `Merge Into` 기능을 활용한 증분 업데이트가 필수적입니다. 특히 읽기 성능과 쓰기 비용 사이의 균형을 위해 MOR 전략을 사용하고, 쓰기 병목을 해소하기 위해 기본 키 기반의 버킷팅을 적용하는 것이 실무적으로 매우 효과적인 접근임을 보여줍니다.

pinterest원문

PinLanding: 멀티모달 (새 탭에서 열림)

Pinterest의 'PinLanding'은 수십억 개의 제품 데이터를 멀티모달 AI를 통해 정교한 쇼핑 컬렉션으로 자동 변환하는 프로덕션 파이프라인입니다. 기존의 수동 큐레이션이나 단순 검색 기록 기반 방식에서 벗어나, 제품의 이미지와 텍스트를 직접 분석하여 사용자의 복잡하고 긴 꼬리형(Long-tail) 검색 의도에 맞는 컬렉션을 생성합니다. 이 시스템은 비전-언어 모델(VLM)을 통한 속성 추출과 CLIP 스타일의 효율적인 임베딩 모델을 결합하여 대규모 데이터셋에서도 정밀도와 확장성을 동시에 확보했습니다. **사용자 쇼핑 의도와 데이터 신호의 특성화** * 사용자의 검색 기록, 자동 완성 상호작용, 필터 사용 패턴을 분석하여 쇼핑 의도의 분포를 파악합니다. * '검은색 칵테일 드레스'와 같은 정형화된 주요 쿼리(Head)뿐만 아니라, '이탈리아 여름 휴가 때 입을 옷'과 같은 서술형 및 대화형 쿼리에 대응하는 것을 목표로 합니다. * 색상, 상황, 스타일, 핏 등 20개 카테고리에 걸친 속성 차원을 정의하여, 수요는 높지만 기존 검색 결과가 부족한 영역을 식별합니다. **VLM과 LLM-as-Judge를 활용한 쇼핑 토픽 정제** * 제품의 이미지와 메타데이터를 비전-언어 모델(VLM)에 입력하여 정규화된 키-값 쌍 형태의 속성을 생성합니다. * 초기 VLM 출력의 너무 구체적이거나 중복된 속성(예: 'boho'와 'bohemian')을 해결하기 위해 빈도 기반 필터링과 임베딩 기반 클러스터링을 수행합니다. * 최종적으로 'LLM-as-judge' 단계를 거쳐 추출된 속성들이 실제 쇼핑 의도와 일치하는지, 의미적으로 일관성이 있는지 평가하여 고품질의 쇼핑 토픽 사전을 구축합니다. **CLIP 스타일 모델을 통한 대규모 속성 할당** * 모든 제품에 VLM을 직접 적용하는 것은 비용이 과다하므로, 이미지-텍스트를 정렬하는 CLIP 스타일의 듀얼 인코더 모델을 별도로 학습시킵니다. * 제품 인코더와 속성 구절 인코더를 통해 각각의 임베딩을 생성하고, 두 벡터 간의 유사도가 임계치를 넘을 때 속성을 할당합니다. * 이 방식은 VLM 대비 연산 비용을 획기적으로 낮추면서도, 제품별 속성 밀도를 높여 더욱 일관된 제품-속성 그래프를 형성합니다. **Ray 및 Spark 기반의 효율적인 배치 추론 및 피드 구축** * 수백만 개의 핀(Pin)과 토픽을 처리하기 위해 Ray 프레임워크를 사용하여 GPU와 CPU 리소스를 독립적으로 확장하며 스트리밍 방식으로 추론을 수행합니다. * CLIP 기반 분류기는 8개의 NVIDIA A100 GPU에서 약 12시간 만에 학습 및 추론을 완료하며, 회당 비용을 약 500달러 수준으로 절감했습니다. * 최종 피드 구성은 Apache Spark를 활용하여 제품과 쇼핑 토픽 간의 속성 유사도를 계산하고, 가중치 기반 스코어링을 통해 관련성 높은 제품들을 컬렉션으로 묶어냅니다. PinLanding 시스템은 AI가 단순한 키워드 매칭을 넘어 제품의 시각적, 맥락적 의미를 깊이 있게 이해할 수 있음을 보여줍니다. 대규모 이커머스 환경에서 사용자에게 개인화되고 탐색 가능한 쇼핑 경험을 제공하려는 기업은 VLM을 통한 '지식 추출'과 CLIP 스타일 모델을 통한 '효율적 확산' 전략을 참고할 가치가 있습니다.

daangn원문

당근 데이터 지도를 그리다: 컬럼 레벨 리니지 구축기 (새 탭에서 열림)

당근마켓(당근) 데이터 가치화팀은 데이터의 흐름을 투명하게 파악하여 신뢰성을 높이기 위해 SQL 파싱 기반의 **컬럼 레벨 데이터 리니지(Column-level Lineage)** 시스템을 구축했습니다. 기존의 테이블 단위 추적으로는 해결하기 어려웠던 연쇄 장애 대응과 민감 정보(PII) 관리 문제를 해결하기 위해, 모든 BigQuery 쿼리 로그를 분석하여 데이터 간의 세부 의존 관계를 시각화했습니다. 이를 통해 당근의 복잡한 데이터 생태계에서 변경 영향도를 정교하게 분석하고 장애 복구 시간을 단축하는 성과를 거두었습니다. ### 데이터 흐름의 불투명성으로 인한 문제점 * **연쇄 실패 대응의 어려움**: 특정 테이블의 파이프라인이 실패했을 때 이를 참조하는 하위 테이블들을 즉각 파악할 수 없어, 수동으로 쿼리를 전수 조사하며 문제를 해결해야 했습니다. * **스키마 변경의 불확실성**: 원천 데이터(MySQL 등)의 컬럼을 삭제하거나 타입을 변경할 때, 해당 컬럼을 사용하는 수많은 파생 테이블 중 어떤 곳에 장애가 발생할지 예측하기 어려웠습니다. * **민감 정보 추적 불가**: PII(개인정보)가 여러 가공 단계를 거치며 어떤 테이블의 어떤 컬럼으로 흘러가는지 파악되지 않아 보안 관리 측면에서 한계가 있었습니다. ### 컬럼 레벨 리니지 도입의 기술적 의사결정 * **테이블 레벨의 한계**: BigQuery의 기본 기능을 통한 테이블 단위 추적은 뷰(View)의 기저 테이블을 정확히 파악하기 어렵고, 세부 컬럼의 변화를 감지하지 못하는 단점이 있었습니다. * **오픈소스(OpenLineage) 대비 효율성**: 다양한 조직이 각기 다른 환경(Airflow, 노트북 등)에서 쿼리를 실행하는 당근의 특성상, 모든 환경에 계측 코드를 심는 방식보다는 중앙화된 BigQuery 로그를 분석하는 방식이 운영 부담이 적다고 판단했습니다. * **SQL 파싱 접근법**: 실행된 모든 SQL의 이력이 남는 `INFORMATION_SCHEMA.JOBS` 뷰를 활용하여, 실행 환경과 관계없이 모든 쿼리로부터 의존성을 추출하는 방식을 채택했습니다. ### 시스템 아키텍처 및 추출 프로세스 * **기술 스택**: 대량의 쿼리 병렬 처리를 위해 **Spark**를 활용하고, SQL 파싱 및 AST(Abstract Syntax Tree) 분석을 위해 **sqlglot** 라이브러리를 사용하며, **Airflow**로 주기적인 추출 프로세스를 자동화했습니다. * **데이터 수집 및 분석**: 모든 GCP 프로젝트에서 쿼리 로그를 수집한 뒤, sqlglot으로 쿼리 구조를 분석하여 `Source Column -> Target Column` 관계를 도출합니다. * **엣지 케이스 처리**: `SELECT *`와 같은 와일드카드 쿼리는 테이블 메타데이터를 결합해 실제 컬럼명으로 확장하고, 복잡한 CTE(Common Table Expressions)나 서브쿼리 내의 의존성도 AST 탐색을 통해 정확하게 추적합니다. ### 데이터 지도를 통한 실질적 변화 * **정교한 영향도 분석**: 특정 컬럼 수정 시 다운스트림에서 이를 참조하는 모든 컬럼을 즉시 확인하여 사전에 장애를 예방할 수 있게 되었습니다. * **거버넌스 강화**: 데이터의 원천부터 최종 활용 단계까지의 흐름을 시각화함으로써 데이터 가계도(Data Genealogy)를 완성하고, 데이터 보안 및 품질 관리 수준을 한 단계 높였습니다. * **운영 효율화**: 장애 발생 시 영향 범위를 데이터 지도를 통해 한눈에 파악함으로써 원인 파악과 복구에 소요되는 리소스를 획기적으로 줄였습니다. 데이터 플랫폼의 규모가 커질수록 수동 관리는 불가능해지므로, 초기부터 SQL 로그를 활용한 자동화된 리니지 체계를 구축하는 것이 중요합니다. 특히 실행 환경이 파편화된 조직일수록 애플리케이션 계측보다는 쿼리 엔진의 로그를 파싱하는 접근법이 빠른 도입과 높은 커버리지를 확보하는 데 유리합니다.

line원문

동적 사용자 분할을 활용한 새로운 A/B 테스트 시스템을 소개합니다 (새 탭에서 열림)

동적 유저 세분화(Dynamic User Segmentation) 기술을 도입한 새로운 A/B 테스트 시스템은 사용자 ID 기반의 단순 무작위 배분을 넘어 특정 속성과 행동 패턴을 가진 정교한 사용자 그룹을 대상으로 실험을 수행할 수 있게 합니다. 이 시스템은 타겟팅 엔진과 테스트 할당 로직을 분리하여 데이터 기반의 의사결정 범위를 개인화된 영역까지 확장하며, 서비스 품질 향상과 리소스 최적화라는 두 가지 목표를 동시에 달성합니다. 결과적으로 개발자와 마케터는 복잡한 사용자 시나리오에 대해 더욱 정확하고 신뢰할 수 있는 실험 데이터를 얻을 수 있습니다. ### 기존 A/B 테스트 방식과 고도화의 필요성 * **무작위 배분의 특징**: 일반적인 시스템은 사용자 ID를 해싱하여 실험군과 대조군으로 무작위 할당하며, 구현이 쉽고 선택 편향(Selection Bias)을 줄일 수 있다는 장점이 있습니다. * **타겟팅의 한계**: 전체 사용자를 대상으로 하는 일반적인 테스트에는 적합하지만, '오사카에 거주하는 iOS 사용자'처럼 특정 조건을 충족하는 집단만을 대상으로 하는 정교한 실험에는 한계가 있습니다. * **고도화된 시스템의 목적**: 사용자 세그먼트를 동적으로 정의함으로써, 서비스의 특정 기능이 특정 사용자 층에게 미치는 영향을 정밀하게 측정하기 위해 도입되었습니다. ### 유저 세분화를 위한 타겟팅 시스템 아키텍처 * **데이터 파이프라인**: HDFS에 저장된 사용자 정보(UserInfo), 모바일 정보(MobileInfo), 앱 활동(AppActivity) 등의 빅데이터를 Spark를 이용해 분석하고 처리합니다. * **세그먼트 연산**: Spark의 RDD 기능을 활용하여 합집합(Union), 교집합(Intersect), 차집합(Subtract) 등의 연산을 수행하며, 이를 통해 복잡한 사용자 조건을 유연하게 조합할 수 있습니다. * **데이터 저장 및 조회**: 처리된 결과는 `{user_id}-{segment_id}` 형태의 키-값 쌍으로 Redis에 저장되어, 실시간 요청 시 매우 낮은 지연 시간으로 해당 사용자의 세그먼트 포함 여부를 확인합니다. ### 효율적인 실험 관리와 할당 프로세스 * **설정 관리(Central Dogma)**: 실험의 설정값은 오픈 소스 설정 저장소인 Central Dogma를 통해 관리되며, 이를 통해 코드 수정 없이 실시간으로 실험 설정을 변경하고 동기화할 수 있습니다. * **할당 로직(Test Group Assigner)**: 클라이언트의 요청이 들어오면 할당기는 Central Dogma에서 실험 정보를 가져오고, Redis를 조회하여 사용자가 타겟 세그먼트에 속하는지 확인한 후 최종 실험군을 결정합니다. * **로그 및 분석**: 할당된 그룹 정보는 로그 스토어에 기록되어 사후 분석 및 대시보드 시각화의 기초 자료로 활용됩니다. ### 주요 활용 사례 및 향후 계획 * **콘텐츠 및 위치 추천**: 특정 사용자 세그먼트에 대해 서로 다른 머신러닝(ML) 모델의 성능을 비교하여 최적의 추천 알고리즘을 선정합니다. * **마케팅 및 온보딩**: 구매 빈도가 낮은 '라이트 유저'에게만 할인 쿠폰 효과를 테스트하거나, '신규 가입자'에게만 온보딩 화면의 효과를 측정하여 불필요한 비용을 줄이고 효율을 높입니다. * **플랫폼 확장성**: 향후에는 LY Corporation 내의 다양한 서비스로 플랫폼을 확장하고, 실험 생성부터 결과 분석까지 한 곳에서 관리할 수 있는 통합 어드민 시스템을 구축할 계획입니다. 이 시스템은 실험 대상자를 정교하게 선별해야 하는 복잡한 서비스 환경에서 데이터의 신뢰도를 높이는 데 매우 효과적입니다. 특히 마케팅 비용 최적화나 신규 기능의 타겟 검증이 필요한 팀이라면, 단순 무작위 할당 방식보다는 유저 세그먼트 기반의 동적 타겟팅 시스템을 구축하거나 활용하는 것을 권장합니다.

naver원문

네이버 TV (새 탭에서 열림)

네이버의 실시간 거래 리포트 시스템은 대규모 데이터를 다양한 조건으로 빠르게 조회하기 위해 Apache Iceberg와 StarRocks의 Materialized View를 핵심 기술로 활용합니다. 단순히 데이터를 적재하는 수준을 넘어, 데이터의 최신성(Freshness)과 저지연(Low-Latency) 응답 속도, 그리고 시스템 확장성을 동시에 확보하는 것이 이번 기술 여정의 핵심 결론입니다. 이를 통해 복잡한 다차원 필터링이 필요한 비즈니스 환경에서도 사용자에게 즉각적인 분석 결과를 제공하는 데이터 레이크하우스 아키텍처를 구현했습니다. **실시간 거래 리포트의 기술적 도전 과제** * 대규모로 발생하는 거래 데이터를 실시간에 가깝게 수집하면서도, 사용자가 원하는 다양한 검색 조건에 즉각 응답해야 하는 성능적 요구사항이 있었습니다. * 데이터의 양이 방대해짐에 따라 기존의 단순 조회 방식으로는 응답 속도가 저하되는 문제가 발생했으며, 데이터의 신선도와 쿼리 성능 사이의 트레이드오프를 해결해야 했습니다. * 다차원 필터링과 집계 연산이 빈번한 리포트 특성상, 인덱싱 최적화와 리소스 효율성을 동시에 고려한 설계가 필요했습니다. **Iceberg와 StarRocks를 활용한 저지연 쿼리 전략** * **Apache Iceberg 기반 데이터 관리**: 데이터 레이크의 스토리지 포맷으로 Iceberg를 채택하여 ACID 트랜잭션을 보장하고, 대규모 데이터셋에 대한 효율적인 스키마 진화와 파티션 관리를 수행합니다. * **StarRocks의 구체화 뷰(Materialized View) 도입**: Iceberg에 저장된 원본 데이터를 직접 조회하는 대신, StarRocks의 Materialized View를 활용해 자주 사용되는 쿼리 결과를 미리 연산하여 저장함으로써 조회 속도를 비약적으로 향상시켰습니다. * **증분 업데이트 및 동기화**: 실시간으로 유입되는 데이터를 Materialized View에 효율적으로 반영하기 위해 Spark와 StarRocks 간의 연동 최적화를 진행하여 데이터의 최신성을 유지합니다. **아키텍처 구성 요소 및 운영 최적화** * **Spark**: 대용량 거래 데이터의 가공 및 Iceberg 테이블로의 수집을 담당하는 컴퓨팅 엔진으로 활용됩니다. * **StarRocks**: 고성능 OLAP 엔진으로서 Iceberg 외부에 위치하며, Materialized View를 통해 복잡한 조인(Join)과 집계(Aggregation) 쿼리를 가속화합니다. * **확장성 확보**: 데이터 노드와 컴퓨팅 리소스를 분리하여 운영함으로써 트래픽 증가에 유연하게 대응할 수 있는 구조를 설계했습니다. 대용량 실시간 분석 시스템을 구축할 때 Apache Iceberg만으로는 쿼리 성능의 한계가 있을 수 있으므로, StarRocks와 같은 고성능 OLAP 엔진의 구체화 뷰를 결합하는 레이크하우스 전략이 효과적입니다. 특히 데이터의 최신성이 중요한 금융 및 거래 리포트 분야에서 이와 같은 기술 조합은 인프라 비용을 절감하면서도 사용자 경험을 극대화할 수 있는 강력한 대안이 됩니다.

netflix원문

스케일링 뮤즈: 넷플릭스가 조단위 데이터에서 데이터 기반 창의적 인사이트를 제공하는 방법 | 넷플릭스 기술 블로그 | 넷플릭스 기술 블로그 (새 탭에서 열림)

넷플릭스의 내부 데이터 분석 플랫폼인 'Muse'는 수조 건 규모의 데이터를 분석하여 홍보용 미디어(아트웍, 영상 클립)의 효과를 측정하고 창작 전략을 지원합니다. 급증하는 데이터 규모와 복잡한 다대다(Many-to-Many) 필터링 요구사항을 해결하기 위해, 넷플릭스는 HyperLogLog(HLL) 스케치와 인메모리 기술인 Hollow를 도입하여 데이터 서빙 레이어를 혁신했습니다. 이를 통해 데이터 정확도를 유지하면서도 수조 행의 데이터를 실시간에 가깝게 처리할 수 있는 고성능 OLAP 환경을 구축했습니다. ### 효율적인 고유 사용자 집계를 위한 HLL 스케치 도입 * **근사치 계산을 통한 성능 최적화:** 고유 사용자 수(Distinct Count)를 계산할 때 발생하는 막대한 리소스 소모를 줄이기 위해 Apache Datasketches의 HLL 기술을 도입했습니다. 약 0.8%~2%의 미세한 오차를 허용하는 대신 집계 속도를 비약적으로 높였습니다. * **단계별 스케치 생성:** Druid 데이터 수집 단계에서 '롤업(Rollup)' 기능을 사용해 데이터를 사전 요약하고, Spark ETL 과정에서는 매일 생성되는 HLL 스케치를 기존 데이터와 병합(hll_union)하여 전체 기간의 통계를 관리합니다. * **데이터 규모 축소:** 수개월에서 수년 치의 데이터를 전수 비교하는 대신, 미리 생성된 스케치만 결합하면 되므로 데이터 처리량과 저장 공간을 획기적으로 절감했습니다. ### Hollow를 활용한 인메모리 사전 집계 및 서빙 * **초저지연 조회 구현:** 모든 쿼리를 Druid에서 처리하는 대신, 자주 사용되는 '전체 기간(All-time)' 집계 데이터는 넷플릭스의 오픈소스 기술인 'Hollow'를 통해 인메모리 방식으로 서빙합니다. * **Spark와 마이크로서비스의 연계:** Spark 작업에서 미리 계산된 HLL 스케치 집계 데이터를 Hollow 데이터셋으로 발행하면, Spring Boot 기반의 마이크로서비스가 이를 메모리에 로드하여 밀리초(ms) 단위의 응답 속도를 제공합니다. * **조인(Join) 병목 해결:** 복잡한 시청자 성향(Audience Affinity) 필터링과 같은 다대다 관계 연산을 메모리 내에서 처리함으로써 기존 아키텍처의 한계를 극복했습니다. ### 데이터 검증 및 아키텍처 현대화 * **신뢰성 보장:** 아키텍처 변경 전후의 데이터 정합성을 확인하기 위해 내부 디버깅 도구를 활용하여 사전/사후 데이터를 정밀하게 비교 검증했습니다. * **기술 스택 고도화:** React 프런트엔드와 GraphQL 레이어, 그리고 gRPC 기반의 Spring Boot 마이크로서비스 구조를 통해 확장성 있는 시스템을 구축했습니다. * **분석 역량 강화:** 이를 통해 단순한 대시보드를 넘어 이상치 감지(Outlier Detection), 미디어 간 성과 비교, 고급 필터링 등 사용자들의 고도화된 분석 요구를 수용할 수 있게 되었습니다. 대규모 OLAP 시스템을 설계할 때 모든 데이터를 실시간으로 전수 계산하기보다는, HLL과 같은 확률적 자료구조와 Hollow 기반의 인메모리 캐싱을 적절히 조합하는 것이 성능 최적화의 핵심입니다. 특히 수조 건 규모의 데이터에서는 완벽한 정확도와 성능 사이의 트레이드오프를 전략적으로 선택하는 것이 시스템의 유연성을 결정짓습니다.