Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Создает определяемую пользователем функцию таблицы PyArrow (UDTF). Эта функция предоставляет собственный интерфейс PyArrow для определяемых пользователем объектов, где метод eval получает PyArrow RecordBatches или Arrays и возвращает итератор таблиц PyArrow или RecordBatches. Это позволяет выполнять векторные вычисления без затрат на обработку строк.
Синтаксис
from pyspark.sql import functions as dbf
@dbf.arrow_udtf(returnType=<returnType>)
class MyUDTF:
def eval(self, ...):
...
Параметры
| Параметр | Тип | Description |
|---|---|---|
cls |
classнеобязательный |
Класс обработчика определяемых пользователем функций таблицы Python. |
returnType |
pyspark.sql.types.StructType или str, необязательно |
Возвращаемый тип определяемой пользователем функции таблицы. Это значение может быть либо объектом StructType, либо строкой типа структуры в формате DDL. |
Примеры
UDTF с входными данными PyArrow RecordBatch:
import pyarrow as pa
from pyspark.sql.functions import arrow_udtf
@arrow_udtf(returnType="x int, y int")
class MyUDTF:
def eval(self, batch: pa.RecordBatch):
# Process the entire batch vectorized
x_array = batch.column('x')
y_array = batch.column('y')
result_table = pa.table({
'x': x_array,
'y': y_array
})
yield result_table
df = spark.range(10).selectExpr("id as x", "id as y")
MyUDTF(df.asTable()).show()
UDTF с входными данными массива PyArrow:
@arrow_udtf(returnType="x int, y int")
class MyUDTF2:
def eval(self, x: pa.Array, y: pa.Array):
# Process arrays vectorized
result_table = pa.table({
'x': x,
'y': y
})
yield result_table
MyUDTF2(lit(1), lit(2)).show()
Замечание
- Метод eval должен принимать PyArrow RecordBatches или Массивы в качестве входных данных
- Метод eval должен выдавать таблицы PyArrow или RecordBatches в качестве выходных данных.