【发布时间】:2018-05-29 09:38:28
【问题描述】:
我有一个用 Java 编写并使用 Spark 2.1 的 Spark 流应用程序。我正在使用KafkaUtils.createDirectStream 来读取来自 Kafka 的消息。我正在为 kafka 消息使用 kryo 编码器/解码器。我在 Kafka 属性-> key.deserializer、value.deserializer、key.serializer、value.deserializer
中指定了这个
当 Spark 以微批次拉取消息时,使用 kryo 解码器成功解码消息。但是我注意到 Spark 执行器创建了一个新的 kryo 解码器实例,用于解码从 kafka 读取的每条消息。我通过将日志放入解码器构造函数中检查了这一点
这对我来说似乎很奇怪。每个消息和每个批次不应该使用相同的解码器实例吗?
我从 kafka 读取的代码:
JavaInputDStream<ConsumerRecord<String, Class1>> consumerRecords = KafkaUtils.createDirectStream(
jssc,
LocationStrategies.PreferConsistent(),
ConsumerStrategies.<String, Class1>Subscribe(topics, kafkaParams));
JavaPairDStream<String, Class1> converted = consumerRecords.mapToPair(consRecord -> {
return new Tuple2<String, Class1>(consRecord.key(), consRecord.value());
});
【问题讨论】:
标签: java apache-spark apache-kafka spark-streaming kryo