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

Azure Synapse PySpark处理5亿条JSON数据遇2G序列化限制报错

解决Azure Synapse PySpark序列化超过2G的报错

问题描述

在Azure Synapse的PySpark Notebook中,从API读取5亿条JSON记录并写入Azure Blob时,触发以下报错:

ValueError                                Traceback (most recent call last)
Cell In [17], line 45
 41             total_results.append((parsed_data))
 42             ##print(total_results)
 43 
 44 ##RDD Spark creation 
 ---> 45 rdd = spark.sparkContext.parallelize(total_results)
 46 df = spark.read.option('multiLine','true').json(rdd)
 48 #Create temporary view on dataframe                                                                                    

  File /opt/spark/python/lib/pyspark.zip/pyspark/context.py:686, in 
SparkContext.parallelize(self, c, numSlices)
683     assert self._jvm is not None
684     return self._jvm.PythonParallelizeServer(self._jsc.sc(), numSlices)
---> 686 jrdd = self._serialize_to_jvm(c, serializer, reader_func, createRDDServer)
687 return RDD(jrdd, self, serializer)

File /opt/spark/python/lib/pyspark.zip/pyspark/context.py:729, in 
SparkContext._serialize_to_jvm(self, data, serializer, reader_func, createRDDServer)
727 try:
728     try:
 ---> 729         serializer.dump_stream(data, tempFile)
730     finally:
731         tempFile.close()

File /opt/spark/python/lib/pyspark.zip/pyspark/serializers.py:224, in 
BatchedSerializer.dump_stream(self, iterator, stream)
223 def dump_stream(self, iterator, stream):
---> 224     self.serializer.dump_stream(self._batched(iterator), stream)

File /opt/spark/python/lib/pyspark.zip/pyspark/serializers.py:146, in 
FramedSerializer.dump_stream(self, iterator, stream)
144 def dump_stream(self, iterator, stream):
145     for obj in iterator:
---> 146         self._write_with_length(obj, stream)

File /opt/spark/python/lib/pyspark.zip/pyspark/serializers.py:160, in 
FramedSerializer._write_with_length(self, obj, stream)
158     raise ValueError("serialized value should not be None")
159 if len(serialized) > (1 << 31):
 ---> 160     raise ValueError("can not serialize object larger than 2G")
161 write_int(len(serialized), stream)
162 stream.write(serialized)

 ValueError: can not serialize object larger than 2G

原代码逻辑是将所有JSON数据存入列表,再创建RDD转换为DataFrame,过滤后写入Parquet:

rdd = spark.sparkContext.parallelize(total_results)
df = spark.read.option('multiLine','true').json(rdd)

#Create temporary view on dataframe                                                                                    
df.createOrReplaceTempView('filter_view')                                                                           
                                                                                                                      
#SQL query to filter on deleteddate value                                                                        
df_filter=spark.sql("""select * from filter_view where DeletedDate is null""")

df_filter.coalesce(800).write.format("parquet").save(stagingpath,mode="overwrite")

报错原因

原代码的核心问题是将5亿条数据全部加载到Driver节点的内存中(存入total_results列表),再通过parallelize转换为RDD时,需要将整个列表序列化传递给JVM,而Spark的序列化限制单个对象不超过2G,直接触发报错。同时这种方式会导致Driver内存溢出,完全违背了Spark的分布式计算设计。

解决方案

1. 采用分布式分批拉取数据

不要在Driver端存储全量数据,而是让Executor节点并行分批拉取API数据,避免Driver成为瓶颈:

# 配置API分页参数(根据实际API调整)
page_size = 10000  # 每页拉取的记录数
total_pages = 50000  # 总页数(根据5亿条计算得出)

# 创建分页编号的RDD,让每个Executor处理部分页码
page_rdd = spark.sparkContext.parallelize(range(1, total_pages + 1), numSlices=200)

# 定义单页数据拉取函数(需根据实际API调整)
def fetch_single_page(page_num):
    import requests
    # 替换为你的API请求逻辑,包含分页参数
    api_url = f"https://your-api-endpoint?page={page_num}&size={page_size}"
    response = requests.get(api_url)
    response.raise_for_status()  # 处理请求错误
    return response.json()

# 分布式处理每个分区的页码,拉取数据并转换为Row对象
from pyspark.sql import Row
def process_partition(page_iterator):
    for page_num in page_iterator:
        page_data = fetch_single_page(page_num)
        # 将JSON转为Spark Row对象,方便后续创建DataFrame
        yield [Row(**item) for item in page_data]

# 创建分布式DataFrame
raw_df = spark.createDataFrame(page_rdd.mapPartitions(process_partition))

# 过滤数据(直接用DataFrame API更高效,无需临时视图)
df_filter = raw_df.filter(raw_df.DeletedDate.isNull())

# 写入Parquet:用repartition均匀分配数据,提升写入性能
df_filter.repartition(800).write.format("parquet").save(stagingpath, mode="overwrite")

2. 额外优化建议

  • API请求防护:添加重试机制和请求延时,避免触发API限流
  • 网络配置:确保Synapse的Executor节点能访问目标API(如配置VNet集成)
  • 避免coalesce:coalesce是减少分区,可能导致数据倾斜,写入大文件时用repartition均匀拆分数据
  • 资源配置:根据数据量调整Synapse Spark池的节点数和内存配置,避免Executor OOM

内容的提问来源于stack exchange,提问作者Arun.K

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 17:01:13