SpiralTrain
Exercises › Block 1 · Appendix A4

Block 1 · Appendix A4

Concurrency

Starter notebook
07-concurrency-starter
Fabric path
/lakehouse/default/Files/data/solutions/07.concurrency/

Open the notebook 07-concurrency-starter in your workspace. It has one cell per step, with the imports and a comment where your code goes. The notebook runs on Fabric's driver, a Linux machine, and every thread and process runs there.

Do one part at a time, 1a, 1b and so on, and compare with the expected output before you go on. Timings in a notebook are noisy : run a cell twice and use the second number, and compare the shape of the results, not the exact seconds. The solution is in 07-concurrency-solution, under every part on the exercise site, and at the back of the exercises PDF.

Step 1Threads

1a

Write a function worker(name) that prints name starting, sleeps for a second and prints name done. Start five threads that run it, as worker-0 to worker-4. Do not join them. Run the cell.

Hint 1

The slide Simple Thread shows how to create and start a Thread with target and args. args takes a tuple.

Check your output
worker-0 starting
worker-1 starting
worker-2 starting
worker-3 starting
worker-4 starting
worker-0 done
worker-1 done
worker-2 done
worker-3 done
worker-4 done

The five starting lines come together, and the done lines follow a second later. The order within each group can differ, and a notebook can show the done lines under the output of the next cell you run.

Show solutionHide solution
python
def worker(name: str) -> None:
    print(f"{name} starting")
    time.sleep(1)
    print(f"{name} done")


threads = [threading.Thread(target=worker, args=(f"worker-{i}",)) for i in range(5)]
for t in threads:
    t.start()

1b

Add print("all done") after the loop that starts the threads, and run the cell again. Look at where the line appears.

Check your output
worker-0 starting
worker-1 starting
worker-2 starting
worker-3 starting
worker-4 starting
all done
worker-0 done
worker-1 done
worker-2 done
worker-3 done
worker-4 done

all done comes before the done lines, and can come before some of the starting lines. The lines can interleave.

Show solutionHide solution
python
for t in threads:
    t.start()
print("all done")

1c

Make the main thread wait for the workers : join every thread before the print. Measure the time from before the threads start to after the print with time.perf_counter, and print it in the message.

Hint 1

The slide Joining Threads shows join(). It has to be called on every thread, so it needs a loop of its own after the start loop.

Check your output
worker-0 starting
...
worker-4 done
all done after 1.00 s

all done comes last, and the time is about one second, not five.

Show solutionHide solution
python
start = time.perf_counter()
threads = [threading.Thread(target=worker, args=(f"worker-{i}",)) for i in range(5)]
for t in threads:
    t.start()
for t in threads:
    t.join()
print(f"all done after {time.perf_counter() - start:.2f} s")

Step 2Race and lock

2a

Start from a new cell. Write a shared counter that starts at 0, and a function increment(times) that adds one to it times times in three separate steps : read counter into a local value, call time.sleep(0), then store value + 1 in counter. The function changes a module variable, so it needs global counter.

Hint 1

time.sleep(0) hands the processor to the other thread between the read and the write. Without it the race is still there, but Python switches threads so rarely inside a tight loop that it seldom shows.

Show solutionHide solution
python
counter = 0


def increment(times: int) -> None:
    global counter
    for _ in range(times):
        value = counter
        time.sleep(0)
        counter = value + 1

2b

Reset counter to 0, start two threads that each run increment(10_000), join them and print counter. Do this three times in a loop.

Hint 1

The total should be 20 000. Reset counter to 0 at the top of every round, or the totals build on each other.

Check your output
10000
10002
10001

Your totals differ from run to run. All of them are far below 20 000.

Show solutionHide solution
python
for _run in range(3):
    counter = 0
    threads = [threading.Thread(target=increment, args=(10_000,)) for _ in range(2)]
    for t in threads:
        t.start()
    for t in threads:
        t.join()
    print(counter)

2c

Write increment_plain(times) with counter += 1 and no sleep, and run it in two threads three times the way you did in 2b. Write a comment that says whether the program is now correct or only lucky, and why.

