arrow_udtf

PyArrow yerel kullanıcı tanımlı tablo işlevi (UDTF) oluşturur. Bu işlev, UDF'ler için PyArrow yerel arabirimi sağlar; burada değerlendirme yöntemi PyArrow RecordBatches veya Diziler alır ve PyArrow Tablolarının veya RecordBatches'in Yineleyicisini döndürür. Bu, satır satır işleme yükü olmadan gerçek vektörleştirilmiş hesaplamayı etkinleştirir.

Sözdizimi

from pyspark.sql import functions as dbf

@dbf.arrow_udtf(returnType=<returnType>)
class MyUDTF:
    def eval(self, ...):
        ...

Parametreler

Parametre Türü Description
cls classopsiyonel Python kullanıcı tanımlı tablo işlev işleyici sınıfı.
returnType pyspark.sql.types.StructType veya strveya isteğe bağlı Kullanıcı tanımlı tablo işlevinin dönüş türü. Değer bir StructType nesnesi veya DDL biçimli yapı türü dizesi olabilir.

Örnekler

PyArrow RecordBatch girişi ile UDTF:

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()

PyArrow Dizisi girişleri ile UDTF:

@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()

Uyarı

  • Değerlendirme yöntemi giriş olarak PyArrow RecordBatches veya Dizileri kabul etmelidir
  • Değerlendirme yöntemi, çıkış olarak PyArrow Tablolarını veya RecordBatches'i vermelidir