SKILL.md
Workload Balancing Skill
Distribute work efficiently across parallel workers to maximize throughput and minimize completion time.
Workflow
- Characterize the workload (uniform vs. variable task times)
- Identify bottlenecks (stragglers, uneven distribution)
- Select balancing strategy based on workload characteristics
- Implement partitioning and scheduling logic
- Monitor and adapt to runtime conditions
Load Balancing Decision Tree
What's the workload characteristic?
Uniform task times:
├── Known count → Static partitioning (equal chunks)
├── Streaming input → Round-robin distribution
└── Large items → Size-aware partitioning
Variable task times:
├── Predictable variance → Weighted distribution
├── Unpredictable → Dynamic scheduling / work stealing
└── Long-tail distribution → Work stealing + time limits
Resource constraints:
├── Memory-bound workers → Memory-aware assignment
├── Heterogeneous workers → Capability-based routing
└── Network costs → Locality-aware placement
Balancing Strategies
Strategy 1: Static Chunking (Uniform Workloads)
Best for: predictable, similar-sized tasks
from concurrent.futures import ProcessPoolExecutor
import numpy as np
def static_balanced_process(items, num_workers=4):
"""Divide work into equal chunks upfront."""
chunks = np.array_split(items, num_workers)
with ProcessPoolExecutor(max_workers=num_workers) as executor:
results = list(executor.map(process_chunk, chunks))
return [item for chunk_result in results for item in chunk_result]
Strategy 2: Dynamic Task Queue (Variable Workloads)
Best for: unpredictable task durations
from concurrent.futures import ProcessPoolExecutor, as_completed
from queue import Queue
def dynamic_balanced_process(items, num_workers=4):
"""Workers pull tasks dynamically as they complete."""
results = []
with ProcessPoolExecutor(max_workers=num_workers) as executor:
# Submit one task per worker initially
futures = {executor.submit(process_item, item): item
for item in items[:num_workers]}
pending = list(items[num_workers:])
while futures:
done, _ = wait(futures, return_when=FIRST_COMPLETED)
for future in done:
results.append(future.result())
del futures[future]
# Submit next task if available
if pending:
next_item = pending.pop(0)
futures[executor.submit(process_item, next_item)] = next_item
return results
