쿼리 실행

쿼리 실행

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 모드를 쓰려면 collectengine="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!

더 알아보기 (Learn more)

  • Polars lazy API의 기본 개념은 Lazy API 문서에서 확인할 수 있어요.
  • 쿼리 플랜과 최적화 과정을 눈으로 보려면 쿼리 플랜 문서를 참고하세요.