异步消息 #
简介 #
异步消息体系采用PubSub即发布/订阅模式允许微服务彼此使用消息进行通信。消息的生产者将消息发送到Topic(话题)无需关心什么应用程序将收到他们的消息。类似地,消费者订阅该主题并接收消息,而无需知道产生这些消息的服务是什么。Sidecar承担中间消息代理负责将消息从消息生产者转移到对该消息感兴趣的所有消费者。当需要将微服务彼此分离时,这种模式特别有用。
概览 #
下图展示了服务A通过SDK或者直接的API调用sidecar的/v1.0/publish接口发布消息。订阅消息的服务对应的sidecar通过调用服务的接口推送订阅消息,服务的具体接口是从上述订阅设置部分所声明的方式在sidecar中设置或者从服务的对应接口拉取。

功能亮点 #
- 至少一次消费,Sidecar确保该消息将至少一次发送给每个消费程序。
- 消息生存时间TTL。发布消息时可以为每个消息设置超时时间,这意味着如果超时未从PubSub组件中读取该消息,则该消息将被丢弃。从而让业务无需关心过期消息的处理。
- 竞争消费者,同一个unique id的消费者互为竞争关系,在推送订阅消息时只有一个实例会收到。不同unique id意味着不同消费者,在消息推送时会推送给所有的消费者。
- 在社区基础上自研完成了对redis cluster作为消息中间件进行支持。
- 多种消息中间件支持。sidecar对上层业务提供统一的API规则和统一的消费订阅能力,业务无需关心不同的消息中间件细节。
- 全链路追踪支持。异步消息的发布和消费可以在同一条trace中展示。
名词解释 #
- pubsubname 在使用异步消息时所使用的的引擎名称,是引擎的唯一标识。这里的引擎概念是异步消息体系中承载信息的中间件载体。在设计上允许服务使用多个引擎,通过pubsubname标识对不同的引擎的使用。在msp中创建了一个默认的redis 集群模式的引擎,中台的普通服务可以直接使用此引擎就可以。
- Topic 如果有过消息队列的使用,对这个概念会很熟悉。Topic或者说主题是一条消息的队列,异步消息通过Topic来分类消息,不同Topic间的数据是相互隔离的。一个Topic会有发布消息的服务和订阅消息的服务,这两个都可以是多个。即服务A、B…可以在TopicXXX上发布消息,服务D、E…可以在TopicXXX上订阅消息。当然前提是在msp上配置了这种规则。
- metadata 对于go服务和PHP服务是不会感知此概念,已经在sdk内部进行了封装。如果是其他服务需要使用异步消息功能,且用到ttl功能时,需要在metadata中附带ttl信息。
- route 在后文的订阅Topic配置中出现此定义。代表当前Topic有消息时,sidecar会回调用服务此接口,将订阅的消息推送给服务。其值为接口路由path,不带"/“前缀。
平台使用 #
消息引擎 #
对于正常业务来说,不会涉及到添加引擎操作,如果是中台普通服务使用可以跳过添加引擎的设置。
消息的传递需要载体,也就是所说的消息引擎,我们目前提供四种类型的消息引擎,分别是 kafka、 rabbitmq 、redis 以及 rediscluster。同时,为了便于开发者添加引擎,我们针对四种引擎抽象了对应的四种模板
创建引擎 #

注:为方便中台使用,已事先创建好底层驱动为redis集群的引擎,业务方无需再次创建
消息组件 #
一个消息组件就是一个服务在指定引擎下的所有发布订阅规则集合。
创建组件 #
引擎创建完毕之后,可针对开发者自己的服务创建对应的消息组件


查看组件 #
已经添加的消息组件都可以在列表页面展示

消息遥测 #
针对已经接入异步消息服务的用户,msp平台提供了基于主题以及服务多个视角的遥测功能,可以在任意时刻看到某一服务、某一主题的各种状态,便于业务方了解当前消息处理的健康度。



发布介绍 #
gosdk 和 phpsdk #
对于go语言和php语言,已经将发布订阅的功能集成进SDK,对于使用方来说就是一个函数调用。例如对于go服务来说一个调用示例如下:
client, _ = gosdk.GetNewClient(_header)
res, err := client.PublishEvent("pubsub", "topic", []byte(`{"message":"aaaaa"}`), &RuntimeConfig{TTLInSecond: 13})
var resMap = make(map[string]interface{})
json.Unmarshal(res, &resMap)
if err == nil && resMap["state"] == 1 {
//消息发布成功...
}
对于php来说一个调用示例如下(phpsdk版本>=0.5.30):
$params = $_POST;
$topic = $params['topic'];
$pubsubName = $params['pubsub_name'];
//通过sdk发送消息
$client = Client::getInstance();
$client->publishEvent($pubsubName, $topic, $data);
通过 HTTP API 调用 #
POST <sidecarAddress>/v1.0/publish/<pubsubname>/<topic>[?<metadata>]
sidecarAddress是服务的sidecar访问地址,通常是通过sdk获取的。pubsubname是组件名称topic是主题metadata是额外的控制字段,query格式的传参格式,当前只有一个可以设置的值metadata.ttlInSeconds代表以秒为单位消息的有效时间,没有过期时间就不需要设置此字段。
返回码定义 #
返回码定义根据中台通用的规则进行了适配。返回的格式内容如下:
{
"state": int // 1代表成功,其他代表失败
"msg": string // 成功会是"success",如果是失败,则为失败原因描述,如"topic XXX is not allowed for unique id XXX"
}