사용자 정의 Python 함수
사용자 정의 Python 함수 (User-defined Python functions)
Polars 표현식은 꽤 강력하고 유연해서, 다른 라이브러리보다 커스텀 Python 함수가 필요할 일이 훨씬 적어요. 그래도 표현식의 상태를 제3자 라이브러리에 넘기거나, 블랙박스 함수를 Polars 데이터에 적용해야 하는 상황이 있을 수 있어요. 이 섹션에서는 그런 일을 가능하게 하는 두 API를 써 볼게요.
출처: 공식문서
참고: 커스텀 Python 함수를 작성하기 전에, 더 나은 성능과 통합을 제공하는 Polars 플러그인을 먼저 고려해 보세요. 커스텀 표현식은 표현식 플러그인 문서를, 커스텀 데이터 소스는 I/O 플러그인 문서를 확인하세요.
이 부분에서 쓰게 될 두 API는 다음과 같아요.
map_elements:Series안의 각 값에 함수를 개별적으로 호출한다.map_batches: 항상 전체Series를 함수에 넘긴다.
map_elements()로 개별 값 처리 (Processing individual values with map_elements())
가장 단순한 경우부터 시작할게요. Series 안의 각 값을 개별적으로 처리하고 싶다고 해 봅시다. 먼저 데이터를 만들게요.
from numba import float64, guvectorize, int64
import numpy as np
import math
import warnings
import polars as pl
from polars.exceptions import PolarsInefficientMapWarning
warnings.simplefilter("ignore", PolarsInefficientMapWarning)
df = pl.DataFrame(
{
"keys": ["a", "a", "b", "b"],
"values": [10, 7, 1, 23],
}
)
print(df)
shape: (4, 2)
┌──────┬────────┐
│ keys ┆ values │
│ --- ┆ --- │
│ str ┆ i64 │
╞══════╪════════╡
│ a ┆ 10 │
│ a ┆ 7 │
│ b ┆ 1 │
│ b ┆ 23 │
└──────┴────────┘
각 개별 값에 math.log()를 호출할 거예요.
def my_log(value):
return math.log(value)
out = df.select(pl.col("values").map_elements(my_log, return_dtype=pl.Float64))
print(out)
shape: (4, 1)
┌──────────┐
│ values │
│ --- │
│ f64 │
╞══════════╡
│ 2.302585 │
│ 1.94591 │
│ 0.0 │
│ 3.135494 │
└──────────┘
동작은 하지만 map_elements()에는 두 가지 문제가 있어요.
- 개별 항목에만 적용 가능: 보통은 개별 항목 하나하나가 아니라
Series전체에 대해 동작해야 하는 계산이 필요할 때가 많아요. - 성능 오버헤드: 각 항목을 개별 처리하고 싶더라도, 항목마다 함수를 호출하는 것은 느려요. 그 많은 추가 함수 호출이 많은 오버헤드를 더합니다.
먼저 첫 번째 문제를 해결하고, 그다음 두 번째 문제를 어떻게 푸는지 볼게요.
map_batches()로 Series 전체 처리 (Processing a whole Series with map_batches())
Series 전체의 내용에 대해 커스텀 함수를 실행하고 싶다고 해 봅시다. 시연 목적으로 Series의 평균과 각 값의 차이를 계산하고 싶다고 할게요.
map_batches() API를 쓰면 이 함수를 전체 Series 또는 group_by()의 개별 그룹에 실행할 수 있어요.
def diff_from_mean(series):
# This will be very slow for non-trivial Series, since it's all Python
# code:
total = 0
for value in series:
total += value
mean = total / len(series)
return pl.Series([value - mean for value in series])
# Apply our custom function to a full Series with map_batches():
out = df.select(pl.col("values").map_batches(diff_from_mean, return_dtype=pl.Float64))
print("== select() with UDF ==")
print(out)
# Apply our custom function per group:
print("== group_by() with UDF ==")
out = df.group_by("keys").agg(
pl.col("values").map_batches(diff_from_mean, return_dtype=pl.Float64)
)
print(out)
== select() with UDF ==
shape: (4, 1)
┌────────┐
│ values │
│ --- │
│ f64 │
╞════════╡
│ -0.25 │
│ -3.25 │
│ -9.25 │
│ 12.75 │
└────────┘
== group_by() with UDF ==
shape: (2, 2)
┌──────┬───────────────┐
│ keys ┆ values │
│ --- ┆ --- │
│ str ┆ list[f64] │
╞══════╪═══════════════╡
│ a ┆ [1.5, -1.5] │
│ b ┆ [-11.0, 11.0] │
└──────┴───────────────┘
사용자 정의 함수로 빠른 연산 (Fast operations with user-defined functions)
순수 Python 구현의 문제는 느리다는 거예요. 일반적으로 빠른 결과를 원한다면 호출하는 Python 코드의 양을 최소화하고 싶을 거예요. 속도를 최대화하려면 컴파일된 언어로 작성된 함수를 쓰는 게 좋습니다. 숫자 계산을 위해 Polars는 NumPy가 정의한 "ufunc"와 "generalized ufunc"라는 한 쌍의 인터페이스를 지원해요. 전자는 각 항목에 개별적으로 실행되고, 후자는 전체 NumPy 배열을 받아 더 유연한 연산을 허용합니다.
NumPy와 SciPy 같은 다른 라이브러리에는 Polars와 함께 쓸 수 있는 미리 작성된 ufunc가 들어 있어요. 예를 들어:
out = df.select(pl.col("values").map_batches(np.log, return_dtype=pl.Float64))
print(out)
shape: (4, 1)
┌──────────┐
│ values │
│ --- │
│ f64 │
╞══════════╡
│ 2.302585 │
│ 1.94591 │
│ 0.0 │
│ 3.135494 │
└──────────┘
여기서 map_batches()를 쓸 수 있는 이유는 numpy.log()가 개별 항목과 전체 NumPy 배열 모두에서 실행될 수 있기 때문이에요. 즉 Python 호출이 하나뿐이고 이후 모든 처리는 빠른 저수준 언어에서 일어나므로, 원래 예시보다 훨씬 빨리 실행됩니다.
예시: Numba로 빠른 커스텀 함수 (Example: A fast custom function using Numba)
NumPy가 제공하는 미리 작성된 함수는 유용하지만, 우리의 목표는 나만의 함수를 작성하는 거예요. 예를 들어 위의 diff_from_mean() 예시의 빠른 버전을 원한다고 해 봅시다. Python으로 이를 쓰는 가장 쉬운 방법은 Numba를 쓰는 것인데, (부분집합의) Python으로 커스텀 함수를 작성하면서 컴파일된 코드의 이점을 얻을 수 있게 해 줍니다.
특히 Numba는 @guvectorize라는 데코레이터를 제공해요. 이 데코레이터는 Python 함수를 빠른 기계 코드로 컴파일해 generalized ufunc를 만들며, Polars가 그걸 쓸 수 있게 해 줍니다.
다음 예시에서 diff_from_mean_numba()는 임포트 시점에 빠른 기계 코드로 컴파일되는데, 여기에 조금 시간이 걸릴 거예요. 그 뒤부터는 모든 함수 호출이 빠르게 실행됩니다. Series는 함수에 넘겨지기 전에 NumPy 배열로 변환됩니다.
# This will be compiled to machine code, so it will be fast. The Series is
# converted to a NumPy array before being passed to the function. See the
# Numba documentation for more details:
# https://numba.readthedocs.io/en/stable/user/vectorize.html
@guvectorize([(int64[:], float64[:])], "(n)->(n)")
def diff_from_mean_numba(arr, result):
total = 0
for value in arr:
total += value
mean = total / len(arr)
for i, value in enumerate(arr):
result[i] = value - mean
out = df.select(
pl.col("values").map_batches(diff_from_mean_numba, return_dtype=pl.Float64)
)
print("== select() with UDF ==")
print(out)
out = df.group_by("keys").agg(
pl.col("values").map_batches(diff_from_mean_numba, return_dtype=pl.Float64)
)
print("== group_by() with UDF ==")
print(out)
== select() with UDF ==
shape: (4, 1)
┌────────┐
│ values │
│ --- │
│ f64 │
╞════════╡
│ -0.25 │
│ -3.25 │
│ -9.25 │
│ 12.75 │
└────────┘
== group_by() with UDF ==
shape: (2, 2)
┌──────┬───────────────┐
│ keys ┆ values │
│ --- ┆ --- │
│ str ┆ list[f64] │
╞══════╪═══════════════╡
│ a ┆ [1.5, -1.5] │
│ b ┆ [-11.0, 11.0] │
└──────┴───────────────┘
generalized ufunc 호출 시 결측 데이터 불가 (Missing data is not allowed when calling generalized ufuncs)
diff_from_mean_numba() 같은 사용자 정의 함수에 넘겨지기 전에 Series는 NumPy 배열로 변환돼요. 불행히도 NumPy 배열에는 결측 데이터라는 개념이 없어요. 원본 Series에 결측 데이터가 있다면 결과 배열이 실제로는 Series와 일치하지 않게 됩니다.
항목별로 결과를 계산한다면 이는 문제가 되지 않아요. 예를 들어 numpy.log()는 각 개별 값에 개별적으로 호출되므로, 그 결측 값들은 계산을 바꾸지 않아요. 하지만 사용자 정의 함수의 결과가 Series의 여러 값에 의존한다면, 결측 값에서 정확히 무슨 일이 일어나야 하는지는 명확하지 않아요.
따라서 @guvectorize로 데코레이션된 Numba 함수 같은 generalized ufunc를 호출할 때, 결측 데이터가 있는 Series를 넘기려 하면 Polars는 오류를 던집니다. 결측 데이터를 없애려면 커스텀 함수 호출 전에 채우거나 제거하면 됩니다.
여러 열 값 결합 (Combining multiple column values)
사용자 정의 함수에 여러 열을 넘기고 싶다면 Struct를 쓸 수 있어요. Struct는 다른 섹션에서 자세히 다룹니다. 기본 아이디어는 여러 열을 Struct로 결합하면 함수가 그 열들을 다시 뽑아낼 수 있다는 거예요.
# Add two arrays together:
@guvectorize([(int64[:], int64[:], float64[:])], "(n),(n)->(n)")
def add(arr, arr2, result):
for i in range(len(arr)):
result[i] = arr[i] + arr2[i]
df3 = pl.DataFrame({"values_1": [1, 2, 3], "values_2": [10, 20, 30]})
out = df3.select(
# Create a struct that has two columns in it:
pl.struct(["values_1", "values_2"])
# Pass the struct to a lambda that then passes the individual columns to
# the add() function:
.map_batches(
lambda combined: add(
combined.struct.field("values_1"), combined.struct.field("values_2")
),
return_dtype=pl.Float64,
)
.alias("add_columns")
)
print(out)
shape: (3, 1)
┌─────────────┐
│ add_columns │
│ --- │
│ f64 │
╞═════════════╡
│ 11.0 │
│ 22.0 │
│ 33.0 │
└─────────────┘
스트리밍 계산 (Streaming calculations)
전체 Series를 사용자 정의 함수에 넘기는 데는 비용이 있어요. 내용이 NumPy 배열로 복사되므로 메모리를 많이 쓸 수 있지요. map_batches에 is_elementwise=True 인자를 쓰면 결과를 함수에 스트리밍할 수 있는데, 이는 모든 값을 한 번에 받지 않을 수도 있다는 뜻이에요.
참고:
is_elementwise인자는 잘못 설정하면 틀린 결과를 낼 수 있어요.is_elementwise=True로 설정한다면 함수가 실제로 요소별로 동작하는지(예: "각 값의 로그를 계산") 확인하세요. 예를 들어 우리 예시 함수diff_from_mean()은 요소별로 동작하지 않아요.
반환 타입 (Return types)
커스텀 Python 함수는 종종 블랙박스예요. Polars는 함수가 무엇을 하는지, 무엇을 반환할지 모릅니다. 그래서 반환 데이터 타입은 자동으로 추론됩니다. Polars는 첫 번째 비-null 값을 기다렸다가 그 값으로 결과 Series의 타입을 결정해요.
Python 타입에서 Polars 데이터 타입으로의 매핑은 다음과 같아요.
int->Int64float->Float64bool->Booleanstr->Stringlist[tp]->List[tp](내부 타입은 같은 규칙으로 추론됨)dict[str, [tp]]->structAny->object(이것은 항상 피해야 해요)
Rust 타입은 다음과 같이 매핑돼요.
i32또는i64->Int64f32또는f64->Float64bool->BooleanString또는str->StringVec<tp>->List[tp](내부 타입은 같은 규칙으로 추론됨)
추론된 타입을 덮어쓰고 싶다면 map_batches에 return_dtype 인자를 넘길 수 있습니다.