Repository: kylin Updated Branches: refs/heads/2.x-staging ed6d971e9 -> 44b020199
KYLIN-1435: Relax the checking for PartitionMetadata and logger the error code Signed-off-by: honma <[email protected]> Project: http://git-wip-us.apache.org/repos/asf/kylin/repo Commit: http://git-wip-us.apache.org/repos/asf/kylin/commit/87cc5094 Tree: http://git-wip-us.apache.org/repos/asf/kylin/tree/87cc5094 Diff: http://git-wip-us.apache.org/repos/asf/kylin/diff/87cc5094 Branch: refs/heads/2.x-staging Commit: 87cc5094c5a72107182a5854df7b3d377409d4dd Parents: ed6d971 Author: yangzhong <[email protected]> Authored: Tue Feb 23 17:37:11 2016 +0800 Committer: honma <[email protected]> Committed: Tue Feb 23 18:26:34 2016 +0800 ---------------------------------------------------------------------- .../java/org/apache/kylin/source/kafka/KafkaStreamingInput.java | 5 ++++- .../java/org/apache/kylin/source/kafka/util/KafkaUtils.java | 5 ++++- 2 files changed, 8 insertions(+), 2 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/kylin/blob/87cc5094/source-kafka/src/main/java/org/apache/kylin/source/kafka/KafkaStreamingInput.java ---------------------------------------------------------------------- diff --git a/source-kafka/src/main/java/org/apache/kylin/source/kafka/KafkaStreamingInput.java b/source-kafka/src/main/java/org/apache/kylin/source/kafka/KafkaStreamingInput.java index ee5a555..bcde47b 100644 --- a/source-kafka/src/main/java/org/apache/kylin/source/kafka/KafkaStreamingInput.java +++ b/source-kafka/src/main/java/org/apache/kylin/source/kafka/KafkaStreamingInput.java @@ -123,7 +123,10 @@ public class KafkaStreamingInput implements IStreamingInput { private Broker getLeadBroker() { final PartitionMetadata partitionMetadata = KafkaRequester.getPartitionMetadata(kafkaClusterConfig.getTopic(), partitionId, replicaBrokers, kafkaClusterConfig); - if (partitionMetadata != null && partitionMetadata.errorCode() == 0) { + if (partitionMetadata != null) { + if (partitionMetadata.errorCode() != 0){ + logger.warn("PartitionMetadata errorCode: "+partitionMetadata.errorCode()); + } replicaBrokers = partitionMetadata.replicas(); return partitionMetadata.leader(); } else { http://git-wip-us.apache.org/repos/asf/kylin/blob/87cc5094/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaUtils.java ---------------------------------------------------------------------- diff --git a/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaUtils.java b/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaUtils.java index 96b7fa7..d496041 100644 --- a/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaUtils.java +++ b/source-kafka/src/main/java/org/apache/kylin/source/kafka/util/KafkaUtils.java @@ -33,7 +33,10 @@ public final class KafkaUtils { public static Broker getLeadBroker(KafkaClusterConfig kafkaClusterConfig, int partitionId) { final PartitionMetadata partitionMetadata = KafkaRequester.getPartitionMetadata(kafkaClusterConfig.getTopic(), partitionId, kafkaClusterConfig.getBrokers(), kafkaClusterConfig); - if (partitionMetadata != null && partitionMetadata.errorCode() == 0) { + if (partitionMetadata != null) { + if (partitionMetadata.errorCode() != 0){ + logger.warn("PartitionMetadata errorCode: "+partitionMetadata.errorCode()); + } return partitionMetadata.leader(); } else { return null;
