What you will be able to do
- Build a DataFrame from a Python list of tuples, with an inferred or explicit schema
- Turn a DataFrame into a list of Row objects with collect, take, head, first and tail, and read values from each Row
- Iterate rows with toLocalIterator or foreach, and explain the memory difference from collect
- Decide when toPandas is safe to use
1.From a Python list to a DataFrame
Conversion works in both directions, and the simplest starting point is a Python list. spark.createDataFrame takes the data as a list of tuples, one tuple per row, plus a schema. The schema can be just a list of column names.
df_children = spark.createDataFrame(
data = [("Mikhail", 15), ("Zaky", 13), ("Zoya", 8)],
schema = ['name', 'age'])If you pass only names, Spark infers each column's type from the Python values. To fix the types yourself, pass a StructType built from StructFields. Each field gives a name, a data type such as StringType() or IntegerType(), and a nullable flag. You import these from pyspark.sql.types.
Checkpoint 1 of 6· Check yourself
You build a DataFrame with schema = ['name', 'age'] and no types. What decides the type of the age column?
With only column names, the types are inferred automatically. A StructType is needed only when you want to set them yourself.
“Notice in the output that the data types of columns of df_children are automatically inferred.”Source: docs.databricks.com
Sources1
2.From a DataFrame to a list: collect and Row objects
Going the other way, collect() returns every record in the DataFrame as a Python list of Row objects, for example [Row(age=14, name='Tom'), ...]. You can filter or select first, and only the rows that remain are collected. Each Row lets you read a field by name with row["name"], and row.asDict() turns it into a plain dictionary. Combined with a list comprehension, that turns a DataFrame into an ordinary Python list of values or dicts.
rows = df.collect()
[row["name"] for row in rows]
# ['Tom', 'Alice', 'Bob']
[row.asDict() for row in rows]
# [{'age': 14, 'name': 'Tom'}, {'age': 23, 'name': 'Alice'}, {'age': 16, 'name': 'Bob'}]That convenience has a cost. The DataFrame lives on the cluster, but the list lives in one Python process, the driver. The reference warns that collect should only be used when the result is expected to be small, because all of the data is loaded into the driver's memory.
Checkpoint 2 of 6· Check yourself
Why does the reference say collect() should only be used on small results?
collect gathers every record into the driver as a Python list, so a large DataFrame can exceed the driver's memory.
“as all the data is loaded into the driver's memory.”Source: docs.databricks.com
Checkpoint 3 of 6· Exam question
A pipeline needs to hand a downstream JSON serializer a Python list of dictionaries, one per row, mapping each column name to its value. Given: ```python data = df.collect() records = [___ for row in data] ``` Which expression completes the list comprehension?
Correct answer: A — `row.asDict()`
- A. `Row.asDict()` converts a single `Row` into a Python dictionary keyed by column name, so mapping it across `data` yields exactly the list of per-row dictionaries the serializer needs.
- B. `toPandas()` is a `DataFrame` method that builds a pandas DataFrame from the whole dataset; an individual `Row` object has no `toPandas` method and calling it raises an `AttributeError`.
- C. A `Row` object has no `schema` attribute; schema information lives on the parent `DataFrame`, so this raises an `AttributeError` rather than producing a dictionary of values.
- D. `select` is a `DataFrame` transformation that returns a new `DataFrame`; a `Row` object has no `select` method, so this call fails with an `AttributeError`.
- E. `rdd` and `collect` are attributes of a `DataFrame`, not of an individual `Row`, so calling them on a single row raises an `AttributeError` instead of returning a dictionary.
Sources2
3.Bringing back only some rows: take, head, first, tail
Often you only need a few rows. take(num) returns the first num rows as a list of Row, or all of them if the DataFrame is smaller. tail(num) does the same from the end. The reference warns that tail moves data into the driver, and a very large num can crash the driver with OutOfMemoryError. head and first look alike but return different things.
df.head()
# Row(age=2, name='Alice')
df.head(1)
# [Row(age=2, name='Alice')]
df.head(0)
# []first() returns the first row as a single Row. If the DataFrame is empty it returns None instead of raising an error, so code that calls first() on a possibly empty result should check for None.
| Call | Returns | Note |
|---|---|---|
| collect() | list of Row (all records) | All data loaded into driver memory |
| take(num) | list of Row (first num) | All records if fewer than num exist |
| tail(num) | list of Row (last num) | Very large num can cause OutOfMemoryError on the driver |
| head() / head(n) | Row / list of Row | n defaults to 1 |
| first() | Row | None if the DataFrame is empty |
Checkpoint 4 of 6· Check yourself
A filter leaves the DataFrame empty. What does df.first() return?
first() returns the first row if there is one and None otherwise. An empty list is what head(0) or take on an empty DataFrame would give.
“First row if DataFrame is not empty, otherwise None.”Source: docs.databricks.com
Sources3
4.Iterating rows: toLocalIterator and foreach
To loop over every row without building the whole list at once, use toLocalIterator(). It returns an iterator over all the rows. Its memory footprint is the key difference from collect: the iterator uses about as much memory as the largest partition, not the whole DataFrame. With prefetchPartitions=True, Spark fetches the next partition before it is needed, which can use the memory of the two largest partitions.
Checkpoint 5 of 6· Fill the gap
Which method returns an iterator whose memory use is bounded by the largest partition?
list(df. ? ())toLocalIterator returns an iterator over the rows and holds roughly one partition at a time. collect returns the full list at once. Wrapping the iterator in list() gives the same rows.
Source: docs.databricks.comSometimes you want to run a function on each row rather than loop yourself. foreach(f) applies f to every Row of the DataFrame. f takes one parameter, the row being processed, and returns nothing. foreachPartition(f) is the per-partition version: f receives an iterator of the rows in one partition.
def func(person):
print(person.name)
df.foreach(func)toLocalIterator. It returns an iterator of rows that your code consumes. foreach takes a function that returns None (its signature is Callable[[Row], None]) and gives no rows back.
5.Converting to pandas with toPandas
toPandas() returns the contents of the DataFrame as a pandas.DataFrame, with a pandas index (0, 1, ...) next to the Spark columns. It has the same limitation as collect: use it only when the result is expected to be small, because all the data is loaded into the driver's memory. It also needs pandas to be installed. The reference says that using it with spark.sql.execution.arrow.pyspark.enabled=True is experimental.
Checkpoint 6 of 6· Match them up
Match each call to the shape of what it gives you
Tap a term, then the definition that fits it.
collect and toPandas both load everything into the driver, as a list or a pandas DataFrame. toLocalIterator streams the rows one partition at a time. foreach only runs a function on each row.
“This method should only be used if the resulting Pandas pandas.DataFrame is expected to be small”Source: docs.databricks.com
Sources7
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.toLocalIterator() is just collect() under another name, so it needs enough driver memory for the whole DataFrame.Why is that wrong?
The iterator uses about as much memory as the largest partition, or the two largest with prefetching. It does not hold the whole DataFrame at once.
Covered in Iterating rows: toLocalIterator and foreach
2.df.head() and df.head(1) return the same thing.Why is that wrong?
Without n, head returns a single Row. With n, it returns a list of Row, even when n is 1.
Covered in Bringing back only some rows: take, head, first, tail
3.tail(num) is always safe because it only returns the last few rows.Why is that wrong?
tail moves data into the driver process, and a very large num can crash the driver with OutOfMemoryError.
Covered in Bringing back only some rows: take, head, first, tail
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://docs.databricks.com/aws/en/pyspark/basicsOfficial docs
“Schemas are defined using the StructType which is made up of StructFields”
↩︎ From a Python list to a DataFrame“Notice in the output that the data types of columns of df_children are automatically inferred.”
↩︎ Checkpoint - 2.
“Returns all the records in the DataFrame as a list of Row.”
↩︎ From a DataFrame to a list: collect and Row objects“as all the data is loaded into the driver's memory.”
↩︎ Checkpoint - 3.
“Will return this number of records or all records if the DataFrame contains less than this number of records.”
↩︎ Bringing back only some rows: take, head, first, tail - 4.
“With prefetch it may consume up to the memory of the 2 largest partitions.”
↩︎ Iterating rows: toLocalIterator and foreach“The iterator will consume as much memory as the largest partition in this DataFrame.”
↩︎ Exam trap 1 - 5.
“A function that accepts one parameter which will receive each row to process.”
↩︎ Iterating rows: toLocalIterator and foreach - 6.https://docs.databricks.com/aws/en/pyspark/reference/classes/dataframe/foreachPartitionOfficial docs
“A function that accepts one parameter which will receive each partition to process.”
↩︎ Iterating rows: toLocalIterator and foreach - 7.
“Usage with spark.sql.execution.arrow.pyspark.enabled=True is experimental.”
↩︎ Converting to pandas with toPandas“This method should only be used if the resulting Pandas pandas.DataFrame is expected to be small”
↩︎ Checkpoint
Also cited
“If n is missing, return a single Row.”
↩︎ Exam trap 2“If n is missing, return a single Row.”
↩︎ Prediction“can crash the driver process with OutOfMemoryError”
↩︎ Exam trap 3“First row if DataFrame is not empty, otherwise None.”
↩︎ Checkpoint