diff --git a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/MqttSource.java b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/MqttSource.java index 144a1e6cc4..8ace035b7a 100644 --- a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/MqttSource.java +++ b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/MqttSource.java @@ -115,7 +115,7 @@ public void messageArrived(String topic, MqttMessage message) throws Exception { headerMap.put("record.messageId", String.valueOf(message.getId())); headerMap.put("record.qos", String.valueOf(message.getQos())); byte[] recordValue = message.getPayload(); - mqttMessagesQueue.offer(new DefaultMessage(recordValue, headerMap), 1, TimeUnit.SECONDS); + mqttMessagesQueue.offer(new DefaultMessage(recordValue, headerMap)); }