Azure Synapse列名含空格致Databricks缓存DataFrame报错的解决办法
问题:Databricks读取带空格列名的Synapse表缓存失败
场景与代码
使用Databricks将Azure Synapse表数据读取至Spark DataFrame,代码如下:
df = spark.read \ .format("com.databricks.spark.sqldw") \ .option("url", sql_dw_connection_string) \ .option("tempDir", temp_dir_url) \ .option("enableServicePrincipalAuth", "true") \ .option("query", query) \ .load() df.cache()
错误信息
缓存DataFrame时触发以下错误:
com.databricks.spark.sqldw.SqlDWSideException: Azure Synapse Analytics failed to execute the JDBC query produced by the connector. Make sure column names do not include any invalid characters such as ';' or whitespace Underlying SQLException(s): - com.microsoft.sqlserver.jdbc.SQLServerException: 110802;An internal DMS error occurred that caused this operation to fail. Details: Exception: Microsoft.SqlServer.DataWarehouse.DataMovement.Common.ExternalAccess.HdfsAccessException, Message: Error occurred while accessing HDFS external file[/sqldw-staging/2023-03-07/07-14-49-790/f030341a-12d7-4ae1-b5e0-81b510a9a9e1/QID68547779_20230307_71451_0.parq.snappy][0]: Java exception raised on call to HdfsBridge_CreateRecordWriter. Java exception message: HdfsBridge::createRecordWriter - Unexpected error encountered when creating a record writer: IllegalArgumentException: field ended by ';': expected ';' but got 'CODE' at line 9: optional binary ATC CODE [ErrorCode = 110802] [SQLState = S0001]
问题原因
Synapse表中存在列名ATC CODE(包含空格),导致连接器生成的JDBC查询在DMS解析时出错。项目要求不能替换列名中的空格,需寻找其他解决方案。
可行解决方案
方案1:在Synapse中创建带引号列名的视图
在Azure Synapse中执行SQL创建视图,用方括号包裹带空格的列名,确保连接器能正确识别:
CREATE VIEW vw_target_table AS SELECT [ATC CODE] AS [ATC CODE], -- 列出其他所有需要的列 col1, col2 FROM target_table
之后在Databricks的query参数中查询该视图即可:
query = "SELECT * FROM vw_target_table"
方案2:修改查询语句,给带空格列名加方括号
直接在query参数的SQL语句中,给带空格的列名加上方括号,避免解析错误:
query = """ SELECT [ATC CODE], col1, col2 FROM target_table """
这样连接器生成的JDBC查询会正确处理带空格的列名,绕过DMS的解析问题。
方案3:临时重命名列缓存,再恢复原列名
读取时临时将带空格的列名替换为无空格名称,缓存完成后再改回原列名,既满足缓存需求又保留原列名:
# 读取时临时重命名带空格的列 df_temp = spark.read \ .format("com.databricks.spark.sqldw") \ .option("url", sql_dw_connection_string) \ .option("tempDir", temp_dir_url) \ .option("enableServicePrincipalAuth", "true") \ .option("query", "SELECT [ATC CODE] AS ATC_CODE, col1, col2 FROM target_table") \ .load() # 缓存临时DataFrame df_temp.cache() # 恢复原列名 df = df_temp.withColumnRenamed("ATC_CODE", "ATC CODE")
内容的提问来源于stack exchange,提问作者Sarath Subramanian
相关产品推荐
相关产品推荐

