版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领
文档简介
20|在分布式环境中,队列、栅栏和STM 14:01Go并发编程实战14:01 已经带你认识了基于etcd实现的Leader、互斥锁和读写锁,今天,来学习下基于etcd的分布式队列、栅栏和STM。只要你学过计算机算法和数据结构相关的知识,队列这种数据结构你一定不陌生,它是一 第12中,专门讲到了一种叫做lock- 下基于etcd key的操作并且提供事务功能的STM(SoftwareTransactionalMemory,软件事务内 并不是从零开始实现一个分布式队列,而是站在etcd的肩膀上,利用etcd提供的功能实现分布式队列。etcd集群的可用性由etcd集群的者来保证, 可以把这些通通交给etcd的运维人员,把 就来了解下etcd提供的分布式队列。etcd通NewQueue,etcd和这个队列的名字,就可以了。代码如下:11func,keyPrefixstring)这个队列只有两个方法,分别是出队和入队,队列中的元素是字符串类型签名如下所示:////func(q*Queue)Enqueue(valstring)func(q*Queue)Dequeue()(string,Dequeue 在接下来讲的例子中,你就可以启动两个节点,一个节点往队列中放入元素,一个节点从队列中取出元素,看看是否能正常取出来。etcd的分布式队列是一种多读多写的队列,下面来借助代码,看一下如何实现分布式队列。 aue>,将一个元素入队,输入pp代1package23import recipe 1518var21
=flag.String("addr"," ","etcdaddresses")queueName=flag.String("name","my-test-queue","queuename")funcmain()//解析etcdendpoints:=strings.Split(*addr,//创建etcdcli,err v3.Config{Endpoints:iferr!=nildefer创建/q:=recipe.NewQueue(cli,consolescanner:=bufio.NewScanner(os.Stdin)forconsolescanner.Scan(){action:=consolescanner.Text()items:=strings.Split(action,"")switchitems[0]{case"push"iflen(items)!=2fmt.Println("mustsetvaluetopush")q.Enqueue(items[1])//case"pop"v,errq.Dequeue()//iferr!={fmt.Println(v)case"quit","exit"://fmt.Println("unknown除了刚刚说的分布式队列,etcd还提供了优先级队列(PriorityQueue)它的用法和队列类似,也提供了出队和入队的操作,只不过,在入队的时候,除了需要把一个值加入到队列, 还需要提供uint6一个,作值的级,级高的元素会优先出队。代1package23import recipe 1619var22
=flag.String("addr"," ","etcdaddresses")queueName=flag.String("name","my-test-queue","queuename")funcmain()//解析etcdendpoints:=strings.Split(*addr,//创建etcdcli,err v3.Config{Endpoints:iferr!=nil defer//创建/q:=recipe.NewPriorityQueue(cli,//从命令 consolescanner:=forconsolescanner.Scan()action:=items:=strings.Split(action,"switchitems[0]case"push":iflen(items)!=3fmt.Println("mustsetvalueandprioritytopush")}pr,err:=strconv.Atoi(items[2 iferr!=nilfmt.Println("mustsetuint16aspriority")}q.Enqueue(items[1],uint16(pr))//case"pop":v,err:=q.Dequeue()//iferr!=nil{}fmt.Println(v)//case"quit","exit"://fmt.Println("unknown}}}etcd中,如果有这类需求的话,你就可以选择用etcd实现。 第17讲中, 学习了循环栅栏CyclicBarrier,它和 第6讲的标准库中的WaitGroup,本质上是同一类并发原语,都是等待同一组goroutine同时执行,或者是等待同一组goroutine都完成。在分布式环境中,Barrier:分布式栅栏。如果持有Barrier的节点释放了它,所有等待这个Barrier的节 的数量,当这些数量的节点都Enter或者Leave的时候,这个栅栏就会放开。所以,先来学习下分布式Barrier分布式Barrier的创建很简单,你只需要提供etcd的 和Barrier的名字就可以了,11func ,keystring)BarrierHold、Release和Wait,funcfunc(b*Barrier)Hold()errorfunc(b*Barrier)Release()errorfunc(b*Barrier)Wait()errorHoldBarrierBarrierWaitReleaseBarrier,也就是打开栅栏。如果使用了这个方法,所有被阻塞Wait方阻塞当前的调用者,直到这个Barrier被release。如果这个栅栏不存在, 就来借助一个例子,来看看Barrier该怎么你可以在一个终端中运行这个程序,执行"hold""eae在另外一个终端中运行这个程序,不断调用"wait"方法,看看是否能正常地跳出阻塞继续代1package23import recipe 1518var21
=flag.String("addr"," ","etcdaddressesbarrierName=flag.String("name","my-test-queue","barriername")funcmain()//解析etcdendpoints:=strings.Split(*addr, //创建etcd cli,err v3.Config{Endpoints: iferr!=nil defer //创建/ b:=recipe.NewBarrier(cli, //从命令 consolescanner.Scan()action:=items:=strings.Split(action,"switchitems[0]case"hold"://持有这个case"release"://释放这个case"wait"://等待barrierfmt.Println("aftercase"quit","exit":fmt.Println("unknown}}65etcd还提供了另外一种栅栏,叫做DoubleBarrier,这也是一种非常有用的栅栏。这个栅栏初始化的时候需要提供一个计数count,如下所示:11funcNewDoubleBarrier(s*concurrency.Session,keystring,countint)Enter和Leavefuncfunc(b*DoubleBarrier)Enter()errorfunc(b*DoubleBarrier)Leave()当调用者调用Enter时,会被阻塞住,直到一共有count(初始化这个栅栏的时候设定的值)个节点调用了Enter,这count个被阻塞的节点才能继续执行。所以,你可以利用它同理,如果你想让一组节点在同一个时刻完成任务,就可以调用Leave方法。节点调用LeavecountLeave再来看一下DoubleBarrier的使用例子。你可以起两个节点,同时执行Enter方法,看看这两个节点是不是先阻塞,才继续执行。然后,你再执行Leave方法,也观察一代1package23import recipe 1619var23
addr=flag.String("addr"," ","etcdaddressesbarrierName=flag.String("name","my-test-doublebarrier","barriername")count=flag.Int("c",2,"")funcmain()//解析etcdendpoints:=strings.Split(*addr,//创建etcdcli,err v3.Config{Endpoints:iferr!=nil defer//创建s1,err:=concurrency.NewSession(cli)iferr!=nil{}defer创建/b:=recipe.NewDoubleBarrier(s1,//consolescanner:=consolescanner.Scan()action:=consolescanner.Text()items:=strings.Split(action,"")switchitems[0]{case"enter"://持有这个case"leave"://释放这个barriercase"quit","exit":fmt.Println("unknown}}} 在第17讲学习的循环栅栏,控制的是同一个进程中的不同goroutine的执行,而分布式栅栏和计数型栅栏控制的是不同节点、不同进程的执行etcd 在学习STM之前 要先了解一下etcd的事务以及它的问题etcd提供了在一个事务中对多个key的更新功能,这一组key的操作要么全部成功,要么全部失败。etcd的事务实现方式是基于CAS方式实现的,融合了Get、Put和Deleteetcd的事务操作如下,分为条件块、成功块和失败块,条件块用来检测事务是否成功,如果成功,就执行Then(...),如果失败,就执行Else(...):11Txn().If(cond1,cond2,...).Then(op1,op2,...,).Else(op1’,op2’, 从账户from向账户to转123123456789funcdoTxnXfer(etcd//代,from,tostring,amountuint)(bool,error)getresp,err:=etcd.Txn(ctx.TODO()).Then(OpGet(from),OpGet(to)).Commit()iferr!=nil{returnfalse,}fromKV:=getresp.Responses[0].GetRangeResponse().Kvs[0]toKV:=getresp.Responses[1].GetRangeResponse().Kvs[1]fromV,toV:=toUInt64(fromKV.Value),toUint64(toKV.Value)iffromV<amount{returnfalse,fmt.Errorf(“insufficient}//txn“=”,fromKV.ModRevision),“=”,toKV.ModRevision))txn=OpPut(from,fromUint64(fromV-amount)),OpPut(to,fromUint64(toV+amount))putresp,erriferr!={return}}returnputresp.Succeeded, 可以看到,虽然可以利用etcd实现事务操作,但是逻辑还是比etcdAPI做STM的操作,提供了更加便利的方法。下 来看一看STM怎么用要使用STM,apply11applyfunc(STM)这个方法包含一个STMkeySTM提供了4Get、Put、Receive和DeletetypetypeSTMinterface{Get(key...string)stringPut(key,valstring,opts...v3.OpOption)Rev(keystring)int64Del(key6使用etcdSTM的时候, 只需要定义一个apply方法,比如说转账方法exchange,然后通过concurrency.NewSTM(cli,exchange),就可以完成转账事务的执行了。STM咋用呢 下面这个例子创建了5个,然后随机选择一些账号两两转账。在转账的时候,要把源账号一半的钱要转给目标账号。这个例子启动了10个goroutine去执行这些事务,每个goroutine要完成100个事务。 代1package23import 1619var21
addr=flag.String("addr", ","etcd24funcmain() //解析etcdendpoints:=strings.Split(*addr,cli,err v3.Config{Endpoints:iferr!=nil defer//设置5个账户,每个账号都有100元,总共500totalAccounts:=fori:=0;i<totalAccounts;i++k:=fmt.Sprintf("accts/%d",if_,err=cli.Put(context.TODO(),k,"100");err!=nil //STMexchange:=func(stmconcurrency.STM)error//from,to:= otalAccounts), iffrom==to//return fromK,toK:=fmt.Sprintf("accts/%d",from),fmt.Sprintf("accts/%d",fromV,toV:=stm.Get(fromK),fromInt,toInt:=0,fmt.Sscanf(fromV,"%d", oV,"%d",//xfer:=fromInt/fromInt,toInt=fromInt-xfer,//stm.Put(fromK,fmt.Sprintf("%d",stm.Put(toK,fmt.Sprintf("%d",return //启动10个goroutinevarwgfori:=0;i<10;i++gofunc()
温馨提示
- 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
- 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
- 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
- 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
- 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
- 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
- 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。
最新文档
- 2026天津市肿瘤医院上合办辅助岗位招聘笔试备考题库及答案详解
- 2026年河南高职单招统一模拟试题及参考答案(原版真题风格)
- 2026年遂昌县带编教师招聘笔试备考试题及答案解析
- 2026年东源县带编教师招聘考试备考题库及答案解析
- 2026年江达县带编教师招聘考试备考题库及答案解析
- 2026年屏山县带编教师招聘笔试参考题库及答案解析
- 2026年汪清县带编教师招聘笔试备考试题及答案解析
- 2026年农安县带编教师招聘考试备考试题及答案解析
- 2026年云和县带编教师招聘考试模拟试题及答案解析
- 2026年乐东黎族自治县带编教师招聘笔试备考试题及答案解析
- 2026年秋浙美版新教材小学美术五年级上册教学计划及进度表
- 2026年秋季新教材浙美版小学美术六年级上册(全册)教案(附目录p94)
- 高一数学 开学第一课 课件-2026-2027学年高一上学期数学人教A版必修第一册
- 2026年安徽矾花源景区运营管理有限公司(筹) 招聘14人考试备考题库及答案详解
- 云南省公路工程竣工文件编制及立卷归档实 用手册
- 2026秋小学湘艺版音乐三年级上册(新教材)教学计划附教学进度表
- 2026下半年上海杨浦区卫健系统事业单位专业技术人员招聘93人笔试题库附答案详解【预热题】
- 长江产业投资集团招聘笔试题目及答案解析
- (2026年秋)人教PEP版五年级上册英语教案
- 新版部编人教版六年级上册道德与法治(课件)第1课 法律是什么
- 2026年二级建造师继续教育试题加答案
评论
0/150
提交评论