如何避免pyspark.sql.SparkSession.sql为查询添加引号?
解决方案
首先要明确:Spark SQL的spark.sql()方法中,传入的参数仅用于值替换(比如WHERE子句里的条件值),不能用来替换表名、列名这类SQL标识符——Spark会自动对这些参数加引号转义,这就是你遇到语法错误的原因。另外你的代码里直接传递DataFrame给my_table参数也是错误的,SQL里的表名需要是注册过的临时视图,不是DataFrame对象。
下面提供几种可行的解决办法:
方法1:使用DataFrame API(推荐)
完全避免手写SQL字符串,用Spark的DataFrame API实现逻辑,可读性和安全性都更高:
from pyspark.sql.functions import col value_column = "Speed" my_table = self.spark.read.format("delta").load(file_path) result = my_table.select( col("ReadingTime").alias("time"), col(value_column).alias("value") ).distinct()
方法2:Python字符串格式化(适用于必须写SQL的场景)
如果一定要用SQL语句,可在Python层面用字符串格式化处理标识符(列名、表名),同时确保参数是可信的(避免SQL注入风险)。注意先将DataFrame注册为临时视图:
value_column = "Speed" my_table = self.spark.read.format("delta").load(file_path) # 注册临时视图 my_table.createOrReplaceTempView("temp_my_table") # 使用f-string格式化标识符 result = self.spark.sql( f""" SELECT DISTINCT temp_my_table.ReadingTime as time, temp_my_table.{value_column} as value FROM temp_my_table """ )
如果需要同时传递值参数(比如过滤条件),可以混合使用:
speed_threshold = 100 result = self.spark.sql( f""" SELECT DISTINCT temp_my_table.ReadingTime as time, temp_my_table.{value_column} as value FROM temp_my_table WHERE temp_my_table.{value_column} > :threshold """, threshold=speed_threshold )
方法3:使用Spark的标识符参数化(Spark 3.2+)
Spark 3.2及以上版本支持用$语法引用标识符,但需要开启相关配置:
# 开启标识符参数化支持 self.spark.conf.set("spark.sql.legacy.parser.quotedRegexColumnNames", "true") value_column = "Speed" my_table.createOrReplaceTempView("temp_my_table") result = self.spark.sql( """ SELECT DISTINCT temp_my_table.ReadingTime as time, temp_my_table.$value_column as value FROM temp_my_table """, value_column=value_column )
不过这种方式兼容性不如前两种,建议优先用前两种方法。
内容的提问来源于stack exchange,提问作者R D
相关产品推荐
相关产品推荐

