diff --git a/cloud-common/cloud-common-kafka/pom.xml b/cloud-common/cloud-common-kafka/pom.xml index 6bd82ac..b515a5b 100644 --- a/cloud-common/cloud-common-kafka/pom.xml +++ b/cloud-common/cloud-common-kafka/pom.xml @@ -23,5 +23,11 @@ kafka-clients 3.0.0 + + + + com.muyu + cloud-common-core + diff --git a/cloud-modules/cloud-modules-carData/pom.xml b/cloud-modules/cloud-modules-carData/pom.xml index 7e7c80b..dd2d5f9 100644 --- a/cloud-modules/cloud-modules-carData/pom.xml +++ b/cloud-modules/cloud-modules-carData/pom.xml @@ -89,7 +89,6 @@ com.muyu cloud-common-kafka - 3.6.3 org.apache.iotdb diff --git a/cloud-modules/cloud-modules-protocolparsing/src/main/java/com/muyu/mqtt/Demo.java b/cloud-modules/cloud-modules-protocolparsing/src/main/java/com/muyu/mqtt/Demo.java index c30d6db..fa5f209 100644 --- a/cloud-modules/cloud-modules-protocolparsing/src/main/java/com/muyu/mqtt/Demo.java +++ b/cloud-modules/cloud-modules-protocolparsing/src/main/java/com/muyu/mqtt/Demo.java @@ -1,14 +1,14 @@ -package com.muyu.mqtt; +package com.muyu.server.mqtt; import com.alibaba.fastjson.JSONObject; -import com.muyu.common.kafka.constants.KafkaConstants; import com.muyu.domain.CarMessage; -import com.muyu.service.CarMessageService; +import com.muyu.server.service.CarMessageService; import jakarta.annotation.PostConstruct; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.eclipse.paho.client.mqttv3.*; +import org.springframework.stereotype.Component; import javax.annotation.Resource; import java.util.ArrayList; @@ -16,7 +16,7 @@ import java.util.List; import java.util.concurrent.CompletableFuture; -//@Component +@Component public class Demo { @Resource private CarMessageService service; @@ -24,10 +24,11 @@ public class Demo { private KafkaProducer kafkaProducer; @PostConstruct public void test() { + String topic = "vehicle"; String content = "Message from MqttPublishSample"; int qos = 2; - String broker = "tcp://60.204.221.52:1883"; + String broker = "tcp://106.15.136.7:1883"; String clientId = "JavaSample"; try { @@ -48,7 +49,7 @@ public class Demo { @Override public void messageArrived(String s, MqttMessage mqttMessage) throws Exception { - List list= service.selectCarMessage(1); + List list= service.selectCarMessageList(1,2); String str = new String( mqttMessage.getPayload() ); System.out.println(str); String[] test = str.split(" "); @@ -81,7 +82,7 @@ public class Demo { results[i] = futures.get(i).get(); } String jsonString = JSONObject.toJSONString( results ); - ProducerRecord producerRecord = new ProducerRecord<>( KafkaConstants.KafkaTopic, jsonString); + ProducerRecord producerRecord = new ProducerRecord<>( "carJsons", jsonString); kafkaProducer.send(producerRecord); } // 接收信息 diff --git a/pom.xml b/pom.xml index 081c4fa..58f348a 100644 --- a/pom.xml +++ b/pom.xml @@ -42,6 +42,7 @@ 5.8.27 4.1.0 2.4.1 + 3.6.3 @@ -218,6 +219,13 @@ ${muyu.version} + + + com.muyu + cloud-common-kafka + ${kafka.version} + + com.muyu