Lesson 9 / 25

Spark SQL and Temporary Views

Mix SQL and the DataFrame API freely using views.

SQL and DataFrames are the same engine

df.createOrReplaceTempView("orders") registers a DataFrame as a temporary view you can query with spark.sql("SELECT ..."). The result is another DataFrame, so you can switch between SQL and the DataFrame API wherever each is clearer. Both are compiled by the same Catalyst optimizer into the same kind of plan, so there is normally no performance difference. SQL is convenient for analysts and for complex aggregations; the DataFrame API is easier to compose in code, test and refactor.

Read, query, write

Spark reads and writes many formats; Parquet with partitioning is the usual choice for analytics.

Three steps: read, query, write.
Figure 3.1 — Read, query and write.

SQL on a DataFrame, run

I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). Revenue per city, highest first. Delhi 430.0 comes first; Mumbai has a null revenue (its only amount is null) and sorts last in a descending order here.

orders.createOrReplaceTempView("orders")
spark.sql("SELECT city, SUM(amount) AS revenue FROM orders GROUP BY city ORDER BY revenue DESC").show()

Output:

+------+-------+
|  city|revenue|
+------+-------+
| delhi|  430.0|
|  pune|  320.0|
|mumbai|   NULL|
+------+-------+

Name views clearly

Temporary views live only for the current session. For tables shared across jobs, use a catalog table (for example Hive metastore, Unity Catalog or an Iceberg/Delta catalog).

Quick check: What does `spark.sql("SELECT ...")` return?

  • A pandas file
  • A string
  • A DataFrame
  • Nothing
Answer

A DataFrame — SQL results are DataFrames and can continue into further transformations.