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
相关产品推荐
相关产品推荐

