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
- URL格式:
- 微软官方驱动(mssql-jdbc/sqljdbc41):
- URL格式:
jdbc:sqlserver://REPLCLPRD01\REPL;databaseName=PRD01_IPS;encrypt=false(如果不需要加密的话) - 驱动类:
com.microsoft.sqlserver.jdbc.SQLServerDriver
- URL格式:
- 打印
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
相关产品推荐
相关产品推荐

