arrow_udtf

Создает определяемую пользователем функцию таблицы 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 в качестве выходных данных.