diff --git a/collector/pom.xml b/collector/pom.xml
index 2d93db59718..fea17e3844d 100644
--- a/collector/pom.xml
+++ b/collector/pom.xml
@@ -206,6 +206,12 @@
+
+
+ com.hivemq
+ hivemq-mqtt-client
+ 1.3.3
+
diff --git a/collector/src/main/java/org/apache/hertzbeat/collector/collect/mqtt/MqttCollectImpl.java b/collector/src/main/java/org/apache/hertzbeat/collector/collect/mqtt/MqttCollectImpl.java
new file mode 100644
index 00000000000..14f8e18145a
--- /dev/null
+++ b/collector/src/main/java/org/apache/hertzbeat/collector/collect/mqtt/MqttCollectImpl.java
@@ -0,0 +1,231 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hertzbeat.collector.collect.mqtt;
+
+import com.hivemq.client.mqtt.MqttVersion;
+import com.hivemq.client.mqtt.datatypes.MqttQos;
+import com.hivemq.client.mqtt.mqtt3.Mqtt3AsyncClient;
+import com.hivemq.client.mqtt.mqtt3.Mqtt3Client;
+import com.hivemq.client.mqtt.mqtt3.Mqtt3ClientBuilder;
+import com.hivemq.client.mqtt.mqtt3.message.connect.connack.Mqtt3ConnAck;
+import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient;
+import com.hivemq.client.mqtt.mqtt5.Mqtt5Client;
+import com.hivemq.client.mqtt.mqtt5.Mqtt5ClientBuilder;
+import com.hivemq.client.mqtt.mqtt5.message.connect.connack.Mqtt5ConnAck;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.function.Consumer;
+import java.util.stream.Collectors;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.hertzbeat.collector.collect.AbstractCollect;
+import org.apache.hertzbeat.collector.constants.CollectorConstants;
+import org.apache.hertzbeat.collector.dispatch.DispatchConstants;
+import org.apache.hertzbeat.common.constants.CommonConstants;
+import org.apache.hertzbeat.common.entity.job.Metrics;
+import org.apache.hertzbeat.common.entity.job.protocol.MqttProtocol;
+import org.apache.hertzbeat.common.entity.message.CollectRep;
+import org.apache.hertzbeat.common.entity.message.CollectRep.MetricsData.Builder;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.util.Assert;
+import org.springframework.util.StopWatch;
+
+/**
+ * collect mqtt metrics
+ */
+public class MqttCollectImpl extends AbstractCollect {
+
+ public MqttCollectImpl() {
+ }
+
+ private static final Logger logger = LoggerFactory.getLogger(MqttCollectImpl.class);
+
+ @Override
+ public void preCheck(Metrics metrics) throws IllegalArgumentException {
+ MqttProtocol mqttProtocol = metrics.getMqtt();
+ Assert.hasText(mqttProtocol.getHost(), "MQTT protocol host is required");
+ Assert.hasText(mqttProtocol.getPort(), "MQTT protocol port is required");
+ Assert.hasText(mqttProtocol.getProtocolVersion(), "MQTT protocol version is required");
+ }
+
+ @Override
+ public void collect(Builder builder, long monitorId, String app, Metrics metrics) {
+ MqttProtocol mqtt = metrics.getMqtt();
+ String protocolVersion = mqtt.getProtocolVersion();
+ MqttVersion mqttVersion = MqttVersion.valueOf(protocolVersion);
+ if (mqttVersion == MqttVersion.MQTT_3_1_1) {
+ collectWithVersion3(metrics, builder);
+ } else if (mqttVersion == MqttVersion.MQTT_5_0) {
+ collectWithVersion5(metrics, builder);
+ }
+ }
+
+ @Override
+ public String supportProtocol() {
+ return DispatchConstants.PROTOCOL_MQTT;
+ }
+
+ /**
+ * collecting data of MQTT 5
+ */
+ private void collectWithVersion5(Metrics metrics, Builder builder) {
+ MqttProtocol mqttProtocol = metrics.getMqtt();
+ Map