☰
Pulsar消息队列与STM32环境监测:从MQTT到分层存储的完整数据链路
2026/9/29 3:54:09 网站建设 项目流程

COSCon‘25 现场最吸睛的标语之一,就是 Pulsar 社区挂出的那句 "Make MQ Great Again!"。说实话,第一眼看到我多少有点想笑:消息队列这种存在了二十多年的"老技术",凭什么喊出"再次伟大"?但等我在 COSCon‘25 x Pulsar Developer Day 2025 泡了一整天,这个想法彻底变了。MQ 不但没有过时,反而在云原生、边缘计算和智能硬件的夹击下长出了新形态。这篇文章不打算写成官方通稿,我就以一个现场开发者的视角,聊聊这场活动的氛围、Pulsar 的几个关键技术点,以及它和最近大家到处搜的 STM32 环境监测系统(DHT11、BH1750、MQ-2、OLED)之间,一条真实可复现的 MQ 数据链路。

如果你正在选型消息队列,或者你手上刚好有一套 STM32 环境监测的板子不知道怎么把数据往上送,这篇文章应该都值得看完。前半部分偏架构理解,后半部分偏实际操作,我尽量把"为什么这样做"也讲清楚。

1. 会场速写:一场把 MQ 当主角的开发者聚会

先聊聊场子。COSCon 本身是国内开源圈一年一度的大聚会,主题跨度非常大,从大前端到操作系统、从 AI 到硬件都有涉猎。而 Pulsar Developer Day 这种垂直技术日叠加在 COSCon 里,本身就是一个很有意思的信号:消息队列领域已经意识到,只在小圈子里自嗨是不够的,必须去跟更广泛的开源生态碰面。

现场的人群也印证了这点。除了常见的后端工程师和架构师,我注意到有相当多从事 IoT/嵌入式开发的人——很多人带着板卡、传感器模块进出,聊的是"数据怎么传上来"。想想也合理:设备端数据要进云端,中间总得有个消息管道,MQ 天然就是这个环节的主角。

1.1 为什么把 COSCon 和 Pulsar Day 放在一起

两个活动叠在一起,主办方的意图其实很明确:Pulsar 作为一个云原生消息平台,它的使用者不止来自传统互联网后端,还有大量物联网、数据集成、实时数仓场景,而 COSCon 恰好能覆盖这些人群。反过来,开源社区也需要一个"大场子"来吸引更多潜在贡献者。

从现场反馈看,这种玩法效果不错。许多原本只逛硬件展区的开发者,被"消息队列怎么跟传感器数据结合"这类议题吸引进了 Pulsar 专场;而一些后端开发者,也在硬件区第一次摸到了真实的传感器模块。这种双向流动,比各自闷头开技术分享要有意思得多。

1.2 现场大家最关心什么

我粗略记了一下现场聊得最多的话题,基本集中在几类:第一,生产环境从 Kafka 迁移到 Pulsar 的坑;第二,多租户隔离和配额管理怎么做;第三,Pulsar 的存储成本和分层卸载(Tiered Storage)实际效果;第四,MQTT 这类轻量协议怎么和 Pulsar 对接,把整套能力延伸到设备端。

这个顺序很有信息量。前两个问题说明 Pulsar 已经从小众尝鲜进入了规模化生产阶段,大家关心的是"用得起、管得住";后两个问题说明边缘接入正在成为新的增长点。整场听下来我有个很深的感受:消息队列技术栈没有"凉",它只是从单机时代的简单工具,变成了平台化时代的基础设施。

2. MQ 选型新逻辑:Pulsar 在 2025 年的差异化竞争力

在聊 Pulsar 的技术细节之前,先捋一下 2025 年 MQ 选型的大环境。这可能是不少人最纠结的部分:Kafka 用了好多年,说要换;RocketMQ 在 Java 生态里也很成熟;RabbitMQ 在老业务系统里根深蒂固;Pulsar 这几年又持续出现在各种案例分享中。到底怎么选?

我个人的判断是,2025 年的选型逻辑和五年前完全不一样了。五年前大家主要看吞吐量和可用性;现在更多看多租户管理、存储成本、云原生适配和协议接入广度。吞吐量早就不是瓶颈,管理和成本反而成了大问题。

2.1 一张表看清主流 MQ 的差异

