Fan-out/Fan-in并发模式的应用场景
Go语言的goroutine和channel天然适合构建流水线式并发处理架构。Fan-out/Fan-in是最常用的并发模式之一:Fan-out将一个channel的数据分发到多个goroutine并行处理,Fan-in将多个goroutine的输出合并到一个channel。这种模式在日志分析、ETL数据处理、爬虫抓取等IO密集型场景中极为常见。
Fan-out/Fan-in的核心价值是:当单个处理步骤成为瓶颈时,通过水平扩展该步骤的并发度来提升吞吐量,而不需要改变上下游的处理逻辑。上游生产者无感知地发送数据,下游消费者无感知地接收处理结果。
基础Fan-out实现:多Worker并行消费
以下实现将输入channel分发给N个Worker并行处理:
package main
import (
"fmt"
"sync"
)
// Fan-out: 将输入分发到多个Worker
func fanOut(input <-chan int, workerCount int) []<-chan int {
workers := make([]<-chan int, workerCount)
for i := 0; i < workerCount; i++ {
workers[i] = worker(input, i)
}
return workers
}
func worker(input <-chan int, id int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for v := range input {
result := v * v
fmt.Printf("worker %d processed %d -> %d\n", id, v, result)
out <- result
}
}()
return out
}
关键设计:每个Worker从同一个input channel读取,Go的channel本身就是并发安全的,多个goroutine同时从一个channel接收数据时,每条数据只会被一个goroutine获取,天然实现了负载均衡。
Fan-in实现:多路合并器
将多个Worker的输出channel合并为单一输出:
// Fan-in: 合并多个channel到一个
func fanIn(channels ...<-chan int) <-chan int {
merged := make(chan int)
var wg sync.WaitGroup
for _, ch := range channels {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for v := range c {
merged <- v
}
}(ch)
}
go func() {
wg.Wait()
close(merged)
}()
return merged
}
Fan-in使用了sync.WaitGroup来追踪所有输入channel的关闭状态。当所有Worker完成处理并关闭各自的输出channel后,合并channel才会关闭。这种模式确保下游不会丢失数据。
完整的Worker Pool管道架构
将Fan-out/Fan-in结合,构建一个可配置并发度的完整处理管道:
package main
import (
"fmt"
"sync"
"time"
)
func generator(items ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, item := range items {
out <- item
time.Sleep(100 * time.Millisecond)
}
}()
return out
}
func processor(in <-chan int, id int) <-chan string {
out := make(chan string)
go func() {
defer close(out)
for v := range in {
start := time.Now()
time.Sleep(time.Duration(200+v%50) * time.Millisecond)
result := fmt.Sprintf("[worker-%d] input=%d result=%d latency=%v",
id, v, v*v, time.Since(start))
out <- result
}
}()
return out
}
func workerPool(input <-chan int, poolSize int) <-chan string {
workers := make([]<-chan string, poolSize)
for i := 0; i < poolSize; i++ {
workers[i] = processor(input, i)
}
return fanInStr(workers...)
}
func fanInStr(channels ...<-chan string) <-chan string {
merged := make(chan string)
var wg sync.WaitGroup
for _, ch := range channels {
wg.Add(1)
go func(c <-chan string) {
defer wg.Done()
for v := range c {
merged <- v
}
}(ch)
}
go func() {
wg.Wait()
close(merged)
}()
return merged
}
func main() {
input := generator(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
results := workerPool(input, 4)
for r := range results {
fmt.Println(r)
}
}
背压控制与速率限制
当生产速度超过消费速度时,需要引入背压机制防止内存无限增长。Go中实现背压最直接的方式是使用带缓冲channel作为限流器:
// 速率限制器:令牌桶模式
func rateLimiter(input <-chan int, rate time.Duration) <-chan int {
out := make(chan int)
go func() {
defer close(out)
ticker := time.NewTicker(rate)
defer ticker.Stop()
for v := range input {
<-ticker.C
out <- v
}
}()
return out
}
// 并发度限制:信号量模式
func boundedWorker(input <-chan int, maxConcurrent int) <-chan string {
out := make(chan string)
sem := make(chan struct{}, maxConcurrent)
var wg sync.WaitGroup
go func() {
for v := range input {
sem <- struct{}{}
wg.Add(1)
go func(val int) {
defer func() {
<-sem
wg.Done()
}()
result := process(val)
out <- result
}(v)
}
wg.Wait()
close(out)
}()
return out
}
令牌桶模式控制处理频率,信号量模式控制瞬时并发上限。两者组合使用可以精确控制系统的负载特征。信号量的缓冲大小即为最大并发数,超出时新任务会阻塞等待,天然形成背压。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-bing-fa-mo-shi-fanoutfanin-guan-dao-jia-gou-yu/