【发布时间】:2018-08-06 10:30:59
【问题描述】:
我有类似的东西,其中spark 是我的sparkContext。我在sparkContext 中导入了implicits._,所以我可以使用$ 语法:
val df = spark.createDataFrame(Seq(("a", 0L), ("b", 1L), ("c", 1L), ("d", 1L), ("e", 0L), ("f", 1L)))
.toDF("id", "flag")
.withColumn("index", monotonically_increasing_id)
.withColumn("run_key", when($"flag" === 1, $"index").otherwise(0))
df.show
df: org.apache.spark.sql.DataFrame = [id: string, flag: bigint ... 2 more fields]
+---+----+-----+-------+
| id|flag|index|run_key|
+---+----+-----+-------+
| a| 0| 0| 0|
| b| 1| 1| 1|
| c| 1| 2| 2|
| d| 1| 3| 3|
| e| 0| 4| 0|
| f| 1| 5| 5|
+---+----+-----+-------+
我想为run_key 的每个非零块创建另一个具有唯一分组键的列,相当于:
+---+----+-----+-------+---+
| id|flag|index|run_key|key|
+---+----+-----+-------+---|
| a| 0| 0| 0| 0|
| b| 1| 1| 1| 1|
| c| 1| 2| 2| 1|
| d| 1| 3| 3| 1|
| e| 0| 4| 0| 0|
| f| 1| 5| 5| 2|
+---+----+-----+-------+---+
它可以是每次运行的第一个值、每次运行的平均值或某个其他值 - 只要保证它是唯一的,这样我就可以对其进行分组以比较其他值,这并不重要组之间。
编辑:顺便说一句,我不需要保留flag 是0 的行。
【问题讨论】:
标签: scala apache-spark apache-spark-sql apache-spark-2.0