【问题标题】:Class methods as Pyspark udf类方法为 Pyspark udf
【发布时间】:2021-08-18 16:24:42
【问题描述】:

我有以下代码

import numpy as np
import pandas as pd

class MyClass:
    def __init__(self, a: pd.Series):
        self.a = a

    def f(self, b: pd.Series):
        return np.exp(a) + b

我还有一个带有双列 ab 的 Pyspark 数据框。我要跑

df.withColumn('c', MyClass(df['a']).f(df['b']))

当然失败了。如何正确调整MyClass 的代码以使其正常工作。 (请注意,我不能简单地将函数 f 写成 Pyspark 函数。

【问题讨论】:

    标签: python pandas apache-spark pyspark user-defined-functions


    【解决方案1】:

    您可以添加一个 UDF 来包装该类:

    import pyspark.sql.functions as F
    import pandas as pd
    import numpy as np
    
    class MyClass:
        def __init__(self, a: pd.Series):
            self.a = a
        def f(self, b: pd.Series):
            return np.exp(self.a) + b
    
    @F.pandas_udf('float')
    def myClassUDF(a: pd.Series, b: pd.Series) -> pd.Series:
        return MyClass(a).f(b)
    
    df = spark.createDataFrame([[0,1], [0,2]],['a','b'])
    
    df.withColumn('c', myClassUDF('a','b')).show()
    +---+---+---+
    |  a|  b|  c|
    +---+---+---+
    |  0|  1|2.0|
    |  0|  2|3.0|
    +---+---+---+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-06-11
      • 2018-09-12
      • 1970-01-01
      • 2020-10-06
      • 1970-01-01
      • 2021-03-18
      • 2021-09-30
      相关资源
      最近更新 更多