PySpark在AWS ECS写入Parquet时遇Java侧无响应错误(仅大文件失败)
问题描述
在AWS ECS中部署PySpark处理数据,先执行以下数据转换操作:
w = Window.partitionBy(*columns).orderBy(F.desc(column1)) df = input_data_frame.withColumn("rank", F.rank().over(w)).where(col("rank") == 1).drop("rank")
随后尝试将处理后的DataFrame写入指定位置:
data_frame.write.mode("overwrite").partitionBy("column2").save(file_location, format="parquet")
仅处理大文件时失败,抛出Py4J网络相关错误,完整错误栈如下:
ERROR:root:Exception while sending command. Traceback (most recent call last): File "/usr/local/lib/python3.10/site-packages/py4j/java_gateway.py", line 1207, in send_command raise Py4JNetworkError("Answer from Java side is empty") py4j.protocol.Py4JNetworkError: Answer from Java side is empty During handling of the above exception, another exception occurred: Traceback (most recent call last): File "/usr/local/lib/python3.10/site-packages/py4j/java_gateway.py", line 1033, in send_command response = connection.send_command(command) File "/usr/local/lib/python3.10/site-packages/py4j/java_gateway.py", line 1211, in send_command raise Py4JNetworkError( py4j.protocol.Py4JNetworkError: Error while receiving ERROR:py4j.java_gateway:An error occurred while trying to connect to the Java server (127.0.0.1:46777) Traceback (most recent call last): File "/app/kafka-connect-ingest/preprocessing/preprocessing_files_invoker.py", line 42, in main executor.execute() File "/app/kafka-connect-ingest/preprocessing/src/preprocessing/preprocessing_files.py", line 64, in execute file_location = self.write_effective_deltas(effective_delta, tmp_dir) File "/app/kafka-connect-ingest/preprocessing/src/preprocessing/preprocessing_files.py", line 195, in write_effective_deltas data_frame.write.mode("overwrite").partitionBy("dl_delete_flag").save(file_location, format="parquet") File "/usr/local/lib/python3.10/site-packages/pyspark/sql/readwriter.py", line 1109, in save self._jwrite.save(path) File "/usr/local/lib/python3.10/site-packages/py4j/java_gateway.py", line 1304, in __call__ return_value = get_return_value( File "/usr/local/lib/python3.10/site-packages/pyspark/sql/utils.py", line 111, in deco return f(*a, **kw) File "/usr/local/lib/python3.10/site-packages/py4j/protocol.py", line 334, in get_return_value raise Py4JError( py4j.protocol.Py4JError: An error occurred while calling o55.save
当前使用的Spark配置如下:
spark = SparkSession.builder.appName("Data preprocessing").getOrCreate() spark._jsc.hadoopConfiguration().set("fs.s3.canned.acl", "BucketOwnerFullControl") spark.conf.set("spark.sql.parquet.outputTimestampType", "TIMESTAMP_MILLIS") spark.conf.set("spark.sql.parquet.writeLegacyFormat", "True") spark.conf.set("spark.sql.parquet.enableVectorizedReader", "False") spark.conf.set("spark.sql.shuffle.partitions", "1") spark.conf.set("spark.sql.files.maxRecordsPerFile", "40000")
需要解决大文件写入失败的问题,确保文件能成功写入目标位置。
解决方案
1. 调整Spark JVM内存配置(核心修复)
Py4J网络错误大多是Java端Spark Driver/Executor进程因内存不足崩溃导致的,大文件处理时内存压力陡增,需针对性调整内存参数:
- 在
SparkSession初始化时添加Driver和Executor内存配置,根据ECS容器实际资源调整(比如容器内存16G时,可设置Driver内存12G、Executor内存8G) - 增加堆外内存配置,避免JVM因堆外内存不足被终止
修改后的SparkSession示例:
spark = SparkSession.builder.appName("Data preprocessing")\ .config("spark.driver.memory", "12g")\ .config("spark.executor.memory", "8g")\ .config("spark.driver.memoryOverhead", "2g")\ .config("spark.executor.memoryOverhead", "2g")\ .getOrCreate()
2. 优化分区与Shuffle配置
当前spark.sql.shuffle.partitions设为1,大文件处理时单个分区数据量过大,会直接导致内存过载:
- 将
spark.sql.shuffle.partitions调整为合理值(比如64或128,根据数据量和容器CPU核心数适配) - 写入前对DataFrame重分区,拆分过大的分区:
# 按分区列+额外分区数拆分,避免单分区数据量过大 data_frame = data_frame.repartition("column2", 64) data_frame.write.mode("overwrite").partitionBy("column2").save(file_location, format="parquet")
3. 提升ECS容器资源配额
确保ECS任务的CPU和内存分配能支撑大文件处理:
- 增加容器CPU配额(比如从1核提升至4核),提升Spark并行处理能力
- 增加容器内存(比如从4G提升至16G),避免系统OOM Killer终止Spark进程
4. 优化Parquet写入性能
- 若无需兼容旧版本Parquet格式,开启
spark.sql.parquet.enableVectorizedReader提升读写性能,降低内存消耗:
spark.conf.set("spark.sql.parquet.enableVectorizedReader", "True")
- 写入时启用Snappy压缩,减少IO压力和磁盘占用:
data_frame.write.mode("overwrite")\ .partitionBy("column2")\ .option("compression", "snappy")\ .save(file_location, format="parquet")
5. 排查S3写入稳定性
如果错误伴随S3超时,添加S3连接优化配置:
spark._jsc.hadoopConfiguration().set("fs.s3a.connection.maximum", "100") spark._jsc.hadoopConfiguration().set("fs.s3a.connection.timeout", "300000") spark._jsc.hadoopConfiguration().set("fs.s3a.socket.timeout", "300000")
内容的提问来源于stack exchange,提问作者Sagar Javeri
相关产品推荐
相关产品推荐

