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());
