From 2614eef40de6b23ee688968ef1dd50c0bfa48e4d Mon Sep 17 00:00:00 2001 From: trialblazerseee <84778104+trialblazerseee@users.noreply.github.com> Date: Wed, 2 Jul 2025 11:41:59 +0530 Subject: [PATCH] MOSIP-42131 - Multiple ActiveMQ Consumer configuration for activeMQ Signed-off-by: trialblazerseee <84778104+trialblazerseee@users.noreply.github.com> --- .../middleware/stage/AbisMiddleWareStage.java | 6 +- .../messagequeue/AbisMessageQueueImpl.java | 5 +- .../stage/ManualAdjudicationStage.java | 6 +- .../stage/ManualAdjudicationStageTest.java | 2 +- .../verification/stage/VerificationStage.java | 6 +- .../stage/VerificationStageTest.java | 2 +- .../StageHealthCheckHandler.java | 2 +- .../core/queue/impl/MosipActiveMqImpl.java | 60 ++++++++++--------- .../core/spi/queue/MosipQueueManager.java | 2 +- 9 files changed, 56 insertions(+), 35 deletions(-) diff --git a/registration-processor/core-processor/registration-processor-abis-middleware-stage/src/main/java/io/mosip/registartion/processor/abis/middleware/stage/AbisMiddleWareStage.java b/registration-processor/core-processor/registration-processor-abis-middleware-stage/src/main/java/io/mosip/registartion/processor/abis/middleware/stage/AbisMiddleWareStage.java index eedd7c44f46..e24934af344 100644 --- a/registration-processor/core-processor/registration-processor-abis-middleware-stage/src/main/java/io/mosip/registartion/processor/abis/middleware/stage/AbisMiddleWareStage.java +++ b/registration-processor/core-processor/registration-processor-abis-middleware-stage/src/main/java/io/mosip/registartion/processor/abis/middleware/stage/AbisMiddleWareStage.java @@ -163,6 +163,10 @@ public class AbisMiddleWareStage extends MosipVerticleAPIManager { private static final String ABIS_QUEUE_NOT_FOUND = "ABIS_QUEUE_NOT_FOUND"; private static final String TEXT_MESSAGE = "text"; + /** Set the Consumer Count which required to listen and process message parallel. */ + @Value("${mosip.regproc.abis.middleware.activemq.consumer.count:1}") + private Integer consumerCount; + /** * Get all the abis queue details,register listener to outbound queue's */ @@ -189,7 +193,7 @@ public void setListener(Message message) { } } }; - mosipQueueManager.consume(queue, abisQueue.getOutboundQueueName(), listener); + mosipQueueManager.consume(queue, abisQueue.getOutboundQueueName(), listener, consumerCount); } } catch (Exception e) { diff --git a/registration-processor/core-processor/registration-processor-abis/src/main/java/io/mosip/registration/processor/abis/messagequeue/AbisMessageQueueImpl.java b/registration-processor/core-processor/registration-processor-abis/src/main/java/io/mosip/registration/processor/abis/messagequeue/AbisMessageQueueImpl.java index d1caf09ddf0..39190c2c325 100644 --- a/registration-processor/core-processor/registration-processor-abis/src/main/java/io/mosip/registration/processor/abis/messagequeue/AbisMessageQueueImpl.java +++ b/registration-processor/core-processor/registration-processor-abis/src/main/java/io/mosip/registration/processor/abis/messagequeue/AbisMessageQueueImpl.java @@ -75,6 +75,9 @@ public class AbisMessageQueueImpl { /** The is connection. */ boolean isConnection = false; + /** consumer Count */ + private Integer consumerCoint = 1; + /** * Run abis queue. * @@ -96,7 +99,7 @@ public void setListener(Message message) { } }; mosipQueueManager.consume(abisQueueDetails.get(i).getMosipQueue(), - abisQueueDetails.get(i).getInboundQueueName(), listener); + abisQueueDetails.get(i).getInboundQueueName(), listener, consumerCoint); } isConnection = true; diff --git a/registration-processor/core-processor/registration-processor-manual-adjudication-stage/src/main/java/io/mosip/registration/processor/adjudication/stage/ManualAdjudicationStage.java b/registration-processor/core-processor/registration-processor-manual-adjudication-stage/src/main/java/io/mosip/registration/processor/adjudication/stage/ManualAdjudicationStage.java index 18291d13f5c..92b6da8e1ef 100644 --- a/registration-processor/core-processor/registration-processor-manual-adjudication-stage/src/main/java/io/mosip/registration/processor/adjudication/stage/ManualAdjudicationStage.java +++ b/registration-processor/core-processor/registration-processor-manual-adjudication-stage/src/main/java/io/mosip/registration/processor/adjudication/stage/ManualAdjudicationStage.java @@ -153,6 +153,10 @@ public class ManualAdjudicationStage extends MosipVerticleAPIManager { private static final String APPLICATION_JSON = "application/json"; + /** Set the Consumer Count which required to listen and process message parallel. */ + @Value("${mosip.regproc.manual.adjudication.activemq.consumer.count:1}") + private Integer consumerCount; + /** * Deploy stage. */ @@ -169,7 +173,7 @@ public void setListener(Message message) { } }; - mosipQueueManager.consume(queue, mvResponseAddress, listener); + mosipQueueManager.consume(queue, mvResponseAddress, listener, consumerCount); } else { throw new QueueConnectionNotFound(PlatformErrorMessages.RPR_PRT_QUEUE_CONNECTION_NULL.getMessage()); diff --git a/registration-processor/core-processor/registration-processor-manual-adjudication-stage/src/test/java/io/mosip/registration/processor/verification/stage/ManualAdjudicationStageTest.java b/registration-processor/core-processor/registration-processor-manual-adjudication-stage/src/test/java/io/mosip/registration/processor/verification/stage/ManualAdjudicationStageTest.java index f7d06c6fca0..90f6b330ab5 100644 --- a/registration-processor/core-processor/registration-processor-manual-adjudication-stage/src/test/java/io/mosip/registration/processor/verification/stage/ManualAdjudicationStageTest.java +++ b/registration-processor/core-processor/registration-processor-manual-adjudication-stage/src/test/java/io/mosip/registration/processor/verification/stage/ManualAdjudicationStageTest.java @@ -143,7 +143,7 @@ public void setUp() throws java.io.IOException, ApisResourceAccessException, Pac //Mockito.when(env.getProperty(SwaggerConstant.SERVER_SERVLET_PATH)).thenReturn("/registrationprocessor/v1/manualverification"); Mockito.when(mosipConnectionFactory.createConnection(any(), any(), any(), any(), anyList())) .thenReturn(mosipQueue); - Mockito.doReturn(new String("str").getBytes()).when(mosipQueueManager).consume(any(), any(), any()); + Mockito.doReturn(new String("str").getBytes()).when(mosipQueueManager).consume(any(), any(), any(), any()); Mockito.doNothing().when(router).setRoute(any()); Mockito.when(router.post(any())).thenReturn(null); Mockito.when(router.get(any())).thenReturn(null); diff --git a/registration-processor/core-processor/registration-processor-verification-stage/src/main/java/io/mosip/registration/processor/verification/stage/VerificationStage.java b/registration-processor/core-processor/registration-processor-verification-stage/src/main/java/io/mosip/registration/processor/verification/stage/VerificationStage.java index 5a6ed1b15ff..d1dde5c4837 100644 --- a/registration-processor/core-processor/registration-processor-verification-stage/src/main/java/io/mosip/registration/processor/verification/stage/VerificationStage.java +++ b/registration-processor/core-processor/registration-processor-verification-stage/src/main/java/io/mosip/registration/processor/verification/stage/VerificationStage.java @@ -152,6 +152,10 @@ public class VerificationStage extends MosipVerticleAPIManager { @Value("${mosip.regproc.verification.message.expiry-time-limit}") private Long messageExpiryTimeLimit; + /** Set the Consumer Count which required to listen and process message parallel. */ + @Value("${mosip.regproc.verification.activemq.consumer.count:1}") + private Integer consumerCount; + private static final String APPLICATION_JSON = "application/json"; /** @@ -170,7 +174,7 @@ public void setListener(Message message) { } }; - mosipQueueManager.consume(queue, mvResponseAddress, listener); + mosipQueueManager.consume(queue, mvResponseAddress, listener,consumerCount); } else { throw new QueueConnectionNotFound(PlatformErrorMessages.RPR_PRT_QUEUE_CONNECTION_NULL.getMessage()); diff --git a/registration-processor/core-processor/registration-processor-verification-stage/src/test/java/io/mosip/registration/processor/verification/stage/VerificationStageTest.java b/registration-processor/core-processor/registration-processor-verification-stage/src/test/java/io/mosip/registration/processor/verification/stage/VerificationStageTest.java index 58e48d709b2..a40ad5bd53e 100644 --- a/registration-processor/core-processor/registration-processor-verification-stage/src/test/java/io/mosip/registration/processor/verification/stage/VerificationStageTest.java +++ b/registration-processor/core-processor/registration-processor-verification-stage/src/test/java/io/mosip/registration/processor/verification/stage/VerificationStageTest.java @@ -125,7 +125,7 @@ public void setUp() throws java.io.IOException, ApisResourceAccessException, Pac //Mockito.when(env.getProperty(SwaggerConstant.SERVER_SERVLET_PATH)).thenReturn("/registrationprocessor/v1/manualverification"); Mockito.when(mosipConnectionFactory.createConnection(any(), any(), any(), any(), anyList())) .thenReturn(mosipQueue); - Mockito.doReturn(new String("str").getBytes()).when(mosipQueueManager).consume(any(), any(), any()); + Mockito.doReturn(new String("str").getBytes()).when(mosipQueueManager).consume(any(), any(), any(), any()); Mockito.doNothing().when(router).setRoute(any()); Mockito.when(router.post(any())).thenReturn(null); Mockito.when(router.get(any())).thenReturn(null); diff --git a/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/abstractverticle/StageHealthCheckHandler.java b/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/abstractverticle/StageHealthCheckHandler.java index 7edb0ae8188..44985179e47 100644 --- a/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/abstractverticle/StageHealthCheckHandler.java +++ b/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/abstractverticle/StageHealthCheckHandler.java @@ -179,7 +179,7 @@ public void setListener(Message message) { } } }; - mosipQueueManager.consume(mosipQueue, HealthConstant.QUEUE_ADDRESS, listener); + mosipQueueManager.consume(mosipQueue, HealthConstant.QUEUE_ADDRESS, listener, 1); isConsumerStarted = true; } } catch (Exception e) { diff --git a/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/queue/impl/MosipActiveMqImpl.java b/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/queue/impl/MosipActiveMqImpl.java index 4a54c65f1ae..b1120d4ed2f 100644 --- a/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/queue/impl/MosipActiveMqImpl.java +++ b/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/queue/impl/MosipActiveMqImpl.java @@ -25,6 +25,8 @@ import javax.jms.MessageProducer; import javax.jms.Session; import javax.jms.TextMessage; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; /** @@ -182,7 +184,7 @@ public Boolean send(MosipQueue mosipQueue, String message, String address, int m * .lang.Object, java.lang.String) */ @Override - public byte[] consume(MosipQueue mosipQueue, String address, QueueListener object) { + public byte[] consume(MosipQueue mosipQueue, String address, QueueListener object, Integer consumerCount) { regProcLogger.debug(LoggerFileConstant.SESSIONID.toString(), LoggerFileConstant.USERID.toString(), "", "MosipActiveMqImpl::consume()::entry"); @@ -195,34 +197,38 @@ public byte[] consume(MosipQueue mosipQueue, String address, QueueListener objec throw new InvalidConnectionException(PlatformErrorMessages.RPR_MQI_INVALID_CONNECTION.getMessage()); } - if (destination == null) { + if (connection == null) { setup(mosipActiveMq); } - MessageConsumer consumer; - try { - if (session == null) { - regProcLogger.error("Session is null. System will retry to create session"); - setup(mosipActiveMq); - } - destination = session.createQueue(address); - consumer = session.createConsumer(destination); - consumer.setMessageListener(QueueListenerFactory.getListener(mosipQueue.getQueueName(), object)); - } catch (JMSException | NullPointerException e) { - regProcLogger.error("*******CONSUME EXCEPTION *****", "*******CONSUME EXCEPTION *****", - "*******CONSUME EXCEPTION *****", ExceptionUtils.getFullStackTrace(e)); - regProcLogger.error(LoggerFileConstant.SESSIONID.toString(), LoggerFileConstant.REGISTRATIONID.toString(), - "", "MosipActiveMqImpl::consume():: error with error message " - + PlatformErrorMessages.RPR_MQI_UNABLE_TO_CONSUME_FROM_QUEUE.getMessage()); - - if (e instanceof NullPointerException && retryCount > 0) { - regProcLogger.warn(LoggerFileConstant.SESSIONID.toString(), LoggerFileConstant.REGISTRATIONID.toString(), - "", "Could not obtain queue connection. System will retry for "+ retryCount +" more times."); - retryCount = retryCount - 1; - consume(mosipQueue, address, object); - } else { - throw new ConnectionUnavailableException( - PlatformErrorMessages.RPR_MQI_UNABLE_TO_CONSUME_FROM_QUEUE.getMessage()); - } + + ExecutorService executorService = Executors.newFixedThreadPool(consumerCount); + + for(int i = 0; i < consumerCount; i++) { + executorService.submit(() -> { + MessageConsumer consumer; + try { + Session session1 = this.connection.createSession(false, Session.AUTO_ACKNOWLEDGE); + Destination destination1 = session1.createQueue(address); + consumer = session1.createConsumer(destination1); + consumer.setMessageListener(QueueListenerFactory.getListener(mosipQueue.getQueueName(), object)); + } catch (JMSException | NullPointerException e) { + regProcLogger.error("*******CONSUME EXCEPTION *****", "*******CONSUME EXCEPTION *****", + "*******CONSUME EXCEPTION *****", ExceptionUtils.getFullStackTrace(e)); + regProcLogger.error(LoggerFileConstant.SESSIONID.toString(), LoggerFileConstant.REGISTRATIONID.toString(), + "", "MosipActiveMqImpl::consume():: error with error message " + + PlatformErrorMessages.RPR_MQI_UNABLE_TO_CONSUME_FROM_QUEUE.getMessage()); + + if (e instanceof NullPointerException && retryCount > 0) { + regProcLogger.warn(LoggerFileConstant.SESSIONID.toString(), LoggerFileConstant.REGISTRATIONID.toString(), + "", "Could not obtain queue connection. System will retry for "+ retryCount +" more times."); + retryCount = retryCount - 1; + consume(mosipQueue, address, object, consumerCount); + } else { + throw new ConnectionUnavailableException( + PlatformErrorMessages.RPR_MQI_UNABLE_TO_CONSUME_FROM_QUEUE.getMessage()); + } + } + }); } regProcLogger.debug(LoggerFileConstant.SESSIONID.toString(), LoggerFileConstant.USERID.toString(), "", "MosipActiveMqImpl::consume()::exit"); diff --git a/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/spi/queue/MosipQueueManager.java b/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/spi/queue/MosipQueueManager.java index 11437250251..c371d685a3c 100644 --- a/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/spi/queue/MosipQueueManager.java +++ b/registration-processor/registration-processor-core/src/main/java/io/mosip/registration/processor/core/spi/queue/MosipQueueManager.java @@ -59,6 +59,6 @@ public interface MosipQueueManager{ * @param address The address * @return the original message */ - public V consume(T mosipQueue, String address, QueueListener object); + public V consume(T mosipQueue, String address, QueueListener object, Integer consumerCount); }