SpiralTrain
Exercises › Block 2 · Exercise 2

Block 2 · Exercise 2

RDDs

Open the notebook 14-rdd-starter in your workspace and run the setup cell, which starts the Spark session. Every step with code to write has a cell with a comment where your code goes. Steps 1 to 3 use two tiny lists. From Step 4 on, the same operations run on the flights.

Do one part at a time, 1a, 1b and so on, and compare with the expected output before you go on. The solution is in 14-rdd-solution, under every part on the exercise site, and at the back of the exercises PDF.

Step 1Set difference

The starter has two lists, list_A = [1, 2, 3, 4, 5] and list_B = [4, 5, 6, 7, 8], and an RDD of each, rdd_A and rdd_B.

1a

Find the elements of list_A that are not in list_B, with one transformation on the RDDs. Print the result as a sorted list.

Hint 1

The slide Set Transformation Functions lists the transformations that work on two RDDs.

Check your output
[1, 2, 3]
Show solutionHide solution
python
print(sorted(rdd_A.subtract(rdd_B).collect()))

1b

Is the transformation you used narrow or wide? Print toDebugString() of the result and look for a shuffle in it.

Hint 1

The slides Narrow Transformations and Wider Transformations say how to tell them apart. In the lineage, a shuffle shows up as a new indented block.

Hint 2

toDebugString() returns bytes. Call .decode() on it before you print it.

Check your output
(<n>) <RDD>[...] at ... []
 |  <RDD>[...] at ... []
 +-(<n>) <RDD>[...] at ... []
    |  <RDD>[...] at ... []
    ...

One RDD per line, the newest at the top. The RDD names, numbers and addresses vary. A line that starts with +- marks a shuffle boundary.

Show solutionHide solution
python
print(rdd_A.subtract(rdd_B).toDebugString().decode())

subtract is the set difference. It is a wide transformation, so it shuffles.

Step 2Symmetric difference

2a

Find everything that is in one list and not in both : [1, 2, 3, 6, 7, 8]. Use two subtractions and a union, and print the result as a sorted list.

Hint 1

Subtract in both directions, then put the two results together.

Check your output
[1, 2, 3, 6, 7, 8]
Show solutionHide solution
python
both_ways = rdd_A.subtract(rdd_B).union(rdd_B.subtract(rdd_A))

print(sorted(both_ways.collect()))

2b

Print getNumPartitions() of both_ways, of rdd_A.subtract(rdd_B) and of rdd_B.subtract(rdd_A). How do the three compare?

Hint 1

union is narrow. The slide Narrow Transformations says what that means for the data.

Check your output
partitions : <a whole number>
A - B      : <a whole number>
B - A      : <a whole number>

The first number is the sum of the other two. The numbers themselves depend on your session.

Show solutionHide solution
python
print("partitions :", both_ways.getNumPartitions())
print("A - B      :", rdd_A.subtract(rdd_B).getNumPartitions())
print("B - A      :", rdd_B.subtract(rdd_A).getNumPartitions())

union moves no data, so it adds the partition counts of its inputs together.

Step 3Partitions

3a

Print the partition count of rdd_A. Nothing you wrote asked for a number, so Spark chose one. Print sc.defaultParallelism next to it.

Hint 1

An RDD has a method that returns its partition count.

Check your output
partitions chosen : <a whole number>
cores for tasks   : <a whole number>

The two numbers are equal, and depend on your session.

Show solutionHide solution
python
print("partitions chosen :", rdd_A.getNumPartitions())
print("cores for tasks   :", sc.defaultParallelism)

Spark chose the first count from the cores available.

3b

Make an RDD of the same five numbers with exactly three partitions, and print the contents of each partition.

Hint 1

sc.parallelize takes the partition count as a second argument. glom() turns each partition into a list.

Hint 2

glom() is safe here because there are five numbers. Never call it on the flights.

Check your output
with three        : [[1], [2, 3], [4, 5]]

Three lists.

Show solutionHide solution
python
print("with three        :", sc.parallelize(list_A, 3).glom().collect())

