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

如何使用PySpark提取嵌套JSON中的指定字段?

用PySpark提取JSON中的指定字段

我有一个带有如下Schema的JSON文件:

root
|-- context: struct (nullable = true)
|    |-- application: struct (nullable = true)
|    |    |-- version: string (nullable = true)
|    |-- custom: struct (nullable = true)
|    |    |-- dimensions: array (nullable = true)
|    |    |    |-- element: struct (containsNull = true)
|    |    |    |    |-- Activity ID: string (nullable = true)
|    |    |    |    |-- Activity Type: string (nullable = true)
|    |    |    |    |-- Bot ID: string (nullable = true)
|    |    |    |    |-- Channel ID: string (nullable = true)
|    |    |    |    |-- Conversation ID: string (nullable = true)
|    |    |    |    |-- Correlation ID: string (nullable = true)
|    |    |    |    |-- From ID: string (nullable = true)
|    |    |    |    |-- Recipient ID: string (nullable = true)
|    |    |    |    |-- StatusCode: string (nullable = true)
|    |    |    |    |-- Timestamp: string (nullable = true)
|    |-- data: struct (nullable = true)
|    |    |-- eventTime: string (nullable = true)
|    |    |-- isSynthetic: boolean (nullable = true)
|    |    |-- samplingRate: double (nullable = true)
|    |-- device: struct (nullable = true)
|    |    |-- roleInstance: string (nullable = true)
|    |    |-- roleName: string (nullable = true)
|    |    |-- type: string (nullable = true)
|    |-- location: struct (nullable = true)
|    |    |-- city: string (nullable = true)
|    |    |-- clientip: string (nullable = true)
|    |    |-- continent: string (nullable = true)
|    |    |-- country: string (nullable = true)
|    |    |-- province: string (nullable = true)
|    |-- operation: struct (nullable = true)
|    |    |-- id: string (nullable = true)
|    |    |-- parentId: string (nullable = true)
|    |-- session: struct (nullable = true)
|    |    |-- isFirst: boolean (nullable = true)
|-- event: array (nullable = true)
|    |-- element: struct (containsNull = true)
|    |    |-- count: long (nullable = true)
|    |    |-- name: string (nullable = true)
|-- internal: struct (nullable = true)
|    |-- data: struct (nullable = true)
|    |    |-- documentVersion: string (nullable = true)
|    |    |-- id: string (nullable = true)

JSON文件示例:

{
  "event": [
    {
      "name": "Activity",
      "count": 1
    }
  ],
  "internal": {
    "data": {
      "id": "79baca55-d168-11ea-b166-6fc861e9e21c",
      "documentVersion": "1.61"
    }
  },
  "context": {
    "application": {
      "version": "Wed 07/22/2020  5:37:05.58 \r\nUTC (fv-az461) [Build 148886] [Repo Intercom] [Branch prod] [Commit XXX] \r\n[XX 1.6.20-140775]  [XXX 1.3.27-144047] \r\n"
    },
    "data": {
      "eventTime": "2020-07-29T06:55:15.6294636Z",
      "isSynthetic": false,
      "samplingRate": 100
    },
    "cloud": {},
    "device": {
      "type": "PC",
      "roleName": "bc-directline-southindia",
      "roleInstance": "RD0003FF905CCA",
      "screenResolution": {}
    },
    "session": {
      "isFirst": false
    },
    "operation": {
      "id": "XXX",
      "parentId": "|XXXX.c4cd9570_"
    },
    "location": {
      "clientip": "0.0.0.0",
      "continent": "XX",
      "country": "XXX",
      "province": "XXX",
      "city": "XXX"
    },
    "custom": {
      "dimensions": [
        {
          "Timestamp": "XXX"
        },
        {
          "StatusCode": "200"
        },
        {
          "Activity ID": "JoH4veTvChCCnzchOD1Lg-f|0000001"
        },
        {
          "From ID": "XXX"
        },
        {
          "Correlation ID": "|54734cb21ba7f143a72ddd03fc865669.c4cd9570_"
        },
        {
          "Channel ID": "directline"
        },
        {
          "Recipient ID": "XXXX"
        },
        {
          "Bot ID": "XXXX"
        },
        {
          "Activity Type": "message"
        },
        {
          "Conversation ID": "XXX"
        }
      ]
    }
  }
}

需要提取Activity ID、Activity Type、Bot ID、Channel ID、Conversation ID、Correlation ID、From ID、Recipient ID、StatusCode、Timestamp到DataFrame,实现步骤如下:

实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, map_keys, map_values, first, monotonically_increasing_id

# 初始化SparkSession
spark = SparkSession.builder.appName("ExtractTargetFields").getOrCreate()

# 读取JSON文件
df = spark.read.json("path/to/your/json/files")

# 给每个原始行添加唯一ID,确保后续透视后每行对应原始一个JSON对象
df_with_row_id = df.withColumn("row_id", monotonically_increasing_id())

# 展开context.custom.dimensions数组,保留行ID
exploded_df = df_with_row_id.select(
    "row_id",
    explode(col("context.custom.dimensions")).alias("dimension")
)

# 将单个键的结构体转换为字段名和字段值的行
key_value_df = exploded_df.select(
    "row_id",
    map_keys(col("dimension"))[0].alias("field_name"),
    map_values(col("dimension"))[0].alias("field_value")
)

# 按行ID分组,透视字段名为列,提取每个字段对应的值
final_df = key_value_df.groupBy("row_id").pivot("field_name").agg(first("field_value"))

# 移除临时的row_id列
final_df = final_df.drop("row_id")

# 查看结果
final_df.show()

代码说明

  1. 添加行ID:因为每个原始JSON的dimensions是数组,每个元素对应一个字段,添加唯一ID可以保证透视后每行对应原始一个JSON对象,不会出现字段值跨行聚合的问题。
  2. 展开数组:用explode函数将数组拆分为多行,每个数组元素对应一行记录。
  3. 提取键值对:利用map_keys和map_values从单个键的结构体中取出字段名和对应的值。
  4. 透视转列:通过pivot将字段名转为列,再用first聚合每个行ID下的字段值(每个行ID下每个字段只会出现一次)。
  5. 清理临时列:移除用于分组的row_id,得到仅包含目标字段的DataFrame。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 09:40:59