消息队列

消息(Message):

消息本质上是一种数据结构(当然,对象也可以看做是一种特殊的消息),它包含消费者与服务双方都能识别的数据,这些数据需要在不同的进程(机器)之间进行传递,并可能会被多个完全不同的客户端消费

队列(Queue):

是先进先出(FIFO, First-In-First-Out)的线性表,通俗的讲队列就是一群人或者事物按照排好的顺序等待接受服务或者处理

MQ全称(Message Queue)

又名消息队列,是一种异步通讯中间件。可以将它理解成邮局,发送者将消息传递到邮局,然后由邮局帮我们发送给具体的消息接收者(消费者),具体发送过程与时间我们无需关心,它也不会干扰我进行其它事情。

它被广泛的应用与跨平台、跨系统的分布式系统之间,为它们提供高效可靠的异步传输机制.

可以看做是一种异步RPC,把一次RPC变为两次或多次,进行内容转存,再在合适的时机投递出去。消息的发送者和接收者不需要在同一时间与消息队列进行交互,消息在被处理或被删除之前一直存储在队列上。

JMS介绍

JMS: java message service,JAVA平台中面相消息服务的中间件接口.

JMS是一种与厂商无关,用来访问消息收发的中间组件,它类似JDBC. JDBC是java用来访问不同厂商的数据库的集成API接口,各厂商在JDBC标准接口下做自己产品的适配.

严格来讲,JMS并不属于消息协议,而是一种规范,是对AMQP,MQTT,STOMP等协议更高一层的抽象。从使用角度看,JMS和JDBC担任差不多的角色,用户都是根据相应的接口可以和实现了JMS的服务进行通信,进行相关的操作。

AMQPMQTTSTOMP是三种最常见、最流行的基于TCP/IP的消息传递协议。

JMS通常包含如下一些角色:

  • JMS provider:实现了JMS接口的消息中间件,如ActiveMQ

  • JMS client:生产或者消费消息的应用

  • JMS producer/publisher:JMS消息生产者

  • JMS consumer/subscriber:JMS消息消费者

  • JMS message:消息,在各个JMS client传输的对象

  • JMS queue:Provider存放等待被消费的消息的地方

  • JMS topic:一种提供多个订阅者消费消息的一种机制;在MQ中常常被提到,topic模式

消息如何从producer端达到consumer端由message-routing来决定。在JMS中,消息路由非常简单,由producer和consumer链接到同一个queue(p2p)或者topic(pub/sub)来实现消息的路由。JMSconsumer同时支持message selector(消息选择器),通过消息选择器,consumer可以只消费那些通过了selector筛选的消息。

AMQP

AMQP(advanced message queuing protocol) 即高级消息队列协议. 一个提供统一消息服务的应用层标准高级消息队列协议,是应用层协议的一个开放标准,为面向消息的中间件设计。基于此协议的客户端与消息中间件可传递消息,并不受客户端/中间件不同产品,不同的开发语言等条件的限制。

AMQP是一种协议,更准确的说是一种binary wire-level protocol(链接协议),兼容JMS

核心角色如下:

  • Message(消息):消息服务器处理消息的原子单元,包括一个内容头,一组属性和一个内容体。
    消息有优先级,高优先级的消息在等待同一消息队列时会比低优先级的消息先发送,而且当消息必须被丢弃时,低优先级的消息优先被丢弃。
    使用AMQP协议,消息服务器不能修改内容体和内容头,但可以在内容头上添加额外信息。

  • PubLisher(消息生产者):发送消息

  • Consumer(消息消费者):消费消息

  • Broker(消息代理):消息队列服务器,负责接收客户端连接,路由消息。

  • Queue(消息队列):Broker中的一个角色,一个Broker中可以有多个Queue,负责保存消息直到发送给不同的消费者。算是消息的容器。一个消息可以被投入一个或多个队列中,每个队列的消息都会等待消费者连接到这个队列并被取走。

  • Exchange(交换路由):Broker中的一个角色,负责接收生产者发送的消息,并路由给服务器中的队列。可以被理解成一个规则表,指明消息该被投到哪个队列中。

  • Channel(信道):信道是一条独立的双向数据流通道。为了解决操作系统无法承受每秒建立特别多的TCP连接。

