【发布时间】:2020-03-23 11:02:21
【问题描述】:
我正在使用 PySpark,这是我的数据框的一部分 -
cleanData.show(4, False)
+------+-----+--------+----+----+------+------+-----+----+-----+-----+-----+----+-----+-----+------+
|STN |WBAN |YEARMODA|TEMP|DEWP|SLP |STP |VISIB|WDSP|MXSPD|GUST |MAX |MIN |PRCP |SNDP |FRSHTT|
+------+-----+--------+----+----+------+------+-----+----+-----+-----+-----+----+-----+-----+------+
|010080|99999|20100101|23.2|19.0|9999.9|9999.9|7.0 |6.0 |15.9 |999.9|33.8*|14.0|0.00H|999.9|001000|
|010080|99999|20100102|20.5|16.4|9999.9|9999.9|6.2 |20.4|33.0 |40.0 |33.4*|8.6*|0.00G|5.1 |001000|
|010080|99999|20100103|6.9 |-3.7|9999.9|9999.9|7.2 |14.1|21.4 |999.9|9.7* |4.5*|0.04G|5.1 |001000|
|010080|99999|20100104|4.9 |-6.2|9999.9|9999.9|8.7 |13.1|19.4 |999.9|6.8* |3.2*|0.00G|999.9|001000|
+------+-----+--------+----+----+------+------+-----+----+-----+-----+-----+----+-----+-----+------+
only showing top 4 rows
数据框中的一些列,如 MAX 和 MIN 列在几个条目的末尾有一个 *。
我需要找出这两列中的最大值和最小值。由于我熟悉 SQL,所以我使用 Spark SQL 发出查询,但是像 MAX 和 ORDER BY 这样的子句不能正常工作,例如 -
spark.sql("select MAX from weather2010uncleaned where not MAX='9999.9' order by MAX desc").show()
+-----+
| MAX|
+-----+
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
|99.9*|
+-----+
only showing top 20 rows
(注 - 9999.9 表示缺失数据)
我认为这是因为所有列都是string 类型,所以我使用.cast() 将它们转换为float 类型(代码在最后链接的github-gist 中)。
但不知何故,强制转换为浮点数将所有条目替换为 *,并在末尾替换为 NULL。
所以,我知道MAX 列中的最大值在132.8 附近(最后可能带有*),但是当我运行此查询以获得最大值时,我只获取128.8。
spark.sql("select STN, YEARMODA AS DATE, MAX from weather2010 where MAX=(select MAX(MAX) from weather2010 where not MAX='9999.9' and not max='99.99')").show()
# +------+--------+-----+
# | STN| DATE| MAX|
# +------+--------+-----+
# |703830|20100613|128.8|
# +------+--------+-----+
这可能是因为在转换为浮动的过程中最大条目被替换为NULL。
有什么办法吗-
- 在使用
createOrReplaceTempView()创建SQL 视图之前,从DataFrame 本身的条目中删除所有*,或者 - 使用 SQL 能够正确运行字符串类型的
MAX、ORDER BY等,同时还包括末尾带有*的条目,因此不需要强制转换,或者 - 如果无法使用 SQL 执行此操作,请单独使用 DataFrame API,尽管我对 API 不是很熟悉。
我不想在这里混淆问题,所以这是我的要点,其中包含更多关于其中一些操作的代码 sn-ps - gist。
【问题讨论】:
-
不熟悉 Spark,但是在普通的 Pandas 中,您不会删除星号然后转换为数字 dtype 吗?
-
我不确定,我以前没用过 Pandas。我是否将 Spark 数据帧转换为 Pandas 数据帧,然后删除星号,然后将其转换回 Spark DF 之类的?
-
不知道,虽然我认为 Spark DataFrames 有许多与 Pandas 对应的方法相同的方法。
-
所以,我无法从 Spark 转换为 Pandas DF。当我尝试使用像
cleanData.toPandas()这样的内置toPandas()函数时,我得到了这个 -py4j.protocol.Py4JJavaError: An error occurred while calling o98.collectToPython. : java.lang.OutOfMemoryError: GC overhead limit exceeded。 -
你检查了吗,除了转换DataFrame没有别的办法了吗?
标签: python sql apache-spark pyspark