基于Scala的Spark 1.6:如何按key和字符位置分组code列字符?
解决Spark 1.6 + Scala下按Key分组并按位置归类Code字符的问题
嗨,针对你提出的需求——按key分组后,把code列的每个字符按位置分别归类(比如第1位字符归一组、第2位归一组),我来给你详细说下实现步骤,适配Spark 1.6的环境:
第一步:先把Code拆成单个字符的数组
你已经写了这部分代码,正好可以作为基础:
val someDF = Seq( (123, "0000"), (123, "X000"), (123, "C111"), (124, "0000"), (124, "0000"), (124, "C200")).toDF("key", "code") val dfWithSplit = someDF.select($"key", split($"code","").as("code_split"))
这一步把每个code字符串拆成了字符数组,比如"0000"变成["0","0","0","0"],方便后续按位置提取字符。
第二步:按Key分组,聚合每个位置的字符
Spark 1.6的DataFrame API没有像新版本那样丰富的高阶函数,所以我们可以用自定义UDF+分组聚合的方式来实现,或者转换成RDD处理,两种方式我都给你列出来:
方式一:用DataFrame + 自定义UDF(更简洁)
首先定义一个UDF,用来把同一个key下的所有字符数组,按位置合并成对应的字符列表:
import org.apache.spark.sql.functions._ // 这个UDF接收一组字符数组,返回按位置聚合后的数组(每个元素是对应位置的所有字符) val aggregateByPosition = udf { (codeArrays: Seq[Seq[String]]) => // 假设所有code长度一致,取第一个数组的长度作为总位置数 val totalPositions = codeArrays.head.length // 遍历每个位置,收集所有数组在该位置的字符 (0 until totalPositions).map(pos => codeArrays.map(_(pos))) } // 分组后聚合所有字符数组,再用UDF处理 val resultDF = dfWithSplit .groupBy($"key") .agg(collect_list($"code_split").as("all_code_arrays")) .select($"key", aggregateByPosition($"all_code_arrays").as("position_chars"))
运行后查看结果:
resultDF.show(false)
输出会是这样:
+---+------------------------------------+ |key|position_chars | +---+------------------------------------+ |123|[[0, X, C], [0, 0, 1], [0, 0, 1], [0, 0, 1]]| |124|[[0, 0, C], [0, 0, 2], [0, 0, 0], [0, 0, 0]]| +---+------------------------------------+
如果想把每个位置拆成单独的列(比如pos1、pos2),可以继续拆分:
val finalDF = resultDF.select( $"key", $"position_chars"(0).as("pos1"), $"position_chars"(1).as("pos2"), $"position_chars"(2).as("pos3"), $"position_chars"(3).as("pos4") ) finalDF.show(false)
输出就会更直观:
+---+---------+---------+---------+---------+ |key|pos1 |pos2 |pos3 |pos4 | +---+---------+---------+---------+---------+ |123|[0, X, C]|[0, 0, 1]|[0, 0, 1]|[0, 0, 1]| |124|[0, 0, C]|[0, 0, 2]|[0, 0, 0]|[0, 0, 0]| +---+---------+---------+---------+---------+
方式二:转换成RDD处理(更灵活)
如果你更熟悉RDD的操作,也可以转成RDD后手动处理:
val resultRDD = dfWithSplit.rdd .groupBy(row => row.getAs[Int]("key")) .mapValues { rows => // 提取当前key下的所有字符数组 val codeArrays = rows.map(_.getAs[Seq[String]]("code_split")).toList val totalPositions = codeArrays.head.length // 按位置收集字符 (0 until totalPositions).map(pos => codeArrays.map(_(pos))) } // 再转回DataFrame val resultDFFromRDD = resultRDD.toDF("key", "position_chars")
后续拆分列的操作和方式一完全一样。
小提醒
- 上面的代码默认所有
code的长度是一致的(比如示例里的4位),如果你的数据里有长度不一样的code,记得先做过滤或者处理,避免数组索引越界的问题。 - Spark 1.6的UDF要注意类型匹配,确保传入的参数类型和UDF定义的一致哦~
内容的提问来源于stack exchange,提问作者nickfrenchy
相关产品推荐
相关产品推荐

