小Cの已经记不起来的博客

用 Docker 搭建消息队列的完整流程

起因

前阵子手头有个小项目,A 服务生成任务,B 服务慢慢消化。一开始图省事直接 HTTP 互调,结果 B 那边一卡,A 跟着超时,重试逻辑写了一坨还是不痛快。索性上消息队列。本来打算直接在机器上装个 RabbitMQ,结果看了一眼安装文档,还得先装 Erlang,就这依赖关系,算了,还是别折腾宿主机了。

下面是实测可用的搭建方式,需要借助 docker,不然你就慢慢装 Erlang 吧。

拉镜像

我选的是 RabbitMQ,原因很简单:自带管理界面,出了问题一眼能看出来,对新手友好。Kafka 那套不是不能用,就是日常小项目用不上那么大阵仗。

直接拉:

docker pull rabbitmq:3.13-management

拉得慢的话可以配个加速域名,比如 docker pull docker.1ms.run/library/rabbitmq:3.13-management。不过这种加速域名时效性很强,你用的时候没准儿已经不可用了,报错就把前缀删掉,或者换别的加速站。

注意一定要带 -management 后缀,不带就是纯净版,没有管理界面,第一次用纯命令行排错会很痛苦(别问我怎么知道的)。

起容器

docker run -d \
  --name rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  -e RABBITMQ_DEFAULT_USER=admin \
  -e RABBITMQ_DEFAULT_PASS=admin123 \
  -v ~/rabbitmq/data:/var/lib/rabbitmq \
  --restart unless-stopped \
  rabbitmq:3.13-management

几个参数简单说下:5672 是 AMQP 端口,程序连接走它;15672 是管理界面端口,浏览器访问用;数据目录映射到宿主机,容器删了队列和配置也还在;--restart unless-stopped 让它开机自动拉起,服务器重启不用管。

嫌命令长也可以写成 compose,效果一样:

services:
  rabbitmq:
    image: rabbitmq:3.13-management
    container_name: rabbitmq
    ports:
      - "5672:5672"
      - "15672:15672"
    environment:
      RABBITMQ_DEFAULT_USER: admin
      RABBITMQ_DEFAULT_PASS: admin123
    volumes:
      - ./rabbitmq-data:/var/lib/rabbitmq
    restart: unless-stopped

跑完 docker ps 看一眼,STATUS 是 Up 就成了。

进管理界面看看

浏览器打开 http://服务器IP:15672,用刚才设的账号密码登录。第一次进去基本是空的,别慌。可以先手动建个队列试试:切到 Queues 标签,Add queue,名字随便起比如 task_queue,其他默认,点创建。列表里出现它,说明服务本身没毛病。

对了,云服务器记得去安全组放行 15672 和 5672,不然浏览器转半天圈,还以为是自己哪里配错了,其实压根没进去。

写段代码试试

光看界面不算数,跑个小脚本验证。Python 的话先 pip install pika。

生产者:

import pika

conn = pika.BlockingConnection(pika.ConnectionParameters('127.0.0.1', 5672, '/', pika.PlainCredentials('admin', 'admin123')))
ch = conn.channel()
ch.queue_declare(queue='task_queue', durable=True)
ch.basic_publish(exchange='', routing_key='task_queue', body='hello mq')
conn.close()
print('sent')

消费者:

import pika, time

conn = pika.BlockingConnection(pika.ConnectionParameters('127.0.0.1', 5672, '/', pika.PlainCredentials('admin', 'admin123')))
ch = conn.channel()
ch.queue_declare(queue='task_queue', durable=True)

def cb(ch, method, props, body):
    print('got:', body.decode())
    time.sleep(1)
    ch.basic_ack(delivery_tag=method.delivery_tag)

ch.basic_qos(prefetch_count=1)
ch.basic_consume(queue='task_queue', on_message_callback=cb)
print('waiting...')
ch.start_consuming()

先跑消费者再跑生产者,不出问题的话就没有问题了,消费者终端能打出 got: hello mq 就算通了

评论

还没有评论。

发表评论

提交后评论将经过自动审核,审核通过后公开展示。

未在播放