当前位置:   article > 正文

修改spark-streaming源码,应对kafka数据倾斜_kafka分区特别倾斜

kafka分区特别倾斜

我们知道,kafka数据分区与spark数据分区一一对应,如果kafka数据倾斜,势必造成spark数据倾斜,在spark-streming源码类

DirectKafkaInputDStream中,compute方法中有对每批次数据的任务切分,在这里修改一下源码,限制生成的每个task消费的数据量,也就是说,将kafka中数据量偏多的分区的数据切分成多个task进行消费,实现逻辑如下:主要的改动就是增加一个参数,控制每个task消费的最大的消息条数,如果超过则切分成多个任务。不过修改后有一个问题,就是由于kafka consumer多线程下消费有问题,因此需要关闭spark-streaming的consumer缓存功能。
val offsetRanges = untilOffsets.flatMap { case (tp, uo) =>
  var currentFo = currentOffsets(tp)
  var currentUo = currentFo + maxMessagePerTask
  val offsetRangeArray = new ju.LinkedList[OffsetRange]()
  while (uo > currentUo) {
    val or = OffsetRange(tp.topic(), tp.partition(), currentFo, currentUo)
    offsetRangeArray.add(or)
    currentFo = currentUo
    currentUo = currentFo + maxMessagePerTask
  }
  val or = OffsetRange(tp.topic, tp.partition, currentFo, uo)
  offsetRangeArray.add(or)
  offsetRangeArray.asScala
}
声明:本文内容由网友自发贡献,不代表【wpsshop博客】立场,版权归原作者所有,本站不承担相应法律责任。如您发现有侵权的内容,请联系我们。转载请注明出处:https://www.wpsshop.cn/w/我家自动化/article/detail/975888
推荐阅读
相关标签
  

闽ICP备14008679号