生产者发送消息时,必须指定消息要被路由到哪些个消息队列中。当消息到消息队列中,消息队列会尝试将消息传给消费者,如果失败,消息队列会存储消息并等待消费者。如果没有消费者,消息队列将选择性的将消息返回给生产者。如果消息被消费掉,消息队列会删除消息,删除的过程或者是及时的,或者是等到消费者消费结果后才删除的。

RabbitMQ是AMQP消息队列最有名的开源实现,当然RabbitMQ同时还可以通过插件支持STOMP、MQTT等协议接入。Kafka、RocketMQ均使用自定义的协议。

MQTT

MQTT,即消息队列遥测传输(Message Queuing Telemetry Transport)。由IBM开发,现在被广泛用于物联网公司。因为他的特点就是轻量,简单,开放和易于实现。所以它常用于很多计算能力有限、带宽低、网络不可靠的远程通信应用场景。

核心角色如下:

  • Publisher(发布者):消息发布客户端

  • Subscriber(订阅者):消息订阅客户端

  • Broker(消息代理):消息服务器端

  • Application Message(应用消息):指通过网络传输的应用数据,一般包括主题和负载。

  • Topic(主题):应用消息的类型,一般消息发布者会确定消息的主题,订阅者根据自己实际情况选择不同的主题进行消息订阅消费。

  • Payload(负载):消息订阅者具体接收的内容。

MQTT协议是通过交换预定义的MQTT控制报文来通信的,控制报文内容由三部分组成:固定报头,可变报头和消息体。固定报头通过标识不同位的值来确定报文类型,包括发布订阅的一些完成状态等;可变报头的内容根据控制报文类型不同而不同,常作为包的标识符;消息体也是根据不同的消息类型有着不同的内容。

MQTT协议中,客户端和服务端是通过请求-应答模式通信的。客户端发送一条控制报文数据给服务器,服务器再发送一条控制报文数据给客户端。

MQTT在发布消息时,有三种Qos等级:

  • 至多一次(0级)

  • 至少一次(1级)

  • 只有一次(2级)

至多一次等级最低,客户端只需要将消息发出去即可,这种等级很低,用于消息不重要但特别多,为了减轻通信压力,就不顾质量,只看数量了。

至少一次等级中等,客户端要保证发出去的消息至少一次被服务端接收到,所以要收到服务端的回应,否则一直发,这种等级一般用于服务端有幂等处理,所以不怕重复消费,还要保证消息不会丢失。

只有一次等级最高,客户端先发消息过去,然后本地记录一个我已发送,但不确定你是否收到的状态,然后服务端接收到消息后,回给客户端一个我已接收的报文,同时服务端记录一个我不确定你知不知道我已接收的状态,然后客户端收到这个已接收的消息后,就确定服务端收到这个消息了,于是把自己本地记录的已发送未确定的状态删除,同时再给客户端发送一个我已经知道你收到的报文,服务端收到这个报文,也会把自己之前记录的状态删掉,整个一条报文只有一次的通信才算完成,这种等级就比较严格了,但质量上去了,相对低等级的,数量就会相对小些,但可靠就是王道,不多不少才是最好的。

只有一次的发送和确定,其实思想和三次握手差不多,都是两端互相确认的过程,所以会一来一回的。如果传输过程中出现丢包,都会由发送者重发上一条消息。

STOMP

STOMP,即流文本定向消息协议(Streaming Text Orientated Messaging Protocal),是一个相对简单的文本消息传输协议,主要特点就是简单易懂,没有特别多的套路。

核心角色如下:

  • 客户端:既可以是生产者,也可以是消费者

  • 服务端:消息中心

ActiveMQ以及它的下一代实现Apache Apollo,是STOMP协议的典型实现。

使用场景

  • 跨平台

  • 多语言

  • 多项目

  • 解耦

  • 分布式事物

  • 流量控制

  • 最终一致性需求

  • RPC调用,

  • 上下游对接,数据变更通知下属

通信模型

JMS具有两种通信模式:

  • Point-to-Point Messaging Domain (点对点)

  • Publish/Subscribe Messaging Domain (发布/订阅模式)

在JMS API出现之前,大部分产品使用“点对点”和“发布/订阅”中的任一方式来进行消息通讯。JMS定义了这两种消息发送模型的规范,它们相互独立。任何JMS的提供者可以实现其中的一种或两种模型,这是它们自己的选择。JMS规范提供了通用接口保证我们基于JMS API编写的程序适用于任何一种模型。

点对点模型

