【发布时间】:2019-09-09 07:08:02
【问题描述】:
基本问题:
我想将 Spark 数据帧 sdf 的“第一行”复制到另一个 Spark 数据帧 sdfEmpty。
我不明白下面的代码出了什么问题。 因此,我期待一个解决方案和一个解释,在我的最小示例中失败了。
一个简单的例子:
// create a spark data frame
import org.apache.spark.sql._
val sdf = Seq(
(1, "a"),
(12, "b"),
(234, "b")
).toDF("A", "B")
sdf.show()
+---+---+
| A| B|
+---+---+
| 1| a|
| 2| b|
| 3| b|
+---+---+
// create an empty spark data frame to store the row
// declare it as var, such that I can change it later
var sdfEmpty = spark.createDataFrame(sc.emptyRDD[Row], sdf.schema)
sdfEmpty.show()
+---+---+
| A| B|
+---+---+
+---+---+
// take the "first" row of sdf as a spark data frame
val row = sdf.limit(1)
// combine the two spark data frames
sdfEmpty = sdfEmpty.union(row)
作为row 是:
row.show()
+---+---+
| A| B|
+---+---+
| 1| a|
+---+---+
sdfEmpty 的预期结果是:
+---+---+
| A| B|
+---+---+
| 1| a|
+---+---+
但我明白了:
sdfEmpty.show()
+---+---+
| A| B|
+---+---+
| 2| b|
+---+---+
问题: 让我感到困惑的是:使用 val row = sdf.limit(1) 我以为我创建了一个永久/不可更改/定义明确的对象。这样当我打印一次并将其添加到某物时,我会得到相同的结果。
备注:(非常感谢丹尼尔的发言)
我知道在 scala 的分布式世界中,没有明确定义的“第一行”概念。我把它放在那里是为了简单起见,我希望那些在类似事情上苦苦挣扎的人会“不小心”使用“第一”这个词。
我试图实现的是:(在一个简化的例子中) 我有一个包含 2 列 A 和 B 的数据框。A 列是部分有序的,B 列是完全有序的。 我想过滤数据w.r.t。列。所以这个想法是某种分而治之:拆分数据框,这样两列都是完全有序的,而不是像往常一样过滤。 (并进行明显的迭代)
为了实现这一点,我需要选择一个定义明确的行并将日期拆分为 w.r.t。行.A.但正如最小示例所示,我的命令不会产生定义明确的对象。
非常感谢
【问题讨论】:
-
你能分享
println(sdf.rdd.partitions.size)的输出吗? -
@moriarty007 println(sdf.rdd.partitions.size) 的输出是3。
标签: scala apache-spark pyspark apache-spark-sql