如何在 PySpark DataFrame 中替换无穷大

Mic*_*ael 3 python pandas apache-spark apache-spark-sql pyspark

似乎不支持替换无穷大值。我尝试了下面的代码,但它不起作用。还是我错过了什么?

a=sqlContext.createDataFrame([(None, None), (1, np.inf), (None, 2)])
a.replace(np.inf, 10)
Run Code Online (Sandbox Code Playgroud)

或者我必须采取痛苦的路线:将PySpark DataFrame转换为pandas DataFrame,替换无穷大值,然后将其转换回PySpark DataFrame

zer*_*323 7

似乎不支持替换无穷大值。

实际上,它看起来像是 Py4J 错误,而不是replace它本身的问题。请参阅在 Python 和 Java 之间支持 nan/inf。

作为解决方法,您可以尝试 UDF(慢选项):

from pyspark.sql.types import DoubleType
from pyspark.sql.functions import col, lit, udf, when

df = sc.parallelize([(None, None), (1.0, np.inf), (None, 2.0)]).toDF(["x", "y"])

replace_infs_udf = udf(
    lambda x, v: float(v) if x and np.isinf(x) else x, DoubleType()
)

df.withColumn("x1", replace_infs_udf(col("y"), lit(-99.0))).show()

## +----+--------+-----+
## |   x|       y|   x1|
## +----+--------+-----+
## |null|    null| null|
## | 1.0|Infinity|-99.0|
## |null|     2.0|  2.0|
## +----+--------+-----+
Run Code Online (Sandbox Code Playgroud)

或这样的表达:

def replace_infs(c, v):
    is_infinite = c.isin([
        lit("+Infinity").cast("double"),
        lit("-Infinity").cast("double")
    ])
    return when(c.isNotNull() & is_infinite, v).otherwise(c)

df.withColumn("x1", replace_infs(col("y"), lit(-99))).show()

## +----+--------+-----+
## |   x|       y|   x1|
## +----+--------+-----+
## |null|    null| null|
## | 1.0|Infinity|-99.0|
## |null|     2.0|  2.0|
## +----+--------+-----+
Run Code Online (Sandbox Code Playgroud)