RabbitMQ,Kafka与RPC原理,

RabbitMQ,Kafka与RPC原理,

参考连接:

  https://www.rabbitmq.com/getstarted.html

rabbitmq默认端口:5672

 笔记整理:

1# -*- coding: utf-8 -*- 2# __author__ = "maple" 3"""""" 4 5# 1. 你了解的消息队列 6""" 7 - Queue,将数据存储当前服务器的内存. 8 - redis 列表, 9 - rabbitMQ/kafka/zeroMQ(专业做消息队列) 10 11 补充:saltstack 12 - ssh:安装方便,但执行效率慢. 13 - agent:执行效率高(基于消息队列zeroMQ做的RPC). 14""" 15 16# 2. 公司在什么情况下会用消息队列? 17""" 18 任务处理,请求数量太多,需要把消息临时放到某个地方. 19 发布订阅,一旦发布消息,所有订阅者都会收到一条相同的消息. 20 21 应用场景: 22 - 长轮询 23 - 智能玩具调用百度AI接口时,celery + RabbitMQ 24 - 生产者&消费者 25""" 26 27# 3. rabbitMQ安装 28""" 29 服务端: 192.168.19.14 30 安装配置epel源 31 $ rpm -ivh http://dl.fedoraproject.org/pub/epel/6/i386/epel-release-6-8.noarch.rpm 32 安装erlang 33 $ yum -y install erlang 34 安装RabbitMQ 35 $ yum -y install rabbitmq-server 36 启动(无用户名密码): 37 service rabbitmq-server start/stop 38 39 设置用户密码: 40 sudo rabbitmqctl add_user xxx pwd # 设置用户名跟密码 41 # 设置用户为administrator角色 42 sudo rabbitmqctl set_user_tags wupeiqi administrator 43 # 设置权限 44 sudo rabbitmqctl set_permissions -p "/" root ".*" ".*" ".*" #此处的 * 代表对rabbitmq的所有队列都有权限. 45 46 service rabbitmq-server start/stop 47 客户端: 48 pip3 install pika 49""" 50 51# 4. 使用 52""" 53 生产者消费者 54 n VS 1 55 n VS m 56 发布订阅 57 fanout,和exchange关联的所有队列都会接收到信息. 58 direct,关键字精确匹配exchange关联的队列都会接收到信息. 59 topic,关键字模糊匹配exchange关联的队列都会接收到信息. 60""" 61 62# 5. exchange是什么? 63""" 64 消息处理的重建建,可以帮助生成者将相关信息发送到指定相关队列. 65""" 66 67# 6. RPC 68""" 69 前戏: 70-> 去哪儿 -> 首都机场票务中心 71 远程过程调用. 72 -> 去哪儿 任务/结果 首都机场票务中心 73 ... 74"""

rabbitmq:在数据安全方面(保证数据不丢失比较擅长)
1.A向B发了一个请求,B把数据放到rabbitmq里了,然后B去rabbitmq里面取数据去处理,
还没有处理完B挂掉了,rabbitmq可以保证数据不丢.
2.A向B发了一个请求,B还没有把数据放到rabbitmq里了,B挂掉了,此时数据就会丢失,跟rabbitmq还没有关系呢;
3.A向B发了一个请求,B把数据放到rabbitmq里了,B还没有去rabbitmq取数据,此时rabbitmq挂掉了,此时数据不会丢失;

即:客户端挂掉或者服务端挂掉,数据都不会丢失的,

kafka:
在分布式以及消息的传递更快一些,消息的存储/发送要比rabbitmq快,

zeromq:集成在saltstack里面,

minion在监听队列,有数据的时候就去取,然后执行命令,执行结果放到另一个队列里面,

执行效率要比rabbitmq跟kafka都要高,

为什么要使用消息队列:
当前的服务器或者应用处理不了太多的数据,需要先把数据放到某个地方(队列),然后慢慢去处理.

代码: 生产者

1# -*- coding: utf-8 -*- 2 3import pika 4# 无密码 5# connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14')) # 先建一个链接 6 7# 有密码 8credentials = pika.PlainCredentials("root","123") 9connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 10channel = connection.channel() # 创建一个对象 11# 声明一个队列(创建一个队列) 12channel.queue_declare(queue='duilie_1') 13 14channel.basic_publish(exchange='', 15 routing_key='duilie_1', # 消息队列名称 16 body='msg7') # 发送的数据 17connection.close() # 记得关闭链接

