Created
July 26, 2012 17:41
-
-
Save baroquebobcat/3183413 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Index: core/src/main/scala/kafka/consumer/TopicCount.scala | |
| =================================================================== | |
| --- core/src/main/scala/kafka/consumer/TopicCount.scala (revision 1355144) | |
| +++ core/src/main/scala/kafka/consumer/TopicCount.scala (working copy) | |
| @@ -18,12 +18,37 @@ | |
| package kafka.consumer | |
| import scala.collection._ | |
| -import scala.util.parsing.json.JSON | |
| +import scala.util.parsing.json._ | |
| import org.I0Itec.zkclient.ZkClient | |
| import java.util.regex.Pattern | |
| import kafka.utils.{ZKGroupDirs, ZkUtils, Logging} | |
| +private class JSONParser extends Parser { | |
| + defaultNumberParser = {input : String => input.toInt} | |
| + | |
| + // copied the important parse methods from the JSON singleton into an | |
| + // instantiable class. | |
| + def parseRaw(input : String) : Option[JSONType] = | |
| + phrase(root)(new lexical.Scanner(input)) match { | |
| + case Success(result, _) => Some(result) | |
| + case _ => None | |
| + } | |
| + def parseFull(input: String): Option[Any] = | |
| + parseRaw(input) match { | |
| + case Some(data) => Some(resolveType(data)) | |
| + case None => None | |
| + } | |
| + | |
| + def resolveType(input: Any): Any = input match { | |
| + case JSONObject(data) => data.transform { | |
| + case (k,v) => resolveType(v) | |
| + } | |
| + case JSONArray(data) => data.map(resolveType) | |
| + case x => x | |
| + } | |
| +} | |
| + | |
| private[kafka] trait TopicCount { | |
| def getConsumerThreadIdsPerTopic: Map[String, Set[String]] | |
| @@ -60,9 +85,6 @@ | |
| private val BLACKLIST_PATTERN = | |
| Pattern.compile("""!(\p{Digit}+)!(.*)""") | |
| - val myConversionFunc = {input : String => input.toInt} | |
| - JSON.globalNumberParser = myConversionFunc | |
| - | |
| def constructTopicCount(group: String, | |
| consumerId: String, | |
| zkClient: ZkClient) : TopicCount = { | |
| @@ -94,7 +116,8 @@ | |
| else { | |
| var topMap : Map[String,Int] = null | |
| try { | |
| - JSON.parseFull(topicCountString) match { | |
| + val parser: JSONParser = new JSONParser | |
| + parser.parseFull(topicCountString) match { | |
| case Some(m) => topMap = m.asInstanceOf[Map[String,Int]] | |
| case None => throw new RuntimeException("error constructing TopicCount : " + topicCountString) | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| [...] kafka.consumer.TopicCount$.constructTopicCount:39] ERROR: error parsing consumer json string [...] | |
| java.lang.NullPointerException | |
| at scala.util.parsing.combinator.Parsers$NoSuccess.<init>(Parsers.scala:131) | |
| at scala.util.parsing.combinator.Parsers$Failure.<init>(Parsers.scala:158) | |
| at scala.util.parsing.combinator.Parsers$$anonfun$acceptIf$1.apply(Parsers.scala:489) | |
| ... | |
| at scala.util.parsing.combinator.Parsers$$anon$2.apply(Parsers.scala:742) | |
| at scala.util.parsing.json.JSON$.parseRaw(JSON.scala:71) | |
| at scala.util.parsing.json.JSON$.parseFull(JSON.scala:85) | |
| at kafka.consumer.TopicCount$.constructTopicCount(TopicCount.scala:32) | |
| at kafka.consumer.ZookeeperConsumerConnector$ZKRebalancerListener.kafka$consumer$ZookeeperConsumerConnector$ZKRebalancerListener$$getTopicCount(ZookeeperConsumerConnector.scala:422) | |
| at kafka.consumer.ZookeeperConsumerConnector$ZKRebalancerListener.kafka$consumer$ZookeeperConsumerConnector$ZKRebalancerListener$$rebalance(ZookeeperConsumerConnector.scala:460) | |
| at kafka.consumer.ZookeeperConsumerConnector$ZKRebalancerListener$$anonfun$syncedRebalance$1.apply$mcVI$sp(ZookeeperConsumerConnector.scala:437) | |
| at scala.collection.immutable.Range$ByOne$class.foreach$mVc$sp(Range.scala:282) | |
| at scala.collection.immutable.Range$$anon$2.foreach$mVc$sp(Range.scala:265) | |
| at kafka.consumer.ZookeeperConsumerConnector$ZKRebalancerListener.syncedRebalance(ZookeeperConsumerConnector.scala:433) | |
| at kafka.consumer.ZookeeperConsumerConnector$ZKRebalancerListener.handleChildChange(ZookeeperConsumerConnector.scala:375) | |
| at org.I0Itec.zkclient.ZkClient$7.run(ZkClient.java:568) | |
| at org.I0Itec.zkclient.ZkEventThread.run(ZkEventThread.java:71) |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment