Browse Source

KAFKA-1890 Fix bug preventing Mirror Maker from successful rebalance; reviewed by Gwen Shapira and Neha Narkhede

pull/1442/head
Jiangjie Qin 10 years ago committed by Neha Narkhede
parent
commit
8cff9119f8
  1. 6
      core/src/main/scala/kafka/tools/MirrorMaker.scala

6
core/src/main/scala/kafka/tools/MirrorMaker.scala

@ -213,11 +213,11 @@ object MirrorMaker extends Logging with KafkaMetricsGroup { @@ -213,11 +213,11 @@ object MirrorMaker extends Logging with KafkaMetricsGroup {
val customRebalanceListenerClass = options.valueOf(consumerRebalanceListenerOpt)
val customRebalanceListener = {
if (customRebalanceListenerClass != null)
Utils.createObject[ConsumerRebalanceListener](customRebalanceListenerClass)
Some(Utils.createObject[ConsumerRebalanceListener](customRebalanceListenerClass))
else
null
None
}
consumerRebalanceListener = new InternalRebalanceListener(mirrorDataChannel, Some(customRebalanceListener))
consumerRebalanceListener = new InternalRebalanceListener(mirrorDataChannel, customRebalanceListener)
connector.setConsumerRebalanceListener(consumerRebalanceListener)
// create producer threads

Loading…
Cancel
Save