dask-parallel-computing
Parallel/distributed computing for larger-than-RAM data. Components: DataFrames (parallel pandas), Arrays (parallel NumPy), Bags, Futures, Schedulers. Scales laptop to HPC cluster. For single-machine speed use polars; for out-of-core without cluster use vaex.
Install / Use
npx skills add jaechang-hits/SciAgent-Skills --skill dask-parallel-computingInstalls into whichever agent you are using.
SKILL.md
Installable skill definition
Quality Score
Category
Development & EngineeringSupported Platforms
Our assessment of dask-parallel-computing
dask-parallel-computing scores 91/100 on our quality scale, 1171st of 4,619 Development & Engineering skills we index (top 26%).
Its SKILL.md is 16 KB long, well organised into 73 sections with 15 code examples: a thorough specification that gives an agent plenty to work with.
It has 367 GitHub stars, a meaningful sign that others use it.
Maintenance, license and trust
- The repository was last updated 37 days ago, so dask-parallel-computing is actively maintained.
- No license is declared. By default that means all rights are reserved: you can read it, but reusing or redistributing it is not clearly permitted. Ask the author before building on it commercially.
- Its trust signals score 88/100, with 1 caution from licensing, adoption, age or documentation. These come from repository metadata, not a code audit — read the skill file before letting an agent act on it.
Safety scan
No issues foundOur scan of the whole file found no instruction hijacking, hidden characters, credential access, data exfiltration or destructive commands.
Automated pattern scan on 2026-10-05. It catches known dangerous patterns, not every risk — read a skill before letting an agent act on it.
dask-parallel-computing compared with similar skills
All 4 of these similar skills score higher than dask-parallel-computing; compare them before choosing.
| Skill | Score | Stars | Updated | Format |
|---|---|---|---|---|
| dask-parallel-computing (this skill)by jaechang-hits | 91 | 367 | 37d ago | SKILL.md |
| Agent-Reachby Panniantong | 100 | 90.8k | 19d ago | CLAUDE.md |
| headroomby headroomlabs-ai | 100 | 74.4k | today | CLAUDE.md |
| ai-job-searchby MadsLorentzen | 100 | 45.0k | 1d ago | CLAUDE.md |
| claude-howtoby luongnv89 | 100 | 41.7k | 4d ago | CLAUDE.md |
Frequently asked questions
- How do I install dask-parallel-computing?
- Run
npx skills add jaechang-hits/SciAgent-Skills --skill dask-parallel-computing. The install tabs above show the steps for each supported agent. - Which AI agents does dask-parallel-computing work with?
- It is written for Universal, as a SKILL.md file. Other agents that read the same format can often use it too.
- Is dask-parallel-computing safe to use?
- Our scan of the whole file found no instruction hijacking, hidden characters, credential access, data exfiltration or destructive commands. It declares no license and scores 88/100 on trust signals. Skills are instructions an agent will follow, so read the file before installing it and do not approve commands you do not understand.
- Is dask-parallel-computing still maintained?
- The repository was last updated 37 days ago, so dask-parallel-computing is actively maintained.
Skill content
View source on GitHubname: dask-parallel-computing description: "Parallel/distributed computing for larger-than-RAM data. Components: DataFrames (parallel pandas), Arrays (parallel NumPy), Bags, Futures, Schedulers. Scales laptop to HPC cluster. For single-machine speed use polars; for out-of-core without cluster use vaex." license: BSD-3-Clause
Dask — Parallel & Distributed Computing
Overview
Dask is a Python library for parallel and distributed computing that scales familiar pandas/NumPy APIs to larger-than-memory datasets. It provides five main components (DataFrames, Arrays, Bags, Futures, Schedulers) and scales from single-machine multi-core to multi-node HPC clusters.
When to Use
- Processing datasets that exceed available RAM (10 GB–100 TB)
- Parallelizing pandas or NumPy operations across multiple cores
- Processing multiple files efficiently (CSV, Parquet, JSON, HDF5, Zarr)
- Building custom parallel workflows with task dependencies
- Distributing workloads across HPC clusters (SLURM, Kubernetes)
- Streaming/ETL pipelines for unstructured data (logs, JSON records)
- For in-memory single-machine speed: use polars instead
- For out-of-core single-machine analytics: use vaex instead
Prerequisites
pip install dask[complete] # All components
pip install dask[dataframe] # DataFrames only
pip install dask[distributed] # Distributed scheduler + dashboard
pip install dask-jobqueue # HPC cluster integration (SLURM, PBS)
Core API
1. DataFrames — Parallel Pandas
import dask.dataframe as dd
# Read multiple files as a single DataFrame
ddf = dd.read_csv('data/2024-*.csv')
ddf = dd.read_parquet('data/', columns=['id', 'value', 'category'])
# Operations are lazy until .compute()
filtered = ddf[ddf['value'] > 100]
result = filtered.groupby('category').agg({'value': ['mean', 'sum']}).compute()
print(result.shape) # (n_categories, 2)
# Custom operations via map_partitions (preferred over apply)
def normalize_partition(df):
df['norm_value'] = (df['value'] - df['value'].mean()) / df['value'].std()
return df
ddf = ddf.map_partitions(normalize_partition)
# Joins
ddf_merged = ddf.merge(lookup_ddf, on='category', how='left')
# Write results
ddf.to_parquet('output/', engine='pyarrow')
# Repartitioning for optimal chunk sizes
ddf = ddf.repartition(npartitions=20) # By count
ddf = ddf.repartition(partition_size='100MB') # By size
# Index management for sorted operations
ddf = ddf.set_index('timestamp', sorted=True)
# Debugging
print(f"Partitions: {ddf.npartitions}")
print(f"Dtypes: {ddf.dtypes}")
sample = ddf.get_partition(0).compute() # Inspect first partition
2. Arrays — Parallel NumPy
import dask.array as da
import numpy as np
# Create from various sources
x = da.random.random((100000, 1000), chunks=(10000, 1000))
x = da.from_array(np_array, chunks=(10000, 1000))
x = da.from_zarr('large_dataset.zarr')
# Standard operations (lazy)
y = (x - x.mean(axis=0)) / x.std(axis=0) # Normalize
z = da.dot(x.T, x) # Matrix multiply
u, s, v = da.linalg.svd(x) # SVD
# Compute and persist
result = y.mean(axis=0).compute()
print(result.shape) # (1000,)
# Custom operations with map_blocks
def custom_filter(block):
from scipy.ndimage import gaussian_filter
return gaussian_filter(block, sigma=2)
filtered = da.map_blocks(custom_filter, x, dtype=x.dtype)
# Rechunking for different access patterns
x_rechunked = x.rechunk({0: 5000, 1: 500})
# Save to disk
da.to_zarr(y, 'normalized.zarr')
3. Bags — Unstructured Data Processing
import dask.bag as db
import json
# Read unstructured data
bag = db.read_text('logs/*.json').map(json.loads)
# Functional operations
valid = bag.filter(lambda x: x['status'] == 'success')
ids = valid.pluck('user_id')
flat = bag.map(lambda x: x['tags']).flatten()
# Aggregation — use foldby instead of groupby (much faster)
counts = bag.foldby(
key='category',
binop=lambda total, x: total + x['amount'],
initial=0,
combine=lambda a, b: a + b,
combine_initial=0
).compute()
# Convert to DataFrame for structured analysis
ddf = valid.to_dataframe(meta={'user_id': 'str', 'amount': 'float64', 'category': 'str'})
4. Futures — Task-Based Parallelism
from dask.distributed import Client
client = Client() # Local cluster with all cores
print(client.dashboard_link) # http://localhost:8787
# Submit individual tasks (executes immediately, not lazy)
def process(x, param):
return x ** param
future = client.submit(process, 42, param=2)
print(future.result()) # 1764
# Map over many inputs
futures = client.map(process, range(100), param=2)
results = client.gather(futures)
print(len(results)) # 100
# Scatter large data to workers (avoids repeated transfers)
import numpy as np
big_data = np.random.random((10000, 1000))
data_future = client.scatter(big_data, broadcast=True)
# Submit tasks using scattered data
futures = [client.submit(process_chunk, data_future, i) for i in range(10)]
results = client.gather(futures)
# Progressive result processing
from dask.distributed import as_completed
for future in as_completed(futures):
result = future.result()
print(f"Completed: {result}")
# Coordination primitives
from dask.distributed import Lock, Queue, Event
lock = Lock('resource-lock')
with lock:
# Thread-safe operation across workers
pass
client.close()
5. Schedulers & Configuration
import dask
# Global scheduler setting
dask.config.set(scheduler='threads') # Default: GIL-releasing numeric work
dask.config.set(scheduler='processes') # Pure Python, GIL-bound work
dask.config.set(scheduler='synchronous') # Debugging with pdb
# Context manager for temporary change
with dask.config.set(scheduler='synchronous'):
result = computation.compute() # Can use pdb here
# Per-compute override
result = ddf.mean().compute(scheduler='processes')
# Distributed scheduler with resource control
from dask.distributed import Client
client = Client(n_workers=4, threads_per_worker=2, memory_limit='4GB')
print(client.dashboard_link)
# HPC cluster integration
from dask_jobqueue import SLURMCluster
from dask.distributed import Client
cluster = SLURMCluster(
cores=24, memory='100GB',
walltime='02:00:00', queue='regular'
)
cluster.scale(jobs=10) # Request 10 SLURM jobs
client = Client(cluster)
# Adaptive scaling
cluster.adapt(minimum=2, maximum=20)
result = computation.compute()
client.close()
Key Concepts
Component Selection Guide
| Data Type | Component | When to Use | |-----------|-----------|-------------| | Tabular (CSV, Parquet) | DataFrames | Standard pandas-like operations at scale | | Numeric arrays (HDF5, Zarr) | Arrays | NumPy operations, linear algebra, image processing | | Text, JSON, logs | Bags | ETL/cleaning → convert to DataFrame for analysis | | Custom parallel tasks | Futures | Dynamic workflows, parameter sweeps, task dependencies | | Any of above | Schedulers | Control execution backend (threads/processes/distributed) |
Control level: DataFrames/Arrays/Bags = high-level lazy API. Futures = low-level immediate execution.
Lazy Evaluation Model
All DataFrames, Arrays, and Bags build a task graph — nothing executes until .compute() or .persist().
.compute()— execute and return result to local memory.persist()— execute and keep result on workers (for reuse across multiple computations)dask.compute(a, b, c)— compute multiple results in a single pass (shares intermediates)
Chunk Size Strategy
Target: ~100 MB per chunk (or 10 chunks per core in worker memory).
| Chunk Size | Effect | |-----------|--------| | Too large (>1 GB) | Memory overflow, poor parallelization | | Optimal (~100 MB) | Good parallelism, manageable memory | | Too small (<1 MB) | Excessive scheduling overhead |
Example: 8 cores, 32 GB RAM → target ~400 MB per chunk (32 GB / 8 cores / 10).
Scheduler Selection Guide
| Scheduler | Overhead | Best For | GIL |
|-----------|----------|----------|-----|
| threads (default) | ~10 µs/task | NumPy, pandas, scikit-learn | Affected |
| processes | ~10 ms/task | Pure Python, text processing | Not affected |
| synchronous | ~1 µs/task | Debugging with pdb | N/A |
| distributed | ~1 ms/task | Dashboard, clusters, advanced features | Configurable |
Common Workflows
Workflow 1: Multi-File ETL Pipeline
import dask.dataframe as dd
import dask
# Extract: Read all CSV files
ddf = dd.read_csv('raw_data/*.csv', dtype={'amount': 'float64'})
# Transform: Clean and process
ddf = ddf[ddf['status'] == 'valid']
ddf['amount'] = ddf['amount'].fillna(0)
ddf = ddf.dropna(subset=['category'])
# Aggregate
summary = ddf.groupby('category').agg({'amount': ['sum', 'mean', 'count']})
# Load: Save as Parquet (columnar, compressed)
summary.to_parquet('output/summary.parquet')
print(f"Processed {len(ddf)} rows across {ddf.npartitions} partitions")
Workflow 2: Large-Scale Array Processing
import dask.array as da
# Load large scientific dataset
x = da.from_zarr('experiment_data.zarr') # e.g., (50000, 50000) float64
print(f"Shape: {x.shape}, Chunks: {x.chunks}")
# Normalize per-column
x_norm = (x - x.mean(axis=0)) / x.std(axis=0)
# Compute covariance matrix
cov = da.dot(x_norm.T, x_norm) / (x_norm.shape[0] - 1)
# SVD for dimensionality reduction (top-k)
u, s, v = da.linalg.svd_compressed(x_norm, k=50)
# Save results
da.to_zarr(u, 'pca_components.zarr')
print(f"Explained variance (top 5): {(s[:5]**2 / (s**2).sum()).compute()}")
Workflow 3: Unstructured Data to Analysis
This workflow is a simple combination of Bags (Section 3) → DataFrames (Section 1): read JSON logs with Bags, filter/transform, convert to DataFrame for groupby analysis. Each step maps directly to Core API examples above.
Key Parameters
| Parameter | Module | Default | Description |
|-----------|--------|---------|-------------|
| npartitions | DataFrame | auto | Number of partitions (controls parallelism) |
| partition_size | DataFrame | — | Target size per partition (e.g., '100MB') |
| chunks | Array | required | Chunk dimensions (e.g., (10000, 1000)) |
| blocksize | Bag | '128 MiB' | File read block size |
| scheduler | All | 'threads' | Execution backend ('threads', 'processes', 'synchronous') |
| n_workers | Distributed | auto | Number of worker processes |
| threads_per_worker | Distributed | auto | Threads per worker |
| memory_limit | Distributed | auto | Per-worker memory limit (e.g., '4GB') |
| sorted | set_index | False | Whether data is pre-sorted (enables optimizations) |
| meta | map_partitions | — | Output DataFrame/Series structure template |
Best Practices
-
Let Dask handle data loading — Never load data into pandas/numpy first then convert. Use
dd.read_csv()/da.from_zarr()directly. -
Batch compute calls — Use
dask.compute(a, b, c)instead of calling.compute()in loops. Allows sharing intermediates. -
Use
map_partitionsoverapply—ddf.apply(func, axis=1)creates one task per row.ddf.map_partitions(func)creates one task per partition. -
Persist reused intermediates — Call
.persist()on data accessed multiple times, thendelwhen done. -
Use the dashboard —
client.dashboard_linkshows task progress, memory usage, worker states. Essential for diagnosing performance issues. -
Anti-pattern — Excessively large task graphs: If
len(ddf.__dask_graph__())returns millions, increase chunk sizes or usemap_partitions/map_blocksto fuse operations. -
Anti-pattern — Wrong scheduler for workload: Using threads for pure Python text processing (GIL-bound) or processes for NumPy operations (unnecessary serialization overhead).
Common Recipes
Re
Truncated for display — read the full file on GitHub.
Related Skills
Agent-Reach
90.8kGive your AI agent eyes to see the entire internet. Read & search Twitter, Reddit, YouTube, GitHub, Bilibili, XiaoHongShu — one CLI, zero API fees.
headroom
74.4kCompress tool outputs, logs, files, and RAG chunks before they reach the LLM. 20% fewer tokens for coding agents, 60-95% fewer tokens for JSON, same answers. Library, proxy, MCP server.
ai-job-search
45.0kThe job search that runs on your machine. AI job application framework built on Claude Code: evaluate postings, tailor CVs, write cover letters, prep interviews. Fork it and own it.
claude-howto
41.7kA visual, example-driven guide to Claude Code — from basic concepts to advanced agents, with copy-paste templates that bring immediate value.
Languages
Trust signals
From repository metadata: license, adoption, age and documentation. Not a code audit — see the Safety scan above for what the skill file itself contains.
