Dask
01 / 02

DataFrame, Array & Lazy Evaluation

Dask: DataFrame, Array & Lazy Evaluation

Dask scales existing Python data-analysis code (pandas, NumPy, scikit-learn) beyond a single machine's memory/CPU, mirroring familiar APIs while adding parallel/distributed execution underneath -- targeting the gap between "fits in memory, plain pandas is fine" and "needs a full big-data system like Spark".

Dask DataFrame: Built from pandas Partitions

import dask.dataframe as dd

# Many separate files map naturally onto Dask's partitioned model --
# each file becomes a partition, no need to concatenate into one
# giant in-memory pandas DataFrame first
ddf = dd.read_csv('data/*.csv')

# pandas-like API -- but LAZY, builds a task graph instead of
# executing immediately
result = ddf.groupby('category')['amount'].mean()

# .compute() triggers actual execution, returns a concrete,
# in-memory pandas result -- must fit in memory, unlike ddf itself
final = result.compute()

# Escape hatch: run arbitrary pandas code per-partition when the
# built-in API doesn't cover what you need
ddf = ddf.map_partitions(lambda df: df.assign(total=df.a + df.b))

Dask Array: Built from NumPy Chunks

import dask.array as da
import numpy as np

# chunks divides the array into NumPy pieces for parallel processing
arr = da.from_array(large_numpy_array, chunks=(1000, 1000))

# Too-large chunks: memory pressure per worker, limited parallelism
# Too-small chunks: excessive per-task scheduling overhead
# -- a real, workload-dependent tuning decision, same trade-off as
# DataFrame partition sizing

result = (arr + 1).sum(axis=0).compute()

Dask Bag: Unstructured Data

import dask.bag as db
import json

# For data that doesn't fit DataFrame/Array's rectangular shape --
# e.g. a large collection of JSON records or log lines
bag = (
    db.read_text('logs/*.json')
    .map(json.loads)
    .filter(lambda record: record['status'] == 'error')
)
errors = bag.compute()

delayed: Parallelizing Arbitrary Code

import dask

# For custom workflows that don't map onto DataFrame/Array/Bag's
# built-in operations -- explicit task graph from ordinary function calls
@dask.delayed
def load(path): ...

@dask.delayed
def process(data): ...

results = [process(load(f)) for f in file_list]
final = dask.compute(*results)

Keep your own version of these notes — editable, searchable, and organised by your stack.

Start free