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

如何从Spark DataFrame的JSON格式字符串列中提取字段值?

提取Spark DataFrame中JSON字符串列的指定字段

推荐方法:Spark原生分布式处理(适合大数据集)

不要用collect()把数据拉到Driver端处理,Spark提供了原生的JSON解析函数,效率更高且避免内存溢出:

  1. 先定义JSON结构的Schema(根据你的实际JSON字段调整类型)
  2. 使用from_json解析字符串列,再提取目标字段
from pyspark.sql.types import StructType, StructField, StringType
from pyspark.sql.functions import from_json, col

# 定义JSON列的Schema,示例中id为字符串类型,可按需修改
json_schema = StructType([
    StructField("id", StringType(), nullable=True)
])

# 解析JSON列并提取id字段
result_df = df_spark.withColumn("parsed_json", from_json(col(JSON_COLUMN), json_schema)) \
                    .select("parsed_json.id")

# 查看结果
result_df.show()

修复你的collect遍历写法(仅适合小数据集测试)

如果只是小数据测试要用collect(),你之前的问题是没有把JSON字符串解析成Python字典,直接尝试取值导致错误,修改如下:

import json

collet_df = df_spark.collect()
for row in collet_df:
    # 先取出JSON格式的字符串
    json_str = row[JSON_COLUMN]
    # 解析为Python字典
    json_dict = json.loads(json_str)
    # 提取id字段(用get避免键不存在时报错)
    new = json_dict.get("id")
    print(new)

注意:collect()会将整个数据集加载到Driver端内存,仅适合小数据集测试,大数据场景务必用第一种方法。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 12:20:34