Step 4Delays per carrier

A flight counts as delayed when its dep_delay is above 15 minutes. The starter builds delayed_pairs : the flights that are delayed, as (carrier, 1) pairs. It skips cancelled flights, because a cancelled flight has no departure delay.

4a

Reduce delayed_pairs to one count per carrier and print the counts, busiest carrier first. Which carrier has the most delayed flights?

Hint 1

The slide ReduceByKey shows how to add up the values that share a key.

Check your output
<carrier>  <count>
<carrier>  <count>
...

One line per carrier, the counts falling down the column. Which carriers and what counts depend on the dataset size.

Show solutionHide solution
python
per_carrier = delayed_pairs.reduceByKey(lambda a, b: a + b)

for carrier, n in sorted(per_carrier.collect(), key=lambda p: -p[1]):
    print(f"{carrier}  {n:>8}")

Filtering before the pair keeps rows out of the shuffle.

4b

Print toDebugString() of per_carrier. How many shuffles are in the chain?

Hint 1

Look for the same kind of indented block as in 1b.

Show solutionHide solution
python
print(per_carrier.toDebugString().decode())

One shuffle, from reduceByKey.

Step 5Sorting

5a

Sort the counts with the busiest carrier first, as an RDD, and print the top five. sortByKey sorts on the key, so the pairs need a different shape first.

Hint 1

The slide SortByKey sorts on the first element of the pair. Put the count there.

Hint 2

sortByKey(ascending=False) sorts from large to small.

Check your output
<carrier>  <count>
...

Five lines, in the same order as in 4a.

Show solutionHide solution
python
ranked = (per_carrier
    .map(lambda pair: (pair[1], pair[0]))
    .sortByKey(ascending=False))

for n, carrier in ranked.take(5):
    print(f"{carrier}  {n:>8}")

The swap is the RDD idiom for ORDER BY on the value.

5b

Look at the Jobs tab. How many jobs did the cell of 5a produce?

Show solutionHide solution

More than one. sortByKey samples the keys in a job of its own before it sorts, to work out how to split the range fairly, and that shows in the Spark UI as an extra job.

Step 6Delay share

A carrier with more flights has more delays, so the count from Step 4 mostly ranks carriers by size. Work out the share of delayed flights instead, as a percentage, for the carriers with at least ten thousand flights.

6a

Make a second pair RDD, totals, with the total number of flights per carrier.

Hint 1

It is Step 4 without the filter : every flight becomes a (carrier, 1) pair.

Show solutionHide solution
python
totals = flights.rdd.map(lambda row: (row["carrier"], 1)).reduceByKey(lambda a, b: a + b)

6b

Join totals with per_carrier, keep the carriers with at least ten thousand flights, and turn the pair of counts into a percentage. Print the shares, highest first.

Hint 1

After a join, each value is a pair of the two values : (flights, delayed) here. mapValues works on the value only.

Hint 2

Filter on the total before you divide.

Check your output
<carrier>   <share>%
<carrier>   <share>%
...

One line per carrier with at least ten thousand flights, the share falling down the column.

Show solutionHide solution
python
shares = (totals.join(per_carrier)
    .filter(lambda pair: pair[1][0] >= 10000)
    .mapValues(lambda counts: 100 * counts[1] / counts[0]))

rdd_shares = sorted(shares.collect(), key=lambda p: -p[1])
for carrier, share in rdd_shares:
    print(f"{carrier}  {share:6.2f}%")

6c

Does the order change against Step 4? Which carrier comes out worst on the share, and where was it on the count?

Show solutionHide solution

The count ranks the big carriers first, and the share ranks the unpunctual ones first. The carrier at the top of the share list is often not the one at the top of the count list.

Step 7DataFrame check

The starter has the DataFrame version of Step 6 : groupBy, agg, and a share column, in by_df.

7a

Run the cell and compare it with your answer to 6b. Do the two answers agree? Check it in code.

Hint 1

Collect by_df into a dictionary of carrier to share, and compare it with the shares from the RDD version.

Hint 2

Round both to six decimals before you compare, so a floating-point difference does not count.