在点对点通信模式中,应用程序由消息队列,发送方,接收方组成。每个消息都被发送到一个特定的队列,接收者从队列中获取消息。队列保留着消息,直到他们被消费或超时。

  • 每个消息只要一个消费者

  • 发送者和接收者在时间上是没有时间的约束,也就是说发送者在发送完消息之后,不管接收者有没有接受消息,都不会影响发送方发送消息到消息队列中。

  • 发送方不管是否在发送消息,接收方都可以从消息队列中去到消息

  • 接收方在接收完消息之后,需要向消息队列应答成功

发布/订阅模型

在发布/订阅消息模型中,发布者发布一个消息,该消息通过topic传递给所有的客户端。该模式下,发布者与订阅者都是匿名的,即发布者与订阅者都不知道对方是谁。并且可以动态的发布与订阅Topic。Topic主要用于保存和传递消息,且会一直保存消息直到消息被传递给客户端。

  • 一个消息可以传递个多个订阅者(即:一个消息可以有多个接受方)

  • 发布者与订阅者具有时间约束,针对某个主题(Topic)的订阅者,它必须创建一个订阅者之后,才能消费发布者的消息,而且为了消费消息,订阅者必须保持运行的状态。一旦消费者退出,相应的订阅以及尚未处理的消息就会丢失。

  • 为了缓和这样严格的时间相关性,JMS允许订阅者创建一个可持久化的订阅。这样,即使订阅者没有被激活(运行),它也能接收到发布者的消息。

接收消息

在JMS中,消息的产生和消息是异步的。对于消费来说,JMS的消息者可以通过两种方式来消费消息。

同步(Synchronous)

  • 在同步消费信息模式模式中,订阅者/接收方通过调用 receive()方法来接收消息。在receive()方法中,线程会阻塞直到消息到达或者到指定时间后消息仍未到达。

  • 例如:远程调用服务,同步RPC

异步(Asynchronous)

  • 使用异步方式接收消息的话,消息订阅者需注册一个消息监听者,类似于事件监听器,只要消息到达,JMS服务提供者会通过调用监听器的onMessage()递送消息。

  • 客户端主服务不需要等待服务处理消息,简单来说就是不阻塞。

分类

本地队列

本地队列按照功能可划分为初始化队列,传输队列,目标队列和死信队列。初始化队列用作消息触发功能。传输队列只是暂存待传的消息,条件许可的情况下,通过管道将消息传送到其他的队列管理器。目标队列是消息的目的地,可以长期存放消息。如果消息不能送达目标队列,也不能再路由出去,则被自动放入死信队列保存。

别名队列&远程队列

只是一个队列定义,用来指定远端队列管理器的队列。使用了远程队列,程序就不需要知道目标队列的位置。

模型队列

模型队列定义了一套本地队列的属性结合,一旦打开模型队列,队列管理器会按照这些属性动态地创建出一个本地队列。

特点

可靠性传输

是消息中间件的重要特点,对于应用来说,只要成功把数据提交给消息中间件,那么关于数据可靠传输的问题就由消息中间件来负责

不重复传输

不重复传播也就是断点续传的功能,特别适合网络不稳定的环境,节约网络资源

异步性传输

异步性传输是指,接受信息双方不必同时在线,具有脱机能力和安全性

消息驱动

接到消息后主动通知消息接收方

支持事务

应用程序可以把一些数据更新组合成一个工作单元,这些更新通常是逻辑相关的,为了保障数据完整性,所有的更新必须同时成功或者同时失败

编程模型

Connection Factories:连接工厂

创建Connection对象的工厂,针对两种不同的jms消息模型,分别有QueueConnectionFactory和TopicConnectionFactory两种。可以通过JNDI来查找ConnectionFactory对象。客户端使用一个连接工厂对象连接到JMS服务提供者,它创建了JMS服务提供者和客户端之间的连接。JMS客户端(如发送者或接受者)会在JNDI名字空间中搜索并获取该连接。使用该连接,客户端能够与目的地通讯,往队列或话题发送/接收消息。

Destination:消息的目的地

目的地指明消息被发送的目的地以及客户端接收消息的来源。JMS使用两种目的地,队列和话题

Connection:连接

表示在客户端和JMS系统之间建立的链接(对TCP/IP socket的包装)。Connection可以产生一个或多个Session。跟ConnectionFactory一样,Connection也有两种类型:QueueConnection和TopicConnection。

Session:会话

