页面加载中,请稍候
来源:开源Linux发布时间:2023-09-141817浏览
询问 AI
作者:Goland猫
https://juejin.cn/post/7245919919223636023
对于大型的互联网应用程序,如电商平台、社交网络、金融交易平台等,每秒钟都会收到大量的请求。在这些应用程序中,需要使用高效的技术来应对高并发的请求,尤其是在短时间内处理大量的请求,如1分钟百万请求。
同时,为了降低用户的使用门槛和提升用户体验,前端需要实现参数的无感知传递。这样用户在使用时,无需担心参数传递的问题,能够轻松地享受应用程序的服务。
在处理1分钟百万请求时,需要使用高效的技术和算法,以提高请求的响应速度和处理能力。Go语言以其高效性和并发性而闻名,因此成为处理高并发请求的优秀选择。Go中有多种模式可供选择,如基于goroutine和channel的并发模型、使用池技术的协程模型等,以便根据具体应用的需要来选择适合的技术模式。
本文代码参考搬至
https://marksuper.xyz/2021/10/08/handle_million_req/http://marcio.io/2015/07/handling-1-million-requests-per-minute-with-golang/
W1 结构体类型,它有五个成员:
typeW1struct{
WgSend*sync.WaitGroup
Wg*sync.WaitGroup
MaxNumint
Chchanstring
DispatchStopchanstruct{}
}
接下来是 Dispatch 方法,它将任务发送到通道 Ch 中。它通过 for 循环来发送 10 倍于 MaxNum 的任务,每个任务都是一个 goroutine。defer 语句用于在任务完成时减少 WgSend 的计数。select 语句用于在任务分发被中止时退出任务发送。
func(w*W1)Dispatch(jobstring){
w.WgSend.Add(10*w.MaxNum)
fori:=0;i<10*w.MaxNum;i++{
gofunc(iint){
deferw.WgSend.Done()
select{
casew.Ch<-fmt.Sprintf("%d",i):
return
case<-w.DispatchStop:
fmt.Println("退出发送job:",fmt.Sprintf("%d",i))
return
}
}(i)
}
}
然后是StartPool方法,它创建了一个 goroutine 池来处理从通道 Ch 中读取到的任务。
如果通道 Ch 还没有被创建,那么它将被创建。如果计数器 WgSend 还没有被创建,那么它也将被创建。如果计数器 Wg 还没有被创建,那么它也将被创建。
如果通道 DispatchStop 还没有被创建,那么它也将被创建。
for 循环用于创建MaxNum 个 goroutine来处理从通道中读取到的任务。defer 语句用于在任务完成时减少 Wg 的计数。
func(w*W1)StartPool(){
ifw.Ch==nil{
w.Ch=make(chanstring,w.MaxNum)
}
ifw.WgSend==nil{
w.WgSend=&sync.WaitGroup{}
}
ifw.Wg==nil{
w.Wg=&sync.WaitGroup{}
}
ifw.DispatchStop==nil{
w.DispatchStop=make(chanstruct{})
}
w.Wg.Add(w.MaxNum)
fori:=0;i<w.MaxNum;i++{
gofunc(){
deferw.Wg.Done()
forv:=rangew.Ch{
fmt.Printf("完成工作:%s\n",v)
}
}()
}
}
最后是 Stop 方法,它停止任务分发并等待所有任务完成。
它关闭了通道 DispatchStop,等待 WgSend 中的任务发送 goroutine 完成,然后关闭通道 Ch,等待 Wg 中的任务处理 goroutine 完成。
func(w*W1)Stop(){
close(w.DispatchStop)
w.WgSend.Wait()
close(w.Ch)
w.Wg.Wait()
}

typeSubWorkerstruct{
JobChanchanstring
}
子协程,它有一个 JobChan,用于接收任务。
Run:SubWorker 的方法,用于启动一个子协程,从 JobChan 中读取任务并执行。
func(sw*SubWorker)Run(wg*sync.WaitGroup,poolChchanchanstring,quitChchanstruct{}){
ifsw.JobChan==nil{
sw.JobChan=make(chanstring)
}
wg.Add(1)
gofunc(){
deferwg.Done()
for{
poolCh<-sw.JobChan
select{
caseres:=<-sw.JobChan:
fmt.Printf("完成工作:%s\n",res)
case<-quitCh:
fmt.Printf("消费者结束......\n")
return
}
}
}()
}
typeW2struct{
SubWorkers[]SubWorker
Wg*sync.WaitGroup
MaxNumint
ChPoolchanchanstring
QuitChanchanstruct{}
}
Dispatch:W2 的方法,用于从 ChPool 中获取 TaskChan,将任务发送给一个 SubWorker 执行。
func(w*W2)Dispatch(jobstring){
jobChan:=<-w.ChPool
select{
casejobChan<-job:
fmt.Printf("发送任务:%s完成\n",job)
return
case<-w.QuitChan:
fmt.Printf("发送者(%s)结束\n",job)
return
}
}
StartPool:W2 的方法,用于初始化协程池,启动所有子协程并把 TaskChan 存储在 ChPool 中。
func(w*W2)StartPool(){
ifw.ChPool==nil{
w.ChPool=make(chanchanstring,w.MaxNum)
}
ifw.SubWorkers==nil{
w.SubWorkers=make([]SubWorker,w.MaxNum)
}
ifw.Wg==nil{
w.Wg=&sync.WaitGroup{}
}
fori:=0;i<len(w.SubWorkers);i++{
w.SubWorkers[i].Run(w.Wg,w.ChPool,w.QuitChan)
}
}
Stop:W2 的方法,用于停止协程的工作,并等待所有协程结束。
func(w*W2)Stop(){
close(w.QuitChan)
w.Wg.Wait()
close(w.ChPool)
}
DealW2 函数则是整个协程池的入口,它通过 NewWorker 方法创建一个 W2 实例,然后调用 StartPool 启动协程池,并通过 Dispatch 发送任务,最后调用 Stop 停止协程池。
funcDealW2(maxint){
w:=NewWorker(w2,max)
w.StartPool()
fori:=0;i<10*max;i++{
gow.Dispatch(fmt.Sprintf("%d",i))
}
w.Stop()
}

