如何验证JDBC_SESSION_INIT_STATEMENT中的SQL语句是否生效?
问题:Spark-JDBC连接SQL Server时JDBC_SESSION_INIT_STATEMENT未生效的验证方法
尝试通过Spark-JDBC连接SQL Server,使用JDBC_SESSION_INIT_STATEMENT创建临时表,再通过主查询从该临时表下载数据,代码如下:
//df is org.apache.spark.sql.DataFrameReader val s = """select * into #tmp_table from ( SELECT op.ID, | op.Date, | op.DocumentID, | op.Amount, | op.AmountCurr, | op.CurrencyID, | operson.ObjectTypeId AS PersonOT, | op.PersonID, | ocontract.ObjectTypeId AS ContractOT, | op.ContractID, | op.DocNum, | op.MomentCreate, | op.ObjectTypeID, | op.OwnerObjectID |FROM dbo.Operation op With (Index = IX_Operation_Date) --Без хинта временами уходит в скан всей таблицы |LEFT JOIN dbo.Object ocontract ON op.ContractID = ocontract.ID |LEFT JOIN dbo.Object operson ON op.PersonID = operson.ID |WHERE op.Date>='2019-01-01' and op.Date<'2020-01-01' AND 1=1 |) wrap_for_single_connect |OPTION (LOOP JOIN, FORCE ORDER, MAX_GRANT_PERCENT=25)""".stripMargin df .option(JDBCOptions.JDBC_SESSION_INIT_STATEMENT, s) .jdbc( jdbcUrl, "(select * from tempdb.#tmp_table) sub", connectionProps)
运行后报错:
com.microsoft.sqlserver.jdbc.SQLServerException: Invalid object name '#tmp_table'.
怀疑JDBC_SESSION_INIT_STATEMENT未生效,因为故意写错该语句后仍得到相同的“无效对象”错误,需验证该配置中的请求是否正常工作。
解决方案与验证方法
一、验证JDBC_SESSION_INIT_STATEMENT是否执行
1. 数据库端日志追踪
- 使用SQL Server Profiler创建跟踪任务,过滤当前连接的SQL语句,查看初始化语句是否被执行。
- 或预先创建一张日志表:
然后修改初始化语句,在末尾追加日志写入逻辑:CREATE TABLE dbo.DebugLog (LogText VARCHAR(255), LogTime DATETIME DEFAULT GETDATE())
执行Spark任务后查询INSERT INTO dbo.DebugLog (LogText) VALUES ('Session init statement executed')dbo.DebugLog,若有新增记录则说明初始化语句已执行。
2. 语法错误注入测试
将初始化语句修改为明显有语法错误的内容(比如把SELECT写成SELCT),如果Spark抛出SQL语法错误,说明初始化语句在执行;若仍报临时表不存在,则证明初始化语句未被触发,需检查参数配置或版本兼容性。
二、临时表作用域问题排查
本地临时表(#开头)仅在当前会话生效,若Spark-JDBC的初始化语句与主查询不在同一会话,就会出现找不到临时表的情况:
- 尝试替换为全局临时表(
##开头),将初始化语句中的#tmp_table改为##tmp_table,主查询改为(select * from tempdb.##tmp_table) sub,验证是否能访问。
三、代码优化建议
若临时表方案存在会话隔离问题,可直接将逻辑合并为CTE(公共表表达式),避免依赖临时表:
val query = """ WITH tmp_table AS ( SELECT op.ID, op.Date, op.DocumentID, op.Amount, op.AmountCurr, op.CurrencyID, operson.ObjectTypeId AS PersonOT, op.PersonID, ocontract.ObjectTypeId AS ContractOT, op.ContractID, op.DocNum, op.MomentCreate, op.ObjectTypeID, op.OwnerObjectID FROM dbo.Operation op With (Index = IX_Operation_Date) LEFT JOIN dbo.Object ocontract ON op.ContractID = ocontract.ID LEFT JOIN dbo.Object operson ON op.PersonID = operson.ID WHERE op.Date>='2019-01-01' and op.Date<'2020-01-01' ) SELECT * FROM tmp_table OPTION (LOOP JOIN, FORCE ORDER, MAX_GRANT_PERCENT=25) """.stripMargin df.jdbc(jdbcUrl, s"($query) sub", connectionProps)
另外,可检查参数名是否正确,部分旧版本Spark需手动使用字符串参数名"sessionInitStatement"替代JDBCOptions.JDBC_SESSION_INIT_STATEMENT常量,同时确保使用的mssql-jdbc驱动版本与Spark、SQL Server版本兼容。
内容的提问来源于stack exchange,提问作者Gumada Yaroslav
相关产品推荐
相关产品推荐

