package main import ( "fmt" "strconv" "strings" "sync" "time" ) // Stage 1: Generates a stream of numbers as strings. func generateNumbers(count int) <-chan string { out := make(chan string) go func() { defer close(out) for i := 1; i <= count; i++ { out <- strconv.Itoa(i) time.Sleep(50 * time.Millisecond) // Simulate generation time } fmt.Println("Generator: Finished generating numbers.") }() return out } // Stage 2: Transforms strings to uppercase and adds a prefix. func transformStrings(in <-chan string) <-chan string { out := make(chan string) go func() { defer close(out) for s := range in { transformed := "PROCESSED_" + strings.ToUpper(s) out <- transformed time.Sleep(30 * time.Millisecond) // Simulate transformation time } fmt.Println("Transformer: Finished transforming strings.") }() return out } // Stage 3: Parses the transformed string, converts to int, and squares it. func squareNumbers(in <-chan string) <-chan int { out := make(chan int) go func() { defer close(out) for s := range in { // Extract number after "PROCESSED_" parts := strings.Split(s, "_") if len(parts) < 2 { fmt.Printf("Squarer: Invalid format received: %s ", s) continue } numStr := parts[len(parts)-1] num, err := strconv.Atoi(numStr) if err != nil { fmt.Printf("Squarer: Could not parse number from '%s': %v ", numStr, err) continue } out <- num * num time.Sleep(20 * time.Millisecond) // Simulate squaring time } fmt.Println("Squarer: Finished squaring numbers.") }() return out } // Final Stage: Collects and prints the results. func collectResults(in <-chan int, wg *sync.WaitGroup) { defer wg.Done() fmt.Println("Collector: Starting to collect results.") for result := range in { fmt.Printf("Collector: Final Result: %d ", result) } fmt.Println("Collector: Finished collecting results.") } func main() { var wg sync.WaitGroup // Stage 1 numbers := generateNumbers(5) // Stage 2 transformed := transformStrings(numbers) // Stage 3 squared := squareNumbers(transformed) // Final Stage - consumer wg.Add(1) go collectResults(squared, &wg) // Wait for the collector to finish wg.Wait() fmt.Println("Main: Pipeline execution complete.") }