| name | task-distribution |
| description | Distribute work across multiple agents using queue-based load balancing. Use for parallel execution, work distribution, and team coordination. |
Task Distribution
Efficiently distribute work across multiple Claude Code agents using RabbitMQ queues.
Quick Start
Distribute Task (Team Leader)
import AgentOrchestrator from './scripts/orchestrator.js';
const orchestrator = new AgentOrchestrator('team-leader');
await orchestrator.initialize();
await orchestrator.startTeamLeader();
await orchestrator.assignTask({
title: "Implement authentication",
description: "JWT-based auth with refresh tokens",
priority: "high"
});
Receive and Execute Task (Worker)
const orchestrator = new AgentOrchestrator('worker');
await orchestrator.initialize();
await orchestrator.startWorker();
Distribution Strategies
Strategy 1: Round-Robin (Fair Distribution)
await channel.prefetch(1);
Strategy 2: Priority-Based
await orchestrator.assignTask({
title: "Critical bug fix",
priority: "critical"
});
await orchestrator.assignTask({
title: "Refactoring",
priority: "low"
});
Strategy 3: Capability-Based Routing
const task = {
title: "Optimize database queries",
requiredCapability: "database"
};
Strategy 4: Batch Distribution
const tasks = [
{ title: "Process file 1.csv" },
{ title: "Process file 2.csv" },
{ title: "Process file 3.csv" }
];
for (const task of tasks) {
await orchestrator.assignTask(task);
}
Load Balancing
Fair Queue Behavior
Prefetch Control
await channel.prefetch(1);
await channel.prefetch(5);
const prefetch = Math.ceil(availableCPU / workerCount);
await channel.prefetch(prefetch);
Worker Scaling
const queueDepth = await getQueueDepth('agent.tasks');
const workerCount = await getConnectedWorkers();
if (queueDepth / workerCount > 10) {
console.log('⚠️ Queue backing up, start more workers!');
}
Task Lifecycle Management
Task States
const taskStates = {
PENDING: 'queued',
ACTIVE: 'processing',
COMPLETED: 'done',
FAILED: 'error'
};
State Tracking
const taskTracker = new Map();
taskTracker.set(taskId, {
state: 'PENDING',
assignedAt: Date.now()
});
taskTracker.set(taskId, {
state: 'ACTIVE',
workerId: 'worker-01',
startedAt: Date.now()
});
taskTracker.set(taskId, {
state: 'COMPLETED',
result: {...},
completedAt: Date.now(),
duration: completedAt - startedAt
});
Progress Monitoring
await publishStatus({
event: 'task_progress',
taskId,
progress: 0.5,
message: 'Halfway through data processing'
}, 'agent.status.task.progress');
console.log(`Task ${taskId}: 50% complete`);
Retry and Failure Handling
Automatic Retry
await client.consumeTasks('agent.tasks', async (msg, { ack, nack, reject }) => {
const { task } = msg;
try {
await executeTask(task);
ack();
} catch (error) {
if (task.retryCount > 0) {
console.log(`Retrying (${task.retryCount} attempts left)`);
task.retryCount--;
nack(true);
} else {
console.error('Max retries reached');
reject();
}
}
});
Dead Letter Queue
await channel.assertQueue('agent.tasks', {
arguments: {
'x-dead-letter-exchange': 'dlx.tasks',
'x-dead-letter-routing-key': 'failed'
}
});
await channel.consume('dlq.tasks', async (msg) => {
console.error('Task failed permanently:', msg);
await publishStatus({
event: 'task_dead_letter',
task: msg.task,
reason: msg.properties.headers['x-death']
}, 'agent.status.task.failed');
});
Circuit Breaker
let consecutiveFailures = 0;
const failureThreshold = 5;
await client.consumeTasks('agent.tasks', async (msg, { ack, nack }) => {
try {
await executeTask(msg.task);
consecutiveFailures = 0;
ack();
} catch (error) {
consecutiveFailures++;
if (consecutiveFailures >= failureThreshold) {
console.error('Circuit breaker opened!');
await channel.cancel(consumerTag);
await publishStatus({
event: 'circuit_breaker_open',
reason: 'Too many consecutive failures'
}, 'agent.status.error');
} else {
nack(true);
}
}
});
Work Distribution Patterns
Pattern 1: Map-Reduce
const chunks = splitDataIntoChunks(largeDataset);
for (const chunk of chunks) {
await assignTask({
title: `Process chunk ${chunk.id}`,
data: chunk
});
}
const results = [];
await consumeResults('agent.results', async (msg) => {
results.push(msg.result);
if (results.length === chunks.length) {
const finalResult = reduce(results);
console.log('Map-reduce complete:', finalResult);
}
});
Pattern 2: Pipeline
await publishTask({ title: 'Fetch data' }, 'queue.fetch');
await consumeTasks('queue.fetch', async (msg, { ack }) => {
const data = await fetchData();
await publishTask({ title: 'Transform', data }, 'queue.transform');
ack();
});
await consumeTasks('queue.transform', async (msg, { ack }) => {
const transformed = await transform(msg.data);
await loadData(transformed);
ack();
});
Pattern 3: Fan-Out / Fan-In
await consumeTasks('agent.tasks', async (msg, { ack }) => {
const { task } = msg;
if (task.type === 'parallel') {
const subtasks = task.subtasks;
for (const subtask of subtasks) {
await assignTask(subtask);
}
let completed = 0;
await consumeResults('agent.results', (result) => {
completed++;
if (completed === subtasks.length) {
console.log('All parallel tasks complete');
}
});
}
ack();
});
Performance Optimization
Batch Assignment
const tasks = generateTasks(100);
const promises = tasks.map(task => assignTask(task));
await Promise.all(promises);
Prefetching Optimization
const optimalPrefetch = calculatePrefetch({
avgTaskDuration: 2000,
workerCount: 5,
targetLatency: 1000
});
await channel.prefetch(optimalPrefetch);
Connection Pooling
const channels = await createChannelPool(5);
let channelIndex = 0;
for (const task of tasks) {
const channel = channels[channelIndex % channels.length];
await channel.sendToQueue('agent.tasks', task);
channelIndex++;
}
Monitoring and Metrics
Distribution Metrics
const metrics = {
tasksAssigned: 0,
tasksCompleted: 0,
tasksFailed: 0,
avgDuration: 0,
queueDepth: 0,
activeWorkers: 0
};
metrics.tasksAssigned++;
metrics.tasksCompleted++;
metrics.avgDuration = calculateAverage(completionTimes);
setInterval(async () => {
const info = await channel.checkQueue('agent.tasks');
metrics.queueDepth = info.messageCount;
metrics.activeWorkers = info.consumerCount;
}, 10000);
Health Checks
setInterval(async () => {
const health = {
queueDepth: await getQueueDepth(),
workerCount: await getWorkerCount(),
avgProcessingTime: calculateAvg(),
failureRate: calculateFailureRate()
};
if (health.queueDepth > 100) {
alert('High queue depth - scale up workers');
}
if (health.failureRate > 0.1) {
alert('High failure rate - investigate workers');
}
}, 60000);
Best Practices
-
Use Persistent Messages for Critical Tasks
await channel.sendToQueue('queue', msg, { persistent: true });
-
Set Reasonable Retry Limits
task.retryCount = 3;
-
Implement Idempotent Task Handlers
async function handleTask(task) {
if (await isProcessed(task.id)) return;
await process(task);
await markProcessed(task.id);
}
-
Monitor Queue Depth
if (queueDepth > threshold) {
scaleUpWorkers();
}
-
Use Priority for Critical Tasks
task.priority = 'critical';
-
Fair Prefetch for Even Distribution
await channel.prefetch(1);
-
Graceful Shutdown
process.on('SIGTERM', async () => {
await channel.cancel(consumerTag);
await channel.close();
});
Examples
See examples/ directory for:
map-reduce.js - Parallel data processing
pipeline.js - Sequential workflow
fan-out-fan-in.js - Parallel sub-tasks
priority-queue.js - Priority-based execution
retry-dlq.js - Retry and dead letter handling