基于Scala函数展开指定数组提取列的技术实现问询
Scala实现:展开数组并提取Sprint组件列
针对你给出的样本数据,我会分两种常见场景提供Spark Scala的实现方案——毕竟你的components数组可能是结构化JSON对象,也可能是像样本里那样的Sprint对象字符串表示,两种情况的处理方式略有不同:
场景1:components是结构化JSON对象数组
如果你的数据源中components存储的是嵌套的JSON对象(而非字符串),可以直接用Spark内置的explode函数展开数组,再提取嵌套字段。
完整代码实现
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ object ExpandSprintComponents { def main(args: Array[String]): Unit = { // 初始化Spark会话 val spark = SparkSession.builder() .appName("ExpandSprintComponents") .master("local[*]") .getOrCreate() import spark.implicits._ // 模拟你的样本数据(结构化JSON格式) val sampleJson = """[{"expand":"names,schema","centralid":10,"centralloc":"balh","components":[{"id":123,"rapidViewId":321,"state":"CLOSED","name":"Sprint 30 - \"abc\"","startDate":"2018-03-09T16:04:40.666+11:00","endDate":"2018-03-23T16:04:00.000+11:00","completeDate":"2018-03-23T14:12:44.680+11:00","sequence":980},{"id":456,"rapidViewId":654,"state":"CLOSED","name":"Sprint 31 - \"abc\"","startDate":"2018-03-23T14:57:17.889+11:00","endDate":"2018-04-06T14:57:00.000+11:00","completeDate":"2018-04-06T12:30:00.000+11:00","sequence":981}]}]""" // 读取JSON数据生成DataFrame val rawDF = spark.read.json(spark.sparkContext.parallelize(Seq(sampleJson))) // 展开数组并提取目标列 val expandedDF = rawDF .withColumn("sprint", explode($"components")) // 将数组的每个元素拆分为单独行 .select( $"centralid", $"centralloc", $"sprint.id".alias("sprint_id"), $"sprint.rapidViewId".alias("rapid_view_id"), $"sprint.state".alias("sprint_state"), $"sprint.name".alias("sprint_name"), $"sprint.startDate".alias("sprint_start_date"), $"sprint.endDate".alias("sprint_end_date"), $"sprint.completeDate".alias("sprint_complete_date"), $"sprint.sequence".alias("sprint_sequence") ) // 打印结果 expandedDF.show(truncate = false) spark.stop() } }
场景2:components是Sprint对象的字符串数组
如果你的components数组里存的是样本中那种com.atlassian.greenhopper.service.sprint.Sprint@xxxx[id=123,...]的字符串,需要先解析这些字符串提取字段,再展开处理。
完整代码实现
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ object ParseAndExpandSprints { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("ParseAndExpandSprints") .master("local[*]") .getOrCreate() import spark.implicits._ // 模拟你给出的样本数据(字符串形式的Sprint对象) val sampleJson = """[{"expand":"names,schema","centralid":10,"centralloc":"balh","components":["com.atlassian.greenhopper.service.sprint.Sprint@89322d3[id=123,rapidViewId=321,state=CLOSED,name=Sprint 30 - \"abc\",startDate=2018-03-09T16:04:40.666+11:00,endDate=2018-03-23T16:04:00.000+11:00,completeDate=2018-03-23T14:12:44.680+11:00,sequence=980]","com.atlassian.greenhopper.service.sprint.Sprint@42e71215[id=456,rapidViewId=654,state=CLOSED,name=Sprint 31 - \"abc\",startDate=2018-03-23T14:57:17.889+11:00,endDate=2018-04-06T14:57:00.000+11:00,completeDate=2018-04-06T12:30:00.000+11:00,sequence=981]"]}]""" val rawDF = spark.read.json(spark.sparkContext.parallelize(Seq(sampleJson))) // 自定义UDF:解析Sprint字符串,提取键值对为Map val parseSprintStr = udf((sprintStr: String) => { // 提取[]内部的内容 val content = sprintStr.split("\\[")(1).split("\\]")(0) // 用正则分割键值对:匹配不在引号内的逗号,避免拆分name字段里的逗号 val keyValuePairs = content.split(",(?=(?:[^\\\"]*\\\"[^\\\"]*\\\")*[^\\\"]*$)") keyValuePairs.map { pair => val Array(key, value) = pair.split("=", 2) // 去除value两端的引号(如果有) val cleanValue = value.replaceAll("^\\\"|\\\"$", "") key.trim -> cleanValue.trim }.toMap }) // 展开数组、解析字符串并提取列 val expandedDF = rawDF .withColumn("sprint_str", explode($"components")) // 展开数组 .withColumn("sprint_map", parseSprintStr($"sprint_str")) // 解析字符串为Map .select( $"centralid", $"centralloc", $"sprint_map.id".cast(IntegerType).alias("sprint_id"), // 转为整数类型 $"sprint_map.rapidViewId".cast(IntegerType).alias("rapid_view_id"), $"sprint_map.state".alias("sprint_state"), $"sprint_map.name".alias("sprint_name"), $"sprint_map.startDate".alias("sprint_start_date"), $"sprint_map.endDate".alias("sprint_end_date"), $"sprint_map.completeDate".alias("sprint_complete_date"), $"sprint_map.sequence".cast(IntegerType).alias("sprint_sequence") ) expandedDF.show(truncate = false) spark.stop() } }
关键说明
- 对于字符串形式的Sprint对象,我用正则表达式处理了字段分割问题,确保
name字段里的逗号不会被误拆分 - 如果需要调整字段的数据类型,直接在
select时用cast转换即可 - 两种方案都基于Spark Scala,因为展开数组提取列是大数据场景下的常见需求,如果你是在普通Scala(非Spark)环境处理,可以把数组遍历逻辑改成普通集合操作,核心思路类似
内容的提问来源于stack exchange,提问作者user2201536
相关产品推荐
相关产品推荐

