【发布时间】: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