package main
import (
"fmt""sync""time"
)
func worker(id int, jobs <-chan int, results chan<- int, wg *sync.WaitGroup) {
defer wg.Done()
for job := range jobs {
fmt.Printf("worker %d 处理任务 %d\n", id, job)
time.Sleep(100 * time.Millisecond)
results <- job * 2
}
}
func main() {
jobs := make(chan int, 100)
results := make(chan int, 100)
var wg sync.WaitGroup
for w := 1; w <= 5; w++ {
wg.Add(1)
go worker(w, jobs, results, &wg)
}
for j := 1; j <= 20; j++ {
jobs <- j
}
close(jobs)
wg.Wait()
close(results)
for r := range results {
fmt.Println("结果:", r)
}
}
二、支持优雅退出
Go
func workerPool(ctx context.Context, jobs <-chan Job, workers int) {
var wg sync.WaitGroup
for i := 0; i < workers; i++ {
wg.Add(1)
gofunc(id int) {
defer wg.Done()
for {
select {
case <-ctx.Done():
fmt.Printf("worker %d 收到退出信号\n", id)
returncase job, ok := <-jobs:
if !ok {
return
}
process(job)
}
}
}(i)
}
wg.Wait()
}
三、Pipeline 模式
Go
func gen(nums ...int) <-chan int {
out := make(chan int)
gofunc() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
func square(in <-chan int) <-chan int {
out := make(chan int)
gofunc() {
defer close(out)
for n := range in {
out <- n * n
}
}()
return out
}
func main() {
for n := range square(gen(1, 2, 3, 4, 5)) {
fmt.Println(n)
}
}
评论(0)
还没有评论,来说两句吧