对比项KafkaRabbitMQRocketMQPulsar
架构模型存储与计算耦合传统 AMQP 代理存储与计算耦合存储与计算分离
多租户能力偏弱,靠集群隔离弱,靠 vhost 逻辑隔离一般,靠 Topic 隔离原生多租户,命名空间级隔离
典型场景日志、流处理、大数据业务解耦、异步任务Java 生态、事务消息统一消息与流、多团队共享集群
运维复杂度中高(扩容要迁移分区)低中中高,但扩容手段更灵活
存储成本优化依赖磁盘扩容一般一般分层存储,可卸载到对象存储

这张表肯定不完全精确,因为各项目版本和社区迭代都有差异,但大方向是准的。值得一提的是,不少团队在生产环境里其实是混用的:业务解耦用 RabbitMQ 或 RocketMQ,大规模流数据用 Kafka,而 Pulsar 往往出现在"既要队列又要流,还要多团队共用一套集群"的场景。它不是要取代谁,而是填补了一个之前没人做好过的位置。

2.2 Pulsar 适合谁、不适合谁

先说适合的。如果你的公司有多条业务线、多套环境,想共用一套消息集群但又要互相隔离配额;如果你的消息量有明显的波峰波谷,希望靠分层存储降低成本;如果你想同时支持队列模式和流模式、不想维护两套系统——Pulsar 会是一个很值得验证的选项。

再说说不适合的。如果你的场景就是一个单体应用内部异步解耦,消息量也不大,那 RabbitMQ 甚至 Redis Stream 都够用,没必要引入 Pulsar 的运维复杂度。如果团队没有专门的中间件运维能力,跑一个分布式 Pulsar 集群其实挺吃力的。工具没有绝对的好坏,只有合不合适。

3. 现场技术笔记:Pulsar 分层架构里的三个关键设计

聊 Pulsar 绕不开它的分层架构。这部分的现场分享密度很高,我把对自己最有启发的三个设计单独记了出来,每个都讲清楚"为什么它是这么设计的"。

3.1 存储与计算分离:把 Broker 变成"无状态"

Pulsar 最核心的设计,是把服务层(Broker)和存储层(BookKeeper)分开。Broker 负责处理生产和消费请求、管理订阅游标,本身不保存消息数据;所有消息都写入 BookKeeper 集群。

这个设计带来的直接好处是:Broker 可以随时扩容、缩容,甚至故障后拉起新节点,不需要做数据迁移,因为它没有本地状态。你可以把它类比成"厨房和仓库分开":炒菜的大厨不用自己囤货,缺物资了随时调货,换个大厨也不影响仓库里的东西。Kafka 早期被人吐槽最多的扩容要迁移分区数据、耗时还容易出问题,Pulsar 用分层架构从一开始就绕开了这个痛点。

3.2 Segment 为中心的存储与分层卸载

BookKeeper 里,消息被分成一个个 Segment(片段)追加存储,分布到多台 Bookie 节点上。这不仅仅是"分片存储"这么简单,它意味着消息的存储位置可以灵活调度,也天然支持副本冗余和故障恢复。

在此基础上,Pulsar 推出了分层存储(Tiered Storage):把老旧的 Segment 自动卸载到对象存储(比如 S3、OSS、MinIO 这类兼容 S3 的服务)上,Broker 需要消费历史数据时再从对象存储里读回来。

这个设计对成本的影响是实打实的。消息数据往往是越新越热、越老越冷,但传统 MQ 不管冷热统统放在本地磁盘上,容量和成本都很难受。有了分层卸载,热数据留在 BookKeeper 保证低延迟,冷数据进对象存储,存储成本能降一个数量级。现场有分享嘉宾给出的生产数据是:启用分层存储后,大分区集群的存储成本降到原来的三分之一以下。具体数字因场景而异,但这个省钱方向是确定的。

3.3 多租户、命名空间与订阅模型

Pulsar 的话题层级是 Tenant(租户)→ Namespace(命名空间)→ Topic(主题)。租户之间可以做认证、配额、存储隔离,命名空间里可以单独配置消息保留策略、备份策略、限流阈值。

