Repository: kylin
Updated Branches:
  refs/heads/1.4-rc b25ed7b53 -> 7411a98db


KYLIN-1696: add catch exception during sending topic meta request to broker


Project: http://git-wip-us.apache.org/repos/asf/kylin/repo
Commit: http://git-wip-us.apache.org/repos/asf/kylin/commit/df46f522
Tree: http://git-wip-us.apache.org/repos/asf/kylin/tree/df46f522
Diff: http://git-wip-us.apache.org/repos/asf/kylin/diff/df46f522

Branch: refs/heads/1.4-rc
Commit: df46f5224fd28eb63b473e4085c1c3962a060705
Parents: b25ed7b
Author: kyotoYaho <[email protected]>
Authored: Mon May 16 13:21:54 2016 +0800
Committer: kyotoYaho <[email protected]>
Committed: Mon May 16 13:21:54 2016 +0800

----------------------------------------------------------------------
 .../kylin/source/kafka/util/KafkaRequester.java     | 16 ++++++++++++++--
 1 file changed, 14 insertions(+), 2 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/kylin/blob/df46f522/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaRequester.java
----------------------------------------------------------------------
diff --git 
a/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaRequester.java
 
b/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaRequester.java
index b78d30f..e39abc3 100644
--- 
a/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaRequester.java
+++ 
b/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaRequester.java
@@ -101,7 +101,13 @@ public final class KafkaRequester {
             consumer = getSimpleConsumer(broker, 
kafkaClusterConfig.getTimeout(), kafkaClusterConfig.getBufferSize(), 
"topic_meta_lookup");
             List<String> topics = 
Collections.singletonList(kafkaClusterConfig.getTopic());
             TopicMetadataRequest req = new TopicMetadataRequest(topics);
-            TopicMetadataResponse resp = consumer.send(req);
+            TopicMetadataResponse resp;
+            try{
+                resp = consumer.send(req);
+            }catch (Exception e){
+                logger.warn("cannot send TopicMetadataRequest successfully: " 
+ e);
+                continue;
+            }
             final List<TopicMetadata> topicMetadatas = resp.topicsMetadata();
             if (topicMetadatas.size() != 1) {
                 break;
@@ -129,7 +135,13 @@ public final class KafkaRequester {
             consumer = getSimpleConsumer(broker, 
kafkaClusterConfig.getTimeout(), kafkaClusterConfig.getBufferSize(), 
"topic_meta_lookup");
             List<String> topics = Collections.singletonList(topic);
             TopicMetadataRequest req = new TopicMetadataRequest(topics);
-            TopicMetadataResponse resp = consumer.send(req);
+            TopicMetadataResponse resp;
+            try{
+                resp = consumer.send(req);
+            }catch (Exception e){
+                logger.warn("cannot send TopicMetadataRequest successfully: " 
+ e);
+                continue;
+            }
             final List<TopicMetadata> topicMetadatas = resp.topicsMetadata();
             if (topicMetadatas.size() != 1) {
                 logger.warn("invalid topicMetadata size:" + 
topicMetadatas.size());

Reply via email to