LLM WikiAccess-protected knowledge portal

WIKI

Polars 1.x: Lazy 실행·스트리밍·Arrow 네이티브로 Pandas 병목을 없애는 Rust DataFrame 엔진

요약 Pandas는 Python 데이터 분석의 표준이었지만 수 GB를 넘어가면 한계가 명확하다. 모든 연산이 단일 스레드로 동작하고, 중간 결과를 매번 메모리에 구체화하며, Python 객체 오버헤드 때문에 실제 연산보다 메모리 관리에 시간이 더 걸린다. Polars 는 이 문제를 Rust로 구현한 DataFrame 엔진으로 해결한다. Apache Arrow 컬럼형 메모리 포맷을 기반으로, 쿼리를 먼저 분석해 불필요한 연산을

경로human/study/content/database-frontier/132-polars-1x-rust-dataframe-lazy-streaming-arrow-data-engineering.md
카테고리Study
태그#arrow #data #engineering #lazy #mysql #streaming #study

요약

Pandas는 Python 데이터 분석의 표준이었지만 수 GB를 넘어가면 한계가 명확하다. 모든 연산이 단일 스레드로 동작하고, 중간 결과를 매번 메모리에 구체화하며, Python 객체 오버헤드 때문에 실제 연산보다 메모리 관리에 시간이 더 걸린다.

Polars는 이 문제를 Rust로 구현한 DataFrame 엔진으로 해결한다. Apache Arrow 컬럼형 메모리 포맷을 기반으로, 쿼리를 먼저 분석해 불필요한 연산을 제거하고(Lazy 실행), 메모리에 다 올리지 못하는 데이터도 청크 단위로 처리한다(스트리밍 모드). 2024년 7월 v1.0 GA 이후 데이터 플랫폼 엔지니어링 현장에서 Pandas 대체재로 빠르게 자리 잡고 있다.

핵심 요약:


배경: Pandas의 구조적 병목

왜 Pandas가 느린가

Pandas의 성능 제약은 세 가지 구조에서 온다.

  1. 단일 스레드 실행: Python GIL과 Pandas 내부 구조 때문에 대부분의 연산이 단일 CPU 코어를 쓴다. 48코어 서버에서도 groupby·merge는 코어 하나가 일한다.
  1. 중간 결과 구체화: df.query(...).groupby(...).agg(...) 같은 체인은 각 단계마다 새 DataFrame을 메모리에 만든다. 필요 없는 열도 모두 포함된다.
  1. 행 기반 메모리 레이아웃: Pandas의 내부 NumPy 배열은 열 단위로 저장되지 않는 경우가 많고, object dtype이 Python 객체 포인터 배열이라 CPU 캐시 미스가 많다.

컬럼형 포맷의 이점

Apache Arrow는 동일한 열의 값을 연속 메모리에 저장한다. SELECT price WHERE region = 'KR' 같은 연산은 price 열과 region 열 두 개만 읽으면 된다. SIMD 벡터화 명령어로 필터를 한 번에 처리할 수 있고, CPU 프리페처가 효율적으로 동작한다.


Polars 아키텍처

Python / SQL API
논리 계획(Logical Plan)
쿼리 옵티마이저
Predicate Pushdown
Projection Pushdown
CSE 제거
Join 순서 최적화
물리 계획(Physical Plan)
Rust 실행 엔진
멀티스레드 파티션
SIMD 필터·집계
Arrow 컬럼형 I/O
DataFrame (Arrow 기반)
입력 소스
Parquet / CSV / JSON
DuckDB / PostgreSQL
Iceberg (PyIceberg)
S3 / GCS / ADLS
스트리밍 모드
청크 단위 처리
RAM < 데이터 크기
.collect(streaming=True)
Polars 쿼리 실행 파이프라인

Lazy API와 Eager API

Polars는 두 실행 모드를 제공한다.

Eager API: 즉시 실행. 디버깅·탐색에 유용하다.

import polars as pl

df = pl.read_parquet("sales.parquet")
result = df.filter(pl.col("region") == "KR").groupby("product").agg(
    pl.col("revenue").sum()
)

Lazy API: 연산 그래프를 쌓고 .collect()에서 한 번에 실행한다.

result = (
    pl.scan_parquet("sales.parquet")     # LazyFrame 반환
    .filter(pl.col("region") == "KR")
    .groupby("product")
    .agg(pl.col("revenue").sum())
    .sort("revenue", descending=True)
    .limit(10)
    .collect()                           # 여기서 실행
)

Lazy API의 중요한 점은 .collect() 전까지 어떤 데이터도 읽히지 않는다는 것이다. 옵티마이저가 다음을 수행한다.


Expression API

Polars의 표현식(expression) 시스템은 열 연산을 조합 가능한 함수로 추상화한다. Pandas의 apply(lambda ...)와 달리 SIMD 벡터화로 실행된다.