对消息进行操作的接口,可以通过session创建生产者、消费者、消息等。Session 提供了事务的功能,如果需要使用session发送/接收多个消息时,可以将这些发送/接收动作放到一个事务中。

Producter:生产值

消息生产者由Session创建,用于往目的地发送消息。生产者实现MessageProducer接口,我们可以为目的地、队列或话题创建生产者。

Consumer:消费者

消息消费者由Session创建,用于接收被发送到Destination的消息。

MessageListener:消息监听器

消息监听器。如果注册了消息监听器,一旦消息到达,将自动调用监听器的onMessage方法。

常用的消息队列中间件

Kafka:

炙手可热的消息中间件,特点是高吞吐量,是Apach的顶级项目,适合产生大数据的互联网服务的数据收集业务

Apache Kafka它最初由LinkedIn公司基于独特的设计实现为一个分布式的提交日志系统( a distributed commit log),之后成为Apache项目的一部分。号称大数据的杀手锏,谈到大数据领域内的消息传输,则绕不开Kafka,这款为大数据而生的消息中间件,以其百万级TPS的吞吐量名声大噪,迅速成为大数据领域的宠儿,在数据采集、传输、存储的过程中发挥着举足轻重的作用。

优点:

  1. 高吞吐量:即使在非常廉价的机器上,Kafka也能做到每秒处理几十万条消息,而它的延迟最低只有几毫秒。

  2. 低延迟:延迟可以控制在ms以内。

  3. 持久性:Kafka可以将消息直接持久化在普通磁盘上,且磁盘读写性能优异。

  4. 扩展性。Kafka集群支持热扩展,Kaka集群启动运行后,用户可以直接向集群中添加。

  5. 容错性。Kafka会将数据备份到多台服务器节点中,即使当某个服务器节点失效时,Zookeeper将通知生产者和消费者从而使用其他的节点,也不会影响整个系统的功能。

  6. 支持多种客户端语言。Kafka支持Java、.NET、PHP、Python等多种语言。

缺点:

  1. 重复消息:Kafka保证每条消息至少送达一次,虽然几率很小,但一条消息可能被送达多次。

  2. 消息乱序:Kafka某一个固定的Partition内部的消息是保证有序的,如果一个Topic有多个Partition,partition之间的消息送达不保证有序。

  3. 复杂性:Kafka需要Zookeeper的支持,Topic一般需要人工创建,部署和维护比一般MQ成本更高。

  4. topic 从几十到几百个时候,吞吐量会大幅度下降,在同等机器下,Kafka 尽量保证 topic 数量不要过多,如果要支撑大规模的 topic,需要增加更多的机器资源。

RabbitMQ:

行业内比较流行的消息中间件,使用的是Erlang语言开发的开源消息队列系统,RabbitMQ可以与Spring做无缝集成,基于AMPQ协议来实现,RabbitMQ对数据的一致性、稳定性、可靠性要求非常严格、不允许丢任何一条数据,但是效率不如KafKa高.

RabbitMQ 2007年发布,是使用Erlang语言开发的开源消息队列系统,基于AMQP协议来实现。AMQP的主要特征是面向消息、队列、路由(包括点对点和发布/订阅)、可靠性、安全。AMQP协议更多用在企业系统内,对数据一致性、稳定性和可靠性要求很高的场景,对性能和吞吐量的要求还在其次。

优点:

  1. 持久化:RabbitMQ可以保证所在的服务器宕机后,消息不会丢失。

  2. 高可用:部分机器宕机了还可以继续使用。

  3. 高级功能:如消息重试、死信队列等

缺点:

  1. 首先是RabbitMQ吞吐量比较低,大概在每秒几万的样子,这样像对于大型电商促销秒杀就不能胜任。

  2. 集群线性扩展比较麻烦。

  3. 开发语言是erlang,懂得人不是很多,无法对其改造。

RocketMQ:

阿里巴巴旗下的消息中间件,纯JAVA开发,对分布式事务处理非常好,可以看作是KafKa的升级版

是阿里开源的消息中间件,它是纯Java开发,具有高吞吐量、高可用性、适合大规模分布式系统应用的特点。RocketMQ思路起源于Kafka,但并不是Kafka的一个Copy,它对消息的可靠传输及事务性做了优化,目前在阿里集团被广泛应用于交易、充值、流计算、消息推送、日志流式处理、binglog分发等场景。

