Go ---Go语言高级编程中订阅/发布模型例子解析

简介: Go ---Go语言高级编程中订阅/发布模型例子解析

Go语言高级编程》确实是本好书,我的反应是:很嫉妒,妈的!写的这么 牛逼!

func main() {
  // 一个过期时间为 0.1秒,缓冲区大小为10的发布者
  // 发布者的缓冲区大小决定,订阅者的缓冲区大下
  // 如果发布的主题订阅者没有接受将会阻塞这个订阅者
  // 新发布的主题该订阅者无法在进行接收
  p := NewPublisher(100*time.Millisecond, 10)
  defer p.Close()
  // 添加两个订阅者,一个订阅全部,一个订阅"golang"
  all := p.Subscribe()
  golang := p.SubscribeTopic(
    // 主题过滤规则
    func(v interface{}) bool {
    //  是字符串类型吗?
    //  如果是,那么这个里面包含golang吗?
    //  满足上面的条件才是我这个订阅者想要的
    //  不然返回 false
    //  如果是 nil,就代表只要发送我就要,相当于全部订阅
    if s, ok := v.(string); ok {
      return strings.Contains(s, "golang")
    }
    return false
  })
  p.Publish("hello, world!")
  p.Publish("hello, golang!")
  p.Publish("golang!")
  go func() {
    for msg := range all {
      fmt.Println("all:",msg)
    }
  }()
  go func() {
    for msg := range golang {
      fmt.Println("golang:",msg)
    }
  }()
//  运行一段时间后退出
  time.Sleep(3*time.Second)
}
// 订阅/发布模型
type (
  // 订阅者为一个通道
  subscriber chan interface{}
  //  主题为一个过滤器
  topicFunc  func(v interface{}) bool
)
// Publisher 发布者对象
type Publisher struct {
  m       sync.RWMutex
  buffer      int
  timeout     time.Duration
  subscribers   map[subscriber] topicFunc
}
// 构建一个发布者对象,可以蛇者超时时间和缓存队列长度
func NewPublisher(publishTimeout time.Duration, buffer int) *Publisher {
  return &Publisher{
    buffer: buffer,
    timeout: publishTimeout,
    subscribers: make(map[subscriber]topicFunc),
  }
}
// 添加一个新的订阅者,订阅所有主题
func (p *Publisher) Subscribe() chan interface{} {
  return p.SubscribeTopic(nil)
}
// 添加一个新的订阅者,订阅过滤器筛选后的主题
func (p *Publisher) SubscribeTopic(topic topicFunc) chan interface{} {
  ch := make(chan interface{}, p.buffer)
  p.m.Lock()
  // 给指定的订阅者,加上主题过滤器
  p.subscribers[ch] = topic
  p.m.Unlock()
  return ch
}
// 退出订阅
func (p *Publisher) Evict(sub chan interface{})  {
  p.m.Lock()
  defer p.m.Unlock()
  // 将该订阅者从发布者的信息中删除
  delete(p.subscribers, sub)
  close(sub)
}
// 发布一个主题
func (p *Publisher) Publish(v interface{})  {
  p.m.RLock()
  defer p.m.RUnlock()
  var wg sync.WaitGroup
  for sub, topic := range p.subscribers {
    wg.Add(1)
    go p.sendTopic(sub, topic, v, &wg)
  }
  // 等待主题发送完成
  wg.Wait()
}
// 发送主题,可以容忍一定的超时
func (p *Publisher) sendTopic(
  sub subscriber, topic topicFunc, v interface{}, wg *sync.WaitGroup,
  )  {
  defer wg.Done()
  // 如果该订阅者没有订阅全部,并且发布的主题又不符合主题过滤器
  // 那么直接返回
  if topic != nil && !topic(v) {
    return
  }
  // 一般 time.After 与 select case一同使用
  // 如果在指定时间内我们定义的通道中没有接受到值,
  // 那么将会执行<-time.After(p.timeout)
  // 是用于判断超时的操作
  select {
  case sub <- v:
  case <-time.After(p.timeout):
  }
}
func (p *Publisher) Close()  {
  p.m.Lock()
  defer p.m.Unlock()
  // 循环关闭所有的订阅者通道
  for sub := range p.subscribers {
    delete(p.subscribers, sub)
    close(sub)
  }
}


相关文章
|
12月前
|
Linux Go iOS开发
Go语言100个实战案例-进阶与部署篇:使用Go打包生成可执行文件
本文详解Go语言打包与跨平台编译技巧,涵盖`go build`命令、多平台构建、二进制优化及资源嵌入(embed),助你将项目编译为无依赖的独立可执行文件,轻松实现高效分发与部署。
1757 162
|
11月前
|
算法 Java Go
【GoGin】(1)上手Go Gin 基于Go语言开发的Web框架,本文介绍了各种路由的配置信息;包含各场景下请求参数的基本传入接收
gin 框架中采用的路优酷是基于httprouter做的是一个高性能的 HTTP 请求路由器,适用于 Go 语言。它的设计目标是提供高效的路由匹配和低内存占用,特别适合需要高性能和简单路由的应用场景。
763 4
|
11月前
|
存储 安全 Java
【Golang】(4)Go里面的指针如何?函数与方法怎么不一样?带你了解Go不同于其他高级语言的语法
结构体可以存储一组不同类型的数据,是一种符合类型。Go抛弃了类与继承,同时也抛弃了构造方法,刻意弱化了面向对象的功能,Go并非是一个传统OOP的语言,但是Go依旧有着OOP的影子,通过结构体和方法也可以模拟出一个类。
495 2
|
Cloud Native 安全 Java
Go:为云原生而生的高效语言
Go:为云原生而生的高效语言
761 1
|
Cloud Native Go API
Go:为云原生而生的高效语言
Go:为云原生而生的高效语言
705 0
|
Cloud Native Java Go
Go:为云原生而生的高效语言
Go:为云原生而生的高效语言
537 0
|
Cloud Native Java 中间件
Go:为云原生而生的高效语言
Go:为云原生而生的高效语言
576 0
|
Cloud Native Java Go
Go:为云原生而生的高效语言
Go:为云原生而生的高效语言
1629 0
|
数据采集 Go API
Go语言实战案例:多协程并发下载网页内容
本文是《Go语言100个实战案例 · 网络与并发篇》第6篇,讲解如何使用 Goroutine 和 Channel 实现多协程并发抓取网页内容,提升网络请求效率。通过实战掌握高并发编程技巧,构建爬虫、内容聚合器等工具,涵盖 WaitGroup、超时控制、错误处理等核心知识点。
|
数据采集 编解码 监控
Go语言实战案例:使用channel实现生产者消费者模型
本文是「Go语言100个实战案例 · 网络与并发篇」第4篇,通过实战案例详解使用 Channel 实现生产者-消费者模型,涵盖并发控制、任务调度及Go语言并发哲学,助你掌握优雅的并发编程技巧。

推荐镜像

更多
  • DNS