如何从Spark DataFrame的JSON格式字符串列中提取字段值?
提取Spark DataFrame中JSON字符串列的指定字段
推荐方法:Spark原生分布式处理(适合大数据集)
不要用collect()把数据拉到Driver端处理,Spark提供了原生的JSON解析函数,效率更高且避免内存溢出:
- 先定义JSON结构的Schema(根据你的实际JSON字段调整类型)
- 使用
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
相关产品推荐
相关产品推荐

