你可以试试这个方法
def longestRun(sq: Seq[String]): Int = {
sq.mkString(" ")
.split("UK")
.filter(_.nonEmpty)
.map(_.trim)
.map(s => s.split(" ").length)
.max
}
case class Flights(passengerID: Int, flight: Seq[String], longestRun: Int)
val data = List(
(1,List("UK","IR","AT","UK","CH", "PK")),
(2,List("CG","IR")),
(3,List("CG","IR","SG","BE","UK")),
(4,List("CG","IR","NO","UK","SG","UK","IR","TJ","AT")),
(5,List("CG","IR"))
)
import spark.implicits._
val df = sc.parallelize(data)
.toDF("passengerID", "Flight")
.map(r => Flights(r(0).asInstanceOf[Int],r(1).asInstanceOf[Seq[String]],longestRun(r(1).asInstanceOf[Seq[String]])))
df.show(truncate = false)
/*
+-----------+------------------------------------+----------+
|passengerID|flight |longestRun|
+-----------+------------------------------------+----------+
|1 |[UK, IR, AT, UK, CH, PK] |2 |
|2 |[CG, IR] |2 |
|3 |[CG, IR, SG, BE, UK] |4 |
|4 |[CG, IR, NO, UK, SG, UK, IR, TJ, AT]|3 |
|5 |[CG, IR] |2 |
+-----------+------------------------------------+----------+
*/