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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 00:40:58