What you will be able to do
- Filter rows with filter or where, using a Column condition or a SQL string, and combine conditions with & and |
- Split a string column into an array with split, and predict the effect of the limit argument
- Turn array elements into rows with explode, keep empty and NULL arrays with explode_outer, and work around the one-explode-per-SELECT rule
1.Keeping only the rows you want: filter and where
Column methods change a DataFrame's width. Row operations change its length, and the simplest one keeps only the rows that satisfy a condition. filter(condition) and where(condition) are interchangeable, and both return a new DataFrame. The condition can take two forms: a Column of BooleanType, such as df.age > 3 or col("c_custkey") == 412449, or a string of SQL, such as "age > 3". Both forms give the same rows.
To combine conditions on Columns, use & for AND and | for OR, and put each comparison in its own parentheses, as the reference does. The Databricks guide shows both, for example (col("c_nationkey") == 20) & (col("c_acctbal") > 1000).
df.filter((df.age > 3) & (df.subject == "Physics")).show()
# +---+----+-------+
# |age|name|subject|
# +---+----+-------+
# | 5| Bob|Physics|
# +---+----+-------+Checkpoint 1 of 5· Check yourself
Which argument is NOT a valid condition for DataFrame.filter?
filter accepts a Boolean Column or a SQL expression string. It does not select rows by position.
“A Column of BooleanType or a string of SQL expressions.”Source: docs.databricks.com
Checkpoint 2 of 5· Exam question
DataFrame `customers` has columns `customer_id`, `email`, `phone`, and `signup_date`. A pipeline runs: ```python result = customers.drop("email", "phone") print(result.columns) ``` What is printed?
Correct answer: D — ['customer_id', 'signup_date']
- A. This list still includes email, but drop was called with email as one of the two column names to remove, so email should not remain in the result.
- B. This list contains only the two dropped names and omits the columns that were kept, which is the opposite of what drop returns.
- C. This is the full original schema; drop was called with two column names as arguments, so both of those columns must be removed from the returned DataFrame.
- D. drop("email", "phone") removes both named columns and returns a new DataFrame retaining every other column in its original order, leaving customer_id and signup_date.
- E. This list still includes phone, but drop was called with phone as one of the two column names to remove, so phone should not remain in the result.
2.Splitting a string column into an array with split
Splitting changes a column's shape rather than the number of rows. split(str, pattern, limit) lives in pyspark.sql.functions rather than on the DataFrame, and it returns a Column whose value is an array of the separated strings. You usually wrap it in select or withColumn. The pattern is a Java regular expression, not a literal delimiter. In the example below, '[ABC]' splits on any of the three capital letters. This matters for characters such as . or |, which have special meanings in a regex.
from pyspark.sql import functions as dbf
df = spark.createDataFrame([('oneAtwoBthreeC',)], ['s',])
df.select('*', dbf.split(df.s, '[ABC]')).show()
df.select('*', dbf.split(df.s, '[ABC]', 2)).show()
df.select('*', dbf.split('s', '[ABC]', -2)).show()limit sets how many times the pattern is applied. If limit > 0, the array has at most limit entries, and the last entry holds everything after the last match. If limit <= 0, the pattern is applied as many times as possible and the array can be any size. Recent versions also accept a column or column name for pattern and limit, so each row can use its own values.
Two entries at most, with the last one holding the remainder: 'one' and 'twoBthreeC'. The pattern is applied only once, so the B and C stay inside the second entry.
Checkpoint 3 of 5· Check yourself
What does a limit of -1 passed to split mean?
Any limit of zero or less means no cap: the pattern is applied as many times as possible.
“pattern will be applied as many times as possible, and the resulting array can be of any size.”Source: docs.databricks.com
Sources3
3.Turning array elements into rows: explode and explode_outer
An array column, whether split produced it or it came from the source data, can be turned into rows. explode(col) returns one new row for each element of an array, or for each entry of a map. The new column is called col for arrays, or key and value for maps, unless you set your own names with alias. The input in the reference example has three rows: i=1 with [1, 2, 3, NULL], i=2 with an empty array, and i=3 with a NULL array.
df.select('*', sf.explode('a')).show()
+---+---------------+----+
| i| a| col|
+---+---------------+----+
| 1|[1, 2, 3, NULL]| 1|
| 1|[1, 2, 3, NULL]| 2|
| 1|[1, 2, 3, NULL]| 3|
| 1|[1, 2, 3, NULL]|NULL|
+---+---------------+----+Losing those rows can be a quiet bug. Customers with no orders, for example, disappear from the result. explode_outer behaves the same way, except that an empty or NULL array produces a single row with NULL in the element column, so rows 2 and 3 survive. posexplode_outer does the same and adds a pos column with each element's position.
| Input array | explode | explode_outer | posexplode_outer |
|---|---|---|---|
| [1, 2, 3, NULL] | 4 rows in col | 4 rows in col | 4 rows, pos 0–3 and col |
| [] (empty) | no row | 1 row, col = NULL | 1 row, pos and col NULL |
| NULL | no row | 1 row, col = NULL | 1 row, pos and col NULL |
Checkpoint 4 of 5· Fill the gap
Rows with an empty or NULL array must stay in the result, with NULL in the element column. Which function completes the call?
df.select('*', sf. ? ('a')).show()Unlike explode, explode_outer produces a NULL element for an empty or NULL array, so the row is kept.
Source: docs.databricks.comOne more rule matters: you can use only one explode per SELECT clause. To explode two arrays, chain two select calls and alias each result.
import pyspark.sql.functions as sf
df = spark.sql('SELECT ARRAY(1,2) AS a1, ARRAY(3,4,5) AS a2')
df.select(
'*', sf.explode('a1').alias('v1')
).select('*', sf.explode('a2').alias('v2')).show()Six. Each of the 2 rows from the first explode is exploded again over the 3 elements of a2, so every (v1, v2) pair appears once.
If the array holds structs, explode it, alias the result, and then use select("s.*") to spread the struct's fields into ordinary columns. The reference example df.select(sf.explode('a').alias("s")).select("s.*") turns an array of {a, b} structs into a two-column table with one row per struct.
Checkpoint 5 of 5· Exam question
DataFrame `df` has a column named `cust_nm`. A developer intends to rename `cust_nm` to `customer_name` and runs: ```python df2 = df.withColumnRenamed("customer_name", "cust_nm") ``` What is the result?
Correct answer: B — `df2` ends up identical to `df` because `withColumnRenamed` found no column named `customer_name` to rename.
- A. Spark's rename method matches the first argument against actual column names, not positions, so this description of how the match happens is incorrect.
- B. withColumnRenamed looks for a column matching its first argument; since no column is named customer_name, it silently returns an unchanged DataFrame rather than touching cust_nm.
- C. This method does not validate that the first argument matches a real column, so it never throws an exception when the name is absent; it simply performs a no-op.
- D. withColumnRenamed only relabels an existing column in place; it never fabricates a new column of nulls, so this outcome does not occur.
- E. There is no separate driver-only schema cache that can diverge from the actual data; the returned DataFrame's schema and data are always consistent with each other.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.explode keeps every input row, filling in NULL when the array is empty or NULL.Why is that wrong?
explode produces no row for an empty or NULL array, so those input rows disappear. explode_outer is the function that keeps them with a NULL element.
Covered in Turning array elements into rows: explode and explode_outer
2.Two array columns can be exploded side by side in the same select.Why is that wrong?
Only one explode is allowed per SELECT clause. Chain a second select to explode the other array.
Covered in Turning array elements into rows: explode and explode_outer
3.split's pattern is a literal delimiter string.Why is that wrong?
The pattern is a Java regular expression, so characters with a regex meaning are interpreted as regex, not as literal text.
Covered in Splitting a string column into an array with split
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“A Column of BooleanType or a string of SQL expressions.”
↩︎ Keeping only the rows you want: filter and where - 2.https://docs.databricks.com/aws/en/pyspark/basicsOfficial docs
“To filter rows, use the filter or where method on a DataFrame to return only certain rows.”
↩︎ Keeping only the rows you want: filter and where“For example, & and | enable you to AND and OR conditions, respectively.”
↩︎ Keeping only the rows you want: filter and where - 3.
“a string representing a regular expression. The regex string should be a Java regular expression.”
↩︎ Splitting a string column into an array with split“the resulting array's last entry will contain all input beyond the last matched pattern”
↩︎ Splitting a string column into an array with split“pyspark.sql.Column: array of separated strings.”
↩︎ Splitting a string column into an array with split“a string representing a regular expression. The regex string should be a Java regular expression.”
↩︎ Exam trap 3“pattern will be applied as many times as possible, and the resulting array can be of any size.”
↩︎ Checkpoint - 4.
“Uses the default column name col for elements in the array and key and value for elements in the map unless specified otherwise.”
↩︎ Turning array elements into rows: explode and explode_outer“Only one explode is allowed per SELECT clause.”
↩︎ Turning array elements into rows: explode and explode_outer“Only one explode is allowed per SELECT clause.”
↩︎ Exam trap 2“Returns a new row for each element in the given array or map.”
↩︎ Prediction - 5.
“Unlike explode, if the array/map is null or empty then null is produced.”
↩︎ Turning array elements into rows: explode and explode_outer“Unlike explode, if the array/map is null or empty then null is produced.”
↩︎ Exam trap 1