关于读取MQTT接入Kafka的输出数据的技术问询
使用Kafka Connect实现MQTT与Kafka的桥接及Java消费者实现
我刚好有过用Kafka Connect搭建MQTT到Kafka数据通道的实践经验,结合你提到的用MQTTLENSES发布样本数据、编写Java消费者的需求,整理了完整的实现方案,帮你把整个流程串起来:
一、配置Kafka Connect MQTT源连接器
首先需要搭建Kafka Connect的MQTT源连接器(这里以Confluent官方的MQTT Connector为例),你可以创建一个connect-mqtt-source.properties配置文件,核心配置如下:
name=mqtt-kafka-bridge connector.class=io.confluent.connect.mqtt.MqttSourceConnector tasks.max=1 # MQTT Broker地址 mqtt.server.uri=tcp://your-mqtt-broker-ip:1883 # 要订阅的MQTT主题 mqtt.topics=mqtt/sensor/data # 转发到的Kafka主题 kafka.topic=kafka/sensor/data mqtt.qos=1 # 消息转换器,这里用字符串格式 value.converter=org.apache.kafka.connect.storage.StringConverter key.converter=org.apache.kafka.connect.storage.StringConverter
然后通过REST API启动连接器:
curl -X POST -H "Content-Type: application/json" --data @connect-mqtt-source.properties http://your-connect-host:8083/connectors
二、用MQTTLENSES发布样本数据
打开MQTTLENSES并连接到你的MQTT Broker,选择刚才配置的mqtt/sensor/data主题,发布一段JSON格式的样本数据,比如:
{"deviceId": "temp-sensor-001", "temperature": 26.3, "humidity": 45, "collectTime": 1718052300000}
三、完整Java Kafka消费者代码
你提供的代码片段看起来包含了解密逻辑,我补全了完整的可运行版本,包含消息接收、解密(如果需要)和JSON解析的完整流程:
package test; import java.nio.charset.StandardCharsets; import java.util.Arrays; import java.util.Properties; import javax.crypto.Cipher; import javax.crypto.spec.SecretKeySpec; import org.json.JSONObject; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; public class SensorDataConsumer { // 注意:密钥需要和MQTT消息发布端的加密密钥完全一致,且长度需符合AES要求(16/24/32字节) private static final String ENCRYPT_KEY = "your-16-byte-secret-key"; private static final String CIPHER_ALGORITHM = "AES"; public static void main(String[] args) { // 配置Kafka消费者参数 Properties consumerProps = new Properties(); consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-broker-ip:9092"); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "sensor-data-consumer-group"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 从最早的消息开始消费 try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps)) { // 订阅目标Kafka主题 consumer.subscribe(Arrays.asList("kafka/sensor/data")); System.out.println("Waiting for messages..."); while (true) { // 拉取消息,超时时间1秒 ConsumerRecords<String, String> records = consumer.poll(java.time.Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { System.out.printf("Received message | Offset: %d | Key: %s | Raw Value: %s%n", record.offset(), record.key(), record.value()); try { // 解密消息(如果你的MQTT消息是加密的,否则注释掉这步直接用record.value()) String decryptedData = decryptMessage(record.value(), ENCRYPT_KEY); System.out.println("Decrypted Data: " + decryptedData); // 解析JSON格式的消息内容 JSONObject dataObj = new JSONObject(decryptedData); String deviceId = dataObj.getString("deviceId"); double temp = dataObj.getDouble("temperature"); int humidity = dataObj.getInt("humidity"); long collectTime = dataObj.getLong("collectTime"); System.out.printf("Parsed Sensor Data | Device: %s | Temp: %.1f°C | Humidity: %d%% | Time: %d%n", deviceId, temp, humidity, collectTime); } catch (Exception e) { System.err.println("Failed to process message: " + e.getMessage()); e.printStackTrace(); } } } } } // AES解密方法 private static String decryptMessage(String encryptedBase64, String secretKey) throws Exception { SecretKeySpec keySpec = new SecretKeySpec(secretKey.getBytes(StandardCharsets.UTF_8), CIPHER_ALGORITHM); Cipher cipher = Cipher.getInstance(CIPHER_ALGORITHM); cipher.init(Cipher.DECRYPT_MODE, keySpec); // 先解码Base64,再解密 byte[] encryptedBytes = java.util.Base64.getDecoder().decode(encryptedBase64); byte[] decryptedBytes = cipher.doFinal(encryptedBytes); return new String(decryptedBytes, StandardCharsets.UTF_8); } }
四、关键注意事项
- 确保Kafka Connect、MQTT Broker、Kafka集群三者网络互通,端口开放(MQTT默认1883,Kafka默认9092,Connect默认8083)
- 如果你的MQTT消息没有加密,直接移除解密相关代码,直接使用
record.value()进行JSON解析即可 - 消费者配置的Kafka主题必须和Kafka Connect配置中的
kafka.topic完全一致 - 确保MQTTLENSES发布的消息格式和消费者的JSON解析逻辑匹配,避免解析报错
内容的提问来源于stack exchange,提问作者abacus service
相关产品推荐
相关产品推荐

