PySpark通过JDBC读取PostgreSQL报BigDecimal类型NaN错误的排查与解决
Bad value for type BigDecimal : NaN错误 环境版本
- Spark 3.3.1
- PostgreSQL 15.4
报错场景
通过PySpark从PostgreSQL读取数据并写入HDFS,核心脚本如下:
df = spark.read.jdbc( url = url, table = f'"{schema}"."{table}"', predicates = predicates, properties = properties ) df.write.mode('overwrite').partitionBy('site').parquet(output_path)
报错堆栈摘要
Caused by: org.apache.spark.SparkException: Job aborted due to stage failure: Task 17 in stage 2.0 failed 4 times, most recent failure: Lost task 17.3 in stage 2.0 (TID 173) (xxx.xxx.xxx executor 4): org.postgresql.util.PSQLException: Bad value for type BigDecimal : NaN
at org.postgresql.jdbc.PgResultSet.toBigDecimal(PgResultSet.java:3249)
at org.postgresql.jdbc.PgResultSet.getBigDecimal(PgResultSet.java:433)
at org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils$.$anonfun$makeGetter$3(JdbcUtils.scala:399)
...
Driver stacktrace:
at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2672)
...
问题成因
- 非法数据写入:PostgreSQL的
numeric/decimal类型原生不支持NaN值,但数据源中存在通过非常规手段(如错误插入逻辑、第三方工具)写入的NaN。 - 类型转换冲突:PostgreSQL JDBC驱动尝试将
NaN转换为Java的BigDecimal类型,但BigDecimal不兼容NaN,直接抛出转换错误。 - Spark自动推断缺陷:Spark读取时自动将目标字段推断为
Decimal类型,遇到非法NaN值时无法完成正常解析。
解决方案
方案1:清理数据源中的非法值
从根源解决,直接修复PostgreSQL中的数据:
-- 将NaN替换为NULL(根据业务需求选择合适的替代值) UPDATE "schema"."table" SET target_column = NULL WHERE target_column = 'NaN'; -- 若业务允许,直接删除含NaN的行 DELETE FROM "schema"."table" WHERE target_column = 'NaN';
方案2:自定义Schema强制读取为字符串
手动指定Schema,将目标字段读取为字符串后再处理:
# 替换为实际表结构,将numeric类型字段定义为STRING custom_schema = "id INT, target_column STRING, site STRING, ..." df = spark.read.jdbc( url=url, table=f'"{schema}"."{table}"', predicates=predicates, properties=properties, schema=custom_schema ) # 后续转换:将字符串转为Decimal,NaN转为NULL from pyspark.sql.functions import when, col df = df.withColumn( "target_column", when(col("target_column") == "NaN", None).otherwise(col("target_column").cast("Decimal(38,18)")) )
方案3:JDBC连接参数兼容处理
在连接属性中添加stringtype=unspecified,让驱动将无法解析的数值转为字符串:
properties = { "user": "your_user", "password": "your_password", "driver": "org.postgresql.Driver", "stringtype": "unspecified" # 添加该参数 } df = spark.read.jdbc( url=url, table=f'"{schema}"."{table}"', predicates=predicates, properties=properties ) # 同样需要处理NaN值 df = df.withColumn( "target_column", when(col("target_column") == "NaN", None).otherwise(col("target_column").cast("Decimal(38,18)")) )
方案4:升级依赖版本
尝试升级PostgreSQL JDBC驱动到42.6.0及以上稳定版,或升级Spark到3.4+版本,新版本可能修复了此类类型转换兼容性问题。
内容的提问来源于stack exchange,提问作者willshen

