apache-cassandra

3 개의 포스트

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 같은 각 도메인에 최적화된 처리를 구현할 수 있다.

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

시계열 워크로드를 위한 동적 재파티셔닝

Netflix의 TimeSeries Abstraction은 Cassandra를 이용해 페타바이트 규모의 시계열 데이터를 밀리초 단위로 처리하지만, 시간이 지날수록 커지는 광범위한 파티션이 지연시간과 타임아웃을 유발했다. 기존의 테이블 단위 재파티셔닝은 전체 데이터가 비슷한 문제를 보일 때는 효과적이지만, 일부 ID만 비정상적으로 커지는 경우에는 적합하지 않았다. 이를 해결하기 위해 Netflix는 읽기 경로에서 넓은 파티션을 감지하고, 해당 TimeSeries ID 단위로 비동기 분할한 뒤 읽기 요청을 자동으로 재라우팅하는 동적 파티셔닝 시스템을 구축했다. ## Cassandra와 넓은 파티션의 문제 - Cassandra는 높은 처리량, 낮은 지연시간, 비용 효율성, 운영 성숙도를 이유로 Netflix의 시계열 저장소로 사용된다. - 그러나 이벤트가 시간에 따라 누적되면 하나의 파티션이 지나치게 커질 수 있다. - 일반적인 읽기 지연은 수 밀리초 수준이지만, 넓은 파티션에서는 특히 데이터 끝부분에서 지연시간이 수 초까지 증가한다. - 이로 인해 요청 타임아웃, 높은 CPU 사용률, 가비지 컬렉션 일시정지, 스레드 큐 대기 등이 발생할 수 있다. - 단순히 Cassandra 클러스터를 확장하는 방식은 비용 문제를 해결하지 못하므로 데이터 배치 자체를 개선해야 한다. ## 기존 TimeSeries 파티셔닝 전략 - 데이터를 일정한 시간 단위의 개별 Time Slice로 나누어 파티션 크기를 제한한다. - 시간 기준으로 데이터를 효율적으로 조회하거나 삭제할 수 있으며, 대량의 tombstone을 처리해야 하는 부담도 줄어든다. - 데이터셋 생성 시 사용자가 예상 트래픽과 이벤트 특성을 입력한다. - 프로비저닝 파이프라인은 해당 입력을 바탕으로 Monte Carlo 시뮬레이션을 수행해 인프라와 파티션 설정을 결정한다. ## 사전 설정 방식의 한계 - 초기 단계에서는 실제 운영 트래픽을 정확히 예측하기 어렵다. - 시간이 지나면서 트래픽 패턴, 클라이언트 동작, 제품 요구사항이 변할 수 있다. - 일부 TimeSeries ID만 다른 ID보다 훨씬 많은 이벤트를 받는 데이터 이상치가 존재할 수 있다. - Time Slice마다 다른 파티션 전략을 적용할 수 있지만, 수천 개 데이터셋의 설정을 사람이 직접 조정하는 것은 지속 가능하지 않다. - 따라서 파티션 상태를 관찰하고 자동으로 조정하는 시스템이 필요하다. ## Time Slice 단위 재파티셔닝 - Cassandra의 `nodetool tablehistograms` 등 introspection API를 활용해 파티션 크기 분포를 관찰한다. - 너무 작은 파티션이 많은 과도한 분할(over-partitioning)과 지나치게 큰 파티션을 모두 탐지할 수 있다. - 예를 들어 파티션 크기가 10KB보다 작으면 읽기 증폭과 스레드 큐잉이 커질 수 있다. - 백그라운드 워커가 애플리케이션에 연결된 Time Slice의 파티션 히스토그램을 감시한다. - 파티션 크기가 설정된 목표 밀도에 미달하면 조정 계수를 계산하고, 이후 생성될 Time Slice의 `time_bucket` 간격을 변경한다. - 목표 파티션 크기는 워크로드에 따라 보통 2MiB~10MiB로 설정된다. - 이 방식은 읽기 지연시간과 타임아웃을 줄이는 데 효과가 있었지만, 테이블 전체가 비슷한 문제를 보일 때만 적합하다. - 특정 ID 몇 개만 넓은 경우에는 전체 테이블의 파티션 전략을 바꾸는 것이 불필요하거나 효과적이지 않다. ## 일부 ID만 문제가 될 때의 대응 - **아무것도 하지 않기** - 애플리케이션의 전체 지표에 영향이 없다면 문제를 감수하는 것이 합리적일 수 있다. - **부분 결과 반환** - 요청이 설정된 지연시간 SLO를 넘으면 진행 중인 요청을 중단한다. - 그때까지 수집한 데이터만 반환해, 전체 결과보다 빠른 응답을 우선하는 클라이언트에 적합하다. - **문제 ID 차단** - 테스트나 스팸 데이터처럼 시스템을 불안정하게 만드는 ID를 차단한다. - 다만 정상적이고 중요한 ID가 큰 경우에는 데이터 전체를 처리해야 하므로 이 방법을 사용할 수 없다. ## ID별 동적 파티셔닝 - 동적 파티셔닝은 테이블 전체가 아니라 특정 TimeSeries ID의 넓은 파티션만 자동으로 분할한다. - 비동기 파이프라인은 다음 세 단계로 구성된다. - **감지:** 읽기 경로에서 특정 파티션의 읽기 바이트 수를 추적하고, 임계치를 넘으면 넓은 파티션으로 판단한다. - **계획 및 분할:** 적절한 크기가 되도록 파티션 분할 작업을 계획하고 비동기적으로 실행한다. - **읽기 제공:** 분할이 완료되면 기존 요청 경로를 투명하게 새 파티션으로 재라우팅한다. - 읽기 작업 중 파티션에서 읽은 바이트가 설정된 한도를 초과하면 Kafka로 감지 이벤트를 발행한다. - 이벤트에는 데이터가 속한 Time Slice 테이블, 문제가 된 `time_series_id`, 기존 `time_bucket`, `event_bucket`, 해당 파티션의 쓰기 종료 여부(`immutable`), 버전 등의 정보가 포함된다. - 이 구조를 통해 정상적인 ID에는 영향을 주지 않으면서, 데이터량이 큰 특정 ID만 선택적으로 재구성할 수 있다. ## 실용적인 결론 - 데이터 전체의 분포가 바뀌었다면 Time Slice 단위 자동 재파티셔닝을 적용하는 것이 효율적이다. - 일부 ID만 비대해지는 경우에는 ID별 동적 파티셔닝이 더 적합하다. - 읽기량, 파티션 크기, 지연시간 SLO를 지속적으로 관찰하고 자동화된 감지·분할·재라우팅 체계를 마련하는 것이 핵심이다.

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

