From: sunil unnava Date: Thu, 6 Dec 2018 11:22:32 +0000 (-0500) Subject: Fix for Kafka Consumer is not safe error X-Git-Tag: 1.1.14^0 X-Git-Url: https://gerrit.onap.org/r/gitweb?p=dmaap%2Fmessagerouter%2Fmsgrtr.git;a=commitdiff_plain;h=83746dbc42bad55e52d4bed2617d0d0ca8634cb5 Fix for Kafka Consumer is not safe error Issue-ID: DMAAP-896 Change-Id: I085dbad1248790796e220267cb3e603ecc6c1067 Signed-off-by: sunil unnava --- diff --git a/pom.xml b/pom.xml index affe8f7..c213101 100644 --- a/pom.xml +++ b/pom.xml @@ -14,7 +14,7 @@ 4.0.0 org.onap.dmaap.messagerouter.msgrtr msgrtr - 1.1.13-SNAPSHOT + 1.1.14-SNAPSHOT jar dmaap-messagerouter-msgrtr Message Router - Restful interface built for kafka diff --git a/src/main/java/org/onap/dmaap/dmf/mr/backends/kafka/Kafka011Consumer.java b/src/main/java/org/onap/dmaap/dmf/mr/backends/kafka/Kafka011Consumer.java index 347f625..2ec323e 100644 --- a/src/main/java/org/onap/dmaap/dmf/mr/backends/kafka/Kafka011Consumer.java +++ b/src/main/java/org/onap/dmaap/dmf/mr/backends/kafka/Kafka011Consumer.java @@ -42,8 +42,7 @@ import org.apache.kafka.common.KafkaException; import org.onap.dmaap.dmf.mr.backends.Consumer; import org.onap.dmaap.dmf.mr.constants.CambriaConstants; - - +import com.att.ajsc.filemonitor.AJSCPropertiesMap; import com.att.eelf.configuration.EELFLogger; import com.att.eelf.configuration.EELFManager; @@ -83,6 +82,12 @@ public class Kafka011Consumer implements Consumer { state = Kafka011Consumer.State.OPENED; kConsumer = cc; fKafkaLiveLockAvoider = klla; + + String consumerTimeOut = AJSCPropertiesMap.getProperty(CambriaConstants.msgRtr_prop, + "consumer.timeout"); + if (null != consumerTimeOut) { + consumerPollTimeOut = Integer.parseInt(consumerTimeOut); + } synchronized (kConsumer) { kConsumer.subscribe(Arrays.asList(topic)); } @@ -147,7 +152,7 @@ public class Kafka011Consumer implements Consumer { ExecutorService service = Executors.newSingleThreadExecutor(); service.execute(future); try { - future.get(5, TimeUnit.SECONDS); // wait 1 + future.get(consumerPollTimeOut, TimeUnit.SECONDS); // wait 1 // second } catch (TimeoutException ex) { // timed out. Try to stop the code if possible. @@ -370,6 +375,7 @@ public class Kafka011Consumer implements Consumer { private long offset; private Kafka011Consumer.State state; private KafkaLiveLockAvoider2 fKafkaLiveLockAvoider; + private int consumerPollTimeOut=5; private static final EELFLogger log = EELFManager.getInstance().getLogger(Kafka011Consumer.class); private final LinkedBlockingQueue> fPendingMsgs; diff --git a/version.properties b/version.properties index 3dd3954..c9b51c7 100644 --- a/version.properties +++ b/version.properties @@ -27,7 +27,7 @@ major=1 minor=1 -patch=13 +patch=14 base_version=${major}.${minor}.${patch}