-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathpipeline.go
67 lines (53 loc) · 1.25 KB
/
pipeline.go
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
package main
import "fmt"
func main() {
input := []int{0, 1, 2, 3, 4, 5, 6, 7, 8, 9}
inputChannel := make(chan int)
doneChannel := make(chan struct{})
defer close(doneChannel)
go func() {
defer close(inputChannel)
for _, data := range input {
select {
case <-doneChannel: // trigger done
return
case inputChannel <- data:
}
}
}()
var index int
for result := range stepTwo(doneChannel, stepOne(doneChannel, inputChannel)) {
index++
fmt.Println(fmt.Sprintf("run %d: %d", index, result))
}
}
func stepOne(doneChannel chan struct{}, inputChannel chan int) chan int {
outputChannel := make(chan int)
go func() {
defer close(outputChannel)
for data := range inputChannel {
result := data + 1 // actual operation
select {
case <-doneChannel: // shutting down step one
return
case outputChannel <- result:
}
}
}()
return outputChannel
}
func stepTwo(doneChannel chan struct{}, inputChannel chan int) chan int {
outputChannel := make(chan int)
go func() {
defer close(outputChannel)
for data := range inputChannel {
result := data - 1 // actual operation
select {
case <-doneChannel: // shutting down step two
return
case outputChannel <- result:
}
}
}()
return outputChannel
}