← 返回
未分类 中文

workload-balancing

Optimize workload distribution across workers, processes, or nodes for efficient parallel execution. Use when asked to balance work distribution, improve par...
Optimize workload distribution across workers, processes, or nodes for efficient parallel execution. Use when asked to balance work distribution, improve par...
lnj22
未分类 clawhub v0.1.0 1 版本 100000 Key: 无需
★ 0
Stars
📥 326
下载
💾 0
安装
1
版本
#latest

概述

Workload Balancing Skill

Distribute work efficiently across parallel workers to maximize throughput and minimize completion time.

Workflow

  1. Characterize the workload (uniform vs. variable task times)
  2. Identify bottlenecks (stragglers, uneven distribution)
  3. Select balancing strategy based on workload characteristics
  4. Implement partitioning and scheduling logic
  5. 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

Strategy 3: Work Stealing (Long-Tail Tasks)

Best for: when some tasks take much longer than others

import asyncio
from collections import deque

class WorkStealingPool:
    def __init__(self, num_workers):
        self.queues = [deque() for _ in range(num_workers)]
        self.num_workers = num_workers

    def distribute(self, items):
        """Initial round-robin distribution."""
        for i, item in enumerate(items):
            self.queues[i % self.num_workers].append(item)

    async def worker(self, worker_id, process_fn):
        """Process own queue, steal from others when empty."""
        while True:
            # Try own queue first
            if self.queues[worker_id]:
                item = self.queues[worker_id].popleft()
            else:
                # Steal from busiest queue
                item = self._steal_work(worker_id)
                if item is None:
                    break

            await process_fn(item)

    def _steal_work(self, worker_id):
        """Steal from the queue with most items."""
        busiest = max(range(self.num_workers),
                      key=lambda i: len(self.queues[i]) if i != worker_id else 0)
        if self.queues[busiest]:
            return self.queues[busiest].pop()  # Steal from end
        return None

Strategy 4: Weighted Distribution

Best for: when task costs are known or estimable

def weighted_partition(items, weights, num_workers):
    """Partition items to balance total weight per worker."""
    # Sort by weight descending (largest first fit)
    sorted_items = sorted(zip(items, weights), key=lambda x: -x[1])

    worker_loads = [0] * num_workers
    worker_items = [[] for _ in range(num_workers)]

    for item, weight in sorted_items:
        # Assign to least loaded worker
        min_worker = min(range(num_workers), key=lambda i: worker_loads[i])
        worker_items[min_worker].append(item)
        worker_loads[min_worker] += weight

    return worker_items

Strategy 5: Async Semaphore Balancing (I/O Workloads)

Best for: limiting concurrent I/O operations

import asyncio

async def semaphore_balanced_fetch(urls, max_concurrent=10):
    """Limit concurrent operations while processing queue."""
    semaphore = asyncio.Semaphore(max_concurrent)

    async def bounded_fetch(url):
        async with semaphore:
            return await fetch(url)

    return await asyncio.gather(*[bounded_fetch(url) for url in urls])

Partitioning Strategies

StrategyBest ForImplementation
------------------------------------
Equal chunksUniform tasksnp.array_split(items, n)
Round-robinStreamingitems[i::n_workers]
Size-weightedKnown sizesBin packing algorithm
Hash-basedConsistent routinghash(key) % n_workers
Range-basedSorted/ordered dataContiguous ranges

Handling Stragglers

Techniques to mitigate slow workers:

# 1. Timeout with fallback
from concurrent.futures import TimeoutError

try:
    result = future.result(timeout=30)
except TimeoutError:
    result = fallback_value

# 2. Speculative execution (backup tasks)
async def speculative_execute(task, timeout=10):
    primary = asyncio.create_task(execute(task))
    try:
        return await asyncio.wait_for(primary, timeout)
    except asyncio.TimeoutError:
        backup = asyncio.create_task(execute(task))  # Retry
        done, pending = await asyncio.wait(
            [primary, backup], return_when=asyncio.FIRST_COMPLETED
        )
        for p in pending:
            p.cancel()
        return done.pop().result()

# 3. Dynamic rebalancing
def rebalance_on_straggler(futures, threshold_ratio=2.0):
    """Redistribute work if one worker falls behind."""
    avg_completion = statistics.mean(completion_times)
    for future, worker_id in futures.items():
        if future.running() and elapsed(future) > threshold_ratio * avg_completion:
            # Cancel and redistribute
            remaining_work = cancel_and_get_remaining(future)
            redistribute(remaining_work, fast_workers)

Monitoring Metrics

Track these for balanced execution:

MetricCalculationTarget
-----------------------------
Load imbalancemax(load) / avg(load)< 1.2
Straggler ratiomax(time) / median(time)< 2.0
Worker utilizationbusy_time / total_time> 90%
Queue depth variancestd(queue_lengths)Low

Anti-Patterns

ProblemCauseFix
---------------------
StarvationLarge tasks block queueBreak into subtasks
Thundering herdAll workers wake at onceJittered scheduling
Hot spotsUneven key distributionBetter hash function
Convoy effectWorkers wait on same resourceFine-grained locking
Over-partitioningToo many small tasksBatch small items

Verification Checklist

Before finalizing balanced code:

  • [ ] Work distribution is roughly even (measure completion times)
  • [ ] No starvation (all workers stay busy)
  • [ ] Stragglers are handled (timeout/retry logic)
  • [ ] Overhead is acceptable (partitioning cost vs. task cost)
  • [ ] Results are complete and correct
  • [ ] Resource utilization is high across workers

版本历史

共 1 个版本

  • v0.1.0 当前
    2026-05-07 20:09 安全 安全

安全检测

腾讯云安全 (Keen)

安全,无风险
查看报告

腾讯云安全 (Sanbu)

安全,无风险
查看报告

🔗 相关推荐

ffmpeg-audio-processing

lnj22
提取、标准化、混合和处理音频轨道,进行音频操控与分析
★ 0 📥 431

pdf

lnj22
全面PDF工具,支持文本/表格提取、新PDF创建、合并/拆分文档、表单处理。当Claude需要...
★ 0 📥 428

pdf

lnj22
全面PDF工具,支持文本/表格提取、新PDF创建、合并/拆分文档、表单处理。当Claude需要...
★ 0 📥 454