Udf

Kullanıcı tanımlı bir işlev (UDF) oluşturur.

Sözdizimi

import pyspark.sql.functions as sf

# As a decorator
@sf.udf
def function_name(col):
    # function body
    pass

# As a decorator with return type
@sf.udf(returnType=<returnType>, useArrow=<useArrow>)
def function_name(col):
    # function body
    pass

# As a function wrapper
sf.udf(f=<function>, returnType=<returnType>, useArrow=<useArrow>)

Parametreler

Parametre Türü Description
f function Optional. Tek başına işlev olarak kullanılırsa Python işlevi.
returnType pyspark.sql.types.DataType veya str Optional. Kullanıcı tanımlı işlevin dönüş türü. Değer bir DataType nesnesi veya DDL biçimli bir tür dizesi olabilir. Varsayılan olarak StringType'ı kullanır.
useArrow bool Optional. (de) serileştirmeyi iyileştirmek için Ok kullanıp kullanmayacağınız. Hiçbiri olduğunda Spark yapılandırması "spark.sql.execution.pythonUDF.arrow.enabled" etkin olur.

Örnekler

Örnek 1: Dönüş türüne sahip lambda, dekoratör ve dekoratör kullanarak UDF oluşturma.

from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf

slen = udf(lambda s: len(s), IntegerType())

@udf
def to_upper(s):
    if s is not None:
        return s.upper()

@udf(returnType=IntegerType())
def add_one(x):
    if x is not None:
        return x + 1

df = spark.createDataFrame([(1, "John Doe", 21)], ("id", "name", "age"))
df.select(slen("name").alias("slen(name)"), to_upper("name"), add_one("age")).show()
+----------+--------------+------------+
|slen(name)|to_upper(name)|add_one(age)|
+----------+--------------+------------+
|         8|      JOHN DOE|          22|
+----------+--------------+------------+

Örnek 2: Anahtar sözcük bağımsız değişkenleriyle UDF.

from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf, col

@udf(returnType=IntegerType())
def calc(a, b):
    return a + 10 * b

spark.range(2).select(calc(b=col("id") * 10, a=col("id"))).show()
+-----------------------------+
|calc(b => (id * 10), a => id)|
+-----------------------------+
|                            0|
|                          101|
+-----------------------------+

Örnek 3: Pandas Series türü ipuçlarını kullanarak vektörleştirilmiş UDF.

from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf, col, PandasUDFType
import pandas as pd

@udf(returnType=IntegerType())
def pd_calc(a: pd.Series, b: pd.Series) -> pd.Series:
    return a + 10 * b

pd_calc.evalType == PandasUDFType.SCALAR
spark.range(2).select(pd_calc(b=col("id") * 10, a="id")).show()
+--------------------------------+
|pd_calc(b => (id * 10), a => id)|
+--------------------------------+
|                               0|
|                             101|
+--------------------------------+

Örnek 4: PyArrow Dizisi türü ipuçları kullanılarak vektörleştirilmiş UDF.

from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf, col, ArrowUDFType
import pyarrow as pa

@udf(returnType=IntegerType())
def pa_calc(a: pa.Array, b: pa.Array) -> pa.Array:
    return pa.compute.add(a, pa.compute.multiply(b, 10))

pa_calc.evalType == ArrowUDFType.SCALAR
spark.range(2).select(pa_calc(b=col("id") * 10, a="id")).show()
+--------------------------------+
|pa_calc(b => (id * 10), a => id)|
+--------------------------------+
|                               0|
|                             101|
+--------------------------------+

Örnek 5: Ok için iyileştirilmiş Python UDF (Spark 4.2'den beri varsayılan).

from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf

# Arrow optimization is enabled by default since Spark 4.2
@udf(returnType=IntegerType())
def my_udf(x):
    return x + 1

# To explicitly disable Arrow optimization and use pickle-based serialization:
@udf(returnType=IntegerType(), useArrow=False)
def legacy_udf(x):
    return x + 1

Örnek 6: Belirlenimi olmayan bir UDF oluşturma.

from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf
import random

random_udf = udf(lambda: int(random.random() * 100), IntegerType()).asNondeterministic()