Polars를 사용하여 수십만 건 데이터 처리 속도 향상
요약
본 글은 대용량 데이터 처리 시 발생하는 메모리 부족(OOM) 및 성능 문제를 Polars 라이브러리를 활용하여 해결하는 방법을 제시합니다. 특히 LazyFrame과 Parquet의 지연 실행 기능을 결합하고, 데이터를 배치 단위로 제한하여 RAM 사용량을 최적화하는 것이 핵심입니다.
핵심 포인트
- Polars는 Rust 기반 DataFrame 라이브러리로 OOM 및 성능 문제 해결에 효과적이다.
- LazyFrame은 실행 계획을 유지하며 불필요한 데이터 읽기를 줄여 메모리 효율성을 높인다.
- 데이터를 Parquet에 순차 쓰기하고, 처리 시 필요한 배치 단위만 RAM/GPU에 로드하는 것이 중요하다.
저는 업무 위탁이나 연구 프로젝트에서 수십만 건 규모의 데이터를 활용하여 기계 학습 모델을 개발하고 있습니다. 모델 개발 과정에서는 전처리, 학습 데이터 생성, 평가, 추론 등 여러 단계에서 동일한 데이터를 반복적으로 읽어 들입니다. 공정별로 전체 데이터를 다시 읽으면 RAM을 과도하게 사용(압박)할 뿐만 아니라, 불필요한 열이나 행을 읽는 시간도 누적됩니다.
처음에는 읽어 들인 레코드를 list에 모은 다음 한 번에 DataFrame으로 변환했습니다. 데이터 양이 늘어나자, 변환 전의 Python 객체와 변환 후의 열 데이터가 동시에 RAM을 사용하게 되어 메모리 부족(OOM) 문제가 발생했습니다. 레코드를 작은 단위로 나누어 처리하면 OOM은 피할 수 있었지만, Python에서 한 줄씩 처리하고 필요 없는 열까지 반복적으로 읽게 되면서 이번에는 처리 시간이 문제점이 되었습니다.
이 글에서 알 수 있는 것들
- 순차 쓰기(Sequential Writing)와 scan_parquet을 채택하여 데이터 로딩 시 메모리 부족(OOM) 회피
- 결과를 전체적으로 유지하는 collect의 판단 기준
- collect_batches를 채택하여 GPU 추론에 전달할 데이터를 배치 단위로 제한하기
Polars란?
이 OOM 및 처리 시간 문제를 해결하기 위해 Polars를 사용했습니다. Polars는 구조화된 데이터를 다루는 DataFrame 라이브러리로, Rust로 작성된 코어를 Python에서 사용할 수 있습니다.
LazyFrame은 처리를 즉시 실행하지 않고 실행 계획(Execution Plan)으로 유지하는 메커니즘입니다. 실행 전에 전체 계획을 최적화하기 때문에 불필요한 열이나 행의 읽기를 줄일 수 있습니다. 또한, 스트리밍 실행에 대응하는 처리에서는 데이터 배치 처리가 가능합니다.
참고: Polars User Guide, Polars Lazy API, Polars Streaming
Parquet과 Polars의 지연 실행(Lazy Execution)을 결합하면, 처리 과정에서 발생하는 메모리상의 문제를 단계적으로 줄일 수 있습니다. 첫 번째 대책은 원시 데이터를 list에 모으지 않고 Parquet으로 순차 쓰기하는 것입니다. 하지만 저장 위치만 바꾼다고 해서 로딩 시 전체 데이터가 RAM에 올라온다는 문제는 남아있습니다. LazyFrame으로 필요한 열과 행을 제한하더라도, collect를 하면 결과 전체가 DataFrame이 되므로, 마지막에는 외부 라이브러리나 GPU로 전달할 단위의 제한이 필요합니다.
이 흐름에 맞춰, 전체 데이터는 Parquet에 저장하고 Python 객체, Polars DataFrame, GPU 메모리에는 처리 중인 데이터만 놓는 구성으로 변경했습니다.
이후부터는 이 구성을 이벤트 레코드와 참조 엔티티의 마스터로 보여드리겠습니다. 기반은 실제 연구용 코드이지만, 고유한 데이터 이름이나 식별자는 제거했습니다.
이하 코드는 polars 1.43.2 이상을 전제로 하며, Parquet을 배치 단위로 쓰는 부분에서는 pyarrow 25.0.1 이상도 사용합니다.
데이터 유지 범위를 결정하기
피크 시점의 RAM 사용량을 좌우하는 것은 동시에 보유하는 데이터의 양입니다. 따라서 재사용할 전체 데이터는 Parquet에 두고, RAM과 GPU 메모리에는 다음 처리에 필요한 범위만 읽어 들입니다.
| 보유 위치 | 보유하는 범위 | 주의 사항 |
|---|---|---|
| Python 객체 | 처리 중인 배치 | list나 dict의 관리 비용이 추가됨 |
| ... | ||
| 첫 번째 대상은 DataFrame을 만들기 전의 Python 객체입니다. 이 단계에서 전체 데이터를 유지하면, 후속 처리를 최적화하기 전에 OOM이 발생하기 때문입니다. |
1. 원시 데이터를 Parquet에 순차 쓰기하기
이터레이터(Iterator)로부터 얻은 레코드를 list에 모으면, DataFrame을 만들기 전에 전체 분량의 Python 객체가 RAM에 올라갑니다.
# 피해야 할 예: records가 전체 건수의 dict를 보유함
records = list(iter_events("events.xml"))
pl.DataFrame(records).write_parquet("events.parquet")
참조 목록이나 큰 문자열을 포함할수록, 각 dict가 사용하는 RAM은 늘어납니다. 게다가 DataFrame으로 변환하는 동안에는, 변환 전의 Python 객체와 변환 후의 열 데이터를 이중으로 보유하는 상태입니다.
이 중복을 피하기 위해, 입력은 이터레이터 상태로 읽고 고정된 건수마다 Parquet에 써냅니다.
이 코드에서도 배치(batch) 분량의 Python 객체는 남습니다. 다만, 보유량의 기준은 입력 전체가 아니라 batch_size와 1개 레코드의 크기입니다. 1개 레코드가 클수록 같은 건수라도 RAM을 사용하므로, 문자열이나 참조 목록이 큰 경우에는 batch_size 조정이 필요합니다.
이렇게 하면 원본 데이터 읽기 시 보유량을 배치 분량으로 제한할 수 있습니다. 다만, 해결되는 것은 저장(save) 시의 OOM(Out Of Memory) 문제입니다. 후속 처리에서 Parquet을 전체 건수로 로드하면, 다시 데이터 전체가 RAM에 올라오는 과제가 남아있습니다.
2. LazyFrame으로 필요한 데이터만 읽기
앞 절에서는 전체 건수를 Parquet으로 임시 저장했지만, read_parquet으로 파일 전체를 읽으면 모든 데이터가 DataFrame이 됩니다.
# 소규모 데이터나 대화형 탐색에는 적합하지만, 전체 건수가 구체화됩니다
events = pl.read_parquet("events_raw.parquet")
validation_events = (
...
즉시 DataFrame을 얻을 수 있기 때문에, 중간 결과를 확인하는 탐색 작업이나, 전체 건수가 RAM에 들어오는 처리에 적합합니다.
전체 건수를 RAM에 올리고 싶지 않은 처리에서는 scan_parquet에서 LazyFrame을 만듭니다. 이 경우 처리를 바로 실행하지 않기 때문에, 필요한 열과 행을 줄이는 작업까지 포함하여 최적화할 수 있습니다. 참고: Polars Sources and Sinks
여기서 제외하는 것은 참조 대상의 마스터 레코드가 누락된 이벤트입니다. 결측(missing) 감지에서는 중첩 열(nested column) 전개로 인해 행 수가 늘어나기 때문에, 전개 후 데이터를 DataFrame으로 유지하면 RAM을 압박합니다. 따라서, 탐색부터 저장까지를 LazyFrame 상태로 연결하는 구성을 사용합니다.
import polars as pl
invalid_event_ids = (
pl.scan_parquet("events_raw.parquet")
...
explode는 중첩된 참조 대상을 한 행씩 분리하는 작업입니다. 전개 후에는 행 수가 늘어나기 때문에, 용도를 결측 확인에 국한합니다. 결측이 있는 event_id를 중복 제거하여 원본 이벤트와 anti join을 하면, 전개 결과의 collect는 필요하지 않습니다.
scan_parquet은 데이터 전체 건수를 DataFrame으로 만들지 않고, 실행 계획을 구성하는 진입점입니다. sink_parquet은 그 계획을 실행하고, 결과를 DataFrame에 collect하지 않고 저장합니다. 참고: Polars Sources and Sinks
LazyFrame 상태로 연결된 처리는 실행 전에 전체 계획을 최적화할 대상이 됩니다. projection pushdown은 필요한 열만 읽고, predicate pushdown은 가능한 필터를 로드 측으로 옮기는 최적화입니다. 참고: Polars Lazy API, Polars Optimizations
최적화가 의도대로 작동하는지는 explain으로 실행 전에 확인할 수 있습니다.
validation_events = (
pl.scan_parquet("events.parquet")
.filter(pl.col("partition") == "validation")
...
출력에는 Parquet SCAN, PROJECT, SELECTION 등이 나타납니다. Polars 버전에 따라 표기가 달라지므로, 문자열을 고정한 테스트에는 적합하지 않습니다. 여기서 확인하는 것은 필요한 열만 읽고, SELECTION이 Parquet SCAN에 포함되어 있는지 여부입니다.
이렇게 하면 Parquet에서 읽는 열과 행을 줄일 수 있습니다. 다만, LazyFrame이 제한하는 것은 실행 도중의 불필요한 로드입니다. 마지막에 collect한 결과는 DataFrame이 되므로, 그 크기에 따라 RAM을 압박할 수 있습니다.
3. collect은 결과의 크기로 판단하기
이전 단계의 처리를 LazyFrame으로 하더라도, collect의 결과는 DataFrame입니다. 스트리밍 실행에서 배치 단위가 되는 것은 중간 처리 과정에 한정되며, 반환되는 DataFrame에는 전체 건수가 유지됩니다. 또한, 스트리밍을 지원하지 않는 처리가 인메모리 엔진으로 전환된다는 점에도 주의해야 합니다. 참고: Polars Streaming
따라서 collect는 반환되는 결과가 RAM에 수용되고, 다음 처리가 DataFrame이나 Python 객체를 필요로 하는 위치에서 실행됩니다. 예를 들어, 학습 대상 이벤트가 참조하는 엔티티 ID는 중복 제거 후 건수가 충분히 작을 경우에 한해 구체화합니다.
training_entity_ids = (
pl.scan_parquet("events.parquet")
.filter(pl.col("partition") == "training")
...
이벤트 수가 많더라도, 참조하는 엔티티 ID의 종류가 적으면 중복 제거 후 결과는 작아집니다. 반면, 유니크 수가 크면 training_entity_ids 자체가 RAM을 압박합니다. 이 경우 선택지는 ID를 collect하지 않고 LazyFrame 상태로 조인하거나, 후술할 배치 단위 추출입니다.
ID 목록과 추출 후 레코드가 RAM에 수용되는 경우에는 마스터에서 필요한 2개 열만 추출합니다.
entities_for_training = (
pl.scan_parquet("entity_master.parquet")
.filter(pl.col("entity_id").is_in(training_entity_ids))
...
이 collect가 외부 라이브러리와의 경계입니다. 여기서 Python의 list나 NumPy 배열로 변환하면, DataFrame과 변환된 데이터를 동시에 보유하는 경우가 발생할 수 있습니다. 필요한 여유 공간은 추출 결과뿐만 아니라 변환된 데이터까지 포함하여 판단해야 합니다.
여기까지는 외부 라이브러리에 전달할 데이터가 RAM에 수용되는 경우의 처리입니다. 평가 데이터 전체가 수용되지 않는 경우에는 collect 위치를 늦추는 것만으로는 해결할 수 없으므로, 입력부터 출력까지를 배치 단위로 나눕니다.
4. Python과 GPU 경계를 배치를 이용해 닫기
모델 추론에서는 PyTorch 등에게 Python 객체나 텐서를 전달하므로, LazyFrame 상태로는 GPU에 전달할 수 없습니다.
전체 평가 데이터를 collect하고 참조처 마스터를 모두 dict로 만들면, Python 객체, DataFrame, GPU 텐서를 동시에 보유하는 상태가 됩니다. 이 중첩을 피하기 위해 collect_batches로 평가 이벤트를 추출하고, 그 배치가 참조하는 엔티티만 읽습니다. 배치 추론 결과를 즉시 쓰기(write)하면, 다음 배치로 넘어가기 전에 불필요한 객체를 해제할 수 있습니다.
from pathlib import Path
import polars as pl
import pyarrow.parquet as pq
...
RESULT_SCHEMA는 출력 열, predict는 추론, to_result_table은 Arrow 형식으로의 변환을 담당합니다. 입력 획득부터 출력까지를 루프 안에 배치함으로써, 보유하는 단위는 다음과 같습니다.
- batch는 평가 이벤트의 일부만
- records와 features_by_id는 해당 배치가 참조하는 만큼만
- GPU에 올리는 입력과 중간 텐서는 1배치분만
- 예측값은 1배치씩 쓰기(write)합니다.
chunk_size는 한 번에 반환되는 배치 행의 수입니다. 처리 전체의 메모리 상한을 나타내는 값은 아닙니다. 조정 기준은 RAM과 GPU 메모리 양쪽이며, 참조처의 수나 특징량 크기에 편차가 있다면 같은 행 수라도 필요한 메모리는 달라집니다.
이 구성은 피크 시점의 보유량을 줄이는 동시에, 배치마다 참조처 마스터를 순회합니다. 따라서 마스터를 RAM에 보유할 수 있는 경우에는 처음에 읽고, 보유할 수 없는 경우는 파일 분할이나 캐싱이 선택지입니다.
collect_batches는 사양(仕様)이 안정화되지 않았으며, sink_parquet 등 Polars 내에서 완결되는 출력 처리보다 느립니다. 따라서 Python이나 GPU로 전달하는 처리에 한정하고, Polars만으로 완결되는 처리에는 sink_parquet를 사용합니다. 버전 고정에 더해, 최대 규모의 실제 데이터 배치(batch)를 포함하는 결합 테스트도 필요합니다. 참고: polars.LazyFrame.collect_batches
이 결과로 외부 모델에 전달할 Python 객체, GPU 텐서, 출력은 배치 단위로 수렴됩니다. 다만, 배치를 만들기 위해서는 순회나 변환의 비용도 있으므로, 채택 기준은 데이터 규모와 후속 처리입니다.
Polars를 채택하는 판단 기준
이전 절까지의 구성이 효과적인 경우는 전체 건수를 RAM에 보유할 수 없어 Parquet 필터링이나 결합을 반복해야 할 때입니다. 반면, 데이터가 작거나 처리의 대부분이 Python 측에 있는 경우에는, 지연 실행(lazy execution)이나 배치화로 인해 증가하는 구현상의 복잡성에 비해 얻는 효과가 작습니다. 이 차이가 데이터량과 처리 내용에 따라 사용을 구분하는 기준입니다.
| 상황 | 채택하기 쉬운 이유 | 주의점 |
|---|---|---|
| Parquet를 여러 번 필터링하는 전처리 | LazyFrame이 읽어올 열(column)과 행(row)을 줄일 수 있음 | collect가 너무 빠르면 이점을 잃음 |
| ... | ||
| 처리 대부분을 자체 구현 함수로 하면, Polars가 최적화할 수 있는 범위에서 벗어나기 때문에 고속화를 기대하기 어렵습니다. 특히 map_elements는 Polars의 표현식(expression)보다 느려지기 쉬우므로, 우선순위는 select나 with_columns 안에서 문자열 표현식, 리스트 표현식, 구조체 표현식을 조합하는 구현입니다. 자체 구현 함수는 필요한 행까지 줄인 후에 사용합니다. 참고: polars.Expr.map_elements |
요약
수십만 건의 처리에 있어서, Polars를 채택한다고 해서 OOM(Out-Of-Memory)을 방지할 수 있는 것은 아닙니다. 생데이터 전체를 보유하면 DataFrame을 만들기 전에 RAM을 다 써버리기 때문에, 첫 번째 대책은 Parquet로 순차적 쓰기입니다. 하지만 Parquet를 전체 건수만큼 읽어오면 같은 문제가 재발합니다. 그래서 scan_parquet와 LazyFrame을 통해 필요한 열과 행에만 읽어오는 범위를 줄입니다. 그럼에도 collect의 결과는 RAM에 올라가므로, 실행하는 것은 결과가 충분히 작을 경우나 외부 처리에 전달하기 직전에만 합니다. 외부 처리로 전달할 데이터도 큰 경우에는, collect_batches를 사용해 입력, Python 객체, GPU 텐서, 출력을 배치 내에 가둡니다.
scan_parquet로 불필요한 열과 행의 읽기를 줄이고, 처리를 Polars의 표현식으로 연결하면, Python에서 한 줄씩 처리하는 범위도 줄일 수 있습니다. 게다가 이전 대책으로 남은 메모리상의 경계를 다음 대책으로 좁힐 수 있다는 점이 Parquet를 중심으로 하는 머신러닝 파이프라인에서 Polars를 채택하는 이유입니다.
Discussion

AI 자동 생성 콘텐츠
본 콘텐츠는 Zenn ML의 원문을 AI가 자동으로 요약·번역·분석한 것입니다. 원 저작권은 원저작자에게 있으며, 정확한 내용은 반드시 원문을 확인해 주세요.
원문 바로가기