【发布时间】:2019-12-10 23:13:04
【问题描述】:
我有一项工作需要在分区的 spark 数据帧上运行,过程如下所示:
rdd = sp_df.repartition(n_partitions, partition_key).rdd.mapPartitions(lambda x: some_function(x))
结果是pandas.dataframe 的rdd,
type(rdd) => pyspark.rdd.PipelinedRDD
type(rdd.collect()[0]) => pandas.core.frame.DataFrame
和rdd.glom().collect() 返回结果如下:
[[df1], [df2], ...]
现在我希望将结果转换为spark dataframe,我做的方式是:
sp = None
for i, partition in enumerate(rdd.collect()):
if i == 0:
sp = spark.createDataFrame(partition)
else:
sp = sp.union(spark.createDataFrame(partition))
return sp
但是,结果可能很大,rdd.collect() 可能会超出驱动程序的内存,所以我需要避免collect() 操作。有没有办法解决这个问题?
提前致谢!
【问题讨论】:
-
你可以运行
rdd.toDF()。或者,spark.createDataFrame(rdd) -
@samakart,不是真的,它会导致错误
ValueError: The truth value of a DataFrame is ambiguous. Use a.empty, a.bool(), a.item(), a.any() or a.all().,我猜它只适用于Row。 -
不,您必须提交正确的架构。在你的情况下,它只是无法弄清楚我想的类型。 createDataFrame 有 schema 参数。它还接受类似 sql 的字符串。我还没有在任何地方找到这种语法的文档。但它只是sql标准。
-
@dre-hh 是 this 你在找什么?
-
是和否:)。所以是的,这是将模式提供为 python 类型的一种方法。但 schema 参数也接受更短的 sql dsl 表示法。例如。
createDataFrame(x, schema="uuid_id STRING, url STRING, title STRING")我试图在文档中找到它支持的类型,但这些基本上是 SQL 表示法中 python 数据类型的类比
标签: python apache-spark pyspark