Hint 1

counter += 1 is still a read, an addition and a write.

Check your output
20000
20000
20000
Show solutionHide solution
python
def increment_plain(times: int) -> None:
    global counter
    for _ in range(times):
        counter += 1


for _run in range(3):
    counter = 0
    threads = [threading.Thread(target=increment_plain, args=(10_000,)) for _ in range(2)]
    for t in threads:
        t.start()
    for t in threads:
        t.join()
    print(counter)

# Only lucky : counter += 1 is still three steps, and a thread switch between two of them
# loses an update. Python switches so rarely inside this loop that the total usually comes out right.

2d

Write increment_locked(times) : the version of increment with the read, the sleep(0) and the write inside with lock:, where lock is a threading.Lock(). Run increment and increment_locked in two threads, three times each. For every run print the name of the function, the total and the time it took.

Hint 1

The slide Locks as Context Managers shows with lock:. Everything that has to happen without the other thread interfering goes inside the block.

Hint 2

target.__name__ gives the name of a function.

Check your output
increment         total  10000  0.83 s
increment         total  10001  0.83 s
increment         total  10000  0.81 s
increment_locked  total  20000  2.17 s
increment_locked  total  20000  2.10 s
increment_locked  total  20000  2.07 s

The times depend on the machine, and the first three totals differ from run to run. The locked version is exact and takes longer.

Show solutionHide solution
python
lock = threading.Lock()


def increment_locked(times: int) -> None:
    global counter
    for _ in range(times):
        with lock:
            value = counter
            time.sleep(0)
            counter = value + 1


for target in (increment, increment_locked):
    for _run in range(3):
        counter = 0
        start = time.perf_counter()
        threads = [threading.Thread(target=target, args=(10_000,)) for _ in range(2)]
        for t in threads:
            t.start()
        for t in threads:
            t.join()
        print(f"{target.__name__:17} total {counter:6}  {time.perf_counter() - start:.2f} s")

Step 3CPU-bound work

The starter has count_primes(limit), which counts the primes below limit and does no sleeping and no I/O.

3a

Time four calls of count_primes(200_000) in a row, in one thread. Use time.perf_counter around the four calls, and print the time.

Check your output
one thread   : 0.87 s

The time depends on the machine. The first run of a cell can be slower than the next ones.

Show solutionHide solution
python
start = time.perf_counter()
for _ in range(4):
    count_primes(200_000)
print(f"one thread   : {time.perf_counter() - start:.2f} s")

3b

Time the same four calls spread over four threads, and print the time.

Hint 1

The slide Simple Thread shows how to start a thread. Start all four, then join all four, and stop the clock after the joins.

Check your output
one thread   : 0.87 s
four threads : 0.90 s

The four threads are no faster than the one thread, and can be slower.

Show solutionHide solution
python
start = time.perf_counter()
threads = [threading.Thread(target=count_primes, args=(200_000,)) for _ in range(4)]
for t in threads:
    t.start()
for t in threads:
    t.join()
print(f"four threads : {time.perf_counter() - start:.2f} s")

3c

Write a comment that says why four threads did not help.

Hint 1

Read the slide Global Interpreter Lock.

Show solutionHide solution
python
# The GIL lets only one thread run Python bytecode at a time, so four threads share one core's
# worth of interpreter. Nothing was waiting, so there was never a moment to hand over.

Step 4Process pool

The starter has count_primes again, with os and Pool imported. Keep print out of any function a pool runs : return values, and print them in the cell.

4a

Time four calls of count_primes(200_000) in one process. Then time the same four calls with Pool(4).map. Print both times.

Hint 1

The slide Pool.map and Pool.apply_async shows pool.map with a function and a list of arguments.

Hint 2

[200_000] * 4 is a list of four arguments.

Check your output
one process 0.975 s, pool of four 0.305 s

The times depend on the machine.

Show solutionHide solution
python
start = time.perf_counter()
for _ in range(4):
    count_primes(200_000)
