You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 21:54:32