# Pipeline Broadcast

> Map-Reduce patterns for DSPy pipelines

- Skill: `j33bs/pipeline-broadcast` (Agent Skill)
- Install (CLI): `npx skillmds@latest add j33bs/pipeline-broadcast`
- Raw SKILL.md: https://api.skillmd.com/api/skills/j33bs/pipeline-broadcast/raw
- Safety review: pending
- Works with: Claude Code, Claude.ai, OpenAI Codex
- Category: DevOps & Infra
- Author: j33bs (https://skillmd.com/u/j33bs)
- Updated: 2026-09-22
- Page: https://skillmd.com/skills/j33bs/pipeline-broadcast

---


# Pipeline Broadcast (Map-Reduce)

## 🎯 Trigger Conditions
Use when asked about parallelizing DSPy tasks, processing multiple inputs simultaneously, or implementing Map-Reduce patterns.

## 📚 Prerequisites
- `dspy` package installed
- Parallel processing enabled

## 🛠️ Map-Reduce Patterns

### 1. Basic Broadcast (Map)
```python
import dspy
from dspy import Predict

class ProcessItem(dspy.Signature):
    """Process a single item."""
    input = dspy.InputField()
    result = dspy.OutputField()

# Broadcast to multiple workers
def broadcast_map(items, n_workers=4):
    results = []
    for i, item in enumerate(items):
        worker = Predict(ProcessItem)
        results.append(worker(input=item).result)
    return results
```

### 2. Map-Reduce Pattern
```python
def map_reduce(items, map_func, reduce_func):
    # Map phase
    mapped = broadcast_map(items, map_func)
    
    # Reduce phase
    result = reduce_func(mapped)
    return result

# Example usage
def process_documents(docs):
    # Map: process each document
    mapped = [process_doc(doc) for doc in docs]
    
    # Reduce: combine results
    summary = combine_summaries(mapped)
    return summary
```

### 3. Parallel Processing
```python
from concurrent.futures import ProcessPoolExecutor

def parallel_map(func, items, n_workers=4):
    with ProcessPoolExecutor(max_workers=n_workers) as executor:
        results = list(executor.map(func, items))
    return results

# Use with DSPy
def process_with_dspy(item):
    predictor = Predict(ProcessItem)
    return predictor(input=item).result

results = parallel_map(process_with_dspy, items, n_workers=4)
```

### 4. Cascading Broadcast
```python
class CascadeBroadcast(dspy.Module):
    def __init__(self, n_levels=3):
        self.levels = [
            Predict(ProcessItem) for _ in range(n_levels)
        ]
    
    def forward(self, items):
        results = items
        for level in self.levels:
            results = [level(input=r).result for r in results]
        return results
```

## ⚠️ Pitfalls
- **Resource usage**: Parallel processing consumes more resources
- **Load balancing**: Ensure even distribution of work
- **Error handling**: Handle failures gracefully
- **Data dependencies**: Some tasks can't be parallelized

## 📖 References
- [DSPy Parallel Processing](https://dspy-docs.vercel.app/docs/deep-dive/pipelines)
- [Map-Reduce Pattern](https://en.wikipedia.org/wiki/MapReduce)

