Scala中按乘客ID合并行程国家码为无重复数组的实现
解决方案:构建乘客连续无重复行程国家数组
针对你的需求,核心问题在于未按行程时间排序,且直接拼接from和to列表会导致连续行程的重复(如上一行的to是下一行的from),以及同一行from与to相同时的冗余。以下是两种可行的Spark实现方案:
前置准备:处理日期排序
首先必须将字符串日期转换为日期类型,并按乘客分组后按日期升序排序,确保行程顺序正确:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.DateType import org.apache.spark.sql.expressions.Window // 定义窗口:按乘客分组,按日期排序 val windowSpec = Window.partitionBy("passengerId").orderBy(to_date(col("date"), "dd/MM/yyyy")) // 转换日期格式并添加行号 val rankedData = data .withColumn("date_dt", to_date(col("date"), "dd/MM/yyyy")) .withColumn("rn", row_number().over(windowSpec))
方案一:合并列表后过滤连续重复
先收集所有from和最后一个to,再通过UDF去除连续重复的国家码:
// 聚合获取from列表和最后一个目的地 val aggregated = rankedData.groupBy("passengerId") .agg( collect_list(col("from")).alias("from_list"), last(col("to")).alias("last_to") ) // 定义UDF:过滤连续重复元素 val removeConsecutiveDuplicates = udf((arr: Seq[String]) => arr.foldLeft(Seq.empty[String]) { case (acc, elem) if acc.isEmpty || acc.last != elem => acc :+ elem case (acc, _) => acc } ) // 合并列表并去重,得到最终行程数组 val result = aggregated.withColumn( "countries", removeConsecutiveDuplicates(concat(col("from_list"), array(col("last_to")))) ) // 查看结果 result.select("passengerId", "countries").show(false)
方案二:聚合时直接跳过重复项
利用窗口函数获取前一行的目的地,在收集时直接跳过与前一个元素重复的国家码:
// 获取前一行的目的地 val withPrevTo = rankedData.withColumn( "prev_to", lag(col("to"), 1).over(windowSpec) ) // 收集行程:第一行取from,后续仅当to与前一个目的地不同时才收集 val result = withPrevTo.groupBy("passengerId") .agg( collect_list( when(col("rn") === 1, col("from")) .otherwise(when(col("to") !== coalesce(col("prev_to"), col("from")), col("to"))) ).alias("countries") ) // 过滤收集过程中产生的null值 .withColumn("countries", filter(col("countries"), x => x.isNotNull)) // 查看结果 result.select("passengerId", "countries").show(false)
输出结果
针对你的测试数据,两种方案都会得到正确的无重复连续行程数组:
+-----------+------------------------------------------------+ |passengerId|countries | +-----------+------------------------------------------------+ |5416 |[il, ar, us, au, sg, be, ca, ch, cg, at, no] | +-----------+------------------------------------------------+
内容的提问来源于stack exchange,提问作者Neeraj Bhanot
相关产品推荐
相关产品推荐

