【问题标题】:read output of kafka coming from MQTT读取来自 MQTT 的 kafka 输出
【发布时间】:2018-08-29 05:31:06
【问题描述】:

我正在使用 Kafka 连接进行 MQTT-KAFKA 连接。我正在从 MQTTLENSES 发布示例数据,并在 java 中为 Kafka 编写了消费者代码:

    package test;


    import java.io.UnsupportedEncodingException;
    import java.nio.charset.StandardCharsets;
    import java.util.Arrays;
    import java.util.Base64;
    import java.util.Iterator;
    import java.util.Properties;




    import javax.crypto.Cipher;
    import javax.crypto.spec.SecretKeySpec;

    import org.json.*;

    import org.apache.kafka.common.serialization.StringDeserializer;

    import com.fasterxml.jackson.databind.util.JSONWrappedObject;

    import org.apache.kafka.clients.consumer.ConsumerRecord;
    import org.apache.kafka.clients.consumer.ConsumerRecords;
    import org.apache.kafka.clients.consumer.KafkaConsumer;

    @SuppressWarnings("unused")
    public class ConsumerTest {


      public static void main(String[] args) throws UnsupportedEncodingException {
        System.out.println("consumer123");
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", "group-1");
        props.put("enable.auto.commit", "true");
        props.put("auto.commit.interval.ms", "1000");
        props.put("auto.offset.reset", "earliest");
        props.put("session.timeout.ms", "30000");
        props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");

        @SuppressWarnings("resource")
        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<>(props);
        kafkaConsumer.subscribe(Arrays.asList(MQTT TOPIC));
        while (true) {


             ConsumerRecords<String, String> records = kafkaConsumer.poll(100);




          for (ConsumerRecord<String, String> record : records) {

              System.out.println(records);


              try
              {
                  String record_data = record.value().toString();

                  JSONObject obj = new JSONObject(record_data);
                  String payload = obj.getString("payload");



                  String s = new String(payload.getBytes(), StandardCharsets.UTF_8);


                  System.out.println(s);
                  System.out.println(record_data);
                  System.out.println(record.key());
                  //System.out.println(decryptData(payload));


              }
              catch(Exception je)
              {
                  System.out.println(je.toString());
              }
          }

        }

      }

对于输入“Hello”,它在消费者中打印输出 => {"schema":{"type":"bytes","optional":false},"payload":"SGVsbG8K"}

如何解码传入kafka消费者的payload?

【问题讨论】:

  • 您正在使用 StringDeserializer。在发送到 kafka 期间你是如何序列化的?你的序列化器,反序列化器应该是一样的。

标签: apache-kafka mqtt payload


【解决方案1】:

您收到的负载是 "base64" 格式。您只需解码这个 base64 字符串:

byte[] plainText = Base64.getDecoder().decode(payload);
System.out.println(new String(plainTexts));

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-07-09
    • 2018-01-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-30
    相关资源
    最近更新 更多