多租户能力为什么在 2025 年尤其重要?因为很多公司的消息集群是多个团队共享的。没有租户隔离时,一个团队把 Topic 打到爆,全集群都跟着倒霉;有了租户和命名空间级别的配额管理,各个业务线互相"老死不相往来",运维也不用整天居中协调。

订阅模型也是 Pulsar 的一个记忆点。它同时支持 Exclusive(独占)、Shared(共享)、Failover(故障转移)、Key_Shared(按键共享)四种订阅方式。Exclusive 保证消息严格有序,但只有一个消费者;Shared 允许多消费者负载均衡,但失去全局顺序;Key_Shared 则在保持同 key 消息有序的前提下,把不同 key 分发给不同消费者,这在订单、设备数据这类需要按维度保序的场景里非常实用。

4. 从服务端到传感器:STM32 环境监测里真实发生的 MQ 链路

说完 Pulsar 本身,我想把话题拉回到一个最近热度很高的方向:STM32 环境监测系统。这段时间很多人都在搜 "stm32 环境监测系统 DHT11 BH1750 MQ-2 OLED",恰好我也是做这类东西的,而且我发现这套硬件方案和 MQ 技术栈天然是一对。这里的 MQ 有两个意思:一个是消息队列 Message Queue,一个是气体传感器模块 MQ-2。标题里"Make MQ Great Again"的 MQ,放在嵌入式场景里,刚好可以一语双关。

4.1 一套典型环境监测板子怎么组成

先列一下这套系统最常见的硬件搭配及其数据接口:

模块作用接口输出数据
STM32 主控采集与逻辑处理GPIO/ADC/I2C/USART—
DHT11温湿度测量单总线8bit 湿度整数 + 8bit 湿度小数 + 8bit 温度整数 + 8bit 温度小数 + 8bit 校验和
BH1750光照强度测量I2C1~65535 lx,16bit 光强
MQ-2可燃气体/烟雾检测ADC模拟电压,可换算为气体浓度相关阻值比
OLED(SSD1306)本地显示I2C(通常 0x3C)温湿度、光照、气体报警状态

这套组合非常有代表性:DHT11 便宜但时序敏感,BH1750 是标准 I2C 从设备,MQ-2 走 ADC 模拟量,OLED 则是嵌入式展示的经典外设。把四种不同接口的传感器都跑通,基本就把 STM32 外设操作过了一遍。

4.2 采集端代码骨架

DHT11 的单总线时序很容易踩坑。读一次数据,主机要先拉低总线至少 18ms 触发,然后释放总线,等待 DHT11 响应(80us 低 + 80us 高),之后每位数据按"50us 低电平 + 26~28us 高电平表示 0,70us 高电平表示 1"来解析。用 HAL 库写的话,关键代码长这样:

