PySpark如何将JSON字符串列拆分为id、经纬度独立列
问题根因
from_json函数的返回值是Struct结构体类型,你当前的代码只是把二进制value转成字符串后解析成了结构体,仍然存在value这一列下,没有把结构体内部的字段拆成顶层独立列,所以输出不符合预期。
修正代码
直接在解析完JSON后,展开结构体字段即可,两种实现方式按需选择:
from pyspark.sql.types import StructType, FloatType, StringType from pyspark.sql.functions import from_json, col, trim # 定义JSON对应的schema location_schema = StructType() \ .add("id", StringType()) \ .add("latitude", FloatType()) \ .add("longitude", FloatType()) # 1. 二进制value转字符串,再解析为结构体列 parsed_df = df.selectExpr("CAST(value AS STRING) as raw_value") \ .withColumn("location_data", from_json(trim(col("raw_value"), '"'), location_schema)) # 2. 展开结构体为独立列,二选一即可 # 方式1:逐个指定字段,适合只需要部分字段的场景 result_df = parsed_df.select( col("location_data.id"), col("location_data.latitude"), col("location_data.longitude") ) # 方式2:用.*通配符批量展开结构体所有字段,适合字段多的场景 # result_df = parsed_df.select(col("location_data.*")) # 验证结果 result_df.printSchema() result_df.show()
额外注意点
- 代码里加了
trim处理首尾双引号的逻辑,是因为你提供的示例CSV数据里JSON外层多套了一层转义双引号,如果不处理会导致JSON解析失败返回null,如果实际从Event Hub读取的数据没有这个多余引号,可以去掉trim部分。 - 经纬度对精度要求高的话,可以把schema里的
FloatType替换为DoubleType,减少精度损耗。 - 如果需要保留Kafka/Event Hub自带的分区、偏移量、 ingestion时间戳等元数据字段,在select步骤把对应字段也加上即可。
内容的提问来源于stack exchange,提问作者user14681827
相关产品推荐
相关产品推荐

