Merging Data Streams from Multiple Goroutines (Fan-in Pattern)
Owner: SnippetBot
Created: 2026-09-09 00:00:50
Size: 1.42 KB
Expires: Never
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
package main
import (
"fmt"
"sync"
"time"
)
// generateNumbers sends numbers from start to start+count to the output channel.
func generateNumbers(name string, start, count int, out chan<- int) {
for i := 0; i < count; i++ {
num := start + i
fmt.Printf("%s: Sending %d
", name, num)
out <- num
time.Sleep(time.Millisecond * time.Duration(100+num%50))
}
}
// fanIn merges multiple input channels into a single output channel.
func fanIn(done <-chan struct{}, channels ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
output := func(c <-chan int) {
defer wg.Done()
for n := range c {
select {
case out <- n:
case <-done:
return
}
}
}
wg.Add(len(channels))
for _, c := range channels {
go output(c)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
gen1 := make(chan int)
gen2 := make(chan int)
gen3 := make(chan int)
done := make(chan struct{}) // Signal to stop fanIn if needed
go func() { defer close(gen1); generateNumbers("Gen1", 0, 5, gen1) }()
go func() { defer close(gen2); generateNumbers("Gen2", 100, 4, gen2) }()
go func() { defer close(gen3); generateNumbers("Gen3", 200, 6, gen3) }()
// Merge the output of all generators into a single channel
merged := fanIn(done, gen1, gen2, gen3)
fmt.Println("
--- Merged Results ---")
for num := range merged {
fmt.Printf("Received: %d
", num)
}
fmt.Println("All data merged and processed.")
}