spark.sql执行SQL抛解析异常但%sql魔法命令可正常运行
问题根因
spark.sql() 原生API默认不支持单次调用执行分号分隔的多条SQL语句,这是它和Notebook中%sql魔法命令的核心差异:
%sql魔法命令在提交执行前会自动对SQL文本做预处理,按分号拆分出单条独立SQL,逐条提交给Spark引擎执行,因此拼接多条语句可以正常运行- 原生
spark.sql()会将传入的整段字符串作为单条SQL做语法解析,分号后的第二条语句会被判定为非法语法,直接抛出解析异常
修复方案
方案1:拆分语句分别执行(推荐,兼容性最强)
将建表和加字段的SQL拆分为独立字符串,分两次调用spark.sql()即可,示例代码如下:
# 建表语句 create_sql = """ CREATE OR REPLACE TABLE Invoices ( InvoiceID INT, CustomerID INT, BillToCustomerID INT, OrderID INT, DeliveryMethodID INT, ContactPersonID INT, AccountsPersonID INT, SalespersonPersonID INT, PackedByPersonID INT, InvoiceDate TIMESTAMP, CustomerPurchaseOrderNumber INT, IsCreditNote STRING, CreditNoteReason STRING, Comments STRING, DeliveryInstructions STRING, InternalComments STRING, TotalDryItems INT, TotalChillerItems STRING, DeliveryRun STRING, RunPosition STRING, ReturnedDeliveryData STRING, ConfirmedDeliveryTime TIMESTAMP, ConfirmedReceivedBy STRING, LastEditedBy INT, LastEditedWhen TIMESTAMP ) LOCATION '/mnt/adls/DQD/udl/Invoices/' """ spark.sql(create_sql) # 新增字段语句 add_column_sql = "ALTER TABLE Invoices ADD COLUMN DQ_Check_Op SMALLINT" spark.sql(add_column_sql)
该写法不受Spark版本限制,执行报错时可以快速定位到具体失败的语句,排查成本更低。
方案2:开启多语句执行配置(Spark 3.x及以上支持)
如果不想拆分语句,可以开启Spark内置的多语句执行开关,配置生效后spark.sql()就支持解析分号分隔的多条语句,示例代码如下:
# 开启多语句执行支持 spark.conf.set("spark.sql.multipleStatements.enabled", "true") # 直接传入拼接好的多段SQL执行 multi_sql = """ CREATE OR REPLACE TABLE Invoices (InvoiceID INT, CustomerID INT, BillToCustomerID INT, OrderID INT, DeliveryMethodID INT, ContactPersonID INT, AccountsPersonID INT, SalespersonPersonID INT, PackedByPersonID INT, InvoiceDate TIMESTAMP, CustomerPurchaseOrderNumber INT, IsCreditNote STRING, CreditNoteReason STRING, Comments STRING, DeliveryInstructions STRING, InternalComments STRING, TotalDryItems INT, TotalChillerItems STRING, DeliveryRun STRING, RunPosition STRING, ReturnedDeliveryData STRING, ConfirmedDeliveryTime TIMESTAMP, ConfirmedReceivedBy STRING, LastEditedBy INT, LastEditedWhen TIMESTAMP) LOCATION '/mnt/adls/DQD/udl/Invoices/'; ALTER TABLE Invoices ADD COLUMN DQ_Check_Op SMALLINT """ spark.sql(multi_sql)
注意:该配置在部分低版本Spark、定制化发行版Spark中可能不支持,生产环境优先使用方案1。
内容的提问来源于stack exchange,提问作者SouravA
相关产品推荐
相关产品推荐