前身是MetaQ,是阿里参考Kafka特点研发的一个队列模型的消息中间件,后开源给apache基金会成为了apache的顶级开源项目,具有高性能、高可靠、高实时、分布式特点。

其实在阿里巴巴内部围绕着RocketMQ内核打造了三款产品,分别是MetaQNotifyAliware MQ

这三者分别采用了不同的模型,MetaQ主要使用了拉模型,解决了顺序消息和海量堆积问题;Notify主要使用了推模型,解决了事务消息;而云产品Aliware MQ则是提供了商业化的版本。

RocketMQ天生为金融互联网领域而生,追求高可靠、高可用、高并发、低延迟,是一个阿里巴巴由内而外成功孕育的典范,除了阿里集团上千个应用外,根据我们不完全统计,国内至少有上百家单位、科研教育机构在使用。

RocketMQ优点:

  • 单机吞吐量:十万级QPS往上。

  • 可用性:非常高,分布式架构

  • 消息可靠性:经过参数优化配置,消息可以做到0丢失

    • 生产者的可靠性保证:生产者发送消息后返回SendResult,如果isSuccess返回true,则表示消息已经确认发送到服务器并被服务器接收保存。整个发送过程是一个同步过程。

    • 服务器的可靠性:消息生产者发送的消息,RocketMQ服务收到后在做必要的校验和检查之后马上保存到磁盘,写入成功后返回给生产者。因此可以确认每条发送结果为成功的消息都会被消息服务器写入磁盘。

    • 消费者的可靠性:消费者是一条一条顺序消费的,之后在成功消费一条后才会消费吓一跳。如果在消费某一条消息时失败则会重试消费这条消息,默认为5次,如果超过最大次数仍然无法消费,则将消息保存到本地,后台线程继续重试消费,主线程则会继续往后走,消费队列后面的消息。

  • 功能支持:MQ功能较为完善,还是分布式的,扩展性好

  • 高级功能:如延迟消息、消息回朔等

  • 支持10亿级别的消息堆积,不会因为堆积导致性能下降

  • 源码是java,我们可以自己阅读源码,定制自己公司的MQ,可以掌控

  • 消息持久性RocketMQ收到消息后,会将消息持久化到文件,并利用Linux文件系统内存来提高性能

  • 消息实时性:RocketMQ采取长轮询+PULL模式保证消息的实时性

  • 消息堆积:支持10亿级别的消息堆积,不会因为消息堆积影响性能

  • 天生为金融互联网领域而生,对于可靠性要求很高的场景,尤其是电商里面的订单扣款,以及业务削峰,在大量交易涌入时,后端可能无法及时处理的情况

  • RoketMQ在稳定性上可能更值得信赖,这些业务场景在阿里双11已经经历了多次考验,如果你的业务有上述并发场景,建议可以选择RocketMQ

RocketMQ缺点:

  • 支持的客户端语言不多,目前是java及c++,其中c++不成熟

  • 社区活跃度不是特别活跃那种

  • 没有在 mq 核心中去实现JMS等接口,有些系统要迁移需要修改大量代码

  • 消息重复:对于消费者来说,通过拉取方式将消息保存到本地,消费完再向服务器返回,在网络异常的情况下可能会出现重复。

  • 消息过滤:

    • 服务器端过滤:减少不必要消息传输,但是会增加服务器负担

    • 客户端过滤:根据客户端需求来定制消息,缺点是客户端会收到对它来说没用的消息,如果客户端无法承载这么多消息就会导致故障

Notify

Notify是淘宝自主研发的一套消息服务引擎,是支撑双11最为核心的系统之中的一个,在淘宝和支付宝的核心交易场景中都有大量使用。消息系统的核心作用就是三点:解耦,异步和并行。

notify诞生于2007年,是伴随着淘宝交易系统进入分布式时代产生的。主要用于核心交易链路的异步解耦,以缩短了下单流程的RT,提高了电商业务的扩展能力。

Notify在设计思路上与传统的MQ有一定的不同,他的核心设计理念是

  • 1. 为了消息堆积而设计系统

  • 2. 无单点,可自由扩展的设计

  • Notify,主要面向需要更加安全可靠地交易类场景,无序推模式
    它的核心特性是: 提供事务支持、不保证消息顺序、消息可能会重复、服务器推模型(一般的消息中间件比如ActiveMQ都是基于服务器push,因为实时性高)。

MetaQ

