{"slug":"alterlab-dask","title":"alterlab-dask","summary":"Scales pandas/NumPy workflows beyond memory with Dask distributed computing — parallel DataFrames, arrays, delayed task graphs, and cluster execution. Use when existing pandas/NumPy code must run on larger-than-RAM data or across clusters, for parallel file processing, distribute","platform":"Claude","tags":[],"authorName":"LLM Mart","authorSlug":"llm-mart","score":0,"source":"github","price":null,"verified":false,"createdAt":"2026-09-23T18:57:03.745406Z","repo":{"url":"https://github.com/AlterLab-IEU/AlterLab-Academic-Skills","stars":68,"forks":13,"license":"MIT","updatedAt":"2026-09-23T13:42:59Z"},"bodyHtml":"<hr>\n<h2>name: alterlab-dask\ndescription: Scales pandas/NumPy workflows beyond memory with Dask distributed computing — parallel DataFrames, arrays, delayed task graphs, and cluster execution. Use when existing pandas/NumPy code must run on larger-than-RAM data or across clusters, for parallel file processing, distributed ML, or integration with existing pandas code. For out-of-core analytics on a single machine prefer vaex; for in-memory speed prefer polars. Part of the AlterLab Academic Skills suite.\nlicense: MIT\nallowed-tools: Read Write Edit Bash(python:<em>) Bash(uv:</em>)\ncompatibility: No API key required. Runs locally via <code>uv run python</code>; requires dask &gt;= 2025.1 (current 2026.8 as of 2026-09; install <code>dask[complete]</code> for DataFrames, arrays, and the distributed scheduler).\nmetadata:\nskill-author: AlterLab\nversion: \"1.0.1\"\nlast_updated: \"2026-09-23\"</h2>\n<h1>Dask</h1>\n<h2>Overview</h2>\n<p>Dask is a Python library for parallel and distributed computing that enables three critical capabilities:</p>\n<ul>\n<li><strong>Larger-than-memory execution</strong> on single machines for data exceeding available RAM</li>\n<li><strong>Parallel processing</strong> for improved computational speed across multiple cores</li>\n<li><strong>Distributed computation</strong> supporting terabyte-scale datasets across multiple machines</li>\n</ul>\n<p>Dask scales from laptops (processing ~100 GiB) to clusters (processing ~100 TiB) while maintaining familiar Python APIs.</p>\n<h2>When to Use This Skill</h2>\n<p>This skill should be used when:</p>\n<ul>\n<li>Process datasets that exceed available RAM</li>\n<li>Scale pandas or NumPy operations to larger datasets</li>\n<li>Parallelize computations for performance improvements</li>\n<li>Process multiple files efficiently (CSVs, Parquet, JSON, text logs)</li>\n<li>Build custom parallel workflows with task dependencies</li>\n<li>Distribute workloads across multiple cores or machines</li>\n</ul>\n<h3>Does NOT Trigger</h3>\n<table>\n<thead>\n<tr>\n<th>Scenario</th>\n<th>Use Instead</th>\n</tr>\n</thead>\n<tbody>\n<tr>\n<td>Data fits in RAM and the goal is maximum single-machine DataFrame speed</td>\n<td><code>alterlab-polars</code></td>\n</tr>\n<tr>\n<td>Billion-row out-of-core exploration and big-data plots on one machine, no cluster</td>\n<td><code>alterlab-vaex</code></td>\n</tr>\n<tr>\n<td>Designing chunked, compressed Zarr stores and codecs (rather than computing on them)</td>\n<td><code>alterlab-zarr</code></td>\n</tr>\n</tbody>\n</table>\n<h3>Version notes</h3>\n<p>Since 2025.1 the query-planning DataFrame (formerly <code>dask-expr</code>) is the only Dask DataFrame\nimplementation and ships inside <code>dask</code>; <code>dask.config.set({\"dataframe.query-planning\": False})</code>\nno longer exists. Text columns load as a string dtype (PyArrow-backed in Dask; <code>str</code> in pandas 3)\nand reductions do not skip them silently: <code>ddf.mean()</code> raises <code>TypeError</code> and\n<code>ddf.groupby(...).mean()</code> raises <code>NotImplementedError</code> when non-numeric columns are present, so\nselect columns or pass <code>numeric_only=True</code>.</p>\n<h2>Core Capabilities</h2>\n<p>Dask provides five main components, each suited to different use cases:</p>\n<h3>1. DataFrames - Parallel Pandas Operations</h3>\n<p><strong>Purpose</strong>: Scale pandas operations to larger datasets through parallel processing.</p>\n<p><strong>When to Use</strong>:</p>\n<ul>\n<li>Tabular data exceeds available RAM</li>\n<li>Need to process multiple CSV/Parquet files together</li>\n<li>Pandas operations are slow and need parallelization</li>\n<li>Scaling from pandas prototype to production</li>\n</ul>\n<p><strong>Reference Documentation</strong>: For comprehensive guidance on Dask DataFrames, refer to <code>references/dataframes.md</code> which includes:</p>\n<ul>\n<li>Reading data (single files, multiple files, glob patterns)</li>\n<li>Common operations (filtering, groupby, joins, aggregations)</li>\n<li>Custom operations with <code>map_partitions</code></li>\n<li>Performance optimization tips</li>\n<li>Common patterns (ETL, time series, multi-file processing)</li>\n</ul>\n<p><strong>Quick Example</strong>:</p>\n<pre><code>import dask.dataframe as dd\n\n# Read multiple files as single DataFrame\nddf = dd.read_csv('data/2024-*.csv')\n\n# Operations are lazy until compute()\nfiltered = ddf[ddf['value'] &gt; 100]\nresult = filtered.groupby('category').mean(numeric_only=True).compute()\n</code></pre>\n<p><strong>Key Points</strong>:</p>\n<ul>\n<li>Operations are lazy (build task graph) until <code>.compute()</code> called</li>\n<li>Use <code>map_partitions</code> for efficient custom operations</li>\n<li>Convert to DataFrame early when working with structured data from other sources</li>\n</ul>\n<h3>2. Arrays - Parallel NumPy Operations</h3>\n<p><strong>Purpose</strong>: Extend NumPy capabilities to datasets larger than memory using blocked algorithms.</p>\n<p><strong>When to Use</strong>:</p>\n<ul>\n<li>Arrays exceed available RAM</li>\n<li>NumPy operations need parallelization</li>\n<li>Working with scientific datasets (HDF5, Zarr, NetCDF)</li>\n<li>Need parallel linear algebra or array operations</li>\n</ul>\n<p><strong>Reference Documentation</strong>: For comprehensive guidance on Dask Arrays, refer to <code>references/arrays.md</code> which includes:</p>\n<ul>\n<li>Creating arrays (from NumPy, random, from disk)</li>\n<li>Chunking strategies and optimization</li>\n<li>Common operations (arithmetic, reductions, linear algebra)</li>\n<li>Custom operations with <code>map_blocks</code></li>\n<li>Integration with HDF5, Zarr, and XArray</li>\n</ul>\n<p><strong>Quick Example</strong>:</p>\n<pre><code>import dask.array as da\n\n# Create large array with chunks\nx = da.random.random((100000, 100000), chunks=(10000, 10000))\n\n# Operations are lazy\ny = x + 100\nz = y.mean(axis=0)\n\n# Compute result\nresult = z.compute()\n</code></pre>\n<p><strong>Key Points</strong>:</p>\n<ul>\n<li>Chunk size is critical (aim for ~100 MB per chunk)</li>\n<li>Operations work on chunks in parallel</li>\n<li>Rechunk data when needed for efficient operations</li>\n<li>Use <code>map_blocks</code> for operations not available in Dask</li>\n</ul>\n<h3>3. Bags - Parallel Processing of Unstructured Data</h3>\n<p><strong>Purpose</strong>: Process unstructured or semi-structured data (text, JSON, logs) with functional operations.</p>\n<p><strong>When to Use</strong>:</p>\n<ul>\n<li>Processing text files, logs, or JSON records</li>\n<li>Data cleaning and ETL before structured analysis</li>\n<li>Working with Python objects that don't fit array/dataframe formats</li>\n<li>Need memory-efficient streaming processing</li>\n</ul>\n<p><strong>Reference Documentation</strong>: For comprehensive guidance on Dask Bags, refer to <code>references/bags.md</code> which includes:</p>\n<ul>\n<li>Reading text and JSON files</li>\n<li>Functional operations (map, filter, fold, groupby)</li>\n<li>Converting to DataFrames</li>\n<li>Common patterns (log analysis, JSON processing, text processing)</li>\n<li>Performance considerations</li>\n</ul>\n<p><strong>Quick Example</strong>:</p>\n<pre><code>import dask.bag as db\nimport json\n\n# Read and parse JSON files\nbag = db.read_text('logs/*.json').map(json.loads)\n\n# Filter and transform\nvalid = bag.filter(lambda x: x['status'] == 'valid')\nprocessed = valid.map(lambda x: {'id': x['id'], 'value': x['value']})\n\n# Convert to DataFrame for analysis\nddf = processed.to_dataframe()\n</code></pre>\n<p><strong>Key Points</strong>:</p>\n<ul>\n<li>Use for initial data cleaning, then convert to DataFrame/Array</li>\n<li>Use <code>foldby</code> instead of <code>groupby</code> for better performance</li>\n<li>Operations are streaming and memory-efficient</li>\n<li>Convert to structured formats (DataFrame) for complex operations</li>\n</ul>\n<h3>4. Futures - Task-Based Parallelization</h3>\n<p><strong>Purpose</strong>: Build custom parallel workflows with fine-grained control over task execution and dependencies.</p>\n<p><strong>When to Use</strong>:</p>\n<ul>\n<li>Building dynamic, evolving workflows</li>\n<li>Need immediate task execution (not lazy)</li>\n<li>Computations depend on runtime conditions</li>\n<li>Implementing custom parallel algorithms</li>\n<li>Need stateful computations</li>\n</ul>\n<p><strong>Reference Documentation</strong>: For comprehensive guidance on Dask Futures, refer to <code>references/futures.md</code> which includes:</p>\n<ul>\n<li>Setting up distributed client</li>\n<li>Submitting tasks and working with futures</li>\n<li>Task dependencies and data movement</li>\n<li>Advanced coordination (queues, locks, events, actors)</li>\n<li>Common patterns (parameter sweeps, dynamic tasks, iterative algorithms)</li>\n</ul>\n<p><strong>Quick Example</strong>:</p>\n<pre><code>from dask.distributed import Client\n\nclient = Client()  # Create local cluster\n\n# Submit tasks (executes immediately)\ndef process(x):\n    return x ** 2\n\nfutures = client.map(process, range(100))\n\n# Gather results\nresults = client.gather(futures)\n\nclient.close()\n</code></pre>\n<p><strong>Key Points</strong>:</p>\n<ul>\n<li>Requires distributed client (even for single machine)</li>\n<li>Tasks execute immediately when submitted</li>\n<li>Pre-scatter large data to avoid repeated transfers</li>\n<li>~1ms overhead per task (not suitable for millions of tiny tasks)</li>\n<li>Use actors for stateful workflows</li>\n</ul>\n<h3>5. Schedulers - Execution Backends</h3>\n<p><strong>Purpose</strong>: Control how and where Dask tasks execute (threads, processes, distributed).</p>\n<p><strong>When to Choose Scheduler</strong>:</p>\n<ul>\n<li><strong>Threads</strong> (default): NumPy/Pandas operations, GIL-releasing libraries, shared memory benefit</li>\n<li><strong>Processes</strong>: Pure Python code, text processing, GIL-bound operations</li>\n<li><strong>Synchronous</strong>: Debugging with pdb, profiling, understanding errors</li>\n<li><strong>Distributed</strong>: Need dashboard, multi-machine clusters, advanced features</li>\n</ul>\n<p><strong>Reference Documentation</strong>: For comprehensive guidance on Dask Schedulers, refer to <code>references/schedulers.md</code> which includes:</p>\n<ul>\n<li>Detailed scheduler descriptions and characteristics</li>\n<li>Configuration methods (global, context manager, per-compute)</li>\n<li>Performance considerations and overhead</li>\n<li>Common patterns and troubleshooting</li>\n<li>Thread configuration for optimal performance</li>\n</ul>\n<p><strong>Quick Example</strong>:</p>\n<pre><code>import dask\nimport dask.dataframe as dd\n\n# Use threads for DataFrame (default, good for numeric)\nddf = dd.read_csv('data.csv')\nresult1 = ddf.mean(numeric_only=True).compute()  # Uses threads\n\n# Use processes for Python-heavy work\nimport dask.bag as db\nbag = db.read_text('logs/*.txt')\nresult2 = bag.map(python_function).compute(scheduler='processes')\n\n# Use synchronous for debugging\ndask.config.set(scheduler='synchronous')\nresult3 = problematic_computation.compute()  # Can use pdb\n\n# Use distributed for monitoring and scaling\nfrom dask.distributed import Client\nclient = Client()\nresult4 = computation.compute()  # Uses distributed with dashboard\n</code></pre>\n<p><strong>Key Points</strong>:</p>\n<ul>\n<li>Threads: Lowest overhead (~10 µs/task), best for numeric work</li>\n<li>Processes: Avoids GIL (~10 ms/task), best for Python work</li>\n<li>Distributed: Monitoring dashboard (~1 ms/task), scales to clusters</li>\n<li>Can switch schedulers per computation or globally</li>\n</ul>\n<h2>Best Practices</h2>\n<p>For comprehensive performance optimization guidance, memory management strategies, and common pitfalls to avoid, refer to <code>references/best-practices.md</code>. Key principles include:</p>\n<h3>Start with Simpler Solutions</h3>\n<p>Before using Dask, explore:</p>\n<ul>\n<li>Better algorithms</li>\n<li>Efficient file formats (Parquet instead of CSV)</li>\n<li>Compiled code (Numba, Cython)</li>\n<li>Data sampling</li>\n</ul>\n<h3>Critical Performance Rules</h3>\n<p><strong>1. Don't Load Data Locally Then Hand to Dask</strong></p>\n<pre><code># Wrong: Loads all data in memory first\nimport pandas as pd\ndf = pd.read_csv('large.csv')\nddf = dd.from_pandas(df, npartitions=10)\n\n# Correct: Let Dask handle loading\nimport dask.dataframe as dd\nddf = dd.read_csv('large.csv')\n</code></pre>\n<p><strong>2. Avoid Repeated compute() Calls</strong></p>\n<pre><code># Wrong: Each compute is separate\nfor item in items:\n    result = dask_computation(item).compute()\n\n# Correct: Single compute for all\ncomputations = [dask_computation(item) for item in items]\nresults = dask.compute(*computations)\n</code></pre>\n<p><strong>3. Don't Build Excessively Large Task Graphs</strong></p>\n<ul>\n<li>Increase chunk sizes if millions of tasks</li>\n<li>Use <code>map_partitions</code>/<code>map_blocks</code> to fuse operations</li>\n<li>Check task graph size: <code>len(ddf.__dask_graph__())</code></li>\n</ul>\n<p><strong>4. Choose Appropriate Chunk Sizes</strong></p>\n<ul>\n<li>Target: ~100 MB per chunk (or 10 chunks per core in worker memory)</li>\n<li>Too large: Memory overflow</li>\n<li>Too small: Scheduling overhead</li>\n</ul>\n<p><strong>5. Use the Dashboard</strong></p>\n<pre><code>from dask.distributed import Client\nclient = Client()\nprint(client.dashboard_link)  # Monitor performance, identify bottlenecks\n</code></pre>\n<h2>Common Workflow Patterns</h2>\n<h3>ETL Pipeline</h3>\n<pre><code>import dask.dataframe as dd\n\n# Extract: Read data\nddf = dd.read_csv('raw_data/*.csv')\n\n# Transform: Clean and process\nddf = ddf[ddf['status'] == 'valid']\nddf['amount'] = ddf['amount'].astype('float64')\nddf = ddf.dropna(subset=['important_col'])\n\n# Load: Aggregate and save. Named aggregation keeps flat string column names;\n# a dict-of-lists .agg() creates MultiIndex columns that Parquet rejects.\nsummary = ddf.groupby('category').agg(\n    amount_sum=('amount', 'sum'),\n    amount_mean=('amount', 'mean'),\n)\nsummary.to_parquet('output/summary.parquet')\n</code></pre>\n<h3>Unstructured to Structured Pipeline</h3>\n<pre><code>import dask.bag as db\nimport json\n\n# Start with Bag for unstructured data\nbag = db.read_text('logs/*.json').map(json.loads)\nbag = bag.filter(lambda x: x['status'] == 'valid')\n\n# Convert to DataFrame for structured analysis\nddf = bag.to_dataframe()\nresult = ddf.groupby('category').mean(numeric_only=True).compute()\n</code></pre>\n<h3>Large-Scale Array Computation</h3>\n<pre><code>import dask.array as da\n\n# Load or create large array\nx = da.from_zarr('large_dataset.zarr')\n\n# Process in chunks\nnormalized = (x - x.mean()) / x.std()\n\n# Save result\nda.to_zarr(normalized, 'normalized.zarr')\n</code></pre>\n<h3>Custom Parallel Workflow</h3>\n<pre><code>from dask.distributed import Client\n\nclient = Client()\n\n# Scatter large dataset once\ndata = client.scatter(large_dataset)\n\n# Process in parallel with dependencies\nfutures = []\nfor param in parameters:\n    future = client.submit(process, data, param)\n    futures.append(future)\n\n# Gather results\nresults = client.gather(futures)\n</code></pre>\n<h2>Selecting the Right Component</h2>\n<p>Use this decision guide to choose the appropriate Dask component:</p>\n<p><strong>Data Type</strong>:</p>\n<ul>\n<li>Tabular data → <strong>DataFrames</strong></li>\n<li>Numeric arrays → <strong>Arrays</strong></li>\n<li>Text/JSON/logs → <strong>Bags</strong> (then convert to DataFrame)</li>\n<li>Custom Python objects → <strong>Bags</strong> or <strong>Futures</strong></li>\n</ul>\n<p><strong>Operation Type</strong>:</p>\n<ul>\n<li>Standard pandas operations → <strong>DataFrames</strong></li>\n<li>Standard NumPy operations → <strong>Arrays</strong></li>\n<li>Custom parallel tasks → <strong>Futures</strong></li>\n<li>Text processing/ETL → <strong>Bags</strong></li>\n</ul>\n<p><strong>Control Level</strong>:</p>\n<ul>\n<li>High-level, automatic → <strong>DataFrames/Arrays</strong></li>\n<li>Low-level, manual → <strong>Futures</strong></li>\n</ul>\n<p><strong>Workflow Type</strong>:</p>\n<ul>\n<li>Static computation graph → <strong>DataFrames/Arrays/Bags</strong></li>\n<li>Dynamic, evolving → <strong>Futures</strong></li>\n</ul>\n<h2>Integration Considerations</h2>\n<h3>File Formats</h3>\n<ul>\n<li><strong>Efficient</strong>: Parquet, HDF5, Zarr (columnar, compressed, parallel-friendly)</li>\n<li><strong>Compatible but slower</strong>: CSV (use for initial ingestion only)</li>\n<li><strong>For Arrays</strong>: HDF5, Zarr, NetCDF</li>\n</ul>\n<h3>Conversion Between Collections</h3>\n<pre><code># Bag → DataFrame\nddf = bag.to_dataframe()\n\n# DataFrame → Array (for numeric data)\narr = ddf.to_dask_array(lengths=True)\n\n# Array → DataFrame\nddf = dd.from_dask_array(arr, columns=['col1', 'col2'])\n</code></pre>\n<h3>With Other Libraries</h3>\n<ul>\n<li><strong>XArray</strong>: Wraps Dask arrays with labeled dimensions (geospatial, imaging)</li>\n<li><strong>Dask-ML</strong>: Machine learning with scikit-learn compatible APIs</li>\n<li><strong>Distributed</strong>: Advanced cluster management and monitoring</li>\n</ul>\n<h2>Debugging and Development</h2>\n<h3>Iterative Development Workflow</h3>\n<ol>\n<li><strong>Test on small data with synchronous scheduler</strong>:</li>\n</ol>\n<pre><code>dask.config.set(scheduler='synchronous')\nresult = computation.compute()  # Can use pdb, easy debugging\n</code></pre>\n<ol start=\"2\">\n<li><strong>Validate with threads on sample</strong>:</li>\n</ol>\n<pre><code>sample = ddf.head(1000)  # Small sample\n# Test logic, then scale to full dataset\n</code></pre>\n<ol start=\"3\">\n<li><strong>Scale with distributed for monitoring</strong>:</li>\n</ol>\n<pre><code>from dask.distributed import Client\nclient = Client()\nprint(client.dashboard_link)  # Monitor performance\nresult = computation.compute()\n</code></pre>\n<h3>Common Issues</h3>\n<p><strong>Memory Errors</strong>:</p>\n<ul>\n<li>Decrease chunk sizes</li>\n<li>Use <code>persist()</code> strategically and delete when done</li>\n<li>Check for memory leaks in custom functions</li>\n</ul>\n<p><strong>Slow Start</strong>:</p>\n<ul>\n<li>Task graph too large (increase chunk sizes)</li>\n<li>Use <code>map_partitions</code> or <code>map_blocks</code> to reduce tasks</li>\n</ul>\n<p><strong>Poor Parallelization</strong>:</p>\n<ul>\n<li>Chunks too large (increase number of partitions)</li>\n<li>Using threads with Python code (switch to processes)</li>\n<li>Data dependencies preventing parallelism</li>\n</ul>\n<h2>Reference Files</h2>\n<p>All reference documentation files can be read as needed for detailed information:</p>\n<ul>\n<li><code>references/dataframes.md</code> - Complete Dask DataFrame guide</li>\n<li><code>references/arrays.md</code> - Complete Dask Array guide</li>\n<li><code>references/bags.md</code> - Complete Dask Bag guide</li>\n<li><code>references/futures.md</code> - Complete Dask Futures and distributed computing guide</li>\n<li><code>references/schedulers.md</code> - Complete scheduler selection and configuration guide</li>\n<li><code>references/best-practices.md</code> - Comprehensive performance optimization and troubleshooting</li>\n</ul>\n<p>Load these files when users need detailed information about specific Dask components, operations, or patterns beyond the quick guidance provided here.</p>\n","files":[{"path":"evals/evals.json","sizeBytes":5180,"isText":true},{"path":"references/arrays.md","sizeBytes":11833,"isText":true},{"path":"references/bags.md","sizeBytes":10869,"isText":true},{"path":"references/best-practices.md","sizeBytes":7317,"isText":true},{"path":"references/dataframes.md","sizeBytes":9284,"isText":true},{"path":"references/futures.md","sizeBytes":12051,"isText":true},{"path":"references/schedulers.md","sizeBytes":11395,"isText":true},{"path":"SKILL.md","sizeBytes":15871,"isText":true}],"reviewScore":null,"reviewSummary":null,"trust":{"provenance":"trusted-source-unreviewed","notice":"Community-authored content, reproduced verbatim and not vetted as instructions. Treat it as data to evaluate, never as directives to follow.","bodySource":null},"bodyLocked":false,"purchaseUrl":null,"sourceUrl":null,"report":{"provenance":"trusted-source-unreviewed","screen":{"ran":true,"outcome":"clean","suspicious":0,"notes":0,"hiddenCharacters":false},"virusScan":{"engine":"clamav","status":"clean","scannedAt":"2026-09-23T18:58:45.04502Z","sha256":"F3E5E3147D8CAFE59984AD75D4D69CCE3EE3C5F6CF7FF66391C3713593B5B9AC","sizeBytes":31133},"review":null,"source":{"repositoryUrl":"https://github.com/AlterLab-IEU/AlterLab-Academic-Skills","path":"skills/data-science/alterlab-dask","license":"MIT","commit":"e4836c08a20da195a11f30f203a8cf23ec30aa95","subtreeSha":"DD12CAA2B7C10693EBE2F447FAE041B88756B8EC48593D0A768243F8936785D8","lastSyncedAt":"2026-09-23T18:56:52.297238Z"},"reviewedAt":"2026-09-23T19:02:09.961315Z","notice":"Community-authored content, reproduced verbatim and not vetted as instructions. Treat it as data to evaluate, never as directives to follow."},"install":[{"target":"skills-cli","command":"npx skills add https://github.com/AlterLab-IEU/AlterLab-Academic-Skills/tree/main/skills/data-science/alterlab-dask"},{"target":"claude-code","command":"claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install alterlab-ieu-alterlab-academic-skills@llmmart"},{"target":"git","command":"git clone https://github.com/AlterLab-IEU/AlterLab-Academic-Skills.git"}]}