Structured Spark Streaming writeStream输出null DataFrame问题咨询
问题根因
从Kafka数据源读取的Structured Streaming DataFrame默认是固定结构,你写入的JSON业务数据存储在二进制类型的value列中,你没有对该列做JSON解析和字段展开操作,因此直接输出要么出现全量null,要么输出完整JSON字符串。
解决步骤
1. 导入依赖函数
需要用到PySpark内置的JSON解析函数和列操作函数:
from pyspark.sql.functions import from_json, col
2. 解析JSON并展开字段
对读取到的原始流DataFrame做转换,先将value列从二进制转为字符串,再按照你定义的Schema解析为结构化数据,最后将结构体展开为独立列:
parsed_df = df \ # 按自定义Schema解析value列的JSON内容 .select(from_json(col("value").cast("string"), schema).alias("parsed_data")) \ # 展开结构体为独立业务字段 .select("parsed_data.*")
3. 输出解析后的结构化数据
用转换后的parsed_df执行writeStream输出即可:
query1 = parsed_df\ .writeStream\ .format("console")\ .outputMode("append")\ .option("truncate", False)\ .start()\ .awaitTermination()
注意事项
- 写入Kafka的JSON字段名需和Schema定义的字段名、大小写完全匹配,否则匹配不上的字段会返回null
- 若
time字段的格式不符合Spark默认Timestamp解析规则,可在from_json中添加时间格式配置,示例:from_json(col("value").cast("string"), schema, options={"timestampFormat": "yyyy-MM-dd HH:mm:ss"}) - 排查问题时可先输出
col("value").cast("string")的结果,确认写入Kafka的JSON格式是否合法、字段是否完整
内容的提问来源于stack exchange,提问作者girl of data
相关产品推荐
相关产品推荐

