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

PySpark连接MSSQL Server报NullPointerException问题求助

PySpark连接MSSQL Server抛出NullPointerException的解决方案

问题概述

PySpark新手尝试连接MSSQL Server,脚本能输出正确的连接参数,但调用spark.read.jdbc()时抛出NullPointerException,已尝试mssql-jdbc-6.2.1.jre7.jar、sqljdbc41.jar、jtds-1.3.1-patch.jar三种驱动,问题未解决。

关键信息

输出的连接参数

('Processing table:', u'POL_ACTION_AMEND')
('Table schema:', u'dbo')
('Source_database:', u'PRD01_IPS')
('SQL Query:', '(SELECT TOP 100 * FROM PRD01_IPS.dbo.POL_EVENT_HISTORY)')
('jdbc_uri:', 'jdbc:jtds:sqlserver://REPLCLPRD01\REPL/PRD01_IPS')

核心代码

select_statement_sql_server = "(SELECT TOP 100 * FROM {source_database}.{schema}.{table_name})"

data_df = (
    spark.read.format("jdbc")
        .option("url", source_jdbc_uri)
        .option("user", source_jdbc_user)
        .option("password", source_jdbc_pass)
        .option("driver", source_jdbc_driver)
        .option("dbtable", select_statement_sql_server.format(
            source_database=source_database, schema=schema, table_name='POL_EVENT_HISTORY'
        ))
        .load()
)
data.show()

报错堆栈

('SQL Query:', '(SELECT TOP 100 * FROM PRD01_IPS.dbo.POL_EVENT_HISTORY)')
Traceback (most recent call last):
  File "/pkg/lxd0bigd/Talend_To_Pyspark/PSTAR_IPS/200job_IPS.py", line 136, in <module>
    .option("dbtable", select_statement_sql_server.format(source_database=source_database, schema=schema, table_name='POL_EVENT_HISTORY')) \
  File "/opt/cloudera/parcels/CDH-7.1.8-1.cdh7.1.8.p0.30990532/lib/spark/python/lib/pyspark.zip/pyspark/sql/readwriter.py", line 172, in load
  File "/opt/cloudera/parcels/CDH-7.1.8-1.cdh7.1.8.p0.30990532/lib/spark/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 1257, in __call__
  File "/opt/cloudera/parcels/CDH-7.1.8-1.cdh7.1.8.p0.30990532/lib/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 63, in deco
  File "/opt/cloudera/parcels/CDH-7.1.8-1.cdh7.1.8.p0.30990532/lib/spark/python/lib/py4j-0.10.7-src.zip/py4j/protocol.py", line 328, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling o236.load.
: java.lang.NullPointerException
    at org.apache.spark.sql.execution.datasources.jdbc.JDBCRDD$.resolveTable(JDBCRDD.scala:71)
    at org.apache.spark.sql.execution.datasources.jdbc.JDBCRelation$.getSchema(JDBCRelation.scala:211)
    at org.apache.spark.sql.execution.datasources.jdbc.JdbcRelationProvider.createRelation(JdbcRelationProvider.scala:35)
    at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:332)
    at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:243)
    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:231)
    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:187)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
    at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
    at py4j.Gateway.invoke(Gateway.java:282)
    at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
    at py4j.commands.CallCommand.execute(CallCommand.java:79)
    at py4j.GatewayConnection.run(GatewayConnection.java:238)
    at java.lang.Thread.run(Thread.java:750)

提交命令

spark-submit --master yarn --deploy-mode client \
    --driver-class-path ${DRIVER_AND_JAR_FILE_PATH} \
    --jars ${DRIVER_AND_JAR_FILE_PATH},${JAR_FILE_XML} \
    --conf "spark.dynamicAllocation.enabled=true" \
    --conf "spark.yarn.dist.files=/etc/hive/conf.cloudera.hive/hive-site.xml" \
    ${ROOT_DIR}/200job_IPS.py \
    --hdfspath ${XML_FILE_HDFS} >> ${LOGFILE_DIR}/200job_log.txt 2>&1 

sh /pkg/lxd0bigd/Talend_To_Pyspark/PSTAR_IPS/Ingestion_PSTAR.sh dev \
    /pkg/lxd0bigd/Talend_To_Pyspark/spark-jobs/lib/mssql-jdbc-6.2.1.jre7.jar

排查与解决步骤

1. 强制校验驱动类与URL匹配

NullPointerException出现在resolveTable阶段,大概率是驱动未正确加载或URL格式不兼容:

  • jtds驱动:
    • URL格式:jdbc:jtds:sqlserver://REPLCLPRD01\REPL;databaseName=PRD01_IPS(注意用分号分隔数据库名,不是斜杠)
    • 驱动类:net.sourceforge.jtds.jdbc.Driver
  • 微软官方驱动(mssql-jdbc/sqljdbc41):
    • URL格式:jdbc:sqlserver://REPLCLPRD01\REPL;databaseName=PRD01_IPS;encrypt=false(如果不需要加密的话)
    • 驱动类:com.microsoft.sqlserver.jdbc.SQLServerDriver
  • 打印source_jdbc_driver变量的值,确认没有拼写错误或为空。

2. 确保驱动Jar在集群全局可访问

Yarn模式下,Executor节点需要能获取到驱动Jar:

  • 把驱动Jar上传到HDFS,比如hdfs:///user/yourname/lib/mssql-jdbc-7.4.1.jre8.jar,然后在--jars参数里指定HDFS路径
  • 避免使用本地路径,除非所有Worker节点的相同路径下都有该Jar包
  • 补充配置Executor端的类路径:
    spark-submit --master yarn --deploy-mode client \
        --driver-class-path hdfs:///path/to/driver.jar \
        --jars hdfs:///path/to/driver.jar \
        --conf "spark.executor.extraClassPath=driver.jar" \
        # 其他配置...
    

3. 简化测试,排除SQL语法/权限问题

  • 先跳过子查询,直接读取表:
    data_df = spark.read.format("jdbc")\
        .option("url", source_jdbc_uri)\
        .option("user", source_jdbc_user)\
        .option("password", source_jdbc_pass)\
        .option("driver", source_jdbc_driver)\
        .option("dbtable", "PRD01_IPS.dbo.POL_EVENT_HISTORY")\
        .load()
    
  • 如果能成功加载,再测试子查询;如果失败,用同账号在SQL Server客户端执行SELECT TOP 100 * FROM PRD01_IPS.dbo.POL_EVENT_HISTORY,验证权限是否正常。

4. 匹配驱动与JDK/Spark版本

CDH7.1.8默认用JDK8,你的mssql-jdbc-6.2.1.jre7.jar是针对JDK7的,存在兼容性问题:

  • 替换为mssql-jdbc-7.4.1.jre8.jar或更高版本的JDK8兼容驱动
  • jtds-1.3.1支持JDK8,但要确认你的SQL Server版本(2000-2016)是否在支持范围内

5. 手动验证驱动加载

在代码中添加驱动加载验证,确认驱动能被正确加载:

from py4j.java_gateway import java_import
java_import(spark._jvm, "java.lang.Class")

try:
    spark._jvm.Class.forName(source_jdbc_driver)
    print("✅ 驱动加载成功")
except Exception as e:
    print(f"❌ 驱动加载失败: {str(e)}")

如果加载失败,检查驱动Jar是否正确传递,或者驱动类名是否拼写错误。


内容的提问来源于stack exchange,提问作者Shiva

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 00:45:54