alone = time.perf_counter() - start

start = time.perf_counter()
with Pool(4) as pool:
    results = pool.map(count_primes, [200_000] * 4)
pooled = time.perf_counter() - start
print(f"one process {alone:.3f} s, pool of four {pooled:.3f} s")

4b

Print os.cpu_count() and the speed-up of the pool, the time of one process divided by the time of the pool. Is the speed-up close to the number of cores, or to four?

Hint 1

The pool has four workers, however many cores the machine has.

Check your output
cores on this machine : 16
speed-up : 3.2

Both numbers depend on the machine.

Show solutionHide solution
python
print("cores on this machine :", os.cpu_count())
print(f"speed-up : {alone / pooled:.1f}")
# The pool has four workers, so the speed-up stays below four, however many cores there are.
# Starting the workers and pickling the arguments take part of it away.

4c

Shrink the work to count_primes(1_000). Time one process and the pool again, for both sizes in one loop, and print a line for each.

Check your output
limit  200000 : one process 0.827 s, pool of four 0.344 s
limit    1000 : one process 0.001 s, pool of four 0.027 s

The times depend on the machine. Look at which of the two is faster for the small job.

Show solutionHide solution
python
for limit in (200_000, 1_000):
    start = time.perf_counter()
    for _ in range(4):
        count_primes(limit)
    alone = time.perf_counter() - start

    start = time.perf_counter()
    with Pool(4) as pool:
        results = pool.map(count_primes, [limit] * 4)
    pooled = time.perf_counter() - start
    print(f"limit {limit:>7} : one process {alone:.3f} s, pool of four {pooled:.3f} s")

4d

Call count_primes for 50 000, 100 000 and 150 000 with pool.apply_async, one call per number. Print the results collected with get().

Hint 1

apply_async returns a handle straight away and takes the arguments as a tuple. get() on the handle waits for the result.

Check your output
apply_async : [5133, 9592, 13848]
Show solutionHide solution
python
with Pool(4) as pool:
    handles = [pool.apply_async(count_primes, (limit,)) for limit in (50_000, 100_000, 150_000)]
    print("apply_async :", [h.get() for h in handles])

Step 5Coroutines

A notebook already runs an event loop, so asyncio.run() refuses to start a second one. Write await directly in a cell.

5a

Write a coroutine fetch(name, seconds) that prints name started, waits with asyncio.sleep(seconds) and returns f"{name} took {seconds}s". Await three of them one after another, ("a", 1), ("b", 2) and ("c", 1). Print each result and the total time.

Hint 1

The slide Simple Coroutine shows async def and await asyncio.sleep. The slide Awaiting Coroutines shows a call that waits for the result.

Check your output
a started
a took 1s
b started
b took 2s
c started
c took 1s
one after another : 4.00 s
Show solutionHide solution
python
async def fetch(name: str, seconds: float) -> str:
    print(f"{name} started")
    await asyncio.sleep(seconds)
    return f"{name} took {seconds}s"


start = time.perf_counter()
for name, seconds in (("a", 1), ("b", 2), ("c", 1)):
    print(await fetch(name, seconds))
print(f"one after another : {time.perf_counter() - start:.2f} s")

5b

Run the same three calls together with asyncio.gather, print the list of results and the total time.

Hint 1

The slide Running Coroutines with gather passes the coroutines as separate arguments, and gather returns a list.

Check your output
a started
b started
c started
['a took 1s', 'b took 2s', 'c took 1s']
gather            : 2.00 s
Show solutionHide solution
python
start = time.perf_counter()
print(await asyncio.gather(fetch("a", 1), fetch("b", 2), fetch("c", 1)))
print(f"gather            : {time.perf_counter() - start:.2f} s")

5c

Change fetch so that it raises ValueError(f"{name} failed") when seconds is negative, after the wait. Gather fetch("a", 1) and fetch("bad", -1) with return_exceptions=True, and print the result.

Hint 1

asyncio.sleep with a negative number returns at once, so the bad call fails before the good one finishes.

