SQL → PySpark
A practical tutorial and reference guide for converting SQL queries into PySpark DataFrame operations. Learn the syntax, understand the concepts, and copy production-ready examples.
Learn SQL → PySpark
If you already know SQL, PySpark becomes much easier when you understand how SQL operations map to DataFrame APIs.
Understand the DataFrame
A PySpark DataFrame is similar to a SQL table. Columns, rows, filters, joins and aggregations can all be expressed using DataFrame operations.
Learn Column Expressions
PySpark uses column expressions such as
col(), when(),
lit(), and functions from
pyspark.sql.functions.
Combine Operations
PySpark transformations can be chained together, allowing you to build complex data pipelines in a readable way.
50 SQL → PySpark Mappings
Search for any SQL operation or PySpark function.
SELECT
Select specific columns from a DataFrame.
SELECT name, salary FROM employees;
df.select("name", "salary")
WHERE
Filter rows using a condition.
SELECT * FROM employees WHERE salary > 50000;
df.filter(df.salary > 50000) # or df.where(df.salary > 50000)
DISTINCT
Remove duplicate rows.
SELECT DISTINCT department FROM employees;
df.select("department").distinct()
ORDER BY
Sort rows by one or more columns.
SELECT * FROM employees ORDER BY salary DESC;
from pyspark.sql.functions import desc
df.orderBy(desc("salary"))
GROUP BY
Group records and perform aggregate calculations.
SELECT department, COUNT(*) FROM employees GROUP BY department;
from pyspark.sql.functions import count
df.groupBy("department") \
.agg(count("*").alias("count"))
HAVING
Filter aggregated results after GROUP BY.
SELECT department, COUNT(*) FROM employees GROUP BY department HAVING COUNT(*) > 10;
from pyspark.sql.functions import count
df.groupBy("department") \
.agg(count("*").alias("cnt")) \
.filter("cnt > 10")
COUNT()
Count rows or non-null values.
SELECT COUNT(*) FROM employees;
from pyspark.sql.functions import count
df.select(count("*"))
SUM()
Calculate the sum of a numeric column.
SELECT SUM(salary) FROM employees;
from pyspark.sql.functions import sum
df.select(sum("salary"))
AVG()
Calculate the average of a numeric column.
SELECT AVG(salary) FROM employees;
from pyspark.sql.functions import avg
df.select(avg("salary"))
MIN()
Find the minimum value.
SELECT MIN(salary) FROM employees;
from pyspark.sql.functions import min
df.select(min("salary"))
MAX()
Find the maximum value.
SELECT MAX(salary) FROM employees;
from pyspark.sql.functions import max
df.select(max("salary"))
INNER JOIN
Return only matching rows from both DataFrames.
SELECT * FROM employees e INNER JOIN departments d ON e.dept_id = d.id;
employees.join(
departments,
employees.dept_id == departments.id,
"inner"
)
LEFT JOIN
Keep every row from the left DataFrame and matching rows from the right.
SELECT * FROM employees e LEFT JOIN departments d ON e.dept_id = d.id;
employees.join(
departments,
employees.dept_id == departments.id,
"left"
)
RIGHT JOIN
Keep every row from the right DataFrame.
SELECT * FROM employees e RIGHT JOIN departments d ON e.dept_id = d.id;
employees.join(
departments,
employees.dept_id == departments.id,
"right"
)
FULL JOIN
Keep matching and non-matching rows from both DataFrames.
SELECT * FROM employees e FULL OUTER JOIN departments d ON e.dept_id = d.id;
employees.join(
departments,
employees.dept_id == departments.id,
"full"
)
CROSS JOIN
Produce the Cartesian product of two DataFrames.
SELECT * FROM employees CROSS JOIN departments;
employees.crossJoin(departments)
UNION ALL
Combine two DataFrames without removing duplicates.
SELECT * FROM employees_2025 UNION ALL SELECT * FROM employees_2026;
df_2025.union(df_2026)
UNION
Combine DataFrames and remove duplicate rows.
SELECT * FROM a UNION SELECT * FROM b;
a.union(b).distinct()
CASE WHEN
Create conditional column logic.
SELECT
name,
CASE
WHEN salary >= 100000 THEN 'High'
ELSE 'Low'
END AS level
FROM employees;
from pyspark.sql.functions import when
df.withColumn(
"level",
when(df.salary >= 100000, "High")
.otherwise("Low")
)
COALESCE
Return the first non-null expression.
SELECT COALESCE(phone, 'N/A') FROM customers;
from pyspark.sql.functions import coalesce, lit
df.select(
coalesce("phone", lit("N/A"))
)
IS NULL
SELECT * FROM customers WHERE phone IS NULL;
df.filter(df.phone.isNull())
IS NOT NULL
SELECT * FROM customers WHERE phone IS NOT NULL;
df.filter(df.phone.isNotNull())
IN
SELECT *
FROM employees
WHERE department IN ('IT', 'HR');
df.filter(
df.department.isin("IT", "HR")
)
BETWEEN
SELECT * FROM employees WHERE salary BETWEEN 50000 AND 100000;
df.filter(
df.salary.between(50000, 100000)
)
LIKE
SELECT * FROM customers WHERE name LIKE 'A%';
df.filter(
df.name.like("A%")
)
ILIKE
Case-insensitive pattern matching where supported (Spark 3.3+).
SELECT * FROM customers WHERE name ILIKE 'john%';
df.filter(
df.name.ilike("john%")
)
CONCAT
SELECT CONCAT(first_name, ' ', last_name) FROM customers;
from pyspark.sql.functions import concat, lit
df.select(
concat("first_name", lit(" "), "last_name")
)
SUBSTRING
SELECT SUBSTRING(name, 1, 3) FROM customers;
from pyspark.sql.functions import substring
df.select(
substring("name", 1, 3)
)
LENGTH
SELECT LENGTH(name) FROM customers;
from pyspark.sql.functions import length
df.select(length("name"))
UPPER
SELECT UPPER(name) FROM customers;
from pyspark.sql.functions import upper
df.select(upper("name"))
LOWER
SELECT LOWER(name) FROM customers;
from pyspark.sql.functions import lower
df.select(lower("name"))
TRIM
SELECT TRIM(name) FROM customers;
from pyspark.sql.functions import trim
df.select(trim("name"))
REPLACE
SELECT REPLACE(name, '-', ' ') FROM customers;
from pyspark.sql.functions import regexp_replace
df.select(
regexp_replace("name", "-", " ")
)
SPLIT
SELECT SPLIT(tags, ',') FROM customers;
from pyspark.sql.functions import split
df.select(
split("tags", ",")
)
ROW_NUMBER()
Assign sequential numbers to rows within a window.
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY department
ORDER BY salary DESC
) AS rn
FROM employees;
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number
w = Window.partitionBy("department") \
.orderBy(df.salary.desc())
df.withColumn("rn", row_number().over(w))
RANK()
RANK() OVER (
PARTITION BY department
ORDER BY salary DESC
)
from pyspark.sql.functions import rank
df.withColumn(
"rank",
rank().over(w)
)
DENSE_RANK()
DENSE_RANK() OVER (
ORDER BY salary DESC
)
from pyspark.sql.functions import dense_rank
df.withColumn(
"dense_rank",
dense_rank().over(w)
)
LAG()
Access a value from a previous row.
LAG(salary) OVER (
ORDER BY employee_id
)
from pyspark.sql.functions import lag
df.withColumn(
"previous_salary",
lag("salary").over(w)
)
LEAD()
Access a value from a following row.
LEAD(salary) OVER (
ORDER BY employee_id
)
from pyspark.sql.functions import lead
df.withColumn(
"next_salary",
lead("salary").over(w)
)
SUM() OVER()
Calculate a windowed or running sum.
SUM(salary) OVER (
PARTITION BY department
)
from pyspark.sql.functions import sum
df.withColumn(
"department_total",
sum("salary").over(w)
)
AVG() OVER()
AVG(salary) OVER (
PARTITION BY department
)
from pyspark.sql.functions import avg
df.withColumn(
"department_avg",
avg("salary").over(w)
)
COUNT() OVER()
COUNT(*) OVER (
PARTITION BY department
)
from pyspark.sql.functions import count
df.withColumn(
"department_count",
count("*").over(w)
)
CURRENT_DATE
SELECT CURRENT_DATE;
from pyspark.sql.functions import current_date df.select(current_date())
CURRENT_TIMESTAMP
SELECT CURRENT_TIMESTAMP;
from pyspark.sql.functions import current_timestamp df.select(current_timestamp())
CAST
SELECT CAST(salary AS INT) FROM employees;
df.select(
df.salary.cast("int")
)
DATE_ADD
SELECT DATE_ADD(order_date, 7) FROM orders;
from pyspark.sql.functions import date_add
df.select(
date_add("order_date", 7)
)
DATEDIFF
SELECT DATEDIFF(end_date, start_date) FROM projects;
from pyspark.sql.functions import datediff
df.select(
datediff("end_date", "start_date")
)
DROP DUPLICATES
Remove duplicate rows from a DataFrame.
SELECT DISTINCT * FROM customers;
df.dropDuplicates() # Specific columns: df.dropDuplicates(["email"])
WITH / CTE
Use an intermediate DataFrame or temporary SQL view to represent a CTE.
WITH high_salary AS (
SELECT *
FROM employees
WHERE salary > 100000
)
SELECT *
FROM high_salary;
high_salary = df.filter(
df.salary > 100000
)
high_salary.show()
TEMP VIEW / CTE
Register a DataFrame as a temporary SQL view when you want to continue writing SQL inside Spark.
WITH employees_cte AS (
SELECT *
FROM employees
)
SELECT *
FROM employees_cte;
df.createOrReplaceTempView(
"employees_view"
)
spark.sql("""
SELECT *
FROM employees_view
""")
Practical PySpark Tutorial
The cheat sheet is useful for quick reference. This section explains how to think about converting a complete SQL query into PySpark.
1. Create a PySpark DataFrame
In SQL, you normally work with tables. In PySpark, the equivalent working object is a DataFrame.
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("SQLToPySpark") \
.getOrCreate()
df = spark.read.csv(
"employees.csv",
header=True,
inferSchema=True
)
df.show()
2. SELECT columns
SQL uses SELECT to choose columns. PySpark uses
select().
SELECT name, department, salary FROM employees;
df.select(
"name",
"department",
"salary"
)
3. WHERE filtering
SQL’s WHERE clause becomes filter()
or where() in PySpark.
SELECT * FROM employees WHERE salary > 70000 AND department = 'IT';
df.filter(
(df.salary > 70000) &
(df.department == "IT")
)
& for AND and | for OR.
Put each condition inside parentheses.
4. GROUP BY and aggregation
SQL aggregation becomes groupBy()
followed by agg().
SELECT
department,
COUNT(*) AS employee_count,
AVG(salary) AS avg_salary
FROM employees
GROUP BY department;
from pyspark.sql.functions import count, avg
result = (
df
.groupBy("department")
.agg(
count("*").alias("employee_count"),
avg("salary").alias("avg_salary")
)
)
5. JOIN DataFrames
A SQL JOIN becomes the DataFrame
join() method.
SELECT
e.name,
d.department_name
FROM employees e
JOIN departments d
ON e.dept_id = d.id;
result = employees.join(
departments,
employees.dept_id == departments.id,
"inner"
).select(
employees.name,
departments.department_name
)
6. CASE WHEN
Conditional SQL logic can be implemented with
when() and otherwise().
SELECT
name,
salary,
CASE
WHEN salary >= 100000 THEN 'Senior'
WHEN salary >= 70000 THEN 'Mid'
ELSE 'Junior'
END AS level
FROM employees;
from pyspark.sql.functions import when
result = df.withColumn(
"level",
when(df.salary >= 100000, "Senior")
.when(df.salary >= 70000, "Mid")
.otherwise("Junior")
)
7. Window Functions
Window functions are extremely important in
real-world data engineering. PySpark uses the
Window class to define the window.
SELECT
name,
department,
salary,
ROW_NUMBER() OVER (
PARTITION BY department
ORDER BY salary DESC
) AS rn
FROM employees;
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number
window_spec = (
Window
.partitionBy("department")
.orderBy(df.salary.desc())
)
result = df.withColumn(
"rn",
row_number().over(window_spec)
)
8. Complete SQL → PySpark Example
Let’s convert a complete query containing filtering, grouping, aggregation, ordering and a calculated column.
SELECT
department,
COUNT(*) AS employee_count,
AVG(salary) AS average_salary,
CASE
WHEN AVG(salary) > 80000
THEN 'High'
ELSE 'Normal'
END AS salary_level
FROM employees
WHERE status = 'ACTIVE'
GROUP BY department
HAVING COUNT(*) > 5
ORDER BY average_salary DESC;
from pyspark.sql.functions import (
count,
avg,
when,
col,
desc
)
result = (
df
.filter(col("status") == "ACTIVE")
.groupBy("department")
.agg(
count("*").alias("employee_count"),
avg("salary").alias("average_salary")
)
.filter(col("employee_count") > 5)
.withColumn(
"salary_level",
when(
col("average_salary") > 80000,
"High"
).otherwise("Normal")
)
.orderBy(desc("average_salary"))
)
result.show()
SQL:
SELECT → WHERE → GROUP BY → HAVING → ORDER BY
PySpark:
select() → filter() → groupBy().agg()
→ filter() → orderBy()
Quick Reference
Keep this table open while writing PySpark transformations.
| # | SQL | PySpark | Purpose |
|---|---|---|---|
| 01 | SELECT |
select() |
Select columns |
| 02 | WHERE |
filter() |
Filter rows |
| 03 | DISTINCT |
distinct() |
Remove duplicates |
| 04 | ORDER BY |
orderBy() |
Sort rows |
| 05 | GROUP BY |
groupBy() |
Group rows |
| 06 | HAVING |
filter() |
Filter aggregates |
| 07 | COUNT() |
count() |
Count |
| 08 | SUM() |
sum() |
Sum |
| 09 | AVG() |
avg() |
Average |
| 10 | MIN() |
min() |
Minimum |
| 11 | MAX() |
max() |
Maximum |
| 12 | INNER JOIN |
join(..., "inner") |
Matching rows |
| 13 | LEFT JOIN |
join(..., "left") |
Keep left rows |
| 14 | RIGHT JOIN |
join(..., "right") |
Keep right rows |
| 15 | FULL JOIN |
join(..., "full") |
Keep both sides |
| 16 | CROSS JOIN |
crossJoin() |
Cartesian product |
| 17 | UNION ALL |
union() |
Combine rows |
| 18 | UNION |
union().distinct() |
Combine unique rows |
| 19 | CASE WHEN |
when().otherwise() |
Conditional logic |
| 20 | COALESCE |
coalesce() |
Handle nulls |
| 21 | IS NULL |
isNull() |
Find nulls |
| 22 | IS NOT NULL |
isNotNull() |
Find non-null |
| 23 | IN |
isin() |
Match list |
| 24 | BETWEEN |
between() |
Range filtering |
| 25 | LIKE |
like() |
Pattern match |
| 26 | ILIKE |
ilike() |
Case-insensitive match |
| 27 | CONCAT |
concat() |
Join strings |
| 28 | SUBSTRING |
substring() |
Extract text |
| 29 | LENGTH |
length() |
String length |
| 30 | UPPER |
upper() |
Uppercase |
| 31 | LOWER |
lower() |
Lowercase |
| 32 | TRIM |
trim() |
Remove spaces |
| 33 | REPLACE |
regexp_replace() |
Replace text |
| 34 | SPLIT |
split() |
Split strings |
| 35 | ROW_NUMBER() |
row_number() |
Sequential ranking |
| 36 | RANK() |
rank() |
Ranking |
| 37 | DENSE_RANK() |
dense_rank() |
Dense ranking |
| 38 | LAG() |
lag() |
Previous row |
| 39 | LEAD() |
lead() |
Next row |
| 40 | SUM() OVER() |
sum().over() |
Window sum |
| 41 | AVG() OVER() |
avg().over() |
Window average |
| 42 | COUNT() OVER() |
count().over() |
Window count |
| 43 | CURRENT_DATE |
current_date() |
Current date |
| 44 | CURRENT_TIMESTAMP |
current_timestamp() |
Current timestamp |
| 45 | CAST |
cast() |
Change datatype |
| 46 | DATE_ADD |
date_add() |
Add days |
| 47 | DATEDIFF |
datediff() |
Date difference |
| 48 | DISTINCT |
dropDuplicates() |
Remove duplicates |
| 49 | WITH / CTE |
Intermediate DataFrame | Reusable transformation |
| 50 | WITH / CTE |
createOrReplaceTempView() |
SQL temp view |
🚀 SQL Developer → PySpark Developer
The biggest shift when moving from SQL to PySpark is understanding that you are no longer primarily writing one SQL statement. You are building a sequence of DataFrame transformations.
- Think of a DataFrame as a table.
-
Think of
select()as SQL SELECT. -
Think of
filter()as SQL WHERE. -
Think of
groupBy().agg()as GROUP BY. -
Think of
join()as SQL JOIN. -
Use
Windowfor ranking and row-to-row calculations. - Chain transformations to build readable pipelines.


