Go语言并发模式fan-out fan-in管道架构与Worker Pool实现

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/

(0)
小编小编
上一篇 8小时前
下一篇 8小时前

相关推荐