RabbitMQ 教程 - “Hello world!”
简介
先决条件
本教程假设 RabbitMQ 已 安装 并在 localhost 上的 标准端口 (5672) 上运行。如果您使用不同的主机、端口或凭据,则需要调整连接设置。
哪里寻求帮助
如果您在学习本教程时遇到困难,可以通过 GitHub Discussions 或 RabbitMQ 社区 Discord 联系我们。
RabbitMQ 是一个消息代理:它接收并转发消息。你可以把它想象成一个邮局:当你把想要寄出的邮件放进邮箱时,你可以确信邮递员最终会将邮件送到你的收件人手中。在这个比喻中,RabbitMQ 就是邮箱、邮局和邮递员。
RabbitMQ 和邮局的主要区别在于,它不处理纸张,而是接收、存储和转发数据的二进制数据块——消息。
RabbitMQ 和消息传递通常会使用一些行话。
-
发布(Producing)指的就是发送。发送消息的程序被称为生产者(producer)。
-
队列 是 RabbitMQ 中邮箱的名称。虽然消息在 RabbitMQ 和你的应用程序之间流动,但它们只能存储在队列 中。队列 的大小仅受主机内存和磁盘空间的限制,它本质上是一个大的消息缓冲区。
许多生产者 可以发送消息到同一个队列,许多消费者 可以尝试从同一个队列 接收数据。
这就是我们表示队列的方式
-
消费 的含义与接收类似。消费者 是一个主要等待接收消息的程序。
请注意,生产者、消费者和代理不一定需要驻留在同一台主机上;事实上,在大多数应用程序中它们都不会。一个应用程序也可以同时是生产者和消费者。
Hello World!
(使用 Pika Python 客户端)
在本教程的这一部分,我们将编写两个小程序(使用 Python):一个生产者(发送者)发送单条消息,一个消费者(接收者)接收消息并将其打印出来。这就是消息传递领域的“Hello World”。
在下面的图表中,“P”是我们的生产者,“C”是我们的消费者。中间的框是一个队列——RabbitMQ 代表消费者存储的消息缓冲区。
我们的整体设计如下所示:
生产者将消息发送到“hello”队列。消费者从该队列接收消息。
RabbitMQ 库
RabbitMQ 支持多种协议。本教程使用 AMQP 0-9-1,这是一种开放的、通用的消息传递协议。RabbitMQ 有许多适用于 不同语言 的客户端。在本系列教程中,我们将使用 Pika 1.0.0,这是 RabbitMQ 团队推荐的 Python 客户端。你可以使用
pip包管理工具来安装它。python -m pip install pika --upgrade
现在我们已经安装了 Pika,可以开始编写代码了。
发送
我们的第一个程序 send.py 将向队列发送单条消息。我们要做的第一件事是与 RabbitMQ 服务器建立连接。
#!/usr/bin/env python
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
我们现在已经连接到了本地机器上的代理(Broker)——这就是 localhost 的由来。如果我们想连接到另一台机器上的代理,只需在此处指定其名称或 IP 地址即可。
接下来,在发送之前,我们需要确保接收方队列存在。如果我们向一个不存在的位置发送消息,RabbitMQ 会直接丢弃该消息。让我们创建一个 hello 队列,消息将发送到该队列。
channel.queue_declare(queue='hello', durable=True, arguments={'x-queue-type': 'quorum'})
此时,我们准备好发送消息了。我们的第一条消息将包含字符串 Hello World!,我们希望将其发送到我们的 hello 队列。
在 RabbitMQ 中,消息永远不能直接发送到队列,它总是需要经过一个 交换机(exchange)。不过我们先不纠结于这些细节 —— 你可以在 本教程的第三部分 中阅读更多关于 交换机 的内容。现在我们只需要知道如何使用由空字符串标识的默认交换机。这个交换机很特殊 —— 它允许我们明确指定消息应该发送到哪个队列。队列名称需要在 routing_key 参数中指定。
channel.basic_publish(exchange='',
routing_key='hello',
body='Hello World!')
print(" [x] Sent 'Hello World!'")
在退出程序之前,我们需要确保网络缓冲区已刷新,并且消息确实已传递给 RabbitMQ。我们可以通过优雅地关闭连接来实现这一点。
connection.close()
发送不起作用!
如果这是您第一次使用 RabbitMQ,但没有看到“Sent”消息,您可能会感到困惑,不知道哪里出了问题。也许代理启动时磁盘空间不足(默认需要至少 50 MB 可用空间),因此拒绝接收消息。检查代理的日志文件,看看是否有资源警报已记录,并在必要时降低可用磁盘空间阈值。配置指南将展示如何设置
disk_free_limit。
接收
我们的第二个程序 receive.py 将从队列中接收消息并将其打印在屏幕上。
同样,首先我们需要连接到 RabbitMQ 服务器。与 Rabbit 建立连接的代码与之前相同。
下一步,就像之前一样,确保队列存在。使用 queue_declare 创建队列是幂等的 —— 我们可以根据需要多次运行该命令,但最终只会创建一个队列。
channel.queue_declare(queue='hello', durable=True, arguments={'x-queue-type': 'quorum'})
你可能会问为什么要再次声明队列 —— 我们已经在之前的代码中声明过了。如果我们确定队列已经存在(例如 send.py 程序之前已经运行过),我们就可以省略这一步。但我们还不确定哪个程序会先运行。在这种情况下,在两个程序中重复声明队列是一种良好的实践。
列出队列
您可能希望查看 RabbitMQ 有哪些队列以及其中有多少消息。您可以使用
rabbitmqctl工具(作为特权用户)来完成此操作。sudo rabbitmqctl list_queues在 Windows 上,省略 sudo。
rabbitmqctl.bat list_queues
从队列中接收消息比较复杂。它通过将一个 callback(回调)函数订阅到队列来工作。每当我们收到一条消息,Pika 库就会调用这个 callback 函数。在我们的例子中,该函数会将消息内容打印在屏幕上。
def callback(ch, method, properties, body):
print(f" [x] Received {body}")
接下来,我们需要告诉 RabbitMQ,这个特定的回调函数应该从我们的 hello 队列接收消息。
channel.basic_consume(queue='hello',
auto_ack=True,
on_message_callback=callback)
为了使该命令成功,我们必须确保我们想要订阅的队列存在。幸运的是,我们已经通过 queue_declare 在上面创建了队列,对此非常有把握。
auto_ack 参数将在 稍后 说明。
最后,我们进入一个永无止境的循环,等待数据并在必要时运行回调,并在程序关闭时捕获 KeyboardInterrupt。
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
if __name__ == '__main__':
try:
main()
except KeyboardInterrupt:
print('Interrupted')
try:
sys.exit(0)
except SystemExit:
os._exit(0)
总而言之
send.py (源代码)
#!/usr/bin/env python
import pika
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost'))
channel = connection.channel()
channel.queue_declare(queue='hello', durable=True, arguments={'x-queue-type': 'quorum'})
channel.basic_publish(exchange='', routing_key='hello', body='Hello World!')
print(" [x] Sent 'Hello World!'")
connection.close()
receive.py (源代码)
#!/usr/bin/env python
import pika, sys, os
def main():
connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost'))
channel = connection.channel()
channel.queue_declare(queue='hello', durable=True, arguments={'x-queue-type': 'quorum'})
def callback(ch, method, properties, body):
print(f" [x] Received {body}")
channel.basic_consume(queue='hello', on_message_callback=callback, auto_ack=True)
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
if __name__ == '__main__':
try:
main()
except KeyboardInterrupt:
print('Interrupted')
try:
sys.exit(0)
except SystemExit:
os._exit(0)
现在我们可以在终端中尝试我们的程序。首先,让我们启动消费者,它将持续运行并等待消息投递:
python receive.py
# => [*] Waiting for messages. To exit press CTRL+C
现在在新的终端中启动生产者。生产者程序会在每次运行后自动停止:
python send.py
# => [x] Sent 'Hello World!'
消费者将打印出消息:
# => [*] Waiting for messages. To exit press CTRL+C
# => [x] Received 'Hello World!'
万岁!我们成功通过 RabbitMQ 发送了第一条消息。如你所见,receive.py 程序并没有退出。它将保持准备状态以接收更多消息,你可以通过 Ctrl-C 中断它。
尝试在新的终端中再次运行 send.py。
我们已经学会了如何从命名队列发送和接收消息。现在是时候进入 第 2 部分,构建一个简单的 工作队列(work queue) 了。