现在有 2 个服务,Service A 和 Service B,通过 REST 接口通信;Service A 在某个业务场景下调用 Service B 的接口完成一个计算密集型任务,假设接口为 http://service_b/api/v1/domain;该任务运行时间很长,但 Service A 不想一直阻塞在接口调用上。为了满足 Service A 的要求,通常有 2 种方案:
2. Service A 在请求 Service B 接口时带上 callback uri,比如 http://service_b/api/v1/domain?callbackuri=http://service_a/api/v1/domain,Service B 收到请求后立即返回 HTTP Status 200 Ok,等任务完成后再调用 Service A callback uri 进行通知。
方案 1 须要轮询接口,轮询太频繁会导致资源浪费,间隔太长又会导致任务完成后 Service A 无法及时感知。显然,方案 2 更加高效,因此也被广泛应用。
方案 2 用到的思想就是本文要介绍的观察者模式(Observer Pattern),GoF 对它的定义如下:
Define a one-to-many dependency between objects so that when one object changes state, all its dependents are notified and updated automatically.
我们将观察者称为 Observer,被观察者(或主体)称为 Subject,那么 Subject 和 Observer 是一对多的关系,当 Subject 状态变更时,所有的 Observer 都会被通知到。
在 简单的分布式应用系统(示例代码工程)中,应用之间通过 network 模块来通信,其中通信模型采用观察者模式:
从上图可知,App 直接依赖 http 模块,而 http 模块底层则依赖 socket 模块:
App2
初始化时,先向 http 模块注册一个 request handler
,处理 App1
发送的 http 请求。request handler
转换为 packet handler
注册到 socket 模块上。App 1
发送 http 请求,http 模块将请求转换为 socket packet
发往 App 2
的 socket 模块。App 2
的 socket 模块收到 packet 后,调用 packet handler
处理该报文;packet handler
又会调用 App 2
注册的 request handler
处理该请求。在上述 socket - http - app 三层模型 中,对 socket 和 http,socket 是 Subject,http 是 Observer;对 http 和 app,http 是 Subject,app 是 Observer。
因为在观察者模式的实现上,socket 模块和 http 模块类似,所以,下面只给出 socket 模块的实现:
// demo/network/socket.go
package network
// 关键点1: 定义Observer接口
// SocketListener Socket报文监听者
type SocketListener interface {
// 关键2: 为Observer定义更新处理方法,入参为相关的上下文对象
Handle(packet *Packet) error
}
// Subject接口
// Socket 网络通信Socket接口
type Socket interface {
// Listen 在endpoint指向地址上起监听
Listen(endpoint Endpoint) error
// Close 关闭监听
Close(endpoint Endpoint)
// Send 发送网络报文
Send(packet *Packet) error
// Receive 接收网络报文
Receive(packet *Packet)
// AddListener 增加网络报文监听者
AddListener(listener SocketListener)
}
// 关键点3: 定义Subject对象
// socketImpl Socket的默认实现
type socketImpl struct {
// 关键点4: 在Subject中持有Observer的集合
listeners []SocketListener
}
// 关键点5: 为Subject定义注册Observer的方法
func (s *socketImpl) AddListener(listener SocketListener) {
s.listeners = append(s.listeners, listener)
}
// 关键点6: 当Subject状态变更时,遍历Observers集合,调用它们的更新处理方法
func (s *socketImpl) Receive(packet *Packet) {
for _, listener := range s.listeners {
listener.Handle(packet)
}
}
...
总结实现观察者模式的几个关键点:
SocketListener
接口。Handle
方法,上下文对象为 Packet
。socketImpl
对象。当然,也可以先将 Subject 抽象为接口,比如上述例子中的 Socket
接口,但大多数情况下都不是必须的。listeners
属性。让 Subject 依赖 Observer 接口,能够使 Subject 与具体的 Observer 实现解耦,提升代码的可扩展性。AddListener
方法。Receive
方法。与观察者模式相近的,是发布-订阅模式(Pub-Sub Pattern),很多人会把两者等同,但它们之间还是有些差异。
从前文的观察者模式实现中,我们发现 Subject 持有 Observer 的引用,当状态变更时,Subject 直接调用 Observer 的更新处理方法完成通知。也就是,Subject 知道有哪些 Observer,也知道 Observer 的数量:
在发布-订阅模式中,我们将发布方称为 Publisher,订阅方称为 Subscriber,不同于观察者模式,Publisher 并不直接持有 Subscriber 引用,它们之间通常通过 Broker 来完成解耦。也即,Publisher 不知道有哪些 Subscriber,也不知道 Subscriber 的数量:
发布-订阅模式被广泛应用在消息中间件的实现上,比如 Apache Kafka 基于 Topic 实现了发布-订阅模式,发布方称为 Producer,订阅方称为 Consumer。
下面,我们通过 简单的分布式应用系统(示例代码工程)中的 mq 模块,展示一个简单的发布-订阅模式实现,在该实现中,我们将 Publisher 的 produce 方法和 Subscriber 的 consume 方法都合并到 Broker 中:
// demo/mq/memory_mq.go
// 关键点1: 定义通信双方交互的消息,携带topic信息
// Message 消息队列中消息定义
type Message struct {
topic Topic
payload string
}
// 关键点2: 定义Broker对象
// memoryMq 内存消息队列,通过channel实现
type memoryMq struct {
// 关键点3: Broker中维持一个队列的map,其中key为topic,value为queue,go语言通常用chan实现。
queues sync.Map // key为Topic,value为chan *Message,每个topic单独一个队列
}
// 关键点4: 为Broker定义Produce方法,根据消息中的topic选择对应的queue发布消息
func (m *memoryMq) Produce(message *Message) error {
record, ok := m.queues.Load(message.Topic())
if !ok {
q := make(chan *Message, 10000)
m.queues.Store(message.Topic(), q)
record = q
}
queue, ok := record.(chan *Message)
if !ok {
return errors.New("model's type is not chan *Message")
}
queue <- message
return nil
}
// 关键点5: 为Broker定义Consume方法,根据topic选择对应的queue消费消息
func (m *memoryMq) Consume(topic Topic) (*Message, error) {
record, ok := m.queues.Load(topic)
if !ok {
q := make(chan *Message, 10000)
m.queues.Store(topic, q)
record = q
}
queue, ok := record.(chan *Message)
if !ok {
return nil, errors.New("model's type is not chan *Message")
}
return <-queue, nil
}
客户端使用时,直接调用 memoryMq
的 Produce
方法和 Consume
方法完成消息的生产和消费:
// 发布方
func publisher() {
msg := NewMessage("test", "hello world")
err := MemoryMqInstance().Produce(msg)
assert.Nil(t, err)
}
// 订阅方
func subscriber() {
result, err := MemoryMqInstance().Consume("test")
assert.Nil(err)
assert.Equal(t, "hello world", result.payload)
}
总结实现发布-订阅模式的几个关键点:
Message
对象。memoryMq
对象。queues
属性。Produce
方法。Consume
方法。实现观察者模式和发布-订阅模式时,都会涉及到 Push 模式或 Pull 模式的选取。所谓 Push 模式,指的是 Subject/Publisher 直接将消息推送给 Observer/Subscriber;所谓 Pull 模式,指的是 Observer/Subscriber 主动向 Subject/Publisher 拉取消息:
Push 模式和 Pull 模式的选择,取决于通信双方处理消息的速率大小。
如果 Subject/Publisher 方生产消息的速率要比 Observer/Subscriber 方处理消息的速率小,可以选择 Push 模式,以求得更高效、及时的消息传递;相反,如果 Subject/Publisher 方产生消息的速率要大,就要选择 Pull 模式,由 Observer/Subscriber 方决定消息的消费速率,否则可能导致 Observer/Subscriber 崩溃。
Pull 模式有个缺点,如果当前无消息可处理,将导致 Observer/Subscriber 空轮询,可以采用类似 Kafka 的解决方案:让 Observer/Subscriber 阻塞一定时长,让出 CPU,避免长期无效的 CPU 空转。
观察者模式和发布-订阅模式中的 Subject 和 Broker,通常都会使用 [单例模式] 来确保它们全局唯一。
可以在 [用Keynote画出手绘风格的配图] 中找到文章的绘图方法。
本文由哈喽比特于2年以前收录,如有侵权请联系我们。
文章来源:https://mp.weixin.qq.com/s/RZqZmjWm_NMzgf8nUhv2Xg
京东创始人刘强东和其妻子章泽天最近成为了互联网舆论关注的焦点。有关他们“移民美国”和在美国购买豪宅的传言在互联网上广泛传播。然而,京东官方通过微博发言人发布的消息澄清了这些传言,称这些言论纯属虚假信息和蓄意捏造。
日前,据博主“@超能数码君老周”爆料,国内三大运营商中国移动、中国电信和中国联通预计将集体采购百万台规模的华为Mate60系列手机。
据报道,荷兰半导体设备公司ASML正看到美国对华遏制政策的负面影响。阿斯麦(ASML)CEO彼得·温宁克在一档电视节目中分享了他对中国大陆问题以及该公司面临的出口管制和保护主义的看法。彼得曾在多个场合表达了他对出口管制以及中荷经济关系的担忧。
今年早些时候,抖音悄然上线了一款名为“青桃”的 App,Slogan 为“看见你的热爱”,根据应用介绍可知,“青桃”是一个属于年轻人的兴趣知识视频平台,由抖音官方出品的中长视频关联版本,整体风格有些类似B站。
日前,威马汽车首席数据官梅松林转发了一份“世界各国地区拥车率排行榜”,同时,他发文表示:中国汽车普及率低于非洲国家尼日利亚,每百户家庭仅17户有车。意大利世界排名第一,每十户中九户有车。
近日,一项新的研究发现,维生素 C 和 E 等抗氧化剂会激活一种机制,刺激癌症肿瘤中新血管的生长,帮助它们生长和扩散。
据媒体援引消息人士报道,苹果公司正在测试使用3D打印技术来生产其智能手表的钢质底盘。消息传出后,3D系统一度大涨超10%,不过截至周三收盘,该股涨幅回落至2%以内。
9月2日,坐拥千万粉丝的网红主播“秀才”账号被封禁,在社交媒体平台上引发热议。平台相关负责人表示,“秀才”账号违反平台相关规定,已封禁。据知情人士透露,秀才近期被举报存在违法行为,这可能是他被封禁的部分原因。据悉,“秀才”年龄39岁,是安徽省亳州市蒙城县人,抖音网红,粉丝数量超1200万。他曾被称为“中老年...
9月3日消息,亚马逊的一些股东,包括持有该公司股票的一家养老基金,日前对亚马逊、其创始人贝索斯和其董事会提起诉讼,指控他们在为 Project Kuiper 卫星星座项目购买发射服务时“违反了信义义务”。
据消息,为推广自家应用,苹果现推出了一个名为“Apps by Apple”的网站,展示了苹果为旗下产品(如 iPhone、iPad、Apple Watch、Mac 和 Apple TV)开发的各种应用程序。
特斯拉本周在美国大幅下调Model S和X售价,引发了该公司一些最坚定支持者的不满。知名特斯拉多头、未来基金(Future Fund)管理合伙人加里·布莱克发帖称,降价是一种“短期麻醉剂”,会让潜在客户等待进一步降价。
据外媒9月2日报道,荷兰半导体设备制造商阿斯麦称,尽管荷兰政府颁布的半导体设备出口管制新规9月正式生效,但该公司已获得在2023年底以前向中国运送受限制芯片制造机器的许可。
近日,根据美国证券交易委员会的文件显示,苹果卫星服务提供商 Globalstar 近期向马斯克旗下的 SpaceX 支付 6400 万美元(约 4.65 亿元人民币)。用于在 2023-2025 年期间,发射卫星,进一步扩展苹果 iPhone 系列的 SOS 卫星服务。
据报道,马斯克旗下社交平台𝕏(推特)日前调整了隐私政策,允许 𝕏 使用用户发布的信息来训练其人工智能(AI)模型。新的隐私政策将于 9 月 29 日生效。新政策规定,𝕏可能会使用所收集到的平台信息和公开可用的信息,来帮助训练 𝕏 的机器学习或人工智能模型。
9月2日,荣耀CEO赵明在采访中谈及华为手机回归时表示,替老同事们高兴,觉得手机行业,由于华为的回归,让竞争充满了更多的可能性和更多的魅力,对行业来说也是件好事。
《自然》30日发表的一篇论文报道了一个名为Swift的人工智能(AI)系统,该系统驾驶无人机的能力可在真实世界中一对一冠军赛里战胜人类对手。
近日,非营利组织纽约真菌学会(NYMS)发出警告,表示亚马逊为代表的电商平台上,充斥着各种AI生成的蘑菇觅食科普书籍,其中存在诸多错误。
社交媒体平台𝕏(原推特)新隐私政策提到:“在您同意的情况下,我们可能出于安全、安保和身份识别目的收集和使用您的生物识别信息。”
2023年德国柏林消费电子展上,各大企业都带来了最新的理念和产品,而高端化、本土化的中国产品正在不断吸引欧洲等国际市场的目光。
罗永浩日前在直播中吐槽苹果即将推出的 iPhone 新品,具体内容为:“以我对我‘子公司’的了解,我认为 iPhone 15 跟 iPhone 14 不会有什么区别的,除了序(列)号变了,这个‘不要脸’的东西,这个‘臭厨子’。