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

如何将Spark Dataset中inner_blob与json_blob合并为单个json_blob列

合并Spark Dataset中的两个结构体字段到单个列

针对你的需求,需要将inner_blob的顶层字段合并到json_blob中,同时保留原identifier_id列,以下是两种语言的实现方案:

Scala 实现

方式1:显式指定要合并的字段(适合字段数量少的场景)

如果inner_blob的字段不多,可以直接显式提取需要合并的字段,和原json_blob的字段组合成新的结构体:

import org.apache.spark.sql.functions._

// 假设你的Dataset名为df
val resultDF = df.withColumn(
  "json_blob",
  struct(
    col("json_blob.*"),  // 包含原json_blob的所有字段
    col("inner_blob.name"),
    col("inner_blob.age")
  )
).drop("inner_blob")  // 移除不再需要的inner_blob列

方式2:自动合并所有字段(适合字段数量多的场景)

如果inner_blob字段较多,可通过将结构体转为Map合并,再转回结构体的方式自动合并所有字段(重复Key会被后者覆盖,这里identifier_id值一致不影响):

import org.apache.spark.sql.functions._

val resultDF = df.withColumn(
  "json_blob",
  from_json(
    to_json(map_concat(to_map(col("json_blob")), to_map(col("inner_blob")))),
    col("json_blob").schema  // 复用原json_blob的Schema保证结构正确
  )
).drop("inner_blob")

Python 实现

方式1:显式指定要合并的字段

from pyspark.sql import functions as F

# 假设你的DataFrame名为df
result_df = df.withColumn(
    "json_blob",
    F.struct(
        F.col("json_blob.*"),
        F.col("inner_blob.name"),
        F.col("inner_blob.age")
    )
).drop("inner_blob")

方式2:自动合并所有字段

from pyspark.sql import functions as F

result_df = df.withColumn(
    "json_blob",
    F.from_json(
        F.to_json(F.map_concat(F.to_map(F.col("json_blob")), F.to_map(F.col("inner_blob")))),
        F.col("json_blob").schema
    )
).drop("inner_blob")

说明

  • 两种方式最终都会生成你期望的结果:保留原identifier_id列,新的json_blob包含原json_blob的所有字段,以及inner_blob中的name、age字段。
  • 若json_blob和inner_blob存在同名且值不同的字段,方式2中inner_blob的字段会覆盖json_blob的字段,可根据业务需求调整map的顺序。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 07:37:50