Spark 3.5.0读取PostgreSQL bpchar列执行show()触发OutOfMemoryError问题
Spark 3.5.0(PySpark)读取PostgreSQL bpchar类型列触发OutOfMemoryError问题
问题现象
在Spark 3.5.0的PySpark环境下,读取PostgreSQL中bpchar类型的列并执行show()等操作时,会抛出OutOfMemoryError,该问题在Spark 3.4.1及更低版本中不会出现。
复现步骤
- PostgreSQL建表语句:
CREATE TABLE test.test_table (id bpchar);
- PySpark读取并执行show()代码:
df = spark.read.format("jdbc").option("url", f"jdbc:postgresql://{host_db}:{db_port}/{dbname}").option("driver", "org.postgresql.Driver").option("query", 'select id from test.test_table').option("user", n).option("password", s).load() df.show()
触发的错误信息
ERROR Executor: 阶段1.0中任务0.0(TID 1)出现异常 java.lang.OutOfMemoryError: Requested array size exceeds VM limit at org.apache.spark.unsafe.types.UTF8String.rpad(UTF8String.java:880) at org.apache.spark.sql.catalyst.util.CharVarcharCodegenUtils.readSidePadding(CharVarcharCodegenUtils.java:62) at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source) at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43) at org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43) at org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388) at org.apache.spark.sql.execution.SparkPlan$$Lambda$2449/0x0000000801093040.apply(Unknown Source) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:890) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:890) at org.apache.spark.rdd.RDD$$Lambda$2450/0x0000000801094040.apply(Unknown Source) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93) at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161) at org.apache.spark.scheduler.Task.run(Task.scala:141) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620) at org.apache.spark.executor.Executor$TaskRunner$$Lambda$2410/0x000000080105b840.apply(Unknown Source) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829)
问题原因
Spark 3.5.0对char/bpchar类型的处理逻辑进行了变更,当读取未指定长度的bpchar列时,错误地尝试将其填充到极大的长度,导致申请超出JVM限制的数组,最终触发OutOfMemoryError。该问题已在Spark社区被反馈。
解决方案
- 显式指定bpchar列长度:建表时明确指定bpchar的长度,例如:
CREATE TABLE test.test_table (id bpchar(10));
- 读取时转换列类型:在JDBC查询中将bpchar列显式转换为varchar类型,避免Spark的错误填充逻辑:
df = spark.read.format("jdbc").option("url", f"jdbc:postgresql://{host_db}:{db_port}/{dbname}").option("driver", "org.postgresql.Driver").option("query", 'select cast(id as varchar) from test.test_table').option("user", n).option("password", s).load()
- 降级Spark版本:暂时回退到Spark 3.4.1及以下版本,避开该问题。
- 等待官方修复:关注Spark后续版本更新,该bug预计会在后续版本中修复。
内容的提问来源于stack exchange,提问作者yjsa
相关产品推荐
相关产品推荐

