# msg **Repository Path**: ouo9527/msg ## Basic Information - **Project Name**: msg - **Description**: No description available - **Primary Language**: Java - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2020-04-24 - **Last Updated**: 2020-12-19 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # Getting Started ### Reference Documentation For further reference, please consider the following sections: * [Official Apache Maven documentation](https://maven.apache.org/guides/index.html) * [Spring Boot Maven Plugin Reference Guide](https://docs.spring.io/spring-boot/docs/2.2.6.RELEASE/maven-plugin/) * [Spring for RabbitMQ](https://docs.spring.io/spring-boot/docs/2.2.6.RELEASE/reference/htmlsingle/#boot-features-amqp) * [Spring for Apache Kafka](https://docs.spring.io/spring-boot/docs/2.2.6.RELEASE/reference/htmlsingle/#boot-features-kafka) * [Spring Data Redis (Access+Driver)](https://docs.spring.io/spring-boot/docs/2.2.6.RELEASE/reference/htmlsingle/#boot-features-redis) ### Guides The following guides illustrate how to use some features concretely: * [Messaging with RabbitMQ](https://spring.io/guides/gs/messaging-rabbitmq/) * [Messaging with Redis](https://spring.io/guides/gs/messaging-redis/) ### rabbitmq ####启动 1、以应用方式启动(2种) rabbitmq-server -detached 后台启动 Rabbitmq-server 直接启动,如果你关闭窗口或者需要在改窗口使用其他命令时应用就会停止  关闭:rabbitmqctl stop 2、以服务方式启动(安装完之后在任务管理器中服务一栏能看到RabbtiMq) rabbitmq-service install 安装服务 rabbitmq-service start 开始服务 Rabbitmq-service stop  停止服务 Rabbitmq-service enable 使服务有效 Rabbitmq-service disable 使服务无效 rabbitmq-service help 帮助 当rabbitmq-service install之后默认服务是enable的,如果这时设置服务为disable的话,rabbitmq-service start就会报错。 当rabbitmq-service start正常启动服务之后,使用disable是没有效果的   关闭:rabbitmqctl stop 3、Rabbitmq 管理插件启动,可视化界面 rabbitmq-plugins enable rabbitmq_management 启动 rabbitmq-plugins disable rabbitmq_management 关闭 4、Rabbitmq节点管理方式 #####channel(通道) channel.basicReject(deliveryTag, true);         basic.reject方法拒绝deliveryTag对应的消息,第二个参数是否requeue,true则重新入队列,否则丢弃或者进入死信队列。 该方法reject后,该消费者还是会消费到该条被reject的消息。 channel.basicNack(deliveryTag, false, true);         basic.nack方法为不确认deliveryTag对应的消息,第二个参数是否应用于多消息,第三个参数是否requeue,与basic.reject区别就是同时支持多个消息,可以nack该消费者先前接收未ack的所有消息。nack后的消息也会被自己消费到。 channel.basicRecover(true);         basic.recover是否恢复消息到队列,参数是是否requeue,true则重新入队列,并且尽可能的将之前recover的消息投递给其他消费者消费,而不是自己再次消费。false则消息会重新被投递给自己。 在RabbitMQ中消费者有2种方式获取队列中的消息: a) 一种是通过basic.consume命令,订阅某一个队列中的消息,channel会自动在处理完上一条消息之后,接收下一条消息。(同一个channel消息处理是串行的)。除非关闭channel或者取消订阅,否则客户端将会一直接收队列的消息。 b) 另外一种方式是通过basic.get命令主动获取队列中的消息,但是绝对不可以通过循环调用basic.get来代替basic.consume,这是因为basic.get RabbitMQ在实际执行的时候,是首先consume某一个队列,然后检索第一条消息,然后再取消订阅。如果是高吞吐率的消费者,最好还是建议使用basic.consume。 简单总结一下就是说: consume是只要队列里面还有消息就一直取。 get是只取了队列里面的第一条消息。 因为get开销大,如果需要从一个队列取消息的话,首选consume方式,慎用循环get方式。 流程:client(produc) -》exchange -》routingKey(可无) -》queue 《- client(consumer) 注意: default Exchange不能进行Binding,也不需要进行绑定。 除default Exchange之外,其他任何Exchange都需要和Queue进行Binding,否则无法进行消息路由(转发) Binding的时候,可以设置一个或多个参数,其中参数要特别注意参数类型,如果Routing key中指定的参数类型和消息中指定的参数类型不一致(header Exchange)也不能进行消息转发。 Direct Exchange,Topic Exchange进行Binding的时候,需要指定Routing key Fanout Exchange,Headers Exchange进行Binding的时候,不需要指定Routing key。 ####queue(队列) name, //队列名称 durable: false, //队列是否持久化.false:队列在内存中,服务器挂掉后,队列就没了;true:服务器重启后,队列将会重新生成.注意:只是队列持久化,不代表队列中的消息持久化!!!! exclusive: false, //队列是否专属,专属的范围针对的是连接,也就是说,一个连接下面的多个信道是可见的.对于其他连接是不可见的.连接断开后,该队列会被删除.注意,不是信道断开,是连接断开.并且,就算设置成了持久化,也会删除. autoDelete: true, //如果所有消费者都断开连接了,是否自动删除.如果还没有消费者从该队列获取过消息或者监听该队列,那么该队列不会删除.只有在有消费者从该队列获取过消息后,该队列才有可能自动删除(当所有消费者都断开连接,不管消息是否获取完) arguments: null //队列的配置 arguments一共10个:(策略) Message TTL : 消息生存期(即键"x-message-ttl",单位毫秒)(标志:TTL) Auto expire : 队列生存期(即键"x-expires",单位毫秒)(标志:Exp) Max length : 队列可以容纳的消息的最大条数(即键"x-max-length",超过这个条数,队列头部的消息将会被丢弃.)(标志:Lim) Max length bytes : 队列可以容纳的消息的最大字节数(即键"x-max-length-bytes",超过这个字节数,队列头部的消息将会被丢弃.)(标志:Lim B) Overflow behaviour : 队列中的消息溢出后如何处理(即键"x-overflow",其值"drop-head"(默认)或"reject-publish".队列中的消息溢出时,如何处理这些消息.要么丢弃队列头部的消息,要么拒绝接收后面生产者发送过来的所有消息)(标志:Ovfl) Dead letter exchange : 溢出的消息需要发送到绑定该死信交换机的队列(即键"x-dead-letter-exchange",当队列中的消息的生存期到了,或者因长度限制被丢弃时,消息会被推送到(绑定到)这台交换机(的队列中),而不是直接丢掉. )(标志:DLX) Dead letter routing key : 溢出的消息需要发送到绑定该死信交换机,并且路由键匹配的队列(即键"x-dead-letter-routing-key")(标志:DLK) Maximum priority : 最大优先级(即键"x-max-priority",优先级属性的类型为 byte。发布消息的时候,可以指定消息的优先级,优先级高的先被消费.如果没有设置该参数,那么该队列不支持消息优先级功能.也就是说,就算发布消息的时候传入了优先级的值,也不会起什么作用.)(标志:Pri) Lazy mode : 懒人模式(即键"x-queue-mode",其值枚举类型,如"lazy"。在该模式下的队列会先将交换机推送过来的消息(尽可能多的)保存在磁盘上,以减少内存(RAM的占用.当消费者开始消费的时候才加载到内存中;如果没有设置懒人模式,队列则会直接利用内存缓存,以最快的速度传递消息.)(标志:Args) Master locator : 将队列设置为主位置模式,确定在节点集群上声明队列主节点所在的规则。(即键"x-queue-master-locator") 补充:DLX,官方用词 "rejected" 应该被翻译成 : "丢弃"或者"抛弃",而不是"拒绝",也就是说,这里的 rejected 和 expire 是队列里面的消息的状态.而不是队列的动作. 因此,reject-publish不一定会进入死信队列。 DLK,"direct"(很多叫它"直连模式"或者"路由模式",偏向叫“精确匹配模式”) "topic"("主题模式",偏向叫“模糊匹配模式") Pri,除了给队列优先级外,还可以给message优先级,一般value<=10 使用QueueBuilder创建队列 ####Exchanges(交换机) 交换机只是“中转机器”即转发器,不是队列,因此没有存储能力,若消息发到某台交换机,如果这时,该交换机未绑定任何队列,那么消息会丢失。 有三种情况可能进死信交换机 1)被reject或者nack,并且requeue设置为false 2)消息最大存活时间(TTL)超时 3)消息数量超过最大队列长度 如何保证消息的不丢失(不是100%),三个地方做到持久化。 1)Exchange需要持久化(声明时durable=true来实现的) 2)Queue需要持久化(声明时queue的持久化是通过durable=true来实现的) 3)Message需要持久化(channel.basicPublish方法中BasicProperties参数对象的deliveryMode属性实现,其值:1:nonpersistent 2:persistent。可利用MessageProperties工具) 4)对于consumer端来说,如果这时autoAck=true,那么当consumer接收到相关消息之后,还没来得及处理就crash掉了, 那么这样也算数据丢失,这种情况也好处理,只需将autoAck设置为false(方法定义如下),然后在正确处理完消息之后进行手动ack(channel.basicAck). 5)对于produce端来说,采用事务或confirm模式(搭配重试)、设置渠道缓存大小(connectionFactory.setChannelCacheSize(100);) 其次,关键的问题是消息在正确存入RabbitMQ之后,还需要有一段时间(这个时间很短,但不可忽视)才能存入磁盘之中,RabbitMQ并不是为每条消息都做fsync的处理,可能仅仅保存到cache中而不是物理磁盘上, 在这段时间内RabbitMQ broker发生crash, 消息保存到cache但是还没来得及落盘,那么这些消息将会丢失。那么这个怎么解决呢?首先可以引入RabbitMQ的mirrored-queue即镜像队列,这个相当于配置了副本, 当master在此特殊时间内crash掉,可以自动切换到slave,这样有效的保障了HA, 除非整个集群都挂掉,这样也不能完全的100%保障RabbitMQ不丢消息,但比没有mirrored-queue的要好很多,很多现实生产环境下都是配置了mirrored-queue的。 还有要在producer引入事务机制或者Confirm机制来确保消息已经正确的发送至broker端,有关RabbitMQ的事务机制或者Confirm机制可以参考:RabbitMQ之消息确认机制(事务+Confirm). RabbitMQ的可靠性涉及producer端的确认机制、broker端的镜像队列的配置以及consumer端的确认机制,要想确保消息的可靠性越高,那么性能也会随之而降,鱼和熊掌不可兼得,关键在于选择和取舍。 参考:https://blog.csdn.net/jmdonghao/article/details/76153757 消息什么时候刷到磁盘? 写入文件前会有一个Buffer,大小为1M,数据在写入文件时,首先会写入到这个Buffer,如果Buffer已满,则会将Buffer写入到文件(未必刷到磁盘)。 有个固定的刷盘时间:25ms,也就是不管Buffer满不满,每个25ms,Buffer里的数据及未刷新到磁盘的文件内容必定会刷到磁盘。 每次消息写入后,如果没有后续写入请求,则会直接将已写入的消息刷到磁盘:使用Erlang的receive x after 0实现,只要进程的信箱里没有消息,则产生一个timeout消息,而timeout会触发刷盘操作。 RabbitMQ中与事务机制有关的方法有三个: txSelect(), txCommit()以及txRollback(), txSelect用于将当前channel设置成transaction模式,txCommit用于提交事务, txRollback用于回滚事务,在通过txSelect开启事务之后,我们便可以发布消息给broker代理服务器了,如果txCommit提交成功了, 则消息一定到达了broker了,如果在txCommit执行之前broker异常崩溃或者由于其他原因抛出异常,这个时候我们便可以捕获异常通过txRollback回滚事务了。 try { channel.txSelect(); channel.basicPublish(exchange, routingKey, MessageProperties.PERSISTENT_TEXT_PLAIN, msg.getBytes()); int result = 1 / 0; channel.txCommit(); } catch (Exception e) { e.printStackTrace(); channel.txRollback(); } 采用事务机制实现会降低RabbitMQ的消息吞吐量,为了提高吞吐率采用Confirm模式(异步的) confirm模式:https://blog.csdn.net/u013256816/article/details/55515234 https://blog.csdn.net/qq_32880973/article/details/98949242 事务(AMQP)、confirm解决:反馈保证消息可靠性 问题:消息到底有没有正确到达broker代理服务器?和如果在消息到达broker之前已经丢失的 mandatory 当mandatory标志位设置为true时,如果exchange根据自身类型和消息routeKey无法找到一个符合条件的queue,那么会调用basic.return方法将消息返回给生产者(Basic.Return + Content-Header + Content-Body);当mandatory设置为false时,出现上述情形broker会直接将消息扔掉。 immediate 当immediate标志位设置为true时,如果exchange在将消息路由到queue(s)时发现对于的queue上么有消费者,那么这条消息不会放入队列中。当与消息routeKey关联的所有queue(一个或者多个)都没有消费者时,该消息会通过basic.return方法返还给生产者。