使用PySpark SQLContext.sql()时提示缺少sqlQuery参数的问题求助
问题分析与解决方案
嘿,我来帮你排查这个报错的问题,其实原因主要有两个,咱们一步步来解决:
1. 核心报错原因:SQLContext.sql() 的用法错误
你代码里直接调用 SQLContext.sql(query) 是不对的——sql 并不是 SQLContext 类的静态方法,它需要通过实例化的SQLContext对象或者更推荐的SparkSession对象来调用。
从你的代码来看,你已经在使用 spark.read.jdbc(),说明你已经有了一个 SparkSession 实例(就是变量 spark),直接用它来执行SQL查询就可以了,完全不需要单独调用 SQLContext。
2. 隐藏的SQL逻辑错误:字段名的引号使用错误
你的查询语句里,把带空格的字段名 record date 用单引号包裹了:
SELECT * FROM example WHERE 'record date' BETWEEN '2020-09-05' AND '2020-09-18'
这里的单引号会让Spark把 'record date' 当成一个字符串常量,而不是字段名,即使语法不报错,也查不到正确的数据。正确的做法是用**反引号(`)**包裹带空格的字段名。
修正后的完整代码
url = 'jdbc:mysql://localhost:3306/test_db' df = spark.read.jdbc(url=url, table='test_data', properties=prop) df.createTempView('example') # 修正字段名的引号,用反引号包裹带空格的字段 query = """SELECT * FROM example WHERE `record date` BETWEEN '2020-09-05' AND '2020-09-18'""" # 用SparkSession实例spark来调用sql方法,而不是SQLContext类 df2 = spark.sql(query) df2.collect()
额外说明
如果你的PySpark版本比较旧(1.x),确实需要用SQLContext的话,那你需要先获取实例:
from pyspark.sql import SQLContext sqlContext = SQLContext(spark.sparkContext) df2 = sqlContext.sql(query)
但在PySpark 2.x及之后的版本,SparkSession已经整合了SQLContext的功能,直接用spark.sql()是更简洁规范的写法。
内容的提问来源于stack exchange,提问作者Edward Chang
相关产品推荐
相关产品推荐

