공개 패키지 인덱스를 PySpark로 재구성하고 원천과 대조
RubyGems 공개 인덱스의 버전 기록 181만 건을 PySpark로 재구성하고, 원천 API와 대조해 이상 행을 찾았습니다.
작업 요약
버전 181만 건 재구성 · 원천 대조 1,768개 중 1,767개 일치 · 이상 행 2건 격리
회사에서 하던 원천 대조를 Spark로 다시 풀었습니다
실무에서는 패키지 레지스트리를 수집해 MySQL에 적재하고, 수집 결과를 원천 목록과 대조해 누락을 찾아왔습니다. 처리 엔진이 Spark로 바뀌면 같은 문제에서 무엇이 달라지는지 직접 확인하려고, 공개된 RubyGems 인덱스로 같은 종류의 작업을 PySpark로 구현했습니다.
회사 코드와 데이터는 사용하지 않았습니다. 입력은 누구나 받을 수 있는 공개 파일과 공개 API뿐입니다.
순서가 의미를 갖는 로그
RubyGems의 /versions 파일(약 23MB)은 표처럼 보이지만 실제로는 기록이 순서대로 쌓인 로그입니다. 앞부분에는 기준 시점의 전체 목록이 이름순으로 있고, 그 뒤에 이후 변경이 한 줄씩 붙습니다. 삭제(yank)된 버전은 - 표시로 다시 기록됩니다.
omen 0.1.0,0.2.0,...,0.9.0 <md5> 기준 목록
omen 0.10.0 <md5> 이후 추가
omen -0.10.0 <md5> 삭제
따라서 버전의 현재 상태는 파일 순서상 마지막 기록으로 정해집니다. Spark는 읽어들인 줄의 순서를 스스로 보장하지 않기 때문에, 이 순서를 어떻게 지킬지가 첫 번째 설계 문제였습니다.
처리 방식
- 셔플이 일어나기 전에
zipWithIndex로 줄 번호를 붙여 파일 순서를 보존했습니다. - 기준 목록과 추가 기록의 경계는 이름 정렬이 처음 깨지는 줄로 찾았습니다. 전체 윈도우를 쓰면 한 태스크로 데이터가 몰리기 때문에, 바로 앞 줄과 self-join해서 비교했습니다.
- 버전 목록을
posexplode로 펼치고, (패키지, 버전, 플랫폼)별 마지막 기록을 윈도우로 골라 삭제된 버전을 제외했습니다. - 기준 시점 상태와 현재 상태를
left_anti조인 두 번으로 비교해 추가·삭제된 버전을 뽑았습니다. 배치 방식의 변경분 추출입니다. - 원본 파일은 날짜별 파티션에 한 번만 저장하고, 이후 단계는 모두 그 원본에서 다시 실행할 수 있게 했습니다.
Spark 4는 기본이 ANSI 모드라 배열 범위를 벗어나면 null 대신 오류가 납니다. 플랫폼이 없는 버전에서 이 오류를 만나, 선택 항목은 get()으로 읽도록 고쳤습니다.
원천 대조에서 찾은 것
재구성한 결과를 RubyGems의 패키지별 /info API와 비교했습니다. 대상은 기준 시점 이후 바뀐 패키지 1,468개 전체와 무작위로 고른 기준 목록 패키지 300개입니다.
첫 대조에서 1,768개 중 1,765개가 일치했습니다. 나머지 3개를 하나씩 확인하니, 2개는 원천 파일에서 체크섬 칸이 비어 있는 줄에서 나왔습니다. 파일 전체에서 이런 줄은 두 개뿐이었고, RubyGems도 두 버전을 실제로 제공하지 않았습니다. 체크섬이 없는 줄은 상태에 반영하지 않고 격리 테이블로 따로 남기도록 규칙과 테스트를 추가했습니다.
규칙을 적용한 뒤 1,767개가 일치했습니다. 남은 1개는 파일을 받은 지 3분 뒤에 올라온 버전이라, 처리 오류가 아닌 시점 차이로 분류했습니다.
Spark와 단일 프로세스 비교
같은 계산을 일반 Python으로도 구현해 결과를 맞춰봤습니다. 현재 버전 1,813,218개, 기준 시점 이후 추가 4,247개와 삭제 10개가 두 구현에서 모두 같았습니다.
처리 시간은 로컬 Spark가 약 19.6초, 단일 Python 프로세스가 약 2.9초였습니다. 23MB는 한 프로세스로 충분한 크기라 Spark가 더 느렸습니다. 이 프로젝트에서 Spark로 얻은 것은 속도보다, 데이터가 한 대에 들어가지 않을 때도 유지되는 처리 구조(순서 보존 → 키별 윈도우 → anti-join 비교)입니다.
검증 기준과 한계
- 삭제 후 재등록, 플랫폼별 버전, 경계 탐지, 체크섬 없는 줄 격리를 다루는 테스트 4개를 두었고, Spark 결과가 단일 프로세스 결과와 같은지 매번 비교합니다.
- 원천 대조 결과는 일치·원천에만 있음·우리에게만 있음으로 나눠 기록하고, 자동으로 고치지 않습니다.
- 단일 머신의 로컬 모드(
local[*])에서만 실행했습니다. 클러스터 환경과 Databricks에서는 실행하지 않았습니다. - 수집 → Spark 처리 → 원천 대조로 이어지는 Airflow DAG는 작성했지만, 실제 Airflow 환경에서 실행하지는 않았습니다.