/

消费者:

1# -*- coding: utf-8 -*- 2import pika 3 4credentials = pika.PlainCredentials("root","123") 5connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 6channel = connection.channel() 7 8# 声明一个队列(创建一个队列) 9channel.queue_declare(queue='s13q1') 10 11def callback(ch, method, properties, body): 12 print("消费者接受到了任务: %r" % body) 13 14channel.basic_consume(callback,queue='s13q1',no_ack=True) # 有消息立即执行callback函数,没有消息就夯住了, 15 16channel.start_consuming() # 开始进行消费.

 //

注意:

11.声明一个队列,如果存在就使用,不存在就创建一个; 2 3有消息会立即执行函数,没有消息就会夯住了, 4 5可以有多个生产者,多个消费者(代表多个人处理同一个队列) 6 7默认处理队列,多个消费者的时候,是轮流来取数据的,

ack回复:带回复功能,针对客户端挂掉之后来做的.(客户端挂掉,数据的保留.服务端进行的数据保存)

ack就是来确保数据安全的,从服务端取走数据之后,如果开启回复的功能,需要等客户端回复,服务端才会把数据清掉,

  假如客户端取走数据之后,没来及处理,挂掉了,也就是没有回复服务端,服务端会保留数据的.

1针对消费者进行的改动,生产者不用改动.import pika 2 3credentials = pika.PlainCredentials("root","123") 4connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 5channel = connection.channel() 6 7# 声明一个队列(创建一个队列) 8channel.queue_declare(queue='s13q2') 9 10def callback(ch, method, properties, body): 11 print("消费者接受到了任务: %r" % body) 12 # int('asdfasdf') 13 ch.basic_ack(delivery_tag=method.delivery_tag) # 此行代码代表告诉服务端,数据已经取走了,只要执行到此,服务端的数据就清掉了. 14 15channel.basic_consume(callback,queue='s13q2',no_ack=False) # 此处的no_ack=False,就是开启回复 16 17channel.start_consuming()

//

durable持久化,服务端挂掉,数据的保留.

在生产者处进行的持久化,

1生产者进行的改动,消费者不用改动.import pika 2# 无密码 3# connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14')) 4 5# 有密码 6credentials = pika.PlainCredentials("root","123") 7connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 8channel = connection.channel() 9# 声明一个队列(创建一个队列)-支持持久化 10channel.queue_declare(queue='s13q3',durable=True) # 标红加粗是新加的参数.另外,一个队列如果建立时没有进行持久化,后续加参数是不管用. 11 12channel.basic_publish(exchange='', 13 routing_key='s13q3', # 消息队列名称 14 body='msg7', 15 properties=pika.BasicProperties( 16 delivery_mode = 2, # make message persistent 17 )) # 标红的是新加的参数. 18connection.close()

/

避免排队,充分利用消费者

1生产者不用动,消费者进行修改.import pika 2 3credentials = pika.PlainCredentials("root","123") 4connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 5channel = connection.channel() 6 7# 声明一个队列(创建一个队列) 8channel.queue_declare(queue='s13q1') 9 10def callback(ch, method, properties, body): 11 print("消费者接受到了任务: %r" % body) 12 13channel.basic_qos(prefetch_count=1) # 此处是进行的配置,哪个闲置,哪个就进行额外的多工作. 14channel.basic_consume(callback,queue='s13q1',no_ack=True) 15 16channel.start_consuming()

 /

 发布订阅

1发布订阅的时候,队列还是要随机生成的, 2 3生成一个exchange(exchange只要创建了就不会再创建,),消费者脚本只要运行一个,就生成一个队列

消费者_订阅者

