【问题标题】:Performing different computations conditioned on a column value in a spark dataframe根据 spark 数据框中的列值执行不同的计算
【发布时间】:2019-11-21 10:55:46
【问题描述】:

我有一个包含 2 列 A 和 B 的 pyspark 数据框。我需要根据 A 列的值对 B 的行进行不同的处理。在普通的熊猫中,我可能会这样做:

import pandas as pd
funcDict = {}
funcDict['f1'] = (lambda x:x+1000)
funcDict['f2'] = (lambda x:x*x)
df = pd.DataFrame([['a',1],['b',2],['b',3],['a',4]], columns=['A','B'])
df['newCol'] = df.apply(lambda x: funcDict['f1'](x['B']) if x['A']=='a' else funcDict['f2']
(x['B']), axis=1)

我能想到的在 (py)spark 中做的简单方法是

使用文件

  • 将数据读入数据帧
  • 按 A 列分区并写入单独的文件 (write.partitionBy)
  • 读入每个文件,然后分别处理

否则

使用表达式

  • 将数据读入数据帧
  • 编写一个笨拙的 expr(从可读性/维护的角度)根据列的值有条件地做一些不同的事情
  • 这不会像上面的 pandas 代码看起来那样“干净”

还有什么其他合适的方法来处理这个要求吗?从效率的角度来看,我希望第一种方法更干净,但由于分区-写-读,运行时间更长,而第二种方法从代码的角度来看并没有那么好,并且更难扩展和维护。

更主要的是,您是否会选择使用完全不同的东西(例如消息队列)(尽管存在相对延迟差异)?

编辑 1

基于我对 pyspark 的有限了解,只要处理不是很复杂,用户 pissall (https://stackoverflow.com/users/8805315/pissall) 提出的解决方案就可以工作。如果发生这种情况,如果不求助于 UDF,我不知道该怎么做,UDF 有其自身的缺点。考虑下面的简单示例

# create a 2-column data frame
# where I wish to extract the city 
# in column B differently based on
# the type given in column A
# This requires taking a different 
# substring (prefix or suffix) from column B
df = sparkSession.createDataFrame([
  (1, "NewYork_NY"),
  (2, "FL_Miami"),
  (1, "LA_CA"),
  (1, "Chicago_IL"),
  (2,"PA_Kutztown")
], ["A", "B"])

# create UDFs to get left and right substrings
# I do not know how to avoid creating UDFs
# for this type of processing
getCityLeft = udf(lambda x:x[0:-3],StringType())
getCityRight = udf(lambda x:x[3:],StringType())

#apply UDFs
df = df.withColumn("city", F.when(F.col("A") == 1, getCityLeft(F.col("B"))) \
                            .otherwise(getCityRight(F.col("B"))))

有没有办法以更简单的方式做到这一点而无需使用 UDF?如果我使用expr,我可以这样做,但正如我之前提到的,它看起来并不优雅。

【问题讨论】:

    标签: dataframe pyspark partitioning


    【解决方案1】:

    使用when怎么样?

    import pyspark.sql.functions as F
    
    df = df.withColumn("transformed_B", F.when(F.col("A") == "a", F.col("B") + 1000).otherwise(F.col("B") * F.col("B")))
    

    在更清楚问题后进行编辑

    您可以在_上使用split,并根据您的情况取其第一部分或第二部分。

    这是预期的输出吗?

    df.withColumn("city", F.when(F.col("A") == 1, F.split("B", "_")[0]).otherwise(F.split("B", "_")[1])).show()
    
    +---+-----------+--------+
    |  A|          B|    city|
    +---+-----------+--------+
    |  1| NewYork_NY| NewYork|
    |  2|   FL_Miami|   Miami|
    |  1|      LA_CA|      LA|
    |  1| Chicago_IL| Chicago|
    |  2|PA_Kutztown|Kutztown|
    +---+-----------+--------+
    

    UDF 方法:

    def sub_string(ref_col, city_col):
        # ref_col is the reference column (A) and city_col is the string we want to sub (B)
        if ref_col == 1:
            return city_col[0:-3]
        return city_col[3:]
    
    sub_str_udf = F.udf(sub_string, StringType())
    df = df.withColumn("city", sub_str_udf(F.col("A"), F.col("B")))
    

    另外,请查看:remove last few characters in PySpark dataframe column

    【讨论】:

    • @pisall,我已经编辑了我的问题,以解决为什么我不知道如何使用您提出的解决方案,因为处理变得更加复杂。谢谢。
    • @KS1 您不需要为此使用 UDF。请看看这是否是预期的输出?我尝试了我们的代码并得到了相同的输出
    • 感谢您的关注。也许这个例子不恰当。在这里,有一个 _ 可以拆分,这使它更简单。一般来说,如果函数取决于 A 列中的实际值,也会使其变得困难。尝试在不存在下划线的情况下执行此操作也是我为之苦恼的事情(例如,我不能只推送“看起来像”的东西substring(F.col('B'),3,length(F.col('B')))
    • 我想这会给你一个TypeError。对于使用动态长度字符串,子字符串不能开箱即用。使用pandas_udf 怎么样?我可以帮你把两个udf合并成一个udf
    • 是的。 UDF 是我试图避免的:)。同样,这只是我期望计算的一个非常简化的版本,所以试图看看是否有其他通用方法。谢谢!
    猜你喜欢
    • 2020-08-23
    • 2019-11-07
    • 1970-01-01
    • 1970-01-01
    • 2022-11-11
    • 2015-06-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多