Spark Scala:将扁平化数据转换为指定嵌套结构并按规则聚合
解决方法
首先,我们需要先把相同A和B分组下的所有C数组元素合并成一个完整的数组,再将这个大数组按每10个元素为一组拆分,最后为每个组分配序号。下面是具体的Scala代码实现,我们优先使用Spark内置函数来避免自定义UDF带来的性能开销:
步骤1:合并相同A、B的所有C元素
首先按A和B分组,收集所有C数组并扁平化,得到每个分组的完整元素列表:
import org.apache.spark.sql.functions._ val mergedDF = df.groupBy("A", "B") .agg(flatten(collect_list("C")).alias("all_elements"))
步骤2:拆分大数组为10元素的块并生成序号
利用Spark的内置函数计算分组数、生成块索引,再通过slice函数切割数组:
val chunkSize = 10 // 你需要的块大小,这里是10 val resultDF = mergedDF // 计算完整数组的长度 .withColumn("array_length", size(col("all_elements"))) // 计算需要拆分的块数(向上取整) .withColumn("num_chunks", ceil(col("array_length") / chunkSize)) // 生成从1开始的块序号数组 .withColumn("chunk_index", sequence(lit(1), col("num_chunks"))) // 计算每个块的起始索引(slice函数从1开始计数,所以后面要+1) .withColumn("start_pos", (col("chunk_index") - 1) * chunkSize) // 切割数组得到当前块 .withColumn("C_chunk", slice(col("all_elements"), col("start_pos") + 1, chunkSize)) // 选择最终需要的列 .select("A", "B", "C_chunk", "chunk_index")
步骤3:(可选)展开每个块的元素
如果你的目标格式是要每个元素关联到它所在的块和序号,可以再对C_chunk执行explode:
val finalResultDF = resultDF .withColumn("C_element", explode(col("C_chunk"))) .select("A", "B", "C_element", "C_chunk", "chunk_index")
格式化输出为你需要的字符串格式
如果要生成你示例中的{5, [1,2,...,10] , 1}这种字符串格式,可以用concat函数拼接:
val formattedDF = resultDF.withColumn("formatted_str", concat( lit("{"), col("A"), lit(", "), col("C_chunk"), lit(" , "), col("chunk_index"), lit("}") ))
测试结果
用你提供的输入数据测试,resultDF会输出:
| A | B | C_chunk | chunk_index |
|---|---|---|---|
| 5 | 1 | [1,2,...,10] | 1 |
| 5 | 1 | [11,12,13] | 2 |
| 5 | 2 | [1,2,3,15,16] | 1 |
| 6 | 1 | [1,2,3] | 1 |
| 7 | 3 | [4,5,6,7] | 1 |
这完全符合你需求中对分组和序号的要求。
内容的提问来源于stack exchange,提问作者Noobie93
相关产品推荐
相关产品推荐

