【发布时间】:2019-11-06 00:35:09
【问题描述】:
我的目标是用零替换 PySpark.DataFrame 列中的所有负元素。
输入数据
+------+
| col1 |
+------+
| -2 |
| 1 |
| 3 |
| 0 |
| 2 |
| -7 |
| -14 |
| 3 |
+------+
所需的输出数据
+------+
| col1 |
+------+
| 0 |
| 1 |
| 3 |
| 0 |
| 2 |
| 0 |
| 0 |
| 3 |
+------+
基本上我可以这样做:
df = df.withColumn('col1', F.when(F.col('col1') < 0, 0).otherwise(F.col('col1'))
或者udf可以定义为
import pyspark.sql.functions as F
smooth = F.udf(lambda x: x if x > 0 else 0, IntegerType())
df = df.withColumn('col1', smooth(F.col('col1')))
或
df = df.withColumn('col1', (F.col('col1') + F.abs('col1')) / 2)
或
df = df.withColumn('col1', F.greatest(F.col('col1'), F.lit(0))
我的问题是,哪一种是最有效的方法? Udf 有优化问题,所以绝对不是这样做的正确方法。但我不知道如何比较其他两种情况。一个答案应该是绝对做实验并比较平均运行时间等等。但我想从理论上比较这些方法(和新方法)。
提前谢谢...
【问题讨论】:
-
spark.sql('''select if(col1 < 0, 0, col1) as col1''') -
sql查询中的F.when和if条件有什么区别(在复杂性方面)?
-
How to measure the execution time of a query on Spark 的可能重复项,如下所示:Spark functions vs UDF performance?,请勿使用
udf代替简单的 spark 函数。
标签: python pyspark pyspark-sql pyspark-dataframes