原来是golang里如果方法传递的不是地址,那么就会做一个拷贝,所以这里调用的wg根本就不是一个对象。
传递的地方传递地址就可以了,如果不传递地址,将会出现死锁
godoSomething(i,&wg,ch)
funcdoSomething(indexint,wg*sync.WaitGroup,chchanint){
在这段代码中,poolCh代表工作者池,sw.JobChan代表工作者的工作通道。当一个工作者完成了工作后,它会将工作结果发送到sw.JobChan,此时可以通过case res := <-sw.JobChan:来接收该工作的结果。
在这个代码块中,还需要处理一个退出信号quitCh。因此,第二个case <-quitCh:用于检测是否接收到了退出信号。如果接收到了退出信号,程序将打印出消息并结束。
需要注意的是,这两个case语句是互斥的,只有当工作者完成工作或收到退出信号时,才会进入其中一个语句。因此,这个循环可以保证在工作者完成工作或收到退出信号时退出。
需要读取两次sw.JobChan的原因是:第一次读取用于将工作者的工作通道放回工作者池中,这样其他工作者就可以使用该通道。第二次读取用于接收工作者的工作结果或退出信号。因此,这两次读取是为了确保能够在正确的时刻将工作者的工作通道放回工作者池中并正确地处理工作结果或退出信号。
根据w2的特点 我自己写了一个w2
import(
"fmt"
"sync"
)
typeSubWorkerNewstruct{
JobChanchanstring
}
typeW2Newstruct{
SubWorkers[]SubWorkerNew
Wg*sync.WaitGroup
MaxNumint
ChPoolchanchanstring
QuitChanchanstruct{}
}
funcNewW2(maxNumint)*W2New{
subWorkers:=make([]SubWorkerNew,maxNum)
fori:=0;i<maxNum;i++{
subWorkers[i]=SubWorkerNew{JobChan:make(chanstring)}
}
pool:=make(chanchanstring,maxNum)
fori:=0;i<maxNum;i++{
pool<-subWorkers[i].JobChan
}
return&W2New{
SubWorkers:subWorkers,
Wg:&sync.WaitGroup{},
MaxNum:maxNum,
ChPool:pool,
QuitChan:make(chanstruct{}),
}
}
func(w*W2New)Dispatch(jobstring){
select{
casejobChannel:=<-w.ChPool:
jobChannel<-job
default:
fmt.Println("Allworkersbusy")
}
}
func(w*W2New)StartPool(){
fori:=0;i<w.MaxNum;i++{
gofunc(subWorker*SubWorkerNew){
w.Wg.Add(1)
deferw.Wg.Done()
for{
select{
casejob:=<-subWorker.JobChan:
fmt.Println("processing",job)
case<-w.QuitChan:
return
}
}
}(&w.SubWorkers[i])
}
}
func(w*W2New)Stop(){
close(w.QuitChan)
w.Wg.Wait()
close(w.ChPool)
for_,subWorker:=rangew.SubWorkers{
close(subWorker.JobChan)
}
}
funcmain(){
w:=NewW2(5)
w.StartPool()
fori:=0;i<20;i++{
w.Dispatch(fmt.Sprintf("job%d",i))
}
w.Stop()
}
但是有几个点需要注意
1.没有考虑JobChan通道的缓冲区大小,如果有大量任务被并发分配,容易导致内存占用过高;
2.每个线程都会执行无限循环,此时线程退出的条件是接收到QuitChan通道的信号,可能导致线程的阻塞等问题;
3.Dispatch函数的默认情况下只会输出"All workers busy",而不是阻塞,这意味着当所有线程都处于忙碌状态时,任务会丢失
4.线程池启动后无法动态扩展或缩小。
这个优化版本改了很多次。有一些需要注意的点是,不然会一直死锁
1.使用sync.WaitGroup来确保线程池中所有线程都能够启动并运行;
2.在Stop函数中,先向SubWorker的JobChan中发送一个关闭信号,再等待所有SubWorker线程退出;
3.在Dispatch函数中,将默认情况下的输出改为阻塞等待可用通道;
packagehandle_million_requests
import(
"fmt"
"sync"
"time"
)
typeSubWorkerNewstruct{
Idint
JobChanchanstring
}
typeW2Newstruct{
SubWorkers[]SubWorkerNew
MaxNumint
ChPoolchanchanstring
QuitChanchanstruct{}
Wg*sync.WaitGroup
}
funcNewW2(maxNumint)*W2New{
chPool:=make(chanchanstring,maxNum)
subWorkers:=make([]SubWorkerNew,maxNum)
fori:=0;i<maxNum;i++{
subWorkers[i]=SubWorkerNew{Id:i,JobChan:make(chanstring)}
chPool<-subWorkers[i].JobChan
}
wg:=new(sync.WaitGroup)
wg.Add(maxNum)
return&W2New{
MaxNum:maxNum,
SubWorkers:subWorkers,
ChPool:chPool,
QuitChan:make(chanstruct{}),
Wg:wg,
}
}
func(w*W2New)StartPool(){
fori:=0;i<w.MaxNum;i++{
gofunc(wg*sync.WaitGroup,subWorker*SubWorkerNew){
deferwg.Done()
for{
select{
casejob:=<-subWorker.JobChan:
fmt.Printf("SubWorker%dprocessingjob%s\n",subWorker.Id,job)
time.Sleep(time.Second)//模拟任务处理过程
case<-w.QuitChan:
return
}
}
}(w.Wg,&w.SubWorkers[i])
}
}
func(w*W2New)Stop(){
close(w.QuitChan)
fori:=0;i<w.MaxNum;i++{
close(w.SubWorkers[i].JobChan)
}
w.Wg.Wait()
}
func(w*W2New)Dispatch(jobstring){
select{
casejobChan:=<-w.ChPool:
jobChan<-job
default:
fmt.Println("Allworkersbusy")
}
}
func(w*W2New)AddWorker(){
newWorker:=SubWorkerNew{Id:w.MaxNum,JobChan:make(chanstring)}
w.SubWorkers=append(w.SubWorkers,newWorker)
w.ChPool<-newWorker.JobChan
w.MaxNum++
w.Wg.Add(1)
gofunc(subWorker*SubWorkerNew){
deferw.Wg.Done()
for{
select{
casejob:=<-subWorker.JobChan:
fmt.Printf("SubWorker%dprocessingjob%s\n",subWorker.Id,job)
time.Sleep(time.Second)//模拟任务处理过程
case<-w.QuitChan:
return
}
}
}(&newWorker)
}
func(w*W2New)RemoveWorker(){
ifw.MaxNum>1{
worker:=w.SubWorkers[w.MaxNum-1]
close(worker.JobChan)
w.MaxNum--
w.SubWorkers=w.SubWorkers[:w.MaxNum]
}
}
AddWorker和RemoveWorker,用于动态扩展/缩小线程池。
funcTestW2New(t*testing.T){
pool:=NewW2(3)
pool.StartPool()
pool.Dispatch("task1")
pool.Dispatch("task2")
pool.Dispatch("task3")
pool.AddWorker()
pool.AddWorker()
pool.RemoveWorker()
pool.Stop()
}

当Dispatch函数向ChPool通道获取可用通道时,会从通道中取出一个SubWorker的JobChan通道,并将任务发送到该通道中。而对于SubWorker来说,并没有进行任务的使用次数限制,所以它可以处理多个任务。
在这个例子中,当任务数量比SubWorker数量多时,一个SubWorker的JobChan通道会接收到多个任务,它们会在SubWorker的循环中按顺序依次处理,直到JobChan中没有未处理的任务为止。因此,如果任务数量特别大,可能会导致某些SubWorker的JobChan通道暂时处于未处理任务状态,而其他的SubWorker在执行任务。
在测试结果中,最后三行中出现了多个"SubWorker 0 processing job",说明SubWorker 0的JobChan通道接收了多个任务,并且在其循环中处理这些任务。下面的代码片段显示了这个过程:
// SubWorker 0 的循环部分
for{
select{
casejob:=<-subWorker.JobChan:
fmt.Printf("SubWorker%dprocessingjob%s\n",subWorker.Id,job)
case<-w.QuitChan:
return
}
}

10T 技术资源大放送!包括但不限于:Linux、虚拟化、容器、云计算、网络、Python、Go 等。在开源Linux公众号内回复10T,即可免费获取!
有收获,点个在看
新闻来源:开源Linux,文中所述为作者独立观点,不代表icspec立场。更多精彩资讯请下载icspec App。如对本稿件有异议,请联系微信客服specltkj。
暂无评论哦,快来评论一下吧!
2026-07-02
2026-06-01