淺談Go語(yǔ)言中高效并發(fā)模式
一、并發(fā)編程基礎(chǔ)回顧
Go語(yǔ)言以其輕量級(jí)的Goroutine和簡(jiǎn)潔的并發(fā)控制機(jī)制(如Channel和Select)聞名,非常適合構(gòu)建高并發(fā)、高吞吐量的系統(tǒng)。
1.1 Goroutine
Go語(yǔ)言中的Goroutine是輕量級(jí)線程,由Go運(yùn)行時(shí)調(diào)度。啟動(dòng)一個(gè)Goroutine非常簡(jiǎn)單:
go func() {
fmt.Println("Hello from goroutine")
}()
1.2 Channel
Channel是Go語(yǔ)言中用于Goroutine之間通信的機(jī)制,支持同步與數(shù)據(jù)傳遞:
ch := make(chan int)
go func() { ch <- 42 }()
value := <-ch
fmt.Println(value)
1.3 Select
Select用于在多個(gè)Channel操作中進(jìn)行選擇,常用于超時(shí)控制、多路復(fù)用等場(chǎng)景:
select {
case msg := <-ch1:
fmt.Println("Got from ch1:", msg)
case msg := <-ch2:
fmt.Println("Got from ch2:", msg)
case <-time.After(1 * time.Second):
fmt.Println("Timeout")
}
二、常用并發(fā)模式詳解
2.1 Worker Pool(工作池模式)
適用場(chǎng)景:任務(wù)分發(fā)與并行處理,如爬蟲(chóng)、批量任務(wù)處理。
案例:使用固定數(shù)量的Worker處理任務(wù)隊(duì)列。
package main
import (
"fmt"
"sync"
"time"
)
func worker(id int, jobs <-chan int, results chan<- int, wg *sync.WaitGroup) {
defer wg.Done()
for j := range jobs {
fmt.Printf("Worker %d started job %d\n", id, j)
time.Sleep(time.Second) // 模擬耗時(shí)任務(wù)
fmt.Printf("Worker %d finished job %d\n", id, j)
results <- j * 2
}
}
func main() {
const numJobs = 5
const numWorkers = 3
jobs := make(chan int, numJobs)
results := make(chan int, numJobs)
var wg sync.WaitGroup
// 啟動(dòng)3個(gè)worker
for w := 1; w <= numWorkers; w++ {
wg.Add(1)
go worker(w, jobs, results, &wg)
}
// 發(fā)送5個(gè)任務(wù)
for j := 1; j <= numJobs; j++ {
jobs <- j
}
close(jobs)
// 等待所有worker完成
wg.Wait()
close(results)
// 收集結(jié)果
for result := range results {
fmt.Println("Result:", result)
}
}
2.2 Fan-in / Fan-out(扇入/扇出模式)
適用場(chǎng)景:將多個(gè)輸入合并到一個(gè)輸出(Fan-in),或?qū)⒁粋€(gè)輸入分發(fā)到多個(gè)處理單元(Fan-out)。
案例:
package main
import (
"fmt"
"sync"
"time"
)
func producer(id int, out chan<- int, wg *sync.WaitGroup) {
defer wg.Done()
for i := 0; i < 3; i++ {
fmt.Printf("Producer %d produced %d\n", id, i)
out <- id*10 + i
time.Sleep(200 * time.Millisecond)
}
}
func consumer(in <-chan int, done chan<- bool) {
for value := range in {
fmt.Printf("Consumed: %d\n", value)
}
done <- true
}
func main() {
ch := make(chan int)
done := make(chan bool)
var wg sync.WaitGroup
// 啟動(dòng)3個(gè)producer
for i := 1; i <= 3; i++ {
wg.Add(1)
go producer(i, ch, &wg)
}
// 啟動(dòng)1個(gè)consumer
go consumer(ch, done)
// 等待所有producer完成
go func() {
wg.Wait()
close(ch)
}()
// 等待consumer完成
<-done
fmt.Println("All done!")
}
2.3 Pipeline(流水線模式)
適用場(chǎng)景:數(shù)據(jù)流處理,分階段處理任務(wù),如ETL流程。
案例:模擬數(shù)據(jù)經(jīng)過(guò)多個(gè)階段的處理。
package main
import (
"fmt"
"sync"
)
func generator(nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
out <- n * n
}
}()
return out
}
func main() {
// 構(gòu)建pipeline: generator -> square -> print
input := generator(1, 2, 3, 4, 5)
squares := square(input)
// 打印結(jié)果
for result := range squares {
fmt.Println(result)
}
}
2.4 Timeout / Cancel(超時(shí)與取消模式)
適用場(chǎng)景:控制任務(wù)執(zhí)行時(shí)間,防止阻塞或資源浪費(fèi)。
案例:使用Context實(shí)現(xiàn)超時(shí)控制。
package main
import (
"context"
"fmt"
"time"
)
func longRunningTask(ctx context.Context) {
select {
case <-time.After(3 * time.Second):
fmt.Println("Task completed")
case <-ctx.Done():
fmt.Println("Task canceled:", ctx.Err())
}
}
func main() {
// 帶超時(shí)的context
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
go longRunningTask(ctx)
// 等待足夠長(zhǎng)時(shí)間觀察結(jié)果
time.Sleep(4 * time.Second)
}
2.5 Future / Promise(異步結(jié)果模式)
適用場(chǎng)景:異步執(zhí)行任務(wù)并獲取結(jié)果,常用于HTTP請(qǐng)求、數(shù)據(jù)庫(kù)查詢等。
案例:模擬異步任務(wù)并獲取結(jié)果。
package main
import (
"fmt"
"time"
)
type Future struct {
result chan int
}
func NewFuture() (*Future, func(int)) {
f := &Future{result: make(chan int)}
return f, func(r int) { f.result <- r }
}
func (f *Future) Get() int {
return <-f.result
}
func asyncCompute(complete func(int)) {
go func() {
time.Sleep(2 * time.Second)
complete(42)
}()
}
func main() {
future, complete := NewFuture()
asyncCompute(complete)
fmt.Println("Waiting for result...")
result := future.Get()
fmt.Println("Result:", result)
}
三、進(jìn)階并發(fā)模式
3.1 Or-Done 模式
適用場(chǎng)景:監(jiān)聽(tīng)多個(gè)Channel,任意一個(gè)完成即退出。
package main
import (
"fmt"
"time"
)
func orDone(channels ...<-chan interface{}) <-chan interface{} {
switch len(channels) {
case 0:
return nil
case 1:
return channels[0]
}
orDone := make(chan interface{})
go func() {
defer close(orDone)
switch len(channels) {
case 2:
select {
case <-channels[0]:
case <-channels[1]:
}
default:
select {
case <-channels[0]:
case <-channels[1]:
case <-channels[2]:
case <-orDone(channels[3:]...):
}
}
}()
return orDone
}
func main() {
sig := make(chan interface{})
time.AfterFunc(2*time.Second, func() { close(sig) })
start := time.Now()
<-orDone(sig, time.After(5*time.Second))
fmt.Printf("Done after %v\n", time.Since(start))
}
3.2 Rate Limiting(限流模式)
適用場(chǎng)景:控制請(qǐng)求頻率,防止系統(tǒng)過(guò)載。
package main
import (
"fmt"
"time"
)
func main() {
requests := make(chan int, 5)
for i := 1; i <= 5; i++ {
requests <- i
}
close(requests)
// 限制每200ms處理一個(gè)請(qǐng)求
limiter := time.Tick(200 * time.Millisecond)
for req := range requests {
<-limiter
fmt.Println("Request", req, time.Now())
}
// 突發(fā)限制模式
burstyLimiter := make(chan time.Time, 3)
for i := 0; i < 3; i++ {
burstyLimiter <- time.Now()
}
go func() {
for t := range time.Tick(200 * time.Millisecond) {
burstyLimiter <- t
}
}()
burstyRequests := make(chan int, 5)
for i := 1; i <= 5; i++ {
burstyRequests <- i
}
close(burstyRequests)
for req := range burstyRequests {
<-burstyLimiter
fmt.Println("Bursty request", req, time.Now())
}
}
四、注意點(diǎn)
4.1 使用說(shuō)明
- 每個(gè)示例都是獨(dú)立的Go程序,可以直接保存為
.go文件運(yùn)行 - 所有示例都包含了必要的包導(dǎo)入和錯(cuò)誤處理
- 每個(gè)示例都有詳細(xì)的注釋說(shuō)明關(guān)鍵步驟
- 可以通過(guò)調(diào)整參數(shù)(如worker數(shù)量、超時(shí)時(shí)間等)觀察不同行為
4.2 模式選擇對(duì)比
| 模式 | 適用場(chǎng)景 | 優(yōu)點(diǎn) | 缺點(diǎn) |
|---|---|---|---|
| Worker Pool | 任務(wù)并行處理 | 控制并發(fā)數(shù),防止資源耗盡 | 需要任務(wù)隊(duì)列管理 |
| Fan-in/Fan-out | 數(shù)據(jù)合并/分發(fā) | 提高吞吐量 | 需注意Channel關(guān)閉時(shí)機(jī) |
| Pipeline | 流式處理 | 模塊化、易擴(kuò)展 | 階段多時(shí)延遲高 |
| Timeout/Cancel | 任務(wù)超時(shí)控制 | 避免阻塞 | 需正確使用Context |
| Future/Promise | 異步結(jié)果獲取 | 非阻塞編程 | 結(jié)果獲取需同步 |
| Or-Done | 多事件監(jiān)聽(tīng) | 靈活組合 | 實(shí)現(xiàn)較復(fù)雜 |
| Rate Limiting | 流量控制 | 防止過(guò)載 | 需合理設(shè)置閾值 |
4.3 選擇建議
- 避免共享內(nèi)存,優(yōu)先使用Channel通信
- 使用
sync.WaitGroup等待Goroutine完成 - 合理使用
context管理生命周期 - 注意Channel的關(guān)閉,避免panic
- 使用
select+default實(shí)現(xiàn)非阻塞操作 - 控制Goroutine數(shù)量,防止資源泄漏
到此這篇關(guān)于淺談Go語(yǔ)言中高效并發(fā)模式 的文章就介紹到這了,更多相關(guān)Go語(yǔ)言高效并發(fā)模式內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Go語(yǔ)言基于viper的conf庫(kù)進(jìn)行配置文件解析
在現(xiàn)代軟件開(kāi)發(fā)中,配置文件是不可或缺的一部分,如何高效地將這些格式解析到 Go 結(jié)構(gòu)體中,一直是開(kāi)發(fā)者的痛點(diǎn),下面我們來(lái)看看如何使用conf進(jìn)行配置文件解析吧2025-03-03
Golang學(xué)習(xí)筆記之延遲函數(shù)(defer)的使用小結(jié)
這篇文章主要介紹了Golang學(xué)習(xí)筆記之延遲函數(shù)(defer),小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧2018-12-12
go語(yǔ)言實(shí)現(xiàn)簡(jiǎn)易比特幣系統(tǒng)錢(qián)包的原理解析
這篇文章主要介紹了go語(yǔ)言實(shí)現(xiàn)簡(jiǎn)易比特幣系統(tǒng)錢(qián)包的原理解析,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2021-04-04
Golang設(shè)計(jì)模式中抽象工廠模式詳細(xì)講解
抽象工廠模式用于生成產(chǎn)品族的工廠,所生成的對(duì)象是有關(guān)聯(lián)的。如果抽象工廠退化成生成的對(duì)象無(wú)關(guān)聯(lián)則成為工廠函數(shù)模式。比如本例子中使用RDB和XML存儲(chǔ)訂單信息,抽象工廠分別能生成相關(guān)的主訂單信息和訂單詳情信息2023-01-01
基于Go實(shí)現(xiàn)TCP長(zhǎng)連接上的請(qǐng)求數(shù)控制
在服務(wù)端開(kāi)啟長(zhǎng)連接的情況下,四層負(fù)載均衡轉(zhuǎn)發(fā)請(qǐng)求時(shí),會(huì)出現(xiàn)服務(wù)端收到的請(qǐng)求qps不均勻的情況或是服務(wù)器無(wú)法接受到請(qǐng)求,因此需要服務(wù)端定期主動(dòng)斷開(kāi)一些長(zhǎng)連接,所以本文給大家介紹了基于Go實(shí)現(xiàn)TCP長(zhǎng)連接上的請(qǐng)求數(shù)控制,需要的朋友可以參考下2024-05-05
golang中defer延遲機(jī)制的實(shí)現(xiàn)示例
本文講解了golang特殊的defer延遲機(jī)制,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2025-09-09

