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

Spark中如何将JSON输出的add、firsts字段转为对象数组?

解决Spark DataFrame转JSON时add/firsts字段转为数组的问题

我已经完成了大部分转换逻辑,但当前生成的JSON中add和firsts是单个对象,需要将它们改为对象数组,以支持包含多个元素的场景。

原代码

case class FirstIdentity(docType: String, docNumber: String, pId: String)
case class SecondIdentity(firm: String, code: String, orgType: String,
                              orgNumber: String, typee: String, perms: Seq[String])
case class General(id: Int, pName: String, description: String, add: Seq[SecondIdentity],
                       delete: Seq[String], act: String, firsts: Seq[FirstIdentity])

val someDF = Seq(
      ("0010XR_TYPE_6","0010XR", "222222", "6", "TYPE", "77444478", "6", 123, 1, "PF 1", "name", "description",
      Seq("PERM1", "PERM2"))
    ).toDF("firm", "code", "org_number", "org_type", "type", "doc_number",
           "doc_type", "id", "p_id", "p_name", "name", "description", "perms")

someDF.createOrReplaceTempView("vw_test")

val filter = spark.sql("""
                        select
                            firm, code, org_number, org_type, type, doc_number,
                             doc_type, id, p_id, p_name, name, description, perms
                         from vw_test
                    """)

val group =
      filter.rdd.map(x => {
          (
            x.getInt(x.fieldIndex("id")),
            x.getString(x.fieldIndex("p_name")),
            x.getString(x.fieldIndex("description")),
            SecondIdentity(
              x.getString(x.fieldIndex("firm")),
              x.getString(x.fieldIndex("code")),
              x.getString(x.fieldIndex("org_type")),
              x.getString(x.fieldIndex("org_number")),
              x.getString(x.fieldIndex("type")),
              x.getSeq(x.fieldIndex("perms"))
            ),
            "act",
            FirstIdentity(
              x.getString(x.fieldIndex("doc_number")),
              x.getString(x.fieldIndex("doc_type")),
              x.getInt(x.fieldIndex("p_id")).toString
            )
          )
        })
        .toDF("id", "name", "desc", "add", "actKey", "firsts")
        .groupBy("id", "name", "desc", "add", "actKey", "firsts")
        .agg(collect_list("add").as("null"))
        .drop("null")

group.toJSON.show(false)

当前输出

{
  "id": 123,
  "name": "PF 1",
  "desc": "description",
  "add": {
    "firm": "0010XR_TYPE_6",
    "code": "0010XR",
    "orgType": "6",
    "orgNumber": "222222",
    "typee": "TYPE",
    "perms": [
      "PERM1",
      "PERM2"
    ]
  },
  "actKey": "act",
  "firsts": {
    "docType": "77444478",
    "docNumber": "6",
    "pId": "1"
  }
}

期望输出

{
  "id": 123,
  "name": "PF 1",
  "desc": "description",
  "add": [
    {
      "firm": "0010XR_TYPE_6",
      "code": "0010XR",
      "orgType": "6",
      "orgNumber": "222222",
      "typee": "TYPE",
      "perms": [
        "PERM1",
        "PERM2"
      ]
    },
    {
      "firm": "0010XR_TYPE_6",
      "code": "0010XR",
      "orgType": "5",
      "orgNumber": "11111",
      "typee": "TYPE2",
      "perms": [
        "PERM1",
        "PERM2"
      ]
    }
  ],
  "actKey": "act",
  "firsts": [
    {
      "docType": "77444478",
      "docNumber": "6",
      "pId": "1"
    },
    {
      "docType": "411133",
      "docNumber": "6",
      "pId": "2"
    }
  ]
}

解决方案

问题出在分组逻辑上:你把add和firsts放到了groupBy的分组键中,导致它们无法被聚合为数组。正确的做法是只按需要唯一标识每组的字段(id, name, desc, actKey)分组,然后对add和firsts分别使用collect_list聚合。

修改后的代码

val group =
  filter.rdd.map(x => {
    (
      x.getInt(x.fieldIndex("id")),
      x.getString(x.fieldIndex("p_name")),
      x.getString(x.fieldIndex("description")),
      SecondIdentity(
        x.getString(x.fieldIndex("firm")),
        x.getString(x.fieldIndex("code")),
        x.getString(x.fieldIndex("org_type")),
        x.getString(x.fieldIndex("org_number")),
        x.getString(x.fieldIndex("type")),
        x.getSeq(x.fieldIndex("perms"))
      ),
      "act",
      FirstIdentity(
        x.getString(x.fieldIndex("doc_number")),
        x.getString(x.fieldIndex("doc_type")),
        x.getInt(x.fieldIndex("p_id")).toString
      )
    )
  })
  .toDF("id", "name", "desc", "add", "actKey", "firsts")
  // 仅按唯一标识字段分组
  .groupBy("id", "name", "desc", "actKey")
  // 对add和firsts分别聚合为数组
  .agg(
    collect_list("add").as("add"),
    collect_list("firsts").as("firsts")
  )

group.toJSON.show(false)

关键说明

  1. 分组键调整:去掉add和firsts作为分组条件,确保相同id/name/desc/actKey的记录会被分到同一组。
  2. 聚合操作:使用collect_list分别收集每组中的add和firsts对象,生成数组类型字段。
  3. 多数据测试:如果需要验证多元素场景,可以给someDF添加更多测试数据,比如:
val someDF = Seq(
  ("0010XR_TYPE_6","0010XR", "222222", "6", "TYPE", "77444478", "6", 123, 1, "PF 1", "name", "description", Seq("PERM1", "PERM2")),
  ("0010XR_TYPE_6","0010XR", "11111", "5", "TYPE2", "411133", "6", 123, 2, "PF 1", "name", "description", Seq("PERM1", "PERM2"))
).toDF("firm", "code", "org_number", "org_type", "type", "doc_number", "doc_type", "id", "p_id", "p_name", "name", "description", "perms")

内容的提问来源于stack exchange,提问作者mclc

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 19:05:46