如何用PySpark提取文本文件中的JSON字符串并生成DataFrame新列?
解决Spark 2.4解析CSV中JSON字段并展开成列的问题
你直接用spark.read.csv读取文件时,只会把第三个json字段当成普通字符串处理,没法自动解析内部的JSON结构。针对Spark 2.4和Python 3.6,按以下步骤操作就能把JSON内的字段展开成DataFrame的独立列:
1. 读取CSV文件
先把文件读进来,保留所有原始字段:
df = spark.read.csv("/content/sample_data/file.txt", header=True, quote='"', escape='"')
这里不用inferSchema=True也没关系,因为后续会手动定义JSON的Schema,原始列按字符串读取即可。
2. 定义JSON字段的Schema
Spark需要明确的Schema来解析JSON结构,尤其是里面嵌套的数组stnKpis,得先定义子Schema再定义主Schema:
from pyspark.sql.types import StructType, StructField, StringType, BooleanType, DoubleType, LongType, ArrayType # 定义stnKpis数组的元素结构 stn_kpis_schema = StructType([ StructField("code", StringType(), nullable=True), StructField("value", DoubleType(), nullable=True), StructField("valueCreatedTs", LongType(), nullable=True), StructField("confidence", StringType(), nullable=True) ]) # 定义整个JSON字段的完整结构 json_schema = StructType([ StructField("line", StringType(), nullable=True), StructField("stn", StringType(), nullable=True), StructField("latitude", StringType(), nullable=True), StructField("longitude", StringType(), nullable=True), StructField("isInterchange", BooleanType(), nullable=True), StructField("isIncidentStn", BooleanType(), nullable=True), StructField("stnKpis", ArrayType(stn_kpis_schema), nullable=True) ])
3. 解析JSON字符串为结构体
用from_json函数把json列解析成Spark的结构体类型,生成一个新列:
from pyspark.sql.functions import from_json df_with_parsed_json = df.withColumn("parsed_json", from_json(df["json"], json_schema))
4. 展开结构体字段
把解析后的结构体里的所有字段展开,同时保留原始的pk、line、date列:
# 直接展开结构体的所有字段,同时保留原始列 final_df = df_with_parsed_json.select("pk", "line", "date", "parsed_json.*") # 如果原始的line列和JSON里的line列重名,可重命名避免冲突 final_df = final_df.withColumnRenamed("line", "csv_line")
查看结果
运行final_df.show(truncate=False)就能看到展开后的完整数据,所有JSON内的字段都变成了DataFrame的独立列。
内容的提问来源于stack exchange,提问作者Debiprasad Mishra
相关产品推荐
相关产品推荐