비디오 검색을 위한 멀티모달 인텔리전스 구현 (새 탭에서 열림)

넷플릭스는 방대한 분량의 원본 영상 데이터에서 창작자가 원하는 특정 순간을 신속하게 찾아낼 수 있도록 여러 전문 AI 모델을 결합한 멀티모달(Multimodal) 검색 시스템을 구축했습니다. 이 시스템은 캐릭터, 환경, 대화 등 서로 다른 모델이 생성한 파편화된 신호들을 하나의 통합된 시간축으로 동기화하여 고차원의 문맥 이해와 실시간 검색을 동시에 실현합니다. 결과적으로 수십억 개의 데이터 포인트 속에서도 창작자의 의도에 부합하는 장면을 지연 시간 없이 정확하게 찾아내는 기술적 해결책을 제시합니다. **비디오 검색의 기술적 복잡성과 한계** * **타임라인 통합의 어려움:** 각 모델은 비디오를 서로 다른 간격으로 분석하여 텍스트 레이블이나 벡터 임베딩 등 상이한 형태의 메타데이터를 생성하므로, 이를 하나의 연대기적 지도로 정렬하는 데 막대한 계산 비용이 발생합니다. * **데이터 규모의 폭발:** 2,000시간 분량의 아카이브는 약 2억 1,600만 프레임에 달하며, 이를 여러 모델로 처리할 경우 수십억 개의 레이블과 벡터 데이터가 생성되어 전통적인 데이터베이스로는 처리가 불가능합니다. * **중복 제거와 하이브리드 스코어링:** 시각적으로 유사한 수천 개의 후보 중 최적의 클립을 제안하기 위해, 단순한 수학적 유사도를 넘어 상징적 텍스트 매칭과 의미론적 벡터 검색을 결합한 정교한 랭킹 엔진이 필요합니다. * **제로 프릭션(Zero-Friction) 검색:** 창작 흐름을 방해하지 않기 위해 수십억 개의 레코드를 탐색하면서도 초 단위 미만의 응답 속도를 유지해야 하는 물리적 제약이 존재합니다. **데이터 수집 및 융합 파이프라인 (Ingestion & Fusion)** * **트랜잭션 영속화 (Transactional Persistence):** 고가용성 파이프라인을 통해 수집된 모델의 원본 주석(Annotation)을 Apache Cassandra에 저장합니다. 이 단계에서는 데이터 무결성과 빠른 쓰기 처리량을 최우선으로 하여 모든 모델 출력을 안전하게 확보합니다. * **오프라인 데이터 융합 (Offline Data Fusion):** Apache Kafka를 통해 비동기적으로 실행되며, 파편화된 모델 데이터를 1초 단위의 '시간 버킷(Temporal Buckets)'으로 정규화합니다. 예를 들어 '조이'라는 캐릭터와 '주방'이라는 배경이 겹치는 구간을 하나의 통합 레코드로 병합하여 복합적인 쿼리가 가능하도록 만듭니다. * **실시간 검색 인덱싱:** 융합된 데이터를 Elasticsearch에 인덱싱합니다. 이때 자산 ID와 시간 버킷을 조합한 복합 키(Composite Key)를 사용하여 업서트(Upsert) 방식으로 데이터를 갱신함으로써 데이터 중복을 방지하고 단일 진실 공급원(Single Source of Truth)을 유지합니다. **효율적인 멀티모달 시스템을 위한 제언** 대규모 영상 자산을 관리하는 시스템에서는 원본 데이터를 실시간으로 검색하는 대신, 데이터를 수집-융합-인덱싱 단계로 분리(Decoupling)하여 처리하는 구조가 필수적입니다. 특히 서로 다른 AI 모델의 출력을 공통된 시간 단위(Time Bucketing)로 정규화하여 저장함으로써, 복잡한 다차원 검색 시 발생하는 계산 부하를 오프라인에서 미리 해결하고 사용자에게는 즉각적인 검색 경험을 제공할 수 있습니다.