【发布时间】: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