Module Module 11. Distributed Data and Computing Frameworks
Objective Use Dask DataFrame to run a partitioned groupby-mean on a synthetic dataset, compare it with Pandas, and explain lazy evaluation, scheduling overhead, partitioning, and framework abstractions.
- Python 3 installed.
- Pandas and NumPy installed.
- Dask DataFrame and Distributed installed.
- Git basics:
clone,add,commit,push. - Concepts from Module 11:
- distributed data-processing challenges;
- high-level distributed frameworks;
- Dask DataFrame basics;
- partitions and workers;
- lazy evaluation.
Install the Python dependencies with:
python -m pip install pandas numpy "dask[dataframe]" distributedMany workloads group records by an identifier and compute aggregates such as mean, sum, or count. Dask DataFrame offers a Pandas-like API while representing work as a task graph over partitions that can be scheduled across workers.
This lab deliberately begins with a Pandas DataFrame and converts it with dd.from_pandas because the goal is to compare APIs and execution models on one machine. For genuinely large or distributed datasets, creating one large Pandas object first is usually the wrong ingestion pattern. Production Dask workloads commonly read partitioned data directly with Dask, for example from Parquet or CSV.
Dask is also not expected to beat Pandas for every in-memory workload. Scheduler, serialization, communication, and process overhead can make Dask slower on small or simple datasets. That is a valid result to analyse.
README.mdthis filelab5_dask_aggregation.pystarter with deterministic data generation, Pandas baseline, and a Dask TODOanalysis.mdwhere you record environment information, timings, and answers
General instructions
- Clone your GitHub Classroom repository.
- Install the required libraries.
- Edit
lab5_dask_aggregation.pyto complete the Dask TODO. - Run the Pandas and Dask versions and verify that their results agree.
- Record observations in
analysis.md. - Commit frequently and push before the deadline.
Read generate_sample_dataframe(num_rows) and run_sequential_aggregation(df).
The starter:
- uses a fixed NumPy random seed so runs are reproducible;
- stores IDs as
int32to reduce memory use; - reports the approximate Pandas DataFrame memory footprint;
- performs
groupby("id")["value"].mean()as the Pandas baseline.
The same Pandas DataFrame is reused for the Dask experiment because neither aggregation mutates it. Do not add unnecessary .copy() calls.
Complete run_parallel_dask_aggregation(df, npartitions).
Your implementation must:
- convert the Pandas DataFrame to a Dask DataFrame with the requested number of partitions;
- create the same groupby-mean operation as the Pandas baseline;
- keep that aggregation lazy until
.compute()is called; - call
.compute()and return the resulting Pandas Series.
The function's timer intentionally includes both construction of the Dask collection/task graph from the existing Pandas DataFrame and the actual computation.
The starter keeps these concepts separate:
NUM_WORKERScontrols the number of worker processes in the local Dask cluster;NUM_PARTITIONScontrols how many DataFrame partitions Dask schedules across those workers.
There may be more partitions than workers. Workers execute tasks; partitions describe pieces of the data.
The defaults are conservative so the lab runs on typical student laptops. If your machine has sufficient memory, you may increase NUM_ROWS after obtaining one successful run.
Run:
python lab5_dask_aggregation.pyThe script reports:
- Python, Pandas, and Dask versions;
- row count;
- worker and partition counts;
- DataFrame memory usage;
- Pandas aggregation time;
- Dask cluster startup time;
- Dask aggregation time;
- correctness verification.
The script must finish with:
Verification: Pandas and Dask results match.
Do not interpret timing results if verification fails.
Answer the questions in analysis.md about:
- whether Dask was faster or slower and why;
- workers versus partitions;
- lazy evaluation and the role of
.compute(); - why framework overhead matters;
- why starting from a large Pandas DataFrame is not a scalable ingestion pattern;
- which concerns Dask handles compared with lower-level multiprocessing or MPI;
- why correctness verification is required before performance conclusions.
- Ensure
lab5_dask_aggregation.pyruns successfully and verification passes. - Ensure
analysis.mdcontains your recorded environment, timings, and answers. - Stage:
git add lab5_dask_aggregation.py analysis.md(orgit add .) - Commit:
git commit -m "Complete Lab 5 Dask Aggregation" - Push:
git push origin main(or your default branch) - Verify on GitHub that
lab5_dask_aggregation.pyandanalysis.mdare updated.