Check your output
a started
bad started
gather keeps the others : ['a took 1s', ValueError('bad failed')]
Show solutionHide solution
python
async def fetch(name: str, seconds: float) -> str:
    print(f"{name} started")
    await asyncio.sleep(seconds)
    if seconds < 0:
        raise ValueError(f"{name} failed")
    return f"{name} took {seconds}s"


results = await asyncio.gather(fetch("a", 1), fetch("bad", -1), return_exceptions=True)
print("gather keeps the others :", results)

5d

Run the same two calls in an asyncio.TaskGroup, and catch the error with except* ValueError. Print the exceptions of the group.

Hint 1

The slide Task Groups shows async with asyncio.TaskGroup() as tg and tg.create_task. A TaskGroup raises an exception group, which except* takes apart.

Check your output
a started
bad started
TaskGroup cancels the rest and raises : (ValueError('bad failed'),)

The good call a is cancelled, and its result never arrives.

Show solutionHide solution
python
try:
    async with asyncio.TaskGroup() as tg:
        tg.create_task(fetch("a", 1))
        tg.create_task(fetch("bad", -1))
except* ValueError as group:
    print("TaskGroup cancels the rest and raises :", group.exceptions)

5e

Write a comment that says which of the two you would use for three calls to the same API, and why.

Show solutionHide solution
python
# TaskGroup : when one call fails, the other two results are usually of no use, and the group stops them at once.
# gather with return_exceptions=True suits calls that do not depend on each other and where a partial result is worth having.

Step 6Blocking call

6a

Write a coroutine polite(name) that prints its name and threading.current_thread().name, and then waits with asyncio.sleep(1). Gather three of them and print the total time.

Check your output
a on MainThread
b on MainThread
c on MainThread
asyncio.sleep : 1.00 s

The thread has the same name on every line. The name can differ on Fabric.

Show solutionHide solution
python
async def polite(name: str) -> None:
    print(f"{name} on {threading.current_thread().name}")
    await asyncio.sleep(1)


start = time.perf_counter()
await asyncio.gather(polite("a"), polite("b"), polite("c"))
print(f"asyncio.sleep : {time.perf_counter() - start:.2f} s")

6b

Write blocking(name), the same coroutine with time.sleep(1) in place of asyncio.sleep(1). Gather three of them and print the total time.

Check your output
a on MainThread
b on MainThread
c on MainThread
time.sleep    : 3.00 s
Show solutionHide solution
python
async def blocking(name: str) -> None:
    print(f"{name} on {threading.current_thread().name}")
    time.sleep(1)


start = time.perf_counter()
await asyncio.gather(blocking("a"), blocking("b"), blocking("c"))
print(f"time.sleep    : {time.perf_counter() - start:.2f} s")

6c

Write a comment that says why the total changed.

Hint 1

The slide The Event Loop Cycle shows what a coroutine gives back to the loop, and when.

Show solutionHide solution
python
# time.sleep does not hand control back to the event loop, so the other coroutines cannot start.
# All three share one thread, and they run one after another. One blocking call anywhere on the
# path, or one library without an async version, and the concurrency is gone without an error.

If time permits

  • Put your CPU-bound count_primes inside a coroutine and confirm it stalls the loop as time.sleep did. Then move it off with asyncio.to_thread and then with loop.run_in_executor on a ProcessPoolExecutor, and decide which is right for CPU-bound work.
  • Replace the Pool in Step 4 with concurrent.futures.ProcessPoolExecutor and then with ThreadPoolExecutor. The same code with a one-word change switches between the two models.
  • Give the Step 2 counter a queue.Queue instead of a lock, with one producer thread and one consumer. Nothing is shared, so no lock is needed.
  • Find the if __name__ == "__main__": line in the step4_process_pool.py cell of the solution notebook. Fabric's Linux driver forks the notebook's process, so the pool needs no guard there and a function defined in a cell reaches the workers. On Windows and macOS the workers start fresh and import the main module, and the guard keeps each of them from starting a pool of its own.

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.