Pipeline Broadcast (Map-Reduce)
🎯 Trigger Conditions
Use when asked about parallelizing DSPy tasks, processing multiple inputs simultaneously, or implementing Map-Reduce patterns.
📚 Prerequisites
dspypackage installed- Parallel processing enabled
🛠️ Map-Reduce Patterns
1. Basic Broadcast (Map)
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
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
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
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