Spring boot Using Rabbitmq or kafka by Profile

Viewed 37

I have a problem with implementing the RabbitMQ listener or Kafka consumer by Springboot Profile.

The RabbitMQ library is added to gradle (implementation 'org.springframework.amqp:spring-rabbit:2.4.6')

And the sources of RabbitMQ listener and Kafka consumer are implemented as follows.

Profile Configured to use Kafka not RabbitMQ for the current message queue. I set it to spring.profiles.active = ProfileConfig.SIMUL, but the following error occurs.

Since it was filtered by Profile, RabbitMQ should not be run, but I don't know why this problem occurs.

@Component
@Profile(ProfileConfig.DEDICATE)
@Slf4j
@AllArgsConstructor
public class RabbitMQListener {
  private final ObjectMapper objectMapper;
  private final CollectionProcessor collectionProcessor;
  private final CollectionManagement collectionManagement;
  private final AssetManagement assetManagement;
  private final AssetProcessor assetProcessor;
  private final TraceService traceService;

  @RabbitListener(queues = {"${rabbitmq.queue.name}"})
  public void receiveMessage(final String message) {
    try {
      JSONParser parser = new JSONParser();
      JSONObject json = (JSONObject) parser.parse(message);
      String messageType = json.get("messageType").toString();
      log.debug("Receive Queue  Key={}, Message = {}", messageType, message);
      AsyncType asyncType = AsyncType.valueOf(messageType);
      executeMessage(asyncType, message);
    } catch (JsonProcessingException | IllegalArgumentException | ParseException e) {
      traceService.removeTraceId();
      traceService.printErrorLog(log, "Fail to deal receive message.", e, PrintStackPolicy.ALL);
    }
  }
}

@Service
@Profile({ProfileConfig.DEV, ProfileConfig.PROD, ProfileConfig.SIMUL})
@Slf4j
@AllArgsConstructor
public class KafkaConsumer {

  private final ObjectMapper objectMapper;
  private final CollectionProcessor collectionProcessor;
  private final CollectionManagement collectionManagement;
  private final AssetManagement assetManagement;
  private final AssetProcessor assetProcessor;
  private final TraceService traceService;

