【发布时间】:2019-08-26 18:13:44
【问题描述】:
我正在使用
创建一个数据框 val snDump = table_raw
.applyMapping(mappings = Seq(
("event_id", "string", "eventid", "string"),
("lot-number", "string", "lotnumber", "string"),
("serial-number", "string", "serialnumber", "string"),
("event-time", "bigint", "eventtime", "bigint"),
("companyid", "string", "companyid", "string")),
caseSensitive = false, transformationContext = "sn")
.toDF()
.groupBy(col("eventid"), col("lotnumber"), col("companyid"))
.agg(collect_list(struct("serialnumber", "eventtime")).alias("snetlist"))
.createOrReplaceTempView("sn")
我在 df 中有这样的数据
eventid | lotnumber | companyid | snetlist
123 | 4q22 | tu56ff | [[12345,67438]]
456 | 4q22 | tu56ff | [[12346,67434]]
258 | 4q22 | tu56ff | [[12347,67455], [12333,67455]]
999 | 4q22 | tu56ff | [[12348,67459]]
我想将其分解为我表中的 2 列中的数据,因为我正在做的是
val serialNumberEvents = snDump.select(col("eventid"), col("lotnumber"), explode(col("snetlist")).alias("serialN"), explode(col("snetlist")).alias("eventT"), col("companyid"))
也试过了
val serialNumberEvents = snDump.select(col("eventid"), col("lotnumber"), col($"snetlist.serialnumber").alias("serialN"), col($"snetlist.eventtime").alias("eventT"), col("companyid"))
但事实证明,explode 只能使用一次,并且我在选择中遇到错误,所以我如何使用explode/或其他东西来实现我想要的。
eventid | lotnumber | companyid | serialN | eventT |
123 | 4q22 | tu56ff | 12345 | 67438 |
456 | 4q22 | tu56ff | 12346 | 67434 |
258 | 4q22 | tu56ff | 12347 | 67455 |
258 | 4q22 | tu56ff | 12333 | 67455 |
999 | 4q22 | tu56ff | 12348 | 67459 |
我查看了很多 stackoverflow 线程,但没有一个对我有帮助。可能已经回答了这样的问题,但是我对 scala 的理解非常少,这可能使我不理解答案。如果这是重复的,那么有人可以指导我找到正确的答案。任何帮助表示赞赏。
【问题讨论】:
-
你可以爆炸两次。
-
Only one generator allowed per select是我爆炸两次时遇到的错误
标签: scala apache-spark struct apache-spark-sql