쿼리 실행
쿼리 실행
Polars의 lazy API로 쿼리를 정의하면 Polars는 그 코드를 곧바로 실행하지 않아요. 대신 코드의 각 줄을 내부 쿼리 그래프에 추가하고, 그 그래프를 최적화해 둡니다. 그리고 실행을 요청할 때는 최적화가 끝난 쿼리 그래프를 기본으로 동작시키죠. 이 장에서는 쿼리를 언제, 어떻게 실행하는지 살펴봅니다.
출처: 공식문서
전체 데이터셋에서 실행하기
예를 들어 Reddit 데이터셋에 대한 쿼리를 하나 정의해 볼게요.
q1 = (
pl.scan_csv("docs/assets/data/reddit.csv")
.with_columns(pl.col("name").str.to_uppercase())
.filter(pl.col("comment_karma") > 0)
)
이 코드를 Reddit CSV에 대해 실행한다고 해도 쿼리는 아직 평가되지 않아요. 앞서 본 것처럼 Polars는 각 코드 줄을 내부 쿼리 그래프에 쌓아두고 최적화만 수행하죠.
쿼리에 .collect 메서드를 붙이면 전체 데이터셋을 대상으로 쿼리가 실행됩니다.
(
pl.scan_csv("docs/assets/data/reddit.csv")
.with_columns(pl.col("name").str.to_uppercase())
.filter(pl.col("comment_karma") > 0)
.collect()
)
shape: (14_029, 6)
┌─────────┬───────────────────────────┬─────────────┬────────────┬───────────────┬────────────┐
│ id ┆ name ┆ created_utc ┆ updated_on ┆ comment_karma ┆ link_karma │
│ --- ┆ --- ┆ --- ┆ --- ┆ --- ┆ --- │
│ i64 ┆ str ┆ i64 ┆ i64 ┆ i64 ┆ i64 │
╞═════════╪═══════════════════════════╪═════════════╪════════════╪═══════════════╪════════════╡
│ 6 ┆ TAOJIANLONG_JASONBROKEN ┆ 1397113510 ┆ 1536527864 ┆ 4 ┆ 0 │
│ 17 ┆ SSAIG_JASONBROKEN ┆ 1397113544 ┆ 1536527864 ┆ 1 ┆ 0 │
│ 19 ┆ FDBVFDSSDGFDS_JASONBROKEN ┆ 1397113552 ┆ 1536527864 ┆ 3 ┆ 0 │
│ 37 ┆ IHATEWHOWEARE_JASONBROKEN ┆ 1397113636 ┆ 1536527864 ┆ 61 ┆ 0 │
│ … ┆ … ┆ … ┆ … ┆ … ┆ … │
│ 1229384 ┆ DSFOX ┆ 1163177415 ┆ 1536497412 ┆ 44411 ┆ 7917 │
│ 1229459 ┆ NEOCARTY ┆ 1163177859 ┆ 1536533090 ┆ 40 ┆ 0 │
│ 1229587 ┆ TEHSMA ┆ 1163178847 ┆ 1536497412 ┆ 14794 ┆ 5707 │
│ 1229621 ┆ JEREMYLOW ┆ 1163179075 ┆ 1536497412 ┆ 411 ┆ 1063 │
└─────────┴───────────────────────────┴─────────────┴────────────┴───────────────┴────────────┘
결과를 보면 1,000만 행 중에서 우리가 준 조건(predicate)을 만족하는 행은 14,029개라는 걸 알 수 있어요. 기본 collect는 모든 데이터를 하나의 배치(batch)로 처리합니다. 즉 쿼리의 최대 메모리 사용 시점에 전체 데이터가 가용 메모리에 들어와야 한다는 뜻이에요.
주의 —
LazyFrame객체 재사용
LazyFrame은 쿼리 플랜, 즉 "계산하겠다는 약속"이지 공통 서브플랜을 캐시한다는 보장은 없어요. 그래서 정의한 뒤 서로 다른 다운스트림 쿼리에서 재사용할 때마다 매번 통째로 다시 계산됩니다.group_by처럼 행 순서를 유지하지 않는 연산을LazyFrame에 정의하면, 실행할 때마다 순서도 바뀔 수 있어요. 이런 연산에는maintain_order=True인자를 사용해 순서를 고정하는 게 좋습니다.
메모리보다 큰 데이터 실행하기
데이터가 사용 가능한 메모리보다 크다면 Polars가 streaming 모드로 데이터를 배치 단위로 처리해 줄 수 있어요. streaming 모드를 쓰려면 collect에 engine="streaming" 인자를 넘기기만 하면 됩니다.
(
pl.scan_csv("docs/assets/data/reddit.csv")
.with_columns(pl.col("name").str.to_uppercase())
.filter(pl.col("comment_karma") > 0)
.collect(engine="streaming")
)
일부 데이터셋 실행하기
큰 데이터셋에 대해 쿼리를 작성하거나 최적화하고 검증하는 동안, 매번 전체 데이터를 조회하면 개발 속도가 크게 느려질 수 있어요. 이럴 때는 파티션의 일부만 스캔하거나, 쿼리 앞부분에 .head를 두고 마지막에 .collect를 쓰면 됩니다.
(
pl.scan_csv("docs/assets/data/reddit.csv")
.head(10)
.with_columns(pl.col("name").str.to_uppercase())
.filter(pl.col("comment_karma") > 0)
.collect()
)
shape: (1, 6)
┌─────┬─────────────────────────┬─────────────┬────────────┬───────────────┬────────────┐
│ id ┆ name ┆ created_utc ┆ updated_on ┆ comment_karma ┆ link_karma │
│ --- ┆ --- ┆ --- ┆ --- ┆ --- ┆ --- │
│ i64 ┆ str ┆ i64 ┆ i64 ┆ i64 ┆ i64 │
╞═════╪═════════════════════════╪═════════════╪════════════╪═══════════════╪════════════╡
│ 6 ┆ TAOJIANLONG_JASONBROKEN ┆ 1397113510 ┆ 1536527864 ┆ 4 ┆ 0 │
└─────┴─────────────────────────┴─────────────┴────────────┴───────────────┴────────────┘
단, 데이터 일부에 대한 집계·필터 결과가 전체 데이터에서 얻는 결과와 반드시 같지는 않다는 점은 기억해 두세요.
분기하는 쿼리 (Diverging queries)
쿼리가 어느 한 지점에서 둘 이상으로 분기하는 경우는 아주 흔해요. 이때는 collect_all을 쓰는 걸 권장합니다. 분기된 쿼리들이 서로 같은 원본 LazyFrame을 공유해도, collect_all이 그 공통 부분을 한 번만 실행하도록 보장해 주거든요.
# Some expensive LazyFrame
lf: LazyFrame
lf_1 = lf.select(pl.all().sum())
lf_2 = lf.some_other_computation()
pl.collect_all([lf_1, lf_2]) # this will execute lf only once!