작업 완료

공개 패키지 인덱스를 PySpark로 재구성하고 원천과 대조

RubyGems 공개 인덱스의 버전 기록 181만 건을 PySpark로 재구성하고, 원천 API와 대조해 이상 행을 찾았습니다.

기간개인 프로젝트
대상 저장소개인 저장소(비공개)
역할단독 설계·구현
PySparkPythonAirflow

작업 요약

버전 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는 읽어들인 줄의 순서를 스스로 보장하지 않기 때문에, 이 순서를 어떻게 지킬지가 첫 번째 설계 문제였습니다.

처리 방식

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 비교)입니다.

검증 기준과 한계