1import pika 2 3credentials = pika.PlainCredentials("root","123") 4connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 5channel = connection.channel() 6 7# exchange='m1',exchange(秘书)的名称 8# exchange_type='fanout' , 秘书工作方式将消息发送给所有的队列 9channel.exchange_declare(exchange='m1',exchange_type='fanout') #exchange='m1',是随意起的名字(和生产者的要对应上) fanout就是所有队列都放一份任务,  代指的是消息传送的模式 10 11# 随机生成一个队列, 12result = channel.queue_declare(exclusive=True) # 随机的生成一个队列 13queue_name = result.method.queue # 队列名称 14# 让exchange和queue进行绑定. 15channel.queue_bind(exchange='m1',queue=queue_name) 16 17 18def callback(ch, method, properties, body): 19 print("消费者接受到了任务: %r" % body) 20 21channel.basic_consume(callback,queue=queue_name,no_ack=True) 22 23channel.start_consuming()

/

生产者_发布者

1import pika 2credentials = pika.PlainCredentials("root","123") 3connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 4channel = connection.channel() 5 6channel.exchange_declare(exchange='m1',exchange_type='fanout') # 指定 exchange 7 8channel.basic_publish(exchange='m1', 9 routing_key='', # 此处routing_key是空的, 10 body='gb') 11 12connection.close()

//

/

/

 带关键字的发布订阅(相当于指定对哪个队列进行发布)

 消费者_订阅者1

1import pika 2 3credentials = pika.PlainCredentials("root","123") 4connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 5channel = connection.channel() 6 7# exchange='m1',exchange(秘书)的名称 8# exchange_type='fanout' , 秘书工作方式将消息发送给所有的队列 9channel.exchange_declare(exchange='m2',exchange_type='direct') # 注意exchange_type的类型要变一下 10 11# 随机生成一个队列 12result = channel.queue_declare(exclusive=True) 13queue_name = result.method.queue 14# 让exchange和queque进行绑定. 15channel.queue_bind(exchange='m2',queue=queue_name,routing_key='test1') # 此处指定关键字,就是队列的名字 16channel.queue_bind(exchange='m2',queue=queue_name,routing_key='test2') 17 18 19def callback(ch, method, properties, body): 20 print("消费者接受到了任务: %r" % body) 21 22channel.basic_consume(callback,queue=queue_name,no_ack=True) 23 24channel.start_consuming()

/

消费者_订阅者2

1import pika 2 3credentials = pika.PlainCredentials("root","123") 4connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 5channel = connection.channel() 6 7# exchange='m1',exchange(秘书)的名称 8# exchange_type='fanout' , 秘书工作方式将消息发送给所有的队列 9channel.exchange_declare(exchange='m2',exchange_type='direct') 10 11# 随机生成一个队列 12result = channel.queue_declare(exclusive=True) 13queue_name = result.method.queue 14 15# 让exchange和queque进行绑定. 16channel.queue_bind(exchange='m2',queue=queue_name,routing_key='test1') 17 18 19def callback(ch, method, properties, body): 20 print("消费者接受到了任务: %r" % body) 21 22channel.basic_consume(callback,queue=queue_name,no_ack=True) 23 24channel.start_consuming()

/

生产者_发布者

