Multiple Concurrent Streams
This guide covers patterns for managing multiple concurrent streams using Go's goroutines, channels, and synchronization primitives.
Basic Concurrent Streaming
Multiple Streams in Parallel
package main
import (
"context"
"fmt"
"os"
"sync"
"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)
type StreamResult struct {
ID int
Text string
Error error
}
func streamConcurrently(ctx context.Context, model provider.LanguageModel, prompts []string) []StreamResult {
results := make([]StreamResult, len(prompts))
var wg sync.WaitGroup
for i, prompt := range prompts {
wg.Add(1)
go func(index int, p string) {
defer wg.Done()
var fullText string
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: p,
})
if err != nil {
results[index] = StreamResult{
ID: index,
Error: err,
}
return
}
// Collect all chunks
for chunk := range stream.Chunks() {
fullText += chunk.Text
fmt.Printf("[Stream %d] %s", index, chunk.Text)
}
// Check for streaming errors
if err := stream.Err(); err != nil {
results[index] = StreamResult{
ID: index,
Text: fullText,
Error: err,
}
return
}
results[index] = StreamResult{
ID: index,
Text: fullText,
}
fmt.Printf("\n[Stream %d] Completed\n", index)
}(i, prompt)
}
wg.Wait()
return results
}
func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4o-mini")
prompts := []string{
"Explain AI in one sentence",
"Explain ML in one sentence",
"Explain DL in one sentence",
"Explain NLP in one sentence",
"Explain CV in one sentence",
}
results := streamConcurrently(ctx, model, prompts)
fmt.Println("\n=== Results ===")
for _, result := range results {
if result.Error != nil {
fmt.Printf("Stream %d: ERROR - %v\n", result.ID, result.Error)
} else {
fmt.Printf("Stream %d: %s\n", result.ID, result.Text)
}
}
}
Fan-Out/Fan-In Pattern
Distributing Work and Collecting Results
package main
import (
"context"
"fmt"
"log"
"os"
"sync"
"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)
type Request struct {
ID int
Prompt string
}
type Response struct {
ID int
Text string
Error error
}
// Fan-out: Distribute prompts to workers
func fanOut(ctx context.Context, prompts []Request, numWorkers int) <-chan Request {
requestChan := make(chan Request, len(prompts))
go func() {
defer close(requestChan)
for _, prompt := range prompts {
select {
case requestChan <- prompt:
case <-ctx.Done():
return
}
}
}()
return requestChan
}
// Worker: Process requests
func worker(ctx context.Context, id int, model provider.LanguageModel, requests <-chan Request, responses chan<- Response) {
for req := range requests {
log.Printf("Worker %d processing request %d", id, req.ID)
var fullText string
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: req.Prompt,
})
if err != nil {
responses <- Response{
ID: req.ID,
Error: err,
}
continue
}
for chunk := range stream.Chunks() {
fullText += chunk.Text
}
if err := stream.Err(); err != nil {
responses <- Response{
ID: req.ID,
Text: fullText,
Error: err,
}
continue
}
responses <- Response{
ID: req.ID,
Text: fullText,
}
}
}
// Fan-in: Collect results from workers
func fanIn(ctx context.Context, numWorkers int, model provider.LanguageModel, requests <-chan Request) <-chan Response {
responses := make(chan Response, numWorkers)
var wg sync.WaitGroup
// Start workers
for i := 0; i < numWorkers; i++ {
wg.Add(1)
go func(workerID int) {
defer wg.Done()
worker(ctx, workerID, model, requests, responses)
}(i)
}
// Close responses channel when all workers done
go func() {
wg.Wait()
close(responses)
}()
return responses
}
func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4o-mini")
// Create requests
requests := []Request{
{ID: 1, Prompt: "What is AI?"},
{ID: 2, Prompt: "What is ML?"},
{ID: 3, Prompt: "What is DL?"},
{ID: 4, Prompt: "What is NLP?"},
{ID: 5, Prompt: "What is CV?"},
{ID: 6, Prompt: "What is RL?"},
{ID: 7, Prompt: "What is GANs?"},
{ID: 8, Prompt: "What is transformers?"},
}
// Fan-out to workers
requestChan := fanOut(ctx, requests, 3)
// Fan-in from workers
responseChan := fanIn(ctx, 3, model, requestChan)
// Collect results
results := make(map[int]Response)
for response := range responseChan {
results[response.ID] = response
}
// Print results in order
fmt.Println("\n=== Results ===")
for i := 1; i <= len(requests); i++ {
if resp, ok := results[i]; ok {
if resp.Error != nil {
fmt.Printf("%d: ERROR - %v\n", resp.ID, resp.Error)
} else {
fmt.Printf("%d: %s\n", resp.ID, resp.Text[:min(100, len(resp.Text))])
}
}
}
}
func min(a, b int) int {
if a < b {
return a
}
return b
}
Rate-Limited Concurrent Streaming
Control Concurrency and Rate Limits
package main
import (
"context"
"fmt"
"log"
"os"
"sync"
"time"
"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
"golang.org/x/time/rate"
)
type RateLimitedStreamer struct {
model provider.LanguageModel
limiter *rate.Limiter
semaphore chan struct{}
}
func NewRateLimitedStreamer(model provider.LanguageModel, requestsPerMinute float64, maxConcurrent int) *RateLimitedStreamer {
return &RateLimitedStreamer{
model: model,
limiter: rate.NewLimiter(rate.Limit(requestsPerMinute/60.0), 1),
semaphore: make(chan struct{}, maxConcurrent),
}
}
func (s *RateLimitedStreamer) Stream(ctx context.Context, prompt string, resultChan chan<- string, errorChan chan<- error) {
// Acquire semaphore (limit concurrency)
select {
case s.semaphore <- struct{}{}:
defer func() { <-s.semaphore }()
case <-ctx.Done():
errorChan <- ctx.Err()
return
}
// Wait for rate limiter
if err := s.limiter.Wait(ctx); err != nil {
errorChan <- err
return
}
// Start streaming
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: s.model,
Prompt: prompt,
})
if err != nil {
errorChan <- err
return
}
var fullText string
for chunk := range stream.Chunks() {
fullText += chunk.Text
}
if err := stream.Err(); err != nil {
errorChan <- err
return
}
resultChan <- fullText
}
func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4o-mini")
// Create streamer: 30 requests per minute, max 5 concurrent
streamer := NewRateLimitedStreamer(model, 30, 5)
prompts := []string{
"Explain AI",
"Explain ML",
"Explain DL",
"Explain NLP",
"Explain CV",
"Explain RL",
"Explain GANs",
"Explain transformers",
}
resultChan := make(chan string, len(prompts))
errorChan := make(chan error, len(prompts))
var wg sync.WaitGroup
start := time.Now()
// Start all streams
for i, prompt := range prompts {
wg.Add(1)
go func(index int, p string) {
defer wg.Done()
log.Printf("Starting stream %d", index)
streamer.Stream(ctx, p, resultChan, errorChan)
log.Printf("Completed stream %d", index)
}(i, prompt)
}
// Wait for completion
go func() {
wg.Wait()
close(resultChan)
close(errorChan)
}()
// Collect results
var results []string
var errors []error
for resultChan != nil || errorChan != nil {
select {
case result, ok := <-resultChan:
if !ok {
resultChan = nil
} else {
results = append(results, result)
}
case err, ok := <-errorChan:
if !ok {
errorChan = nil
} else {
errors = append(errors, err)
}
}
}
elapsed := time.Since(start)
fmt.Printf("\n=== Summary ===\n")
fmt.Printf("Total time: %v\n", elapsed)
fmt.Printf("Successful: %d\n", len(results))
fmt.Printf("Errors: %d\n", len(errors))
for i, result := range results {
fmt.Printf("\nResult %d:\n%s\n", i+1, result[:min(100, len(result))])
}
}
func min(a, b int) int {
if a < b {
return a
}
return b
}
Real-Time Stream Aggregation
Combine Multiple Streams in Real-Time
package main
import (
"context"
"fmt"
"log"
"os"
"sync"
"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)
type StreamChunk struct {
SourceID int
Text string
}
func aggregateStreams(ctx context.Context, model provider.LanguageModel, prompts []string) <-chan StreamChunk {
aggregated := make(chan StreamChunk, 100)
var wg sync.WaitGroup
for i, prompt := range prompts {
wg.Add(1)
go func(sourceID int, p string) {
defer wg.Done()
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: p,
})
if err != nil {
log.Printf("Stream %d error: %v", sourceID, err)
return
}
// Forward chunks to aggregated channel
for chunk := range stream.Chunks() {
select {
case aggregated <- StreamChunk{
SourceID: sourceID,
Text: chunk.Text,
}:
case <-ctx.Done():
return
}
}
if err := stream.Err(); err != nil {
log.Printf("Stream %d error: %v", sourceID, err)
}
}(i, prompt)
}
// Close aggregated channel when all streams complete
go func() {
wg.Wait()
close(aggregated)
}()
return aggregated
}
func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4o-mini")
prompts := []string{
"List 3 programming languages",
"List 3 database systems",
"List 3 cloud providers",
}
// Get aggregated stream
aggregated := aggregateStreams(ctx, model, prompts)
// Process chunks as they arrive from any stream
streamTexts := make(map[int]string)
for chunk := range aggregated {
streamTexts[chunk.SourceID] += chunk.Text
fmt.Printf("[Stream %d] %s", chunk.SourceID, chunk.Text)
}
fmt.Println("\n\n=== Final Results ===")
for id, text := range streamTexts {
fmt.Printf("\nStream %d:\n%s\n", id, text)
}
}
Error Handling Across Streams
Robust Error Handling
package main
import (
"context"
"fmt"
"log"
"os"
"sync"
"time"
"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)
type StreamStatus struct {
ID int
Prompt string
Text string
Error error
StartTime time.Time
EndTime time.Time
ChunkCount int
}
func streamWithMonitoring(ctx context.Context, id int, model provider.LanguageModel, prompt string) StreamStatus {
status := StreamStatus{
ID: id,
Prompt: prompt,
StartTime: time.Now(),
}
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: prompt,
})
if err != nil {
status.Error = err
status.EndTime = time.Now()
return status
}
for chunk := range stream.Chunks() {
status.Text += chunk.Text
status.ChunkCount++
// Check for context cancellation
select {
case <-ctx.Done():
status.Error = ctx.Err()
status.EndTime = time.Now()
return status
default:
}
}
status.Error = stream.Err()
status.EndTime = time.Now()
return status
}
func main() {
// Create context with timeout
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4o-mini")
prompts := []string{
"Explain quantum computing",
"Explain artificial intelligence",
"Explain blockchain technology",
"Explain machine learning",
"Explain neural networks",
}
statusChan := make(chan StreamStatus, len(prompts))
var wg sync.WaitGroup
// Start all streams with monitoring
for i, prompt := range prompts {
wg.Add(1)
go func(id int, p string) {
defer wg.Done()
status := streamWithMonitoring(ctx, id, model, p)
statusChan <- status
}(i, prompt)
}
// Wait and close
go func() {
wg.Wait()
close(statusChan)
}()
// Collect and analyze statuses
var statuses []StreamStatus
successCount := 0
errorCount := 0
for status := range statusChan {
statuses = append(statuses, status)
if status.Error != nil {
errorCount++
log.Printf("Stream %d FAILED: %v", status.ID, status.Error)
} else {
successCount++
duration := status.EndTime.Sub(status.StartTime)
log.Printf("Stream %d SUCCESS: %d chunks in %v", status.ID, status.ChunkCount, duration)
}
}
// Print summary
fmt.Println("\n=== Summary ===")
fmt.Printf("Total streams: %d\n", len(statuses))
fmt.Printf("Successful: %d\n", successCount)
fmt.Printf("Failed: %d\n", errorCount)
// Print detailed results
for _, status := range statuses {
duration := status.EndTime.Sub(status.StartTime)
fmt.Printf("\n[Stream %d] %v\n", status.ID, duration)
if status.Error != nil {
fmt.Printf(" Error: %v\n", status.Error)
} else {
fmt.Printf(" Chunks: %d\n", status.ChunkCount)
fmt.Printf(" Text length: %d\n", len(status.Text))
}
}
}
Best Practices
1. Always Use WaitGroups
var wg sync.WaitGroup
for _, prompt := range prompts {
wg.Add(1)
go func(p string) {
defer wg.Done()
// Process stream
}(prompt)
}
wg.Wait()
2. Limit Concurrency
sem := make(chan struct{}, maxConcurrent)
// Acquire: sem <- struct{}{}
// Release: <-sem
3. Use Context for Cancellation
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
4. Handle Errors Gracefully
if err := stream.Err(); err != nil {
log.Printf("Stream error: %v", err)
// Handle error, don't panic
}
5. Monitor Stream Health
log.Printf("Stream %d: %d chunks, %v elapsed", id, chunkCount, time.Since(start))