Repository: kafka Updated Branches: refs/heads/trunk 7a36d3647 -> 80223baa3
HOTFIX: Rename WakeupException in MirrorMaker Author: Guozhang Wang <[email protected]> Reviewers: Ismael Juma Closes #375 from guozhangwang/HFWakeup Project: http://git-wip-us.apache.org/repos/asf/kafka/repo Commit: http://git-wip-us.apache.org/repos/asf/kafka/commit/80223baa Tree: http://git-wip-us.apache.org/repos/asf/kafka/tree/80223baa Diff: http://git-wip-us.apache.org/repos/asf/kafka/diff/80223baa Branch: refs/heads/trunk Commit: 80223baa3b7125fd656ed85b450498a21924ca71 Parents: 7a36d36 Author: Guozhang Wang <[email protected]> Authored: Tue Oct 27 19:32:42 2015 -0700 Committer: Guozhang Wang <[email protected]> Committed: Tue Oct 27 19:32:42 2015 -0700 ---------------------------------------------------------------------- core/src/main/scala/kafka/tools/MirrorMaker.scala | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/kafka/blob/80223baa/core/src/main/scala/kafka/tools/MirrorMaker.scala ---------------------------------------------------------------------- diff --git a/core/src/main/scala/kafka/tools/MirrorMaker.scala b/core/src/main/scala/kafka/tools/MirrorMaker.scala index 3cf754b..79001de 100755 --- a/core/src/main/scala/kafka/tools/MirrorMaker.scala +++ b/core/src/main/scala/kafka/tools/MirrorMaker.scala @@ -32,12 +32,13 @@ import kafka.message.MessageAndMetadata import kafka.metrics.KafkaMetricsGroup import kafka.serializer.DefaultDecoder import kafka.utils.{CommandLineUtils, CoreUtils, Logging} -import org.apache.kafka.clients.consumer.{ConsumerWakeupException, Consumer, ConsumerRecord, KafkaConsumer} +import org.apache.kafka.clients.consumer.{Consumer, ConsumerRecord, KafkaConsumer} import org.apache.kafka.clients.producer.internals.ErrorLoggingCallback import org.apache.kafka.clients.producer.{KafkaProducer, ProducerConfig, ProducerRecord, RecordMetadata} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.serialization.ByteArrayDeserializer import org.apache.kafka.common.utils.Utils +import org.apache.kafka.common.errors.WakeupException import scala.collection.JavaConversions._ import scala.util.control.ControlThrowable @@ -390,7 +391,7 @@ object MirrorMaker extends Logging with KafkaMetricsGroup { } catch { case cte: ConsumerTimeoutException => trace("Caught ConsumerTimeoutException, continue iteration.") - case cwe: ConsumerWakeupException => + case we: WakeupException => trace("Caught ConsumerWakeupException, continue iteration.") } maybeFlushAndCommitOffsets()
