You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 04:15:04