在2011年,Linkedin开源了一款全新的消息队列kafka,是面向高吞吐量设计的系统,主要用于日志和大数据领域。同年,据说是出于兴趣的原因,当时的notify负责人伯岩花了两个礼拜的时间,开发了java版的kafka,取名为metamorphosis(变形记是奥地利作家卡夫卡的名作,算是对kafka的致敬),主要用于集团内部的日志传输业务。

metamorphosis的第一个版本是基本上是参考kafka的架构,每个topic都会对应broker的多个分区文件,一旦topic数量增多,broker的分区文件数也会随着增大,本来高性能的顺序写文件会变成随机写,吞吐量会有较大的下降。2012年9月,誓嘉解决了这个难题,发布metaq 2.0对存储层进行重新设计以满足大型互联网复杂业务的需求,性能不再随着分区数增大而下降。

MetaQ作为一个分布式的消息中间件,需要依赖zookeeper,对于一些规模不大、单机应用的场景,我个人并不是特别支持尝试用MetaQ,因为多一个依赖系统,其实就是多一份风险,在这些简单场景下,可能类似memcacheq、kestrel甚至redis等轻量级MQ就非常合适。而MetaQ一开始就是为大规模分布式系统设计的,如果不当使用,可能没有带来好处,反而多出一堆问题。

ActiveMQ:

是Apache出品,最流行的,能力强劲的开源消息总线。官方社区现在对ActiveMQ 5.x维护越来越少,较少在大规模吞吐的场景中使用。

Apach出品的最流行的老牌开源消息总线,具有丰富的api,小型企业颇受欢迎,但是性能诟病,吞吐量并不高。

优点

  • ActiveMQ采用消息推送方式,所以最适合的场景是默认消息都可在短时间内被消费。数据量越大,查找和消费消息就越慢,消息积压程度与消息速度成反比。

  • activemq可以很好的运行在任何JVM上,而不只是集成到JBoss的应用服务器中;

  • activemq对Spring有很好的支持.直接集成在Spring依赖中.

  • activemq支持跨网络的分布式目的地

缺点

  • 吞吐量低。由于ActiveMQ需要建立索引,导致吞吐量下降。这是无法克服的缺点,只要使用完全符合JMS规范的消息中间件,就要接受这个级别的TPS。

  • 无分片功能。这是一个功能缺失,JMS并没有规定消息中间件的集群、分片机制。而由于ActiveMQ是为企业级开发设计的消息中间件,初衷并不是为了处理海量消息和高并发请求。如果一台服务器不能承受更多消息,则需要横向拆分。ActiveMQ官方不提供分片机制,需要自己实现。

适用场景

对TPS要求比较低的系统,可以使用ActiveMQ来实现,一方面比较简单,能够快速上手开发,另一方面可控性也比较好,还有比较好的监控机制和界面

TPS:是TransactionsPerSecond的缩写,也就是事务数/秒。它是软件测试结果的测量单位。一个事务是指一个客户机向服务器发送请求然后服务器做出反应的过程。客户机在发送请求时开始计时,收到服务器响应后结束计时,以此来计算使用的时间和完成的事务个数。

ZeroMQ

号称最快的消息队列系统,专门为高吞吐量/低延迟的场景开发,在金融界的应用中经常使用,偏重于实时数据通信场景。ZMQ能够实现RabbitMQ不擅长的高级/复杂的队列,但是开发人员需要自己组合多种技术框架,开发成本高。因此ZeroMQ具有一个独特的非中间件的模式,更像一个socket library,你不需要安装和运行一个消息服务器或中间件,因为你的应用程序本身就是使用ZeroMQ API完成逻辑服务的角色。但是ZeroMQ仅提供非持久性的队列,如果down机,数据将会丢失。如:Twitter的Storm中使用ZeroMQ作为数据流的传输。

ZeroMQ套接字是与传输层无关的:ZeroMQ套接字对所有传输层协议定义了统一的API接口。默认支持 进程内(inproc) ,进程间(IPC) ,多播,TCP协议,在不同的协议之间切换只要简单的改变连接字符串的前缀。可以在任何时候以最小的代价从进程间的本地通信切换到分布式下的TCP通信。ZeroMQ在背后处理连接建立,断开和重连逻辑。

ZeroMQ(简称ZMQ)是一个基于消息队列的多线程网络库,其对套接字类型、连接处理、帧、甚至路由的底层细节进行抽象,提供跨越多种传输协议的套接字。

