Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. map vs flatMap

PySpark · Execution Model

map vs flatMap

Easypyspark-45
rddmapflatmapexplode

Question

What is the difference between map and flatMap?

Solution

map turns each input element into exactly one output element. flatMap turns each input element into zero or more output elements, and then flattens them all into one collection.

Word count shows the difference

lines = sc.parallelize(["spark is fast", "spark is lazy"])

lines.map(lambda l: l.split(" ")).collect()
# [['spark', 'is', 'fast'], ['spark', 'is', 'lazy']]   -> 2 elements, each a list

lines.flatMap(lambda l: l.split(" ")).collect()
# ['spark', 'is', 'fast', 'spark', 'is', 'lazy']        -> 6 elements

map keeps the shape: two lines in, two lists out. flatMap takes those lists apart, so you get a flat sequence of words, ready for map(lambda w: (w, 1)) and reduceByKey.

Because flatMap can return an empty list, it can also act as a filter: return [] for rows you want to drop and [x] for rows you keep.

The DataFrame way

You rarely write flatMap on DataFrames. The same job uses split and explode:

from pyspark.sql import functions as F

words = (lines_df
         .select(F.explode(F.split("text", " ")).alias("word"))
         .groupBy("word").count())

split turns the string into an array. explode makes one row per array element. This stays inside Spark's optimizer and avoids running Python for every row, which is why it is much faster than an RDD flatMap with a Python lambda.

Which to use

New code should use DataFrames and explode. Knowing map versus flatMap is still asked because it checks you understand one-to-one versus one-to-many transformations. A likely follow-up is "what is the DataFrame equivalent" and the answer is explode(split(...)).

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext