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

PySpark中将字典列表转为JSON时遇ndarray序列化错误求助

解决Spark pandas_udf中ndarray无法JSON序列化的问题

问题场景

通过Spark分组聚合生成descr列(值为JSON字符串组成的列表),尝试用pandas_udf压缩该列时触发错误:

TypeError: Object of type ndarray is not JSON serializable

原始代码如下:

from pyspark.sql import SparkSession, functions as F

df = spark.createDataFrame(
    [
        (2022,'A1', "cat", 'eng', 3, 56.768639), 
        (2022,'A1', "rabbit", 'eng', 10, 56.768639), 
        (2022, 'A2', "dog", 'eng', 10, 54.114841),
        (2022, 'A2', "mouse", 'eng', 20, 81.114841),
    ],
    ["data",'group', "word", 'lang', 'count', 'value']
)

df2 = df\
    .groupBy('data', 'group', 'lang')\
    .agg(F.collect_list(F.to_json(F.struct(F.col('count'), F.col('value'), F.col('word')))).alias('descr'))

# 报错的pandas_udf
import base64
import gzip
import json
from pyspark.sql.functions import pandas_udf, StringType

@pandas_udf(StringType())
def jsn(lst):
    return lst.apply(lambda lst: base64.b64encode(gzip.compress(json.dumps(lst).encode('utf-8'))).decode("utf-8"))

df3= df2.withColumn('descr2', jsn(F.col('descr')))

错误原因

pandas_udf接收的Series中,列表类型的元素会被自动转为numpy ndarray,而json.dumps无法直接序列化numpy数组类型,因此触发报错。

解决方案

方案1:在udf中将ndarray转回Python列表

修改udf逻辑,先将numpy数组转为原生Python列表再序列化:

@pandas_udf(StringType())
def jsn(lst_series):
    def process_item(arr):
        # 将numpy数组转为Python列表
        python_list = arr.tolist()
        # 序列化并压缩
        json_str = json.dumps(python_list)
        compressed_data = gzip.compress(json_str.encode('utf-8'))
        return base64.b64encode(compressed_data).decode("utf-8")
    return lst_series.apply(process_item)

方案2:优化Spark聚合逻辑(更高效)

直接利用Spark内置函数生成完整的JSON数组,避免后续处理ndarray的麻烦:

# 聚合时直接收集结构体并转为JSON数组
df2 = df\
    .groupBy('data', 'group', 'lang')\
    .agg(F.to_json(F.collect_list(F.struct(F.col('count'), F.col('value'), F.col('word')))).alias('descr'))

# 简化udf,直接压缩已生成的JSON字符串
@pandas_udf(StringType())
def compress_json(json_series):
    def process_str(s):
        compressed_data = gzip.compress(s.encode('utf-8'))
        return base64.b64encode(compressed_data).decode("utf-8")
    return json_series.apply(process_str)

df3 = df2.withColumn('descr2', compress_json(F.col('descr')))

方案2的优势:利用Spark内置函数完成JSON数组的生成,减少了udf中的类型转换开销,性能更优,从根源上避免了ndarray序列化问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 13:01:43