【问题标题】:Decode message Kafka-mqtt解码消息Kafka-mqtt
【发布时间】:2018-08-28 18:38:30
【问题描述】:

我使用 kafka-mqtt 连接器从 mqtt 代理接收关于 kafka 主题的消息。然后我在 spark 上使用 kafka 消费者从 kafka 主题中读取了这些消息。当我打印消息时,这就是结果。如何正确阅读消息?

这是设置消费者和创建流的代码。

SparkConf sparkConf = new SparkConf().setAppName("GestoreSoccorso").setMaster("local[2]");
JavaStreamingContext ssc = new JavaStreamingContext(sparkConf, new Duration(500));

Map<String, Object> kafkaParams = new HashMap<>();
kafkaParams.put("bootstrap.servers", "10.0.4.215:9092");
kafkaParams.put("key.deserializer", StringDeserializer.class);
kafkaParams.put("value.deserializer", StringDeserializer.class);
kafkaParams.put("group.id", "use_a_separate_group_id_for_each_stream");
kafkaParams.put("auto.offset.reset", "earliest");
kafkaParams.put("enable.auto.commit", false);
Collection<String> topics = Arrays.asList("ankioverdrive_v1_events");

JavaInputDStream<ConsumerRecord<String, String>> stream =
        KafkaUtils.createDirectStream(
                ssc,
                LocationStrategies.PreferConsistent(),
                ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams)
        );

然后这是我用来从主题读取消息并打印它们的代码。

stream.foreachRDD(new VoidFunction<JavaRDD<ConsumerRecord<String, String>>>() {
            @Override
            public void call(JavaRDD<ConsumerRecord<String, String>> consumerRecordJavaRDD) throws Exception {
                consumerRecordJavaRDD.foreach(new VoidFunction<ConsumerRecord<String, String>>() {
                    @Override
                    public void call(ConsumerRecord<String, String> stringStringConsumerRecord) throws Exception {
                        String stringa=stringStringConsumerRecord.value();
                        System.out.println(DEBUG+"DATI RICEVUTI -> "+ stringa);

最后是输出

DEBUG: DATI RICEVUTI ->     ܑ�՛Y
SKULL   �    unknown0% 

【问题讨论】:

  • 如何知道 MQTT 数据被序列化为字符串?
  • 我知道是因为我看到了mqtt消息的代码,它是一个json转换成字符串然后发布的。
  • 如果是字符串或 JSON,您将看不到 ܑ�՛Y 字符。消息要么是压缩的、加密的,要么不是 JSON/明文。你能显示kafka-console-consumer的输出吗?
  • 我现在看不到输出控制台,明天可以。我今天早上看到了输出,在消费者控制台上,有很多这种类型的行:x00/x00/.../SKULL/..../unknown/... 我认为这是十六进制。跨度>
  • 看起来像,是的。您的主题包含不完全是 UTF-8 字符串的二进制数据。在 Spark 中,您可以使用ByteArrayDeserializer.class.getName(),但您仍然需要确定如何解码字节,这只能通过知道它们进入 Kafka 的格式来完成

标签: java apache-spark apache-kafka mqtt


【解决方案1】:

您为参数key.deserializervalue.deserializer 传递了错误的值。

代替

kafkaParams.put("key.deserializer", StringDeserializer.class);
kafkaParams.put("value.deserializer", StringDeserializer.class);

你需要通过

kafkaParams.put("key.deserializer", StringDeserializer.class.getName());
kafkaParams.put("value.deserializer", StringDeserializer.class.getName());

或者干脆

kafkaParams.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
kafkaParams.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");

【讨论】:

  • 这是假设 MQTT 数据实际上是一个字符串
  • 我尝试了,但它不起作用。我相信问题在于消息的编码。也就是说,我不知道消息是在 mqtt 在 kafka 主题上发送时编码的。但是消息是一个字符串。有人可以告诉我如何查看已使用的编码吗?以及如何修改它?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-08-28
  • 1970-01-01
  • 2016-12-07
  • 2016-03-06
  • 1970-01-01
相关资源
最近更新 更多