# 여러 집계를 한 번에
result = df.groupby("category").agg([
    pl.col("price").mean().alias("avg_price"),
    pl.col("price").std().alias("std_price"),
    pl.col("qty").sum().alias("total_qty"),
    (pl.col("price") * pl.col("qty")).sum().alias("total_revenue"),
])

# 윈도우 함수 (그룹 내 순위)
df_with_rank = df.with_columns(
    pl.col("revenue")
      .rank("dense", descending=True)
      .over("region")         # 그룹 내 적용
      .alias("rank_in_region")
)

# 문자열 연산
df_clean = df.with_columns(
    pl.col("name").str.strip_chars().str.to_lowercase(),
    pl.col("date").str.to_date("%Y-%m-%d"),
)

over() 함수는 SQL의 PARTITION BY와 동일한 의미다. Pandas에서 transform으로 복잡하게 처리해야 했던 것을 한 줄로 표현한다.


스트리밍 모드

메모리보다 큰 데이터를 처리할 때 .collect(streaming=True)를 사용한다.

# 500GB 로그 파일을 8GB RAM 서버에서 처리
result = (
    pl.scan_parquet("logs/year=2026/**/*.parquet")  # glob 패턴
    .filter(pl.col("level") == "ERROR")
    .groupby("service")
    .agg(pl.col("latency_ms").mean())
    .collect(streaming=True)    # 배치 청크 단위 처리
)

스트리밍 모드에서 Polars는 데이터를 고정 크기 배치(기본 524,288 행)로 나눠 처리한다. groupby·join 같은 연산도 가능하지만 내부적으로 spill-to-disk 메커니즘을 사용한다.

현재 제약: v1.x 스트리밍은 일부 복잡한 연산(일부 윈도우 함수, 중첩된 join)에서 fallback을 사용하거나 지원하지 않을 수 있다. Polars는 .explain(streaming=True) 명령으로 어떤 연산이 스트리밍으로 처리되는지 미리 확인할 수 있다.


SQL 인터페이스

Polars 1.x는 polars.sql 모듈을 통해 표준 SQL을 직접 실행할 수 있다.

ctx = pl.SQLContext({"sales": df_sales, "products": df_products})

result = ctx.execute("""
    SELECT p.name, SUM(s.revenue) AS total_revenue
    FROM sales s
    JOIN products p ON s.product_id = p.id
    WHERE s.date >= '2026-01-01'
    GROUP BY p.name
    ORDER BY total_revenue DESC
    LIMIT 20
""").collect()

SQLContext는 Lazy 계획으로 컴파일되므로 SQL 실행도 옵티마이저의 혜택을 받는다.


데이터 플랫폼 통합

DuckDB와 협력

DuckDB와 Polars는 모두 Arrow 포맷을 사용하므로 제로 복사 교환이 가능하다.

import duckdb

# Polars DataFrame → DuckDB 뷰 (복사 없음)
duckdb.register("orders", df_polars.to_arrow())
result_duck = duckdb.sql("SELECT * FROM orders WHERE total > 1000").arrow()

# DuckDB 결과 → Polars (복사 없음)
df_result = pl.from_arrow(result_duck)

일반적인 패턴: 파일 읽기와 초기 필터링은 Polars Lazy로, 복잡한 집계 SQL은 DuckDB로, 결과 후처리는 다시 Polars로.

PyIceberg (Apache Iceberg)

from pyiceberg.catalog import load_catalog

catalog = load_catalog("glue", **{"type": "glue"})
table = catalog.load_table("prod.orders")

# Polars Lazy로 Iceberg 테이블 스캔
df = pl.from_arrow(
    table.scan()
         .filter("region = 'KR'")
         .select(["order_id", "revenue", "date"])
         .to_arrow()
)

Spark와의 관계

Spark와 Polars는 경쟁 관계이기도 하지만, 실무에서는 역할이 다르다.

Spark와 Polars는 Arrow IPC로 데이터를 교환한다. Spark의 toPandas() 대신 toArrow()pl.from_arrow()를 사용하면 변환 비용을 낮출 수 있다.


Pandas에서 Polars로 마이그레이션

가장 흔한 패턴 변환:

작업PandasPolars
파일 읽기pd.read_parquet(...)pl.scan_parquet(...) (lazy)
필터df[df.col > 0]df.filter(pl.col("col") > 0)
열 추가df["new"] = df["a"] + df["b"]df.with_columns((pl.col("a") + pl.col("b")).alias("new"))
groupby 집계df.groupby("k")["v"].sum()df.groupby("k").agg(pl.col("v").sum())
apply (행별)df["col"].apply(func)df["col"].map_elements(func) (느림) 또는 표현식 조합 (빠름)
날짜 필터df[df["date"] >= "2026-01-01"]df.filter(pl.col("date") >= pl.lit("2026-01-01").str.to_date())

map_elements 주의: Polars의 map_elements(구 apply)는 Python 함수를 행별로 호출하므로 Pandas apply만큼 느리다. 가능하면 Polars 표현식 조합으로 대체해야 SIMD 가속의 혜택을 받는다.


운영 체크리스트


References