【问题标题】:Why does Kafka Direct Stream create a new decoder for every message?为什么 Kafka Direct Stream 为每条消息创建一个新的解码器?
【发布时间】: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


    【解决方案1】:

    如果我们想了解 Spark 如何在内部从 Kafka 获取数据,我们需要查看 KafkaRDD.compute,这是为每个 RDD 实现的方法,它告诉框架如何计算 @987654323 @:

    override def compute(thePart: Partition, context: TaskContext): Iterator[R] = {
      val part = thePart.asInstanceOf[KafkaRDDPartition]
      assert(part.fromOffset <= part.untilOffset, errBeginAfterEnd(part))
      if (part.fromOffset == part.untilOffset) {
        logInfo(s"Beginning offset ${part.fromOffset} is the same as ending offset " +
        s"skipping ${part.topic} ${part.partition}")
        Iterator.empty
      } else {
        new KafkaRDDIterator(part, context)
      }
    }
    

    这里重要的是else 子句,它创建了一个KafkaRDDIterator。这在内部有:

    val keyDecoder = classTag[U].runtimeClass.getConstructor(classOf[VerifiableProperties])
      .newInstance(kc.config.props)
      .asInstanceOf[Decoder[K]]
    
    val valueDecoder = classTag[T].runtimeClass.getConstructor(classOf[VerifiableProperties])
      .newInstance(kc.config.props)
      .asInstanceOf[Decoder[V]]
    

    如您所见,它通过反射为每个底层分区创建一个键解码器和值解码器的实例。这意味着它不是每个消息而是每个Kafka分区生成的。

    为什么要这样实现?我不知道。我假设是因为与 Spark 内部发生的所有其他分配相比,键和值解码器的性能损失应该可以忽略不计。

    如果您分析了您的应用并发现这是一个分配热路径,您可以打开一个问题。否则,我不会担心。

    【讨论】:

    • 研究得很好! #印象深刻
    • @Yuval:我使用的是 Kafka 0.10.x。 Spark 使用缓存的 kafka 消费者(每个执行程序),其中缓存键由消费者 ID、主题 ID、分区 ID 标识。每个 kafka 分区都有一个解码器是有意义的,否则 Spark 将如何并行解码消息。我期望的是,一个新的解码器必须在缓存消费者中的每个分区创建一次,就是这样!我在轻负载下看不到这个问题,但只有在我每秒发送 1000 条消息时才会出现。可能我正在进入“GC”循环。您对如何启用 KafkaRDD 类中的日志记录有任何想法吗?
    • @scorpio Kafka 0.10.x 根本不需要解码器。它返回底层ConsumerRecord,您可以选择如何处理它。您是否正在 map 中创建解码器的实例?
    • @Yuval:添加了相关代码sn-p。我在 kafka 属性中指定键和值解码器。您发现问题了吗?
    猜你喜欢
    • 2017-04-23
    • 2016-05-06
    • 1970-01-01
    • 1970-01-01
    • 2021-07-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多