ZMQ是网络通信中新的一层,介于应用层和传输层之间(按照TCP/IP划分),其是一个可伸缩层,可并行运行,分散在分布式系统间。

ZMQ不是单独的服务,而是一个嵌入式库,它封装了网络通信、消息队列、线程调度等功能,向上层提供简洁的API,应用程序通过加载库文件,调用API函数来实现高性能网络通信。

nanomsg: ZeroMQ作者用C语言新写的消息队列库,解决了ZeroMQ存在的一些问题,改善了一些令人诟病的设计。

特性:

  • 无锁的队列模型:对于跨线程间的交互(用户端和session)之间的数据交换通道pipe,采用无锁的队列算法CAS;在pipe的两端注册有异步事件,在读或者写消息到pipe的时,会自动触发读写事件。

  • 批量处理的算法:对于批量的消息,进行了适应性的优化,可以批量的接收和发送消息。

  • 多核下的线程绑定,无须CPU切换:区别于传统的多线程并发模式,信号量或者临界区,zeroMQ充分利用多核的优势,每个核绑定运行一个工作者线程,避免多线程之间的CPU切换开销。

  • 消息发送端的内存或者磁盘中。不支持持久化。

缺点:

  • 线上找到的资料大部分是原api文档照翻译,用处都不大,看文档还不如直接看官方原文。

  • window下ZeroMQ目前不支持进程间通讯管道,新开个端口吧。

  • 订阅模式没找到数量限制的配置,只能自己写个队列控制链接数量。

  • 可以设置ip限制,设置属性 ZMQ_TCP_ACCEPT_FILTER,配置如 192.168.24.0/22, 代表二进制的ip地址中,前22位必须匹配

  • 令人诟病的 zmq_ctx,以及其设计上有缺陷的

总结

KafKa与RabbitMQ相比较,KafKa更适合高吞吐的处理,但是对数据一致性和稳定性的处理不如RabbitMQ

一般的业务系统要引入 MQ,最早大家都用 ActiveMQ,但是现在确实大家用的不多了,没经过大规模吞吐量场景的验证,社区也不是很活跃,所以大家还是算了吧,我个人不推荐用这个了。

后来大家开始用 RabbitMQ,但是确实 erlang 语言阻止了大量的 Java 工程师去深入研究和掌控它,对公司而言,几乎处于不可控的状态,但是确实人家是开源的,比较稳定的支持,活跃度也高。

不过现在确实越来越多的公司会去用 RocketMQ,确实很不错,毕竟是阿里出品,但社区可能有突然黄掉的风险(目前 RocketMQ 已捐给 Apache,但 GitHub 上的活跃度其实不算高),推荐用 RocketMQ,否则回去老老实实用 RabbitMQ 吧,人家有活跃的开源社区,绝对不会黄。

所以中小型公司,技术实力较为一般,技术挑战不是特别高,用 RabbitMQ 是不错的选择;大型公司,基础架构研发实力较强,用 RocketMQ 是很好的选择。

如果是大数据领域的实时计算、日志采集等场景,用 Kafka 是业内标准的,绝对没问题,社区活跃度很高,绝对不会黄,何况几乎是全世界这个领域的事实性规范。

Notify与MetaQ对比

RocketMQ,kafka对比

功能

消息队列 RocketMQ

Apache RocketMQ (开源)

消息队列 Kafka

Apache Kafka (开源)

安全防护

支持

不支持

支持

不支持

主子账号支持

支持

不支持

支持

不支持

可靠性

- 同步刷盘 - 同步双写 - 超3份数据副本 - 99.99999999%

- 同步刷盘 - 异步刷盘

- 同步刷盘 - 同步双写 - 超3份数据副本 - 99.99999999%

异步刷盘,丢数据概率高

可用性

- 非常好,99.95% - Always Writable

- 非常好,99.95% - Always Writable

横向扩展能力

- 支持平滑扩展 - 支持百万级 QPS

支持

- 支持平滑扩展 - 支持百万级 QPS

支持

Low Latency

支持

不支持

支持

不支持

消费模型

Push / Pull

Push / Pull

Push / Pull

Pull

定时消息

支持(可精确到秒级)

支持(只支持18个固定 Level)

暂不支持

不支持

事务消息

支持

不支持

不支持

不支持

顺序消息

支持

支持

暂不支持

支持

全链路消息轨迹

支持

不支持

暂不支持

