如何将JSON对象数组解析为Spark DataFrame?Databricks结构化流求助
解决Spark结构化流读取JSON数组文件返回null的问题
你遇到的问题根源很明确:你的JSON文件是整个文件封装成一个JSON数组([{...}, {...}, ...]),但Spark的json数据源默认是按「每行一个独立JSON对象」的格式来读取的,它没法直接识别这种整文件数组的结构,所以导致所有字段匹配失败,返回全null。
下面是具体的解决步骤和代码修改方案:
1. 调整Schema定义
首先确保你的json_schema是单个tweet对象的结构,然后我们需要把它包装成数组类型,用来解析整个文件的数组内容:
from pyspark.sql.types import StructType, StructField, StringType, LongType, ArrayType from pyspark.sql.functions import from_json, col, explode # 单个tweet的Schema(根据你的实际数据结构调整字段类型) json_schema = StructType([ StructField("run_stamp", StringType(), nullable=True), StructField("user", StringType(), nullable=True), StructField("id", LongType(), nullable=True), StructField("source", StringType(), nullable=True), StructField("favorite_count", LongType(), nullable=True), StructField("retweet_count", LongType(), nullable=True) ]) # 包装成数组类型,对应文件里的整个JSON数组结构 array_schema = ArrayType(json_schema)
2. 修改流读取逻辑
我们需要先把整个文件的内容读成字符串,再解析成JSON数组,最后用explode把数组中的每个tweet对象展开成单独的行:
tweetstore = "/tmp/tweets/" streamingInputDF = (spark # 先以text格式读取整个文件的内容(每行对应一个完整的JSON数组字符串) .readStream .text(tweetstore) # 将读取到的字符串解析为JSON数组 .withColumn("tweets_array", from_json(col("value"), array_schema)) # 展开数组:数组中的每个元素变成独立的一行,给对象起别名`tweet` .select(explode(col("tweets_array")).alias("tweet")) # 提取`tweet`对象中的各个字段 .select( "tweet.run_stamp", "tweet.user", "tweet.id", "tweet.source", "tweet.favorite_count", "tweet.retweet_count" ) )
3. 后续查询逻辑保持不变
你的后续写入内存表和SQL查询部分不需要修改,直接沿用即可:
streamingQuery = streamingInputDF \ .writeStream \ .format("memory") \ .queryName("tweetstream") \ .outputMode("append")\ .start() streamingDF = spark.sql("select * from tweetstream order by 1 desc")
为什么之前的方法失效?
Spark原生的json数据源默认处理的是行分隔式JSON(Line-delimited JSON),也就是文件中每行都是一个独立的JSON对象:
{"run_stamp": "...", "user": "..."} {"run_stamp": "...", "user": "..."}
而你的文件是把所有tweet对象包裹在一个数组里,Spark会把整个数组当作一个「不符合你定义的单对象Schema」的结构,所以所有字段都无法匹配,最终返回null。
通过先读取文本、解析数组、再展开的方式,就能完美适配你的JSON文件格式啦。
内容的提问来源于stack exchange,提问作者Mark
相关产品推荐
相关产品推荐