uint8_t dht11_read_data(uint8_t *humidity, uint8_t *temperature) { uint8_t data[5] = {0}; // 1. 触发信号:拉低 20ms 再释放 HAL_GPIO_WritePin(DHT11_GPIO_Port, DHT11_Pin, GPIO_PIN_RESET); HAL_Delay(20); HAL_GPIO_WritePin(DHT11_GPIO_Port, DHT11_Pin, GPIO_PIN_SET); // 2. 等待响应:先等总线被拉低,再等被拉高 while (HAL_GPIO_ReadPin(DHT11_GPIO_Port, DHT11_Pin) == GPIO_PIN_SET); while (HAL_GPIO_ReadPin(DHT11_GPIO_Port, DHT11_Pin) == GPIO_PIN_RESET); while (HAL_GPIO_ReadPin(DHT11_GPIO_Port, DHT11_Pin) == GPIO_PIN_SET); // 3. 读取 40 位数据(5 字节) for (int i = 0; i < 40; i++) { while (HAL_GPIO_ReadPin(DHT11_GPIO_Port, DHT11_Pin) == GPIO_PIN_RESET); uint32_t t = 0; while (HAL_GPIO_ReadPin(DHT11_GPIO_Port, DHT11_Pin) == GPIO_PIN_SET) { t++; delay_us(1); } data[i / 8] <<= 1; if (t > 40) data[i / 8] |= 1; // 高电平持续时长区分 0/1 } // 4. 校验和 if ((data[0] + data[1] + data[2] + data[3]) == data[4]) { *humidity = data[0]; *temperature = data[2]; return 1; } return 0; }

注意,delay_us需要自己实现,标准 HAL 只有毫秒级HAL_Delay,可以用 DWT 或 SysTick 做微秒延时。我实测的经验是:DHT11 判断位值时,用"50us 归零区间加高电平时长"来区分 0 和 1,比死等边沿要稳定;另外,两次读取之间至少隔 1 秒,否则 DHT11 会返回旧数据甚至不响应。

MQ-2 读取更简单,本质是 ADC 采样。用 STM32 的 ADC 读引脚电压,然后按数据手册的灵敏度特性做阈值判断:

uint32_t adc_value = 0; HAL_ADC_Start(&hadc1); adc_value = HAL_ADC_GetValue(&hadc1); float voltage = adc_value * 3.3f / 4095.0f; // 12位 ADC,参考电压 3.3V if (voltage > GAS_ALARM_THRESHOLD) { HAL_GPIO_WritePin(BUZZER_GPIO_Port, BUZZER_Pin, GPIO_PIN_RESET); // 拉低触发蜂鸣器 }

MQ-2 上电后需要预热,第一次读数通常偏高,建议开机 60 秒后再进入正式检测逻辑,否则你会在半夜被家里的蜂鸣器吓醒。

4.3 数据上行:为什么不用裸 HTTP 而是 MQ

采集到数据后,下一步是往服务端送。很多新手第一反应是让 STM32 直接 HTTP POST 到后端接口。这个方案在实验室跑通当然可以,但到了真实部署就会暴露问题:设备断网重连时请求直接丢失;服务器接口抖动一次,设备端就得写一整套重试逻辑;多个设备同时上报,服务端容易被冲垮。

换成消息方式就顺滑得多。设备端用 MQTT 这种轻量协议发布消息,比如发布到主题stm32/device-01/sensors:

mosquitto_pub -h broker_host -t "stm32/device-01/sensors" -m '{"temp":26.5,"humidity":60,"light":320,"gas":0.42}'

MQTT 协议的 QoS 等级(0/1/2)可以控制消息可靠性,遗嘱消息(LWT)还能在设备异常掉线时通知服务端。消息到了 Broker 之后,再通过桥接转发到 Pulsar 这样的大规模消息平台做持久化、分析和多团队消费。

这就是我前面说的"MQ 链路":STM32 采集端到 MQTT Broker 是一段,MQTT Bridge 到 Pulsar 是一段,Pulsar 到下游数据分析或告警服务又是一段。每一段都由消息驱动,中间层通过队列天然做了削峰填谷和故障缓冲。

5. 动手试验:本地跑通 Pulsar 与 MQTT 转发桥接

理论聊完了,来点能直接抄作业的东西。我按自己在 Pulsar Developer Day 之后复现的流程写一遍,照着做就能在本地把 Pulsar 跑起来,并且让 MQTT 传感器数据转发进 Pulsar。

5.1 一条命令拉起 Pulsar 单机

前提是机器上装了 Docker。拉镜像直接跑单机模式:

docker run --name pulsar-dev -d \ -p 6650:6650 -p 8080:8080 \ apachepulsar/pulsar:3.3.1 standalone
  • 6650 是客户端连接端口,8080 是管理 REST API 端口。
  • standalone 模式自带一个最小 BookKeeper 和 Broker,足够本地验证。

启动后用管理 API 看一眼状态:

curl http://localhost:8080/pulsar/v2/brokers/health

返回 OK 就说明起来了。我遇到过的情况是 Docker Desktop 内存配额给太低导致容器反复重启,建议至少给 4GB 内存。

5.2 用 Python 跑一次"生产-消费"

装 Python 客户端:

pip install pulsar-client

然后一个文件验证生产:

import pulsar client = pulsar.Client("pulsar://localhost:6650") producer = client.create_producer("persistent://public/default/env-sensor") for i in range(10): producer.send(("sensor-reading-%d" % i).encode("utf-8")) print("完成10条消息生产") client.close()

再开一个进程验证消费:

import pulsar client = pulsar.Client("pulsar://localhost:6650") consumer = client.subscribe( "persistent://public/default/env-sensor", "demo-subscription" ) while True: msg = consumer.receive(timeout_ms=5000) if msg: print("收到:", msg.data().decode("utf-8")) consumer.acknowledge(msg) else: break client.close()

注意,消费端如果是先于生产者订阅的,那它只消费后续新消息;如果生产者先发了消息而当时没有订阅存在,这些消息会保留在 topic 里,订阅一经创建默认从头消费。standalone 模式的默认保留策略是全部保留,这正好让你观察"迟到的订阅者"行为。

5.3 桥接 MQTT:让 STM32 数据流进 Pulsar

本地先起一个 MQTT Broker,我用的是 EMQX,一条 Docker 命令的事:

docker run --name emqx -d -p 1883:1883 -p 8083:8083 emqx/emqx:5.6.0

然后写一个 Python 桥接脚本,监听 MQTT 主题,把消息转发到 Pulsar:

import json import paho.mqtt.client as mqtt import pulsar PULSAR_URL = "pulsar://localhost:6650" MQTT_BROKER = "localhost" MQTT_TOPIC = "stm32/+/sensors" pulsar_client = pulsar.Client(PULSAR_URL) producer = pulsar_client.create_producer( "persistent://public/default/env-sensor-from-mqtt" ) def on_message(client, userdata, msg): payload = msg.payload.decode("utf-8") try: data = json.loads(payload) data["source_topic"] = msg.topic producer.send(json.dumps(data).encode("utf-8")) print("forwarded:", msg.topic, payload) except json.JSONDecodeError: print("忽略非JSON消息:", msg.topic) mqtt_client = mqtt.Client() mqtt_client.on_message = on_message mqtt_client.connect(MQTT_BROKER, 1883, 60) mqtt_client.subscribe(MQTT_TOPIC) mqtt_client.loop_forever()

然后在另一个终端模拟 STM32 设备发布消息:

mosquitto_pub -t "stm32/device-01/sensors" -m '{"temp":26.5,"humidity":60,"light":320,"gas":0.42}'

如果脚本打印 forward 了,再用 5.2 的消费端去读 Pulsar 里的env-sensor-from-mqtttopic,就能看到这条数据。至此,一条从设备到 Pulsar 的完整 MQ 链路就通了。

5.4 我踩过的几个坑

  • Pulsar 客户端默认走 6650,REST API 走 8080,两个端口都要映射,漏一个管理功能就不可用。
  • standalone 模式默认没有启用认证,不用填 token;但如果后面开了认证,所有客户端都要同步配置 token,否则清一色 401。
  • Paho 的on_message回调里千万别放阻塞耗时操作,比如同步写数据库,会造成 MQTT 消息堆积。桥接脚本里转发 Pulsar 用异步 send,问题不大。
  • 如果模拟设备发的 JSON 带 BOM 头,json.loads会报错,记得先decode再.strip('\ufeff')。

6. 散场后的工程复盘与下一步打算

一天逛下来,我最想带回家的其实不是某个具体功能,而是三个判断:第一,消息队列正在从"应用之间的管道"变成"整个系统的事件底座",设备数据、服务数据最终都汇聚到一条统一的消息链路上;第二,Pulsar 的多租户和分层存储,解决的正是规模化之后最头疼的管理和成本问题,这个定位在 2025 年非常能打;第三,嵌入式端和 MQ 技术栈的融合速度比我想象中快,现场能看到不少人已经在用 MQTT 把 STM32 环境监测数据送进云端消息平台。

我自己的下一步计划,是把家里那套 STM32 环境监测板子的数据正式接进 Pulsar。具体做法是:板子上 DHT11、BH1750、MQ-2 每分钟采集一次,OLED 本地刷新显示,数据通过 ESP8266 走 MQTT 发布;本地跑一个 Pulsar standalone 加 EMQX 桥接,Pulsar 里按租户home、命名空间env-sensor建 topic;下游写一个简单的 Python 消费服务,做历史曲线和告警推送。

这套方案放在生产环境里当然还有不少要补的,比如 Pulsar 集群高可用、消息 schema 校验、数据脱敏和权限管理。但对于个人项目和个人成长来说,先把链路跑通、再把每一段原理吃透,比一开始就追求大而全要靠谱得多。

最后再分享一个我做这类项目的小习惯:不管用什么 MQ,消息格式尽早统一成 JSON 并带上version字段。消息队列最怕的不是量,而是消费端解析不了旧消息。一个version字段,能让你在后续演进协议的时候少掉很多头发。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询