如何用PySpark StructType适配含未知长度坐标的API JSON结果?
问题描述
我需要将PySpark StructType格式化为匹配某API返回的JSON结果,该API返回的完整结果示例如下:
{ "type": "Polygon", "coordinates": [ [ [-74.53703811195342, 43.93214895162186], [-74.53765823498132, 43.932511606633376], [-74.53790321887529, 43.933052813967755], ... [-74.53653908240167, 43.93223696479821], [-74.53703811195342, 43.93214895162186] ] ] }
注意:获取JSON前,coordinates元素的长度是未知且动态的。
我的尝试代码如下:
# set up UDF schema = StructType([ StructField("type", StringType()), StructField("coordinates", StringType()) # I just put a random StringType ]) # then after the API query, result_df = request_df \ .withColumn("result", udf_executeRestApi(col("body"))) df = result_df.select([col for col in result_df.columns]) df.show()
执行后,Polygon对象能显示,但格式错误,预定义的type和coordinates字段未被识别。
解决方案
- 修正Schema定义
你之前给coordinates用StringType完全不符合API返回的嵌套数组结构。正确的Schema要把coordinates定义为三层嵌套的浮点数组,对应API返回的数组<数组<数组<double>>>结构(外层数组包含多边形的环,每个环是经纬度坐标对组成的数组,每个坐标对是两个浮点数)。
修改后的Schema代码:
from pyspark.sql.types import StructType, StructField, StringType, ArrayType, DoubleType schema = StructType([ StructField("type", StringType(), nullable=False), StructField("coordinates", ArrayType(ArrayType(ArrayType(DoubleType()))), nullable=False) ])
- 确保UDF返回对应结构
你的UDFudf_executeRestApi必须返回符合上述Schema的Python数据结构(字典+嵌套列表),不能是JSON字符串。如果当前UDF返回的是JSON字符串,需要先解析成Python对象:
def execute_rest_api(body): import requests # 调用API并解析响应为Python字典 resp = requests.post("你的API地址", json=body) return resp.json() # 注册UDF时绑定正确的Schema udf_executeRestApi = udf(execute_rest_api, schema)
- 正确提取嵌套字段
运行DataFrame代码后,就可以直接提取result中的type和coordinates字段了:
result_df = request_df.withColumn("result", udf_executeRestApi(col("body"))) # 展开嵌套字段并移除原始result列 df = result_df.select( "*", col("result.type").alias("polygon_type"), col("result.coordinates").alias("polygon_coordinates") ).drop("result") df.show(truncate=False)
内容的提问来源于stack exchange,提问作者Jane
相关产品推荐
相关产品推荐

