页面加载中,请稍候
来源:嵌入式应用研究院发布时间:2023-05-243094浏览
询问 AI点击上方“嵌入式应用研究院”,选择“置顶/星标公众号”
干货福利,第一时间送达!来源 | golang编程笔记
整理&排版| 嵌入式应用研究院
1 什么是协程协程是与其他函数或方法一起并发运行的函数或方法。Go协程可以看作是轻量级线程,与线程相比,创建一个 Go 协程的成本很小。
1.1 协程与线程的对比调用函数或者方法时,在前面加上关键字go,可以让一个新的Go协程并发地运行。需要注意:
使用示例如下:
packagemain
import(
"fmt"
"time"
)
funcnumbers(){
fori:=1;i<=5;i++{
time.Sleep(250*time.Millisecond)
fmt.Printf("%d",i)
}
}
funcalphabets(){
fori:='a';i<='e';i++{
time.Sleep(400*time.Millisecond)
fmt.Printf("%c",i)
}
}
funcmain(){
gonumbers()//启动协程
goalphabets()//启动协程
//等待子协程允许完毕,后面介绍更高级的信道方式,这里就简单的等待
time.Sleep(3000*time.Millisecond)
fmt.Println("mainterminated")
}
//输出:1a23b4c5demainterminated
下图可以清晰的看到三个协程的运行关系:![]()
信道可以想像成协程之间通信的管道。如同管道中的水会从一端流到另一端,通过使用信道,数据也可以从一端发送,在另一端接收。所有信道都关联了一个类型。信道只能运输这种类型的数据,而运输其他类型的数据都是非法的。chan T表示T类型的信道,使用make函数进行初始化。例如:
a:=make(chanint)2.2 信道的收发
信道旁的箭头方向指定了是发送数据还是接收数据
data:=<-a//读取信道a,保存值到data
a<-data//写入信道a
发送与接收默认是阻塞的。当把数据发送到信道时,程序控制会在发送数据的语句处发生阻塞,直到有其它协程从信道读取到数据,才会解除阻塞。与此类似,当读取信道的数据时,如果没有其它的协程把数据写入到这个信道,那么读取过程就会一直阻塞着。**信道的这种特性能够帮助Go协程之间进行高效的通信,不需要用到其他编程语言常见的显式锁或条件变量。 借助阻塞这个特性,我们可以用一个读操作等待子协程结束,而不是使用sleep:
funchello(donechanbool){2.3 小心死锁
fmt.Println("Helloworldgoroutine")
done<-true//子协程结束,写入数据
}
funcmain(){
done:=make(chanbool)//创建bool信道
gohello(done)
<-done//读操作,一直阻塞直到子协程结束
fmt.Println("mainfunction")
}
使用信道需要考虑的一个重点是死锁。
数据发送方可以关闭信道,通知接收方这个信道不再有数据发送过来。当从信道接收数据时,接收方可以多用一个变量来检查信道是否已经关闭。
funcproducer(chnlchanint){
fori:=0;i<10;i++{
chnl<-i
}
close(chnl)//关闭信道
}
funcmain(){
ch:=make(chanint)
goproducer(ch)
for{
v,ok:=<-ch//判断信道是否关闭
ifok==false{
break
}
fmt.Println("Received",v,ok)
}
}
上面的语句里,如果成功接收信道所发送的数据,那么 ok 等于 true。而如果 ok 等于 false,说明我们试图读取一个关闭的通道。从关闭的信道读取到的值会是该信道类型的零值。 或者我们可以用range遍历信道,代替上面示例中的for循环:
funcmain(){2.5 缓冲信道
ch:=make(chanint)
goproducer(ch)
forv:=rangech{//range可以在信道关闭后自动结束,不用显示的判断
fmt.Println("Received",v)
}
}
上面无缓冲信道的发送和接收过程是阻塞的,读写操作会一直阻塞。我们还可以创建一个有缓冲的信道(Buffered Channel)。只在缓冲已满的情况,才会阻塞向缓冲信道发送数据。同样,只有在缓冲为空的时候,才会阻塞从缓冲信道接收数据。 通过向 make 函数时再传递一个表示容量的参数(指定缓冲的大小,sizeof(type) * capacity),就可以创建缓冲信道。
ch:=make(chantype,capacity)//capacity应该大于0。无缓冲信道的容量默认为0
缓冲区容量和长度的区别:
使用示例如下:
funcwrite(chchanint){2.6 select
fori:=0;i<5;i++{
ch<-i//写入两个值之后缓冲区满,阻塞等待缓冲区空闲
fmt.Println("successfullywrote",i,"toch")
}
close(ch)
}
funcmain(){
ch:=make(chanint,2)//缓冲大小为2
gowrite(ch)
time.Sleep(2*time.Second)
forv:=rangech{
fmt.Println("readvalue",v,"fromch")
time.Sleep(2*time.Second)
}
}
select 语句用于在多个发送/接收信道操作中进行选择。该语法与 switch 类似,所不同的是,这里的每个 case 语句都是信道操作。
使用示例:
funcserver1(chchanstring){3 WaitGroup3.1 WaitGroup的使用
time.Sleep(6*time.Second)
ch<-"fromserver1"
}
funcserver2(chchanstring){
time.Sleep(3*time.Second)
ch<-"fromserver2"
}
funcmain(){
output1:=make(chanstring)
output2:=make(chanstring)
goserver1(output1)
goserver2(output2)
select{//一直阻塞,直到某个信道可用
cases1:=<-output1:
fmt.Println(s1)
cases2:=<-output2:
fmt.Println(s2)
}
}
WaitGroup可以用来等待一批go协程执行结束,类似于C++的std::Thread::join。使用示例如下:
import(3.2 实现一个协程池
"fmt"
"sync"
"time"
)
funcprocess(iint,wg*sync.WaitGroup){//waitgroup参数指针,因为要修改内部的值,不能是值传递
fmt.Println("startedGoroutine",i)
time.Sleep(2*time.Second)
fmt.Printf("Goroutine%dended\n",i)
wg.Done()//子协程结束,调用done减少计数器
}
funcmain(){
no:=3
varwgsync.WaitGroup//定义waitgroup
fori:=0;i<no;i++{
wg.Add(1)//+1,增加计数器
goprocess(i,&wg)
}
wg.Wait()//等待协程结束,计数器为0
fmt.Println("Allgoroutinesfinishedexecuting")
}
基本思路:
代码和解析如下:
packagemain4 协程的同步手段4.1 互斥与Mutex
import(
"fmt"
"math/rand"
"sync"
"time"
)
//定义任务和结果两个结构体
typeJobstruct{
idint
randomnoint
}
typeResultstruct{
jobJob//包含job结构体
sumofdigitsint
}
//创建任务和结果的两个缓冲信道
varjobs=make(chanJob,10)
varresults=make(chanResult,10)
//计算一个整数每一位相加的和
funcdigits(numberint)int{
sum:=0
no:=number
forno!=0{
digit:=no%10
sum+=digit
no/=10
}
time.Sleep(2*time.Second)
returnsum
}
//遍历job信道,计算后每个job的数字并将结果写入reslut信道
funcworker(wg*sync.WaitGroup){
forjob:=rangejobs{
output:=Result{job,digits(job.randomno)}
results<-output
}
wg.Done()
}
//初始化waitgroup,并开启多个协程开始计算
funccreateWorkerPool(noOfWorkersint){
varwgsync.WaitGroup
fori:=0;i<noOfWorkers;i++{
wg.Add(1)
goworker(&wg)
}
wg.Wait()//等待协程池中协程都执行完毕
close(results)
}
//创建job,添加到job信道中
funcallocate(noOfJobsint){
fori:=0;i<noOfJobs;i++{
randomno:=rand.Intn(999)
job:=Job{i,randomno}
jobs<-job
}
close(jobs)
}
//打印所有计算结果
funcresult(donechanbool){
forresult:=rangeresults{
fmt.Printf("Jobid%d,inputrandomno%d,sumofdigits%d\n",result.job.id,result.job.randomno,result.sumofdigits)
}
done<-true
}
funcmain(){
startTime:=time.Now()
noOfJobs:=100
goallocate(noOfJobs)
done:=make(chanbool)
goresult(done)//这里会一直读取result信道,直到信道关闭
noOfWorkers:=10
createWorkerPool(noOfWorkers)
<-done//等待result打印完毕
endTime:=time.Now()
diff:=endTime.Sub(startTime)
fmt.Println("totaltimetaken",diff.Seconds(),"seconds")
}
Mutex用于提供一种加锁机制(Locking Mechanism),可确保在某时刻只有一个协程在临界区运行,以防止出现竞态条件。Mutex可以在sync包内找到。Mutex 定义了两个方法:Lock和Unlock。所有在 Lock 和 Unlock 之间的代码,都只能由一个Go协程执行,于是就可以避免竞态条件。
mutex.Lock()
x=x+1
mutex.Unlock()
使用示例:
//互斥锁保证线程同步4.2 原子操作
packagemain
import(
"fmt"
"sync"
)
vartotalstruct{//全局的结构体变量
sync.Mutex//互斥锁
valueint
}
funcworker(wg*sync.WaitGroup){
deferwg.Done()
fori:=0;i<=100;i++{
total.Lock()//加锁
total.value++
total.Unlock()//解锁
}
}
funcmain(){
varwgsync.WaitGroup
wg.Add(2)
goworker(&wg)
goworker(&wg)
wg.Wait()
fmt.Println(total.value)
}
用互斥锁来保护一个数值型的共享资源,麻烦且效率低下。标准库的sync/atomic包对原子操作提供了丰富的支持:sync/atomic包对基本的数值类型及复杂对象的读写都提供了原子操作的支持。atomic.Value原子对象提供了Load和Store两个原子方法,分别用于加载和保存数据,返回值和参数都是interface{}类型。
//原子操作实现线程同步4.3 阻塞信道
packagemain
import(
"fmt"
"sync"
"sync/atomic"
)
vartotaluint64
funcworker(wg*sync.WaitGroup){
deferwg.Done()
variuint64
fori=0;i<=100;i++{
atomic.AddUint64(&total,1)//原子操作,线程安全的
}
}
funcmain(){
varwgsync.WaitGroup
wg.Add(2)
goworker(&wg)
goworker(&wg)
wg.Wait()
fmt.Println(atomic.LoadUint64(&total))//读取值
}
上面的示例我们也可以用信道来实现互斥(还是推荐实际中使用Mutex),使用大小为1的缓冲信道可以导致可写阻塞,这样其他协程就不能继续执行,只能等待阻塞结束。在并发编程中,对共享资源的正确访问需要精确的控制,在目前的绝大多数语言中,都是通过加锁等线程同步方案来解决这一困难问题,而Go语言却另辟蹊径,它将共享的值通过Channel传递(实际上多个独立执行的线程很少主动共享资源)。在任意给定的时刻,最好只有一个Goroutine能够拥有该资源。
//使用channel实现线程同步
packagemain
import(
"fmt"
"sync"
)
vartotaluint64
funcworker(wg*sync.WaitGroup,chchanbool){
deferwg.Done()
variuint64
fori=0;i<=100;i++{
ch<-true//信道被写入值,其他协程到这一句也想写入值,就会阻塞等待信道可写
total++
<-ch//本协程读取信道,信道空了,其他协程可以写入了
}
}
funcmain(){
ch:=make(chanbool,1)//创建大小为1的chan
varwgsync.WaitGroup
wg.Add(2)
goworker(&wg,ch)
goworker(&wg,ch)
wg.Wait()
fmt.Println(total)//读取值
}
不仅如此,我们还可以通过设置chan的缓存大小来控制最大并发数。
5 常见并发模型5.1 生产者消费者模型通过平衡生产线程和消费线程的工作能力来提高程序的整体处理数据的速度。简单地说,就是生产者生产一些数据,然后放到成果队列中,同时消费者从成果队列中来取这些数据。这样就让生产消费变成了异步的两个过程。当成果队列中没有数据时,消费者就进入饥饿的等待中;而当成果队列中数据已满时,生产者则面临因产品挤压导致CPU被剥夺的下岗问题。 Go可以使用带缓冲区的chan作为成功队列,由不同的协程负责接入和读取,很简单的实现生产者消费者模型:
packagemain5.2 发布订阅模型
import(
"fmt"
"os"
"os/signal"
"syscall"
)
//生产者:生成factor整数倍的序列
funcProducer(factorint,outchan<-int){
fori:=0;;i++{
out<-i*factor//往信道缓冲区内写入数据
}
}
//消费者
funcConsumer(in<-chanint){
forv:=rangein{
fmt.Println(v)//从信道读取数据打印
}
}
funcmain(){
ch:=make(chanint,64)//成果队列,大小为64
//开启了2个Producer生产流水线,分别用于生成3和5的倍数的序列
//然后开启1个Consumer消费者线程,打印获取的结果
goProducer(3,ch)//生成3的倍数的序列
goProducer(5,ch)//生成5的倍数的序列
goConsumer(ch)//消费生成的队列
//Ctrl+C退出
sig:=make(chanos.Signal,1)
signal.Notify(sig,syscall.SIGINT,syscall.SIGTERM)
fmt.Printf("quit(%v)\n",<-sig)
}
发布订阅(publish/subscribe)模型通常被简写为pub/sub模型。在这个模型中,消息生产者成为发布者(publisher),而消息消费者则成为订阅者(subscriber),生产者和消费者是M:N的关系。在传统生产者和消费者模型中,是将消息发送到一个队列中,而发布订阅模型则是将消息发布给一个主题。在发布订阅模型中,每条消息都会传送给多个订阅者。发布者通常不会知道、也不关心哪一个订阅者正在接收主题消息。订阅者和发布者可以在运行时动态添加,是一种松散的耦合关系,这使得系统的复杂性可以随时间的推移而增长。在现实生活中,像天气预报之类的应用就可以应用这个并发模式。 示例代码如下:
//发布订阅模型实现
packagepubsub
import(
"sync"
"time"
)
type(
subscriberchaninterface{}//订阅者为一个管道
topicFuncfunc(vinterface{})bool//主题为一个过滤器
)
//发布者对象
typePublisherstruct{
msync.RWMutex//读写锁,保护订阅者map
bufferint//订阅队列的缓存大小
timeouttime.Duration//发布超时时间
subscribersmap[subscriber]topicFunc//订阅者信息
}
//构建一个发布者对象,可以设置发布超时时间和缓存队列的长度
funcNewPublisher(publishTimeouttime.Duration,bufferint)*Publisher{
return&Publisher{//返回对象指针
buffer:buffer,
timeout:publishTimeout,
subscribers:make(map[subscriber]topicFunc),//创建订阅者map
}
}
//添加一个新的订阅者,订阅全部主题
func(p*Publisher)Subscribe()chaninterface{}{
returnp.SubscribeTopic(nil)
}
//添加一个新的订阅者,订阅过滤器筛选后的主题
func(p*Publisher)SubscribeTopic(topictopicFunc)chaninterface{}{
ch:=make(chaninterface{},p.buffer)
p.m.Lock()
p.subscribers[ch]=topic
p.m.Unlock()
returnch
}
//退出订阅
func(p*Publisher)Evict(subchaninterface{}){
p.m.Lock()
deferp.m.Unlock()//函数退出时解锁
delete(p.subscribers,sub)//根据key删除map中一项
close(sub)//关闭chan
}
//发布一个主题
func(p*Publisher)Publish(vinterface{}){
p.m.RLock()
deferp.m.RUnlock()
varwgsync.WaitGroup
forsub,topic:=rangep.subscribers{
wg.Add(1)
gop.sendTopic(sub,topic,v,&wg)
}
wg.Wait()
}
//关闭发布者对象,同时关闭所有的订阅者管道。
func(p*Publisher)Close(){
p.m.Lock()
deferp.m.Unlock()
forsub:=rangep.subscribers{
delete(p.subscribers,sub)
close(sub)
}
}
//发送主题,可以容忍一定的超时
func(p*Publisher)sendTopic(subsubscriber,topictopicFunc,vinterface{},wg*sync.WaitGroup){
deferwg.Done()
iftopic!=nil&&!topic(v){
return
}
//监听subchan写入成功或超时
select{
casesub<-v:
case<-time.After(p.timeout):
}
}
我们可以选择订阅全部,或指定自定义函数只订阅符合要求的消息,返回chan对象:
all:=p.Subscribe()//添加一个订阅者,订阅全部消息
//添加一个订阅者,只关系有golang字符串的内容
golang:=p.SubscribeTopic(func(vinterface{})bool{
ifs,ok:=v.(string);ok{
returnstrings.Contains(s,"golang")
}
returnfalse
})
新闻来源:嵌入式应用研究院,文中所述为作者独立观点,不代表icspec立场。更多精彩资讯请下载icspec App。如对本稿件有异议,请联系微信客服specltkj。
暂无评论哦,快来评论一下吧!
2026-07-12
2026-07-08