如何从Kafka readStream生成的Parquet文件中读取真实value值?
解决Spark读取Kafka生成的Parquet文件中value列的字节数据问题
你的问题根源在于:Spark从Kafka读取的value列默认是二进制类型(BinaryType),写入Parquet文件后保留了该类型,所以直接读取后看到的是字节数组形式的数据。要获取真实数据,按以下两步处理即可:
步骤1:将字节数组转换为字符串
使用Spark内置函数把字节数组转成可读字符串,两种方式任选:
方式一:用binaryToString函数
from pyspark.sql.functions import binaryToString # 将value列的字节数组转为字符串 str_df = parquet_df.withColumn("value_str", binaryToString("value")) # 查看转换后的结果 str_df.select("value_str").show(truncate=False)
方式二:用cast类型转换
如果你的Spark版本不支持binaryToString,可以用类型强转替代:
str_df = parquet_df.withColumn("value_str", parquet_df["value"].cast("string"))
步骤2:解析JSON字符串(适配你的数据格式)
从字节开头的[22 7B ...]来看,原始数据是JSON格式(22是双引号ASCII码,7B是左大括号),需要把字符串解析为结构化数据:
from pyspark.sql.functions import from_json from pyspark.sql.types import StructType, StructField, StringType, DoubleType # 必须根据你的实际JSON结构定义Schema,以下是示例模板 schema = StructType([ StructField("device_id", StringType(), True), StructField("electric_data", DoubleType(), True), StructField("collect_time", StringType(), True) ]) # 解析JSON字符串为结构化数据 final_df = str_df.withColumn("real_data", from_json("value_str", schema)) # 展开结构化数据到顶层列 final_df = final_df.select("real_data.*") final_df.show()
完整可运行代码
import glob from pyspark.sql.functions import binaryToString, from_json from pyspark.sql.types import StructType, StructField, StringType, DoubleType # 读取Parquet文件 parquet_files = glob.glob("/home/ubuntu/tmp/output" + "/*.parquet") parquet_df = spark.read.parquet(*parquet_files) # 字节转字符串 str_df = parquet_df.withColumn("value_str", binaryToString("value")) # 定义匹配实际数据的Schema schema = StructType([ StructField("device_id", StringType()), StructField("electric_data", DoubleType()), StructField("collect_time", StringType()) ]) # 解析并获取真实数据 final_df = str_df.withColumn("real_data", from_json("value_str", schema)).select("real_data.*") final_df.show()
内容的提问来源于stack exchange,提问作者Burak
相关产品推荐
相关产品推荐