Check your output
+-------+-------+-------+-----+
|carrier|flights|delayed|share|
+-------+-------+-------+-----+
|    ...|    ...|    ...|  ...|
...
+-------+-------+-------+-----+

rdd and dataframe agree : True
Show solutionHide solution
python
by_df.show()

df_shares = {r["carrier"]: r["share"] for r in by_df.collect()}
agree = all(round(share, 6) == round(df_shares[carrier], 6) for carrier, share in rdd_shares)
print("rdd and dataframe agree :", agree)

7b

Which version is shorter, and which one did Catalyst get to optimise? Open the SQL tab of the Spark UI. Does the RDD version appear there at all?

Hint 1

The slide RDDs Under DataFrames says what Catalyst can and cannot see.

Show solutionHide solution

The DataFrame version has the same numbers in about a third of the code, and Catalyst gets to optimise it. The RDD version does not show up in the SQL tab, because there is no query to plan. An RDD runs any Python you give it, and Catalyst can optimise none of it.

Step 8Broadcast and accumulate

8a

Build a dictionary city_by_code from the airports table : airport code to city. Broadcast it as cities, and make an accumulator missing that starts at 0.

Hint 1

airports.select("iata_code", "municipality").collect() gives a list of rows. A dictionary comprehension turns it into the dictionary.

Hint 2

The slides Broadcast Variables and Accumulators show sc.broadcast and sc.accumulator.

Show solutionHide solution
python
city_by_code = {r["iata_code"]: r["municipality"] for r in
                airports.select("iata_code", "municipality").collect()}
cities = sc.broadcast(city_by_code)
missing = sc.accumulator(0)

8b

Name the origin of every flight with the broadcast dictionary, and count the flights per city. Add 1 to the accumulator for every origin that is not in the dictionary, in the same pass. Print the five busiest cities, then the accumulator.

Hint 1

A function passed to map reads cities.value and can call missing.add(1).

Hint 2

Read the accumulator only after the action has run. An accumulator can be added to in a task and read in the driver, never the other way round.

Check your output
<city>                   <count>
<city>                   <count>
... five lines
origins missing from the airports table : <a whole number>

The counts fall down the column.

Show solutionHide solution
python
def name_origin(row):
    city = cities.value.get(row["origin"])
    if city is None:
        missing.add(1)
        city = "unknown"
    return (city, 1)

by_city = flights.rdd.map(name_origin).reduceByKey(lambda a, b: a + b)
busiest = by_city.takeOrdered(5, key=lambda pair: -pair[1])

for city, n in busiest:
    print(f"{city:<24} {n:>8}")
print("origins missing from the airports table :", missing.value)

cities.unpersist()

The broadcast replaces a join with a dictionary lookup, and the accumulator counts the lookups that failed without a second pass over the data.

8c

How many airports are in the broadcast, how many flight origins are missing from it, and which city has the most departures?

Hint 1

len(city_by_code) counts the airports.

Check your output
airports in the broadcast : <a whole number>
origins missing           : <a whole number>
busiest city              : <city>
Show solutionHide solution
python
print("airports in the broadcast :", len(city_by_code))
print("origins missing           :", missing.value)
print("busiest city              :", busiest[0][0])

Step 9Stop session

9a

Click the session indicator and stop the session.

Show solutionHide solution

Nothing to run. A session holds its executors until the idle timeout whether you are running anything or not, and with the whole class on one capacity a forgotten session is why somebody else cannot start theirs.

If time permits

  • Cache per_carrier and time its count() twice. Does the cache pay for itself, and does the Storage tab show what you expect? Steps 5 and 6 both use it.
  • Redo Step 4 with groupByKey().mapValues(len) instead of reduceByKey. The answer is the same. Compare the two toDebugString() outputs and the shuffle write in the Stages tab.
  • Count the delays per route, keyed on (origin, dest) rather than on the carrier. How many routes are there, and what does that do to the shuffle?

Tried it yourself first?

The solution is a spoiler. Work through the hints first : a wrong attempt teaches more than a solution you only read.