如何无需自定义代码从RabbitMQ Stream获取纯MQTT消息
RabbitMQ MQTT扇出模式消息提取问题
目标
使用RabbitMQ实现MQTT消息的扇出模式。
数据流
IoT数据 -> MQTT交换机 -> RabbitMQ Stream -> 多个消费者
部署与验证
本地运行RabbitMQ(5552为Stream端口,1883为MQTT端口):
docker run -it --rm --name rabbitmq -p 5552:5552 -p 15672:15672 -p 5672:5672 -p 1883:1883 \ -e RABBITMQ_SERVER_ADDITIONAL_ERL_ARGS='-rabbitmq_stream advertised_host localhost' \ rabbitmq:4-management docker exec rabbitmq rabbitmq-plugins enable rabbitmq_stream rabbitmq_stream_management docker exec rabbitmq rabbitmq-plugins enable rabbitmq_mqtt
已验证MQTT消息的发布与订阅:
# 发布 ❯ mosquitto_pub -h localhost -p 1883 -t test_mqtt -m "Hello, MQTT" -u guest -P guest # 订阅 ❯ mosquitto_sub -h localhost -p 1883 -t test_mqtt -u guest -P guest Hello, MQTT
问题场景
希望通过订阅RabbitMQ Stream而非MQTT来实现扇出模式,已通过上述mosquitto_pub命令向名为mqtt_stream的Stream发送了3条消息。使用RabbitMQ官方教程中的rstream客户端订阅Stream:
import asyncio from rstream import ( AMQPMessage, Consumer, ConsumerOffsetSpecification, MessageContext, OffsetType, ) STREAM_NAME = "mqtt_stream" # 5GB STREAM_RETENTION = 5000000000 async def receive(): async with Consumer(host="localhost", username="guest", password="guest") as consumer: async def on_message(msg: AMQPMessage, message_context: MessageContext): print("Got message: {} from stream {}".format(msg, message_context.stream)) print("Press control + C to close") await consumer.start() await consumer.subscribe( stream=STREAM_NAME, callback=on_message, offset_specification=ConsumerOffsetSpecification(OffsetType.LAST, None), ) try: await consumer.run() except (KeyboardInterrupt, asyncio.CancelledError): print("Closing Consumer...") return with asyncio.Runner() as runner: runner.run(receive())
运行后控制台输出带二进制前缀的消息:
❯ python receive.py Press control + C to close Got message: b'\x00Sp\xc0\x06\x05B@@AC\x00Sr\xc12\x04\xa3\nx-exchange\xa1\tamq.topic\xa3\rx-routing-key\xa1\ttest_mqtt\x00Su\xa0\x0bHello, MQTT' from stream mqtt_stream
问题
是否无需编写自定义代码,就能从中提取出纯字符串(最终为JSON格式)的消息内容?注:使用rstream直接向Stream发布消息时,相同代码可输出干净的字符串。
内容的提问来源于stack exchange,提问作者Igor L.
相关产品推荐
相关产品推荐

