在 PySpark 环境中执行基础 SQL 查询的步骤是什么?
考察说明
考查候选人是否掌握 PySpark SQL 的基本使用流程和常用 API。
回答思路
- 【回答框架 1】PySpark SQL 允许通过 SQL 语句或 DataFrame API 操作结构化数据。核心是将数据加载为 DataFrame 并注册为临时视图或表,之后使用 spark.sql 方法执行 SQL 查询。
- 【回答框架 2】具体步骤:首先创建 SparkSession,它是 PySpark 的入口。然后使用 spark.read 读取数据(如 parquet、csv、json),返回 DataFrame。接着调用 createOrReplaceTempView 注册临时视图,命名任意表名。最后通过 spark.sql 编写 SELECT 等查询语句。
- 【回答框架 3】例如:spark.sql('SELECT * FROM my_table').show() 即可输出查询结果。查询返回新的 DataFrame,可继续使用 DataFrame API 处理或调用 collect 获取结果。
- 【回答框架 4】注意事项:临时视图生命周期为当前 SparkSession,级别为会话级;使用 createGlobalTempView 可跨会话,但需要以 global_temp 作为数据库前缀访问。
- 【关键点 1】使用 spark.sql 执行 SQL 前需将 DataFrame 注册为临时视图或表。
- 【关键点 2】createOrReplaceTempView 是会话级别的临时表,适合单任务内查询。
- 【关键点 3】SQL 查询结果返回 DataFrame,可继续链式操作或收集数据。
- 【易错点 1】混淆 createTempView 与 createGlobalTempView 的作用域,导致跨会话访问失败。
- 【易错点 2】未将 DataFrame 注册为视图就调用 spark.sql 会报 Table or view not found 错误。
- 【易错点 3】认为 SQL 查询返回的是 RDD 而非 DataFrame,后续误用 RDD 操作。