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

关于读取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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:05:15