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

PySpark应用函数报错:不支持嵌套StructType的Arrow转换

问题场景

在Pandas中通过apply调用API获取响应,存入新列后用json_normalize解析嵌套JSON完全正常,但切换到PySpark pandas(psdf)执行apply时,抛出TypeError: Nested StructType not supported in conversion from Arrow错误,无法完成后续的嵌套结构解析。

报错原因

PySpark pandas的apply方法底层依赖Apache Arrow进行Pandas与Spark数据格式的转换,而当前版本的Arrow转换不支持嵌套StructType(即多层嵌套的字典/对象结构)。当apply返回嵌套字典时,Arrow无法将其正确转换为Spark可识别的类型,从而触发报错。

解决方案

以下两种方法可解决该问题,优先推荐方法一(原生Spark UDF),性能与兼容性更好:

方法一:使用Spark原生UDF替代ps.apply

Spark原生UDF对复杂嵌套类型的支持更完善,步骤如下:

  1. 提前定义API返回数据的StructType schema(需与API响应结构完全匹配)
  2. 编写UDF实现API调用逻辑,返回符合schema的字典
  3. 应用UDF到Spark DataFrame,再通过select直接解析嵌套字段(替代json_normalize)
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, DoubleType, StringType
import requests
import json

# 初始化SparkSession
spark = SparkSession.builder.appName("API_Normalization").getOrCreate()

# 定义API响应的StructType schema(根据实际API返回调整)
api_schema = StructType([
    StructField("location", StructType([
        StructField("lat", DoubleType()),
        StructField("lon", DoubleType()),
        StructField("city", StringType())
    ])),
    StructField("data", StructType([
        StructField("geographies", StringType())
    ]))
])

def call_api(a, b):
    # 替换为实际API调用逻辑
    # resp = requests.get(f"https://your-api.com?param1={a}&param2={b}")
    # return resp.json()
    # 以下为模拟返回数据
    return {
        "location": {"lat": a * 0.1, "lon": b * 0.1, "city": "SampleCity"},
        "data": {"geographies": "SampleRegion"}
    }

# 注册UDF
api_udf = udf(call_api, api_schema)

# 创建Spark DataFrame
sdf = spark.createDataFrame([(1,4), (2,5), (3,6)], ["A", "B"])

# 应用UDF生成API响应列
sdf = sdf.withColumn("api_response", api_udf(sdf["A"], sdf["B"]))

# 解析嵌套字段,等价于pandas的json_normalize
normalized_sdf = sdf.select(
    "A", "B",
    "api_response.location.lat",
    "api_response.location.lon",
    "api_response.location.city",
    "api_response.data.geographies"
)

# 如需转回PySpark pandas DataFrame
normalized_psdf = ps.DataFrame(normalized_sdf)

方法二:ps.apply返回JSON字符串,再解析

如果必须使用PySpark pandas的apply,可先将API响应转为JSON字符串(Arrow支持字符串类型),再通过Spark的from_json函数解析嵌套结构:

import pyspark.pandas as ps
import json
from pyspark.sql.functions import from_json, col
from pyspark.sql.types import StructType, StructField, DoubleType, StringType

def call_api_str(row):
    # 调用API后返回JSON字符串
    resp_data = {
        "location": {"lat": row["A"] * 0.1, "lon": row["B"] * 0.1},
        "data": {"geographies": "SampleRegion"}
    }
    return json.dumps(resp_data)

# 创建PySpark pandas DataFrame
psdf = ps.DataFrame({"A": [1,2,3], "B": [4,5,6]})

# apply返回字符串列
psdf["api_response_str"] = psdf.apply(call_api_str, axis=1)

# 定义API响应schema(同方法一)
api_schema = StructType([
    StructField("location", StructType([
        StructField("lat", DoubleType()),
        StructField("lon", DoubleType())
    ])),
    StructField("data", StructType([
        StructField("geographies", StringType())
    ]))
])

# 转为Spark DataFrame解析字符串
sdf = psdf.to_spark()
sdf = sdf.withColumn("api_response", from_json(col("api_response_str"), api_schema))

# 解析嵌套字段
normalized_sdf = sdf.select("A", "B", "api_response.location.*", "api_response.data.geographies")

# 转回PySpark pandas DataFrame
normalized_psdf = ps.DataFrame(normalized_sdf)
关键注意事项
  • 必须提前明确API返回的结构并定义对应的StructType schema,Spark无法自动推断复杂嵌套类型
  • 避免在ps.apply中直接返回嵌套字典,这是触发Arrow转换错误的核心原因
  • 如果API响应结构不固定,可使用MapType替代StructType,但解析时需手动处理键值对,灵活性较低

内容的提问来源于stack exchange,提问作者kms

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 05:53:27