这有点奇怪。
当我创建一个dataframe,然后用函数pow进行一些转换时,它就能工作了。但当我推动它在现实世界中运行时,它就没有了。在我的虚拟场景中,列的数据类型和实际场景是相同的。
错误
Method pow([class java.lang.Double, class java.lang.Double]) does not exist这起作用(用合成的数据)。
from pyspark.sql.types import StructType,StructField, IntegerType, DoubleType
columns = ["CounterpartID","Year","Month","Day","churnprobability", "deadprobability"]
data = [(1234, 2021,5,12, 0.85,0.6),(1224, 2022,6,12, 0.75,0.6),(1345, 2022,5,13, 0.8,0.2),(234, 2021,7,12, 0.9,0.8), (1654, 2021,7,12, 1.40,20.0), (7548, 2021,7,12, -1.40,20.0), (6582, 2021,7,12, -1.40,20.0)]
schema = StructType([ \
StructField("CounterpartID",IntegerType(),False), \
StructField("Year",IntegerType(),False), \
StructField("Month",IntegerType(),False), \
StructField("Day", IntegerType(), False), \
StructField("churnprobability", DoubleType(), False), \
StructField("deadprobability", DoubleType(), False) \
])
df = spark.createDataFrame(data=data,schema=schema)
df.printSchema()
df.show(truncate=False)
abc=df.withColumn("client_id", f.col("CounterpartID"))\
.withColumn("year", f.col("Year"))\
.withColumn("month", f.col("Month"))\
.withColumn("day", f.col("Day"))\
.withColumn("churn_probability_unit", f.col("churnprobability").cast(IntegerType()))\
.withColumn("churn_probability_nanos", ((f.col("churnprobability") - f.col("churnprobability").cast(IntegerType())) * pow(10,9)).cast(IntegerType()))\
.withColumn("dead_probability_unit", f.col("deadprobability").cast(IntegerType()))\
.withColumn("dead_probability_nanos", (f.col("deadprobability") %1 * pow(10,9)).cast(IntegerType()))\
.select("client_id", "year", "month", "day", "churn_probability_unit", "churn_probability_nanos", "dead_probability_unit","dead_probability_nanos")\
abc.show()但是,在实际的场景(生产作业)中,我没有df,而是有一个真实的dataframe (当然),其中的所有列都具有与上面我的虚拟dataframe相同的数据类型:
例如在这里:

但是,当我执行相同的转换时,它会抱怨无法找到带有双参数的pow函数。这是堆栈跟踪(如下)。我还查找了用于pow的docs 这里,但是它没有提到任何关于数据类型的内容。一些S.O帖子建议双倍应该是可以的。哪里出问题了?当然,我可以将它乘以实际数字,而不是使用pow,但我仍然希望更好地理解这一点。任何问题,我都能帮你回答。
----> 7 .withColumn("churn_probability_nanos", ((f.col("churnprobability") % 1.0) * pow(10,9)).cast(IntegerType()))\
8 .withColumn("dead_probability_unit", f.col("deadprobability").cast(IntegerType()))\
9 .withColumn("dead_probability_nanos", (f.col("deadprobability") %1 * pow(10,9)).cast(IntegerType()))\
/databricks/spark/python/pyspark/sql/functions.py in pow(col1, col2)
737 Returns the value of the first argument raised to the power of the second argument.
738 """
--> 739 return _invoke_binary_math_function("pow", col1, col2)
740
741
/databricks/spark/python/pyspark/sql/functions.py in _invoke_binary_math_function(name, col1, col2)
73 and wraps the result with :class:`~pyspark.sql.Column`.
74 """
---> 75 return _invoke_function(
76 name,
77 # For legacy reasons, the arguments here can be implicitly converted into floats,
/databricks/spark/python/pyspark/sql/functions.py in _invoke_function(name, *args)
57 """
58 jf = _get_get_jvm_function(name, SparkContext._active_spark_context)
---> 59 return Column(jf(*args))
60
61
/databricks/spark/python/lib/py4j-0.10.9-src.zip/py4j/java_gateway.py in __call__(self, *args)
1302
1303 answer = self.gateway_client.send_command(command)
-> 1304 return_value = get_return_value(
1305 answer, self.gateway_client, self.target_id, self.name)
1306
/databricks/spark/python/pyspark/sql/utils.py in deco(*a, **kw)
108 def deco(*a, **kw):
109 try:
--> 110 return f(*a, **kw)
111 except py4j.protocol.Py4JJavaError as e:
112 converted = convert_exception(e.java_exception)
/databricks/spark/python/lib/py4j-0.10.9-src.zip/py4j/protocol.py in get_return_value(answer, gateway_client, target_id, name)
328 format(target_id, ".", name), value)
329 else:
--> 330 raise Py4JError(
331 "An error occurred while calling {0}{1}{2}. Trace:\n{3}\n".
332 format(target_id, ".", name, value))
Py4JError: An error occurred while calling z:org.apache.spark.sql.functions.pow. Trace:
py4j.Py4JException: Method pow([class java.lang.Double, class java.lang.Double]) does not exist
at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:341)
at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:362)
at py4j.Gateway.invoke(Gateway.java:289)
at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
at py4j.commands.CallCommand.execute(CallCommand.java:79)
at py4j.GatewayConnection.run(GatewayConnection.java:251)
at java.lang.Thread.run(Thread.java:748)发布于 2022-08-11 07:37:24
Spark中的函数是在Scala中开发的,因此是强类型的。这意味着函数的签名是函数的一部分。例如:
def foo(bar:str):
...
# AND
def foo(bar:int):
...这两个函数在python中是相同的。它们在Scala中是不同的函数,因为输入类型不同。
在您的示例中,不存在输入为double的pyspark pow。但是,接受columns作为输入的columns存在。它可能在您的第一个示例中起作用,因为您不是在使用pyspark pow,而是使用内置的python pow。
我建议你永远使用F.function_name。许多python函数和pyspark函数具有相同的名称,因此Python将用导入覆盖其内置函数。
您应该简单地将代码更改为:
from pyspark.sql import functions as F
F.pow(F.lit(10), F.lit(9))注意:lit使用输入参数创建一个列。
https://stackoverflow.com/questions/73316653
复制相似问题