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()
/
/
/
/