You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 23:45:43