ActiveMQ Artemis+Telegraf MQTT监听:10万消息已投递但4千滞留队列
MQTT QoS 2消息投递测试异常问题
测试背景
向ActiveMQ Artemis发送10万条持久化MQTT消息(QoS 2),该主题配置两个Telegraf监听器,分别部署在VM 85和VM 86上,监听器负责将消息写入各自服务器的InfluxDB。
测试目标与前置配置
核心目标:即使VM 86离线,所有投递到VM 85的消息也需同步投递到VM 86。
前置配置:
- 两个监听器均使用唯一客户端ID连接代理
- 设置
clean-session = false - 以QoS 2订阅目标主题,确保消息发送时无论监听器是否活跃,订阅关系持续存在
测试步骤
初始状态:两个监听器均未连接代理
- 启动VM 85上的监听器
- 发送10万条测试消息
- 确认消息已全部投递至VM 85的监听器并写入对应InfluxDB
- 启动VM 86上的监听器
- 确认消息已全部投递至VM 86的监听器并写入对应InfluxDB
异常现象
- 两台VM的InfluxDB均收到完整的10万条消息,但VM 86对应的ActiveMQ队列仍显示约4.3千条消息滞留
- 重启VM 86上的监听器,监听器会显示正在写入更多数据,但InfluxDB中的消息总数仍保持10万条(客户端发送的消息包含递增数值和唯一时间戳,理论上不会产生重复记录,InfluxDB重复记录覆盖机制未触发)
疑问点
为何必须重启VM 86的监听器才能清空对应队列?
相关配置与代码
Telegraf未尝试的参数
## Maximum messages to read from the broker that have not been written by an ## output. For best throughput set based on the number of metrics within ## each message and the size of the output's metric_batch_size. ## ## For example, if each message from the queue contains 10 metrics and the ## output metric_batch_size is 1000, setting this to 100 will ensure that a ## full batch is collected and the write is triggered immediately without ## waiting until the next flush_interval. # max_undelivered_messages = 1000
根据输出日志判断,批量大小默认值为1000,但输出前可读取的最大消息数似乎更大——重启后监听器会输出4.3千条消息,可这些消息实际已被输出过,此现象无法解释。
MQTT客户端发送代码
package abc; import java.time.Instant; import org.eclipse.paho.client.mqttv3.MqttClient; import org.eclipse.paho.client.mqttv3.MqttConnectOptions; import org.eclipse.paho.client.mqttv3.MqttException; import org.eclipse.paho.client.mqttv3.MqttMessage; import org.eclipse.paho.client.mqttv3.MqttSecurityException; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; import com.influxdb.client.domain.WritePrecision; import com.influxdb.client.write.Point; public class MqttPublishSample { public static void main(String[] args) throws MqttSecurityException, MqttException, InterruptedException { String broker = "tcp://localhost:1883"; String clientId = "JavaSample"; MemoryPersistence persistence = new MemoryPersistence(); int qos = 2; int start = Integer.parseInt(args[0]); int end = Integer.parseInt(args[1]); String topic = args[2]; if (topic == null) { topic = "testtopic/999"; } System.out.println("start: " + start + ", end: " + end + ", topic: " + topic + " qos: " + qos); MqttClient sampleClient = new MqttClient(broker, clientId, persistence); MqttConnectOptions connOpts = new MqttConnectOptions(); connOpts.setCleanSession(false); connOpts.setUserName("admin"); connOpts.setPassword("xxxxxxx".toCharArray()); System.out.println("Connecting to broker: " + broker); sampleClient.connect(connOpts); System.out.println("Connected"); for (int i = start; i <= end; i++) { // print out every 100 if (i%100 == 0) { System.out.println("i: " + i); } try { Point point = Point.measurement("temperature").addTag("machine", "unit43").addField("external", i) .time(Instant.now(), WritePrecision.NS); String content = point.toLineProtocol(); MqttMessage message = new MqttMessage(content.getBytes()); message.setQos(qos); sampleClient.publish(topic, message); Thread.sleep(10); } catch (MqttException me) { System.out.println("reason " + me.getReasonCode()); System.out.println("msg " + me.getMessage()); System.out.println("loc " + me.getLocalizedMessage()); System.out.println("cause " + me.getCause()); System.out.println("excep " + me); me.printStackTrace(); } } sampleClient.disconnect(); System.out.println("Disconnected"); } }
VM 85上的Telegraf配置
############################################################################### # INPUT PLUGINS # ############################################################################### [[inputs.mqtt_consumer]] servers = ["tcp://127.0.0.1:1883"] ## Topics that will be subscribed to. topics = [ "testtopic/#", ] ## The message topic will be stored in a tag specified by this value. If set ## to the empty string no topic tag will be created. # topic_tag = "topic" ## When using a QoS of 1 or 2, you should enable persistent_session to allow ## resuming unacknowledged messages. qos = 2 persistent_session = true ## If unset, a random client ID will be generated. client_id = "InfluxData_on_86_listen_local" ## Username and password to connect MQTT server. username = "admin" password = "xxxxxx" data_format = "influx" [[inputs.mqtt_consumer]] servers = ["tcp://10.102.11.86:1883"] ## Topics that will be subscribed to. topics = [ "testtopic/#", ] ## The message topic will be stored in a tag specified by this value. If set ## to the empty string no topic tag will be created. # topic_tag = "topic" ## When using a QoS of 1 or 2, you should enable persistent_session to allow ## resuming unacknowledged messages. qos = 2 persistent_session = true ## If unset, a random client ID will be generated. client_id = "InfluxData_on_86_listen_85" ## Username and password to connect MQTT server. username = "admin" password = "xxxx" data_format = "influx" ############################################################################### # OUTPUT PLUGINS # ############################################################################### [[outputs.influxdb_v2]] ## The URLs of the InfluxDB cluster nodes. ## ## Multiple URLs can be specified for a single cluster, only ONE of the ## urls will be written to each interval. urls = ["http://127.0.0.1:8086"] ## Token for authentication. token = "xxxx" ## Organization is the name of the organization you wish to write to. organization = "xxxx" # ## Destination bucket to write into. bucket = "events"
内容的提问来源于stack exchange,提问作者Jack BeNimble
相关产品推荐
相关产品推荐