1import pika 2credentials = pika.PlainCredentials("root","123") 3connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 4channel = connection.channel() 5 6channel.exchange_declare(exchange='m2',exchange_type='direct') 7 8channel.basic_publish(exchange='m2', 9 routing_key='test1', # 此处有哪个关键字, 绑定关键字的消费者就能拿到数据 10 body='xxxx') 11 12connection.close()

/

 关键字的模糊匹配

 生产者_发布

1import pika 2credentials = pika.PlainCredentials("root","123") 3connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 4channel = connection.channel() 5 6channel.exchange_declare(exchange='m3',exchange_type='topic') # 注意:exchange_type的类型又变了 7 8channel.basic_publish(exchange='m3', 9 routing_key='old.alex.py', # 队列名称 10 body='x1') 11 12connection.close()

/

消费者_订阅者_1

1import pika 2 3credentials = pika.PlainCredentials("root","123") 4connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 5channel = connection.channel() 6 7# exchange='m1',exchange(秘书)的名称 8# exchange_type='fanout' , 秘书工作方式将消息发送给所有的队列 9channel.exchange_declare(exchange='m3',exchange_type='topic') 10 11# 随机生成一个队列 12result = channel.queue_declare(exclusive=True) 13queue_name = result.method.queue 14# 让exchange和queque进行绑定. 15channel.queue_bind(exchange='m3',queue=queue_name,routing_key='old.*') # *是匹配一个单词,#是可以匹配多个单词的, 16 17 18def callback(ch, method, properties, body): 19 print(method.routing_key) 20 print("消费者接受到了任务: %r" % body) 21 22channel.basic_consume(callback,queue=queue_name,no_ack=True) 23 24channel.start_consuming()

/

消费者_订阅者_2

1import pika 2 3credentials = pika.PlainCredentials("root","123") 4connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 5channel = connection.channel() 6 7# exchange='m1',exchange(秘书)的名称 8# exchange_type='fanout' , 秘书工作方式将消息发送给所有的队列 9channel.exchange_declare(exchange='m3',exchange_type='topic') 10 11# 随机生成一个队列 12result = channel.queue_declare(exclusive=True) 13queue_name = result.method.queue 14# 让exchange和queque进行绑定. 15channel.queue_bind(exchange='m3',queue=queue_name,routing_key='old.#') 16 17 18def callback(ch, method, properties, body): 19 print("消费者接受到了任务: %r" % body) 20 21channel.basic_consume(callback,queue=queue_name,no_ack=True) 22 23channel.start_consuming()

/

 rabbitmq实现的rpc

rpc:远程过程调用

1rpc:远程过程调用. 2 3例子: 4就是去哪网,携程网,都去票务中心订票, 5他们把任务放到一个队列里面(q1队列),同时建立一个自己的队列(q_qn,q_xc)票务中心监控队列(q1),有任务就去处理, 6处理完,把对应的结果放到对应的队列里面(q_qn,q_xc,), 7去哪网,携程网,在往队列里面放任务时, 就把自己创建的队列带过去了,票务中心知道自己该往哪个队列里面放结果.

/

去哪网

1import pika 2import uuid 3 4class FibonacciRpcClient(object): 5 def __init__(self): 6 credentials = pika.PlainCredentials("root", "123") 7 self.connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14', credentials=credentials)) 8 self.channel = self.connection.channel() 9 10 # 随机生成一个消息队列(用于接收结果) 11 result = self.channel.queue_declare(exclusive=True) 12 self.callback_queue = result.method.queue 13 14 # 监听消息队列中是否有值返回,如果有值则执行 on_response 函数(一旦有结果,则执行on_response) 15 self.channel.basic_consume(self.on_response, no_ack=True,queue=self.callback_queue) 16 17 def on_response(self, ch, method, props, body): 18 if self.corr_id == props.correlation_id: 19 self.response = body 20 21 def call(self, n): 22 self.response = None 23 self.corr_id = str(uuid.uuid4()) 24 25 # 去哪网 给 票务中心 发送一个任务: 任务id = corr_id / 任务内容 = '30' / 用于接收结果的队列名称 26 self.channel.basic_publish(exchange='', 27 routing_key='rpc_queue', # 票务中心接收任务的队列名称 28 properties=pika.BasicProperties( 29 reply_to = self.callback_queue, # 用于接收结果的队列 30 correlation_id = self.corr_id, # 任务ID 31 ), 32 body=str(n)) 33 while self.response is None: 34 self.connection.process_data_events() 35 36 return self.response 37 38fibonacci_rpc = FibonacciRpcClient() 39 40response = fibonacci_rpc.call(50) 41print('返回结果:',response)

/

票务中心

1import pika 2credentials = pika.PlainCredentials("root","123") 3connection = pika.BlockingConnection(pika.ConnectionParameters('192.168.19.14',credentials=credentials)) 4channel = connection.channel() 5 6# 票务中心监听任务队列 7channel.queue_declare(queue='rpc_queue') 8 9def on_request(ch, method, props, body): 10 n = int(body) 11 response = n + 100 12 # props.reply_to 要放结果的队列. 13 # props.correlation_id 任务 14 ch.basic_publish(exchange='', 15 routing_key=props.reply_to, 16 properties=pika.BasicProperties(correlation_id= props.correlation_id), 17 body=str(response)) 18 ch.basic_ack(delivery_tag=method.delivery_tag) 19 20channel.basic_qos(prefetch_count=1) 21channel.basic_consume(on_request, queue='rpc_queue') 22channel.start_consuming()

/

/

/

/

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )