Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
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