Pyspark从Postgres同步超1.5亿条数据到S3存Parquet性能优化求助
PySpark同步Postgres 1.5亿大表到S3性能优化建议
先修复原脚本基础语法问题:
username、password变量未赋值,jdbcUrl字符串拼接逻辑存在语法错误,修正后再应用以下优化方案。
1. 修复JDBC读取瓶颈(核心问题)
原脚本未配置并行读取,全表1.5亿条数据走单连接拉取,是任务失败的核心原因,优化如下:
- 启用JDBC分片并行读取,利用表的
created分区键作为分片列,将读取任务拆分为多个并行子任务,避免单连接压力过大:
def read_data(self): conpropererties = self.connection_props()[0] url = self.connection_props()[1] # 补充Postgres驱动配置,开启谓词下推 conpropererties["driver"] = "org.postgresql.Driver" conpropererties["pushDownPredicate"] = "true" # 先查询分片列的上下界 bound_query = f"(select min({self.partition_key}) as min_val, max({self.partition_key}) as max_val from {self.database}.{self.schema}.{self.table}) as bound" bound_row = spark.read.jdbc(url=url, table=bound_query, properties=conpropererties).collect()[0] lower_bound = bound_row["min_val"] upper_bound = bound_row["max_val"] # 并行读取,numPartitions根据集群资源和Postgres负载调整,建议20~50 df = spark.read.jdbc( url=url, table=f"{self.database}.{self.schema}.{self.table}", properties=conpropererties, partitionColumn=self.partition_key, lowerBound=lower_bound, upperBound=upper_bound, numPartitions=30 ) return df
- 移除不合理的
coalesce(200)操作:单连接读取的DataFrame初始只有1个分区,coalesce不会触发shuffle,无法实现分区拆分,完全无效反而会导致后续数据倾斜。
2. 写入S3阶段优化
- 提前配置S3写入优化参数,减少写入开销:
# 初始化Spark时配置以下参数 spark.conf.set("spark.hadoop.fs.s3a.fast.upload", "true") spark.conf.set("spark.hadoop.fs.s3a.multipart.size", "104857600") # 100M分片上传 spark.conf.set("spark.sql.parquet.compression.codec", "snappy") # 开启snappy压缩,降低文件体积 spark.conf.set("spark.sql.shuffle.partitions", "200") # 调整shuffle分区数匹配资源
- 控制分区文件大小:写入前添加
repartition(self.partition_key)操作,避免生成大量小文件,单Parquet文件大小控制在128M~1G区间性能最优。 - 如果
created字段是时间戳类型,建议先转换为日期字段做粗分区,避免生成过多细碎分区,降低S3管理开销。
3. 稳定性优化
- 控制JDBC并行度不超过Postgres实例最大连接数的1/3,避免同步任务打垮业务库。
- 1.5亿条数据可按时间切片分批同步,比如按天维度拉取写入,避免单次任务处理数据量过大导致超时。
- 调整Spark资源配置:Executor内存建议设置为8G16G,单Executor核心数设置为24,Executor总数根据集群总资源调整。
内容的提问来源于stack exchange,提问作者Tanya Kapoor
相关产品推荐
相关产品推荐