  @KafkaListener(topics = {"${aws.kafka.topic}"})
  public void consume(@Payload String message, @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) String key){
    try {
      JSONParser parser = new JSONParser();
      JSONObject json = (JSONObject) parser.parse(message);
      String messageType = json.get("messageType").toString();
      log.debug("Receive Queue  Key={}, Message = {}", messageType, message);
      AsyncType asyncType = AsyncType.valueOf(messageType);
      executeMessage(asyncType, message);
    } catch (JsonProcessingException | IllegalArgumentException | ParseException e) {
      traceService.removeTraceId();
      traceService.printErrorLog(log, "Fail to deal receive message.", e, PrintStackPolicy.ALL);
    }
  }
}
WARN  22-09-08 14:41:23 Rabbit health check failed - [RMI TCP Connection(5)-xxx.xxx.xxx.xxx] [RabbitHealthIndicator:94]
org.springframework.amqp.AmqpConnectException: java.net.ConnectException: Connection refused
    at org.springframework.amqp.rabbit.support.RabbitExceptionTranslator.convertRabbitAccessException(RabbitExceptionTranslator.java:61)
    at org.springframework.amqp.rabbit.connection.AbstractConnectionFactory.createBareConnection(AbstractConnectionFactory.java:602)
    at org.springframework.amqp.rabbit.connection.CachingConnectionFactory.createConnection(CachingConnectionFactory.java:724)
    at org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils.createConnection(ConnectionFactoryUtils.java:252)
    at org.springframework.amqp.rabbit.core.RabbitTemplate.doExecute(RabbitTemplate.java:2175)
    at org.springframework.amqp.rabbit.core.RabbitTemplate.execute(RabbitTemplate.java:2148)
    at org.springframework.amqp.rabbit.core.RabbitTemplate.execute(RabbitTemplate.java:2128)
    at org.springframework.boot.actuate.amqp.RabbitHealthIndicator.getVersion(RabbitHealthIndicator.java:49)
    at org.springframework.boot.actuate.amqp.RabbitHealthIndicator.doHealthCheck(RabbitHealthIndicator.java:44)
    at org.springframework.boot.actuate.health.AbstractHealthIndicator.health(AbstractHealthIndicator.java:82)
    at org.springframework.boot.actuate.health.HealthIndicator.getHealth(HealthIndicator.java:37)
    at org.springframework.boot.actuate.health.HealthEndpoint.getHealth(HealthEndpoint.java:77)
    at org.springframework.boot.actuate.health.HealthEndpoint.getHealth(HealthEndpoint.java:40)
    at org.springframework.boot.actuate.health.HealthEndpointSupport.getContribution(HealthEndpointSupport.java:130)
    at org.springframework.boot.actuate.health.HealthEndpointSupport.getAggregateContribution(HealthEndpointSupport.java:141)
    at org.springframework.boot.actuate.health.HealthEndpointSupport.getContribution(HealthEndpointSupport.java:126)
    at org.springframework.boot.actuate.health.HealthEndpointSupport.getHealth(HealthEndpointSupport.java:95)
    at org.springframework.boot.actuate.health.HealthEndpointSupport.getHealth(HealthEndpointSupport.java:66)
    at org.springframework.boot.actuate.health.HealthEndpoint.health(HealthEndpoint.java:71)
    at org.springframework.boot.actuate.health.HealthEndpoint.health(HealthEndpoint.java:61)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:568)
    at org.springframework.util.ReflectionUtils.invokeMethod(ReflectionUtils.java:282)
    at org.springframework.boot.actuate.endpoint.invoke.reflect.ReflectiveOperationInvoker.invoke(ReflectiveOperationInvoker.java:74)
    at org.springframework.boot.actuate.endpoint.annotation.AbstractDiscoveredOperation.invoke(AbstractDiscoveredOperation.java:60)
    at org.springframework.boot.actuate.endpoint.jmx.EndpointMBean.invoke(EndpointMBean.java:122)
    at org.springframework.boot.actuate.endpoint.jmx.EndpointMBean.invoke(EndpointMBean.java:97)
    at java.management/com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.invoke(DefaultMBeanServerInterceptor.java:814)
    at java.management/com.sun.jmx.mbeanserver.JmxMBeanServer.invoke(JmxMBeanServer.java:802)
    at java.management.rmi/javax.management.remote.rmi.RMIConnectionImpl.doOperation(RMIConnectionImpl.java:1472)
    at java.management.rmi/javax.management.remote.rmi.RMIConnectionImpl$PrivilegedOperation.run(RMIConnectionImpl.java:1310)
    at java.management.rmi/javax.management.remote.rmi.RMIConnectionImpl.doPrivilegedOperation(RMIConnectionImpl.java:1405)
    at java.management.rmi/javax.management.remote.rmi.RMIConnectionImpl.invoke(RMIConnectionImpl.java:829)
    at java.base/jdk.internal.reflect.GeneratedMethodAccessor219.invoke(Unknown Source)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:568)
    at java.rmi/sun.rmi.server.UnicastServerRef.dispatch(UnicastServerRef.java:360)
    at java.rmi/sun.rmi.transport.Transport$1.run(Transport.java:200)
    at java.rmi/sun.rmi.transport.Transport$1.run(Transport.java:197)
    at java.base/java.security.AccessController.doPrivileged(AccessController.java:712)
    at java.rmi/sun.rmi.transport.Transport.serviceCall(Transport.java:196)
    at java.rmi/sun.rmi.transport.tcp.TCPTransport.handleMessages(TCPTransport.java:587)
    at java.rmi/sun.rmi.transport.tcp.TCPTransport$ConnectionHandler.run0(TCPTransport.java:828)
    at java.rmi/sun.rmi.transport.tcp.TCPTransport$ConnectionHandler.lambda$run$0(TCPTransport.java:705)
    at java.base/java.security.AccessController.doPrivileged(AccessController.java:399)
    at java.rmi/sun.rmi.transport.tcp.TCPTransport$ConnectionHandler.run(TCPTransport.java:704)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
    at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: java.net.ConnectException: Connection refused
    at java.base/sun.nio.ch.Net.pollConnect(Native Method)
    at java.base/sun.nio.ch.Net.pollConnectNow(Net.java:672)
    at java.base/sun.nio.ch.NioSocketImpl.timedFinishConnect(NioSocketImpl.java:549)
    at java.base/sun.nio.ch.NioSocketImpl.connect(NioSocketImpl.java:597)
    at java.base/java.net.SocksSocketImpl.connect(SocksSocketImpl.java:327)
    at java.base/java.net.Socket.connect(Socket.java:633)
    at com.rabbitmq.client.impl.SocketFrameHandlerFactory.create(SocketFrameHandlerFactory.java:60)
    at com.rabbitmq.client.ConnectionFactory.newConnection(ConnectionFactory.java:1223)
    at com.rabbitmq.client.ConnectionFactory.newConnection(ConnectionFactory.java:1173)
    at org.springframework.amqp.rabbit.connection.AbstractConnectionFactory.connectAddresses(AbstractConnectionFactory.java:640)
    at org.springframework.amqp.rabbit.connection.AbstractConnectionFactory.connect(AbstractConnectionFactory.java:615)
    at org.springframework.amqp.rabbit.connection.AbstractConnectionFactory.createBareConnection(AbstractConnectionFactory.java:565)
    ... 49 common frames omitted
0 Answers
Related