不支持

消息堆积能力

百亿级别 不影响性能

百亿级别 影响性能

百亿级别 不影响性能

影响性能

消息堆积查询

支持

支持

支持

不支持

消息回溯

支持

支持

支持

不支持

消息重试

支持

支持

暂不支持

不支持

死信队列

支持

支持

不支持

不支持

性能(常规)

非常好 百万级 QPS

非常好 十万级 QPS

非常好 百万级 QPS

非常好 百万级 QPS

性能(万级 Topic 场景)

非常好 百万级 QPS

非常好 十万级 QPS

非常好 百万级 QPS

性能(海量消息堆积场景)

非常好 百万级 QPS

非常好 十万级 QPS

非常好 百万级 QPS

开发语言

Java

Erlang

Java

C

客户端支持语言

Java、C、 C++、 Python、 PHP、 Perl、.net 等

Java、C、 C++、 Python、 PHP、 Perl、.net 等

Java C++(不成熟)

python、 java、 php、.net 等

事务

支持

不支持

支持

不支持

集群

支持

支持

支持

不支持

负载均衡

支持

支持

支持

不支持

主流的消息队列对比

特性

ActiveMQ

RabbitMQ

RocketMQ

Kafka

资料文档

资料数量多

资料数量多

资料数量少,建议去官网上看

资料数量中等

开发语言

Java

Erlang

Java

Scala

支持的协议

OpenWire、STOMP、REST、XMPP、AMQP

AMQP

自己定义的一套…

自己定义的一套…(基于TCP)

消息存储

内存、磁盘、数据库。支持少量堆积。

内存、磁盘。支持少量堆积

磁盘。支持大量堆积。

内存、磁盘、数据库。支持大量堆积

消息事务

支持

支持。客户端将信道设置为事务模式,只有当消息被RabbitMQ接收,事务才能提交成功,否则在捕获异常后进行回滚。使用事务会使得性能有所下降

支持

支持

负载均衡

支持负载均衡。基于zookeeper实现负载均衡。

负载均衡的支持不好

支持

支持

单机吞吐量

万级,比 RocketMQ、Kafka 低一个数量级

同 ActiveMQ

10 万级,支撑高吞吐

10 万级,高吞吐,一般配合大数据类的系统来进行实时数据计算、日志采集等场景

topic 数量对吞吐量的影响

topic 可以达到几百/几千的级别,吞吐量会有较小幅度的下降,这是 RocketMQ 的一大优势,在同等机器下,可以支撑大量的 topic

topic 从几十到几百个时候,吞吐量会大幅度下降,在同等机器下,Kafka 尽量保证 topic 数量不要过多,如果要支撑大规模的 topic,需要增加更多的机器资源

时效性

ms 级

微秒级,这是 RabbitMQ 的一大特点,延迟最低

ms 级

延迟在 ms 级以内

可用性

高,基于主从架构实现高可用

高,基于主从架构实现高可用

非常高,分布式架构

非常高,分布式,一个数据多个副本,少数机器宕机,不会丢失数据,不会导致不可用

消息可靠性

有较低的概率丢失数据

基本不丢

经过参数优化配置,可以做到 0 丢失

同 RocketMQ

功能支持

MQ 领域的功能极其完备

基于 erlang 开发,并发能力很强,性能极好,延时很低

MQ 功能较为完善,还是分布式的,扩展性好

功能较为简单,主要支持简单的 MQ 功能,在大数据领域的实时计算以及日志采集被大规模使用

顺序消息

不支持

不支持

支持

支持

消息回溯

不支持

不支持

支持指定分区offset位置的回溯

支持指定分区offset位置的回溯

消息重试

不支持

不支持,但是可以利用消息确认机制实现。

支持

不支持,但是可以利用消息确认机制实现

事务对比

kafka 开始事务时保证了 事务内的发送消息动作要不然一起成功,要不然一起失败,作用域单体应用.

rabiit事务也是保证发送消息同时成功失败.

rocketMQ的事务使用了两阶段提交, 保证本地数据库事务和消息一起成功或失败.

与seata事务相比, 是柔性事务,保持最终一致性. 优点是解耦,缺点有延迟.

文章作者: 刘同学
本文链接:
版权声明: 本站所有文章除特别声明外,均采用 CC BY-NC-SA 4.0 许可协议。转载请注明来自 刘同学的小站
中间件 rocketMQ
喜欢就支持一下吧