Skip to content

Instantly share code, notes, and snippets.

@BekirUzun
Last active May 21, 2026 12:00
Show Gist options
  • Select an option

  • Save BekirUzun/797550f6387583f50d82aebcc740e44c to your computer and use it in GitHub Desktop.

Select an option

Save BekirUzun/797550f6387583f50d82aebcc740e44c to your computer and use it in GitHub Desktop.
Simple RabbitMQ health indicator with min consumer and max message count

RabbitMQ Custom Health Check

Simple RabbitMQ health indicator with min consumer and max message count.

Usage

import lombok.RequiredArgsConstructor;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.core.RabbitAdmin;
import org.springframework.boot.health.contributor.HealthIndicator;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@RequiredArgsConstructor
@Configuration
public class RabbitQueueHealth {

    private static final int MIN_CONSUMER_COUNT = 2;
    private static final int MAX_MESSAGE_COUNT = 100000;

    private final Queue someQueue;
    private final Queue otherQueue;
    private final RabbitAdmin rabbitAdmin;

    @Bean
    public HealthIndicator rabbitMqQueueHealthIndicator() {
        return new RabbitQueueCheckHealthIndicator(rabbitAdmin)
                .addQueueCheck(someQueue, MAX_MESSAGE_COUNT, MIN_CONSUMER_COUNT)
                .addQueueCheck(otherQueue, Integer.MAX_VALUE, MIN_CONSUMER_COUNT)
    }
}
package com.project.configuration;
import lombok.Builder;
import lombok.RequiredArgsConstructor;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.core.RabbitAdmin;
import org.springframework.boot.health.contributor.Health;
import org.springframework.boot.health.contributor.HealthIndicator;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
@RequiredArgsConstructor
public class RabbitQueueCheckHealthIndicator implements HealthIndicator {
private final RabbitAdmin rabbitAdmin;
private final List<QueueCheck> queueChecks = new ArrayList<>();
private static final String QUEUE_NOT_FOUND = "Queue not found";
public RabbitQueueCheckHealthIndicator addQueueCheck(Queue queue, int maxMessageCount, int minConsumerCount) {
var queueCheck = QueueCheck.builder()
.queue(queue)
.maxMessageCount(maxMessageCount)
.minConsumerCount(minConsumerCount)
.build();
queueChecks.add(queueCheck);
return this;
}
@Override
public Health health() {
var healthBuilder = Health.up();
boolean healthy = true;
for (var queueCheck : queueChecks) {
if (!checkQueueHealth(healthBuilder, queueCheck)) {
healthy = false;
}
}
if (!healthy) {
healthBuilder.down();
}
return healthBuilder.build();
}
private boolean checkQueueHealth(Health.Builder healthBuilder, QueueCheck queueCheck) {
var queueName = queueCheck.queue().getName();
var queueProperties = rabbitAdmin.getQueueProperties(queueName);
if (queueProperties == null) {
healthBuilder.withDetail(queueName, QUEUE_NOT_FOUND);
return false;
}
int messageCount = toInt(queueProperties.get(RabbitAdmin.QUEUE_MESSAGE_COUNT));
int consumerCount = toInt(queueProperties.get(RabbitAdmin.QUEUE_CONSUMER_COUNT));
var details = new HashMap<String, Object>();
details.put("messageCount", messageCount);
details.put("consumerCount", consumerCount);
details.put("maxMessageCount", queueCheck.maxMessageCount());
details.put("minConsumerCount", queueCheck.minConsumerCount());
healthBuilder.withDetail(queueName, details);
return messageCount <= queueCheck.maxMessageCount() &&
consumerCount >= queueCheck.minConsumerCount();
}
private int toInt(Object value) {
if (value instanceof Number number) {
return number.intValue();
}
return 0;
}
@Builder
private record QueueCheck(Queue queue, int maxMessageCount, int minConsumerCount) {
}
}
package com.project.configuration;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.QueueBuilder;
import org.springframework.amqp.rabbit.core.RabbitAdmin;
import java.util.Properties;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.InstanceOfAssertFactories.map;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
class RabbitQueueCheckHealthIndicatorTest {
@InjectMocks
private RabbitQueueCheckHealthIndicator rabbitQueueCheckHealthIndicator;
@Mock
private RabbitAdmin rabbitAdmin;
@Test
void it_should_report_up_when_message_and_consumer_counts_are_within_thresholds() {
Queue queue1 = QueueBuilder.durable("queue.name").build();
Queue queue2 = QueueBuilder.durable("queue.name2").build();
rabbitQueueCheckHealthIndicator
.addQueueCheck(queue1, 10, 1)
.addQueueCheck(queue2, Integer.MAX_VALUE, 1);
when(rabbitAdmin.getQueueProperties(queue1.getName())).thenReturn(buildQueueProperties(5, 2));
when(rabbitAdmin.getQueueProperties(queue2.getName())).thenReturn(buildQueueProperties(0, 1));
var health = rabbitQueueCheckHealthIndicator.health();
assertThat(health).isNotNull();
assertThat(health.getStatus().getCode()).isEqualTo("UP");
assertThat(health.getDetails())
.containsKey(queue1.getName())
.containsKey(queue2.getName());
assertThat(health.getDetails().get(queue1.getName()))
.asInstanceOf(map(String.class, Object.class))
.containsEntry("messageCount", 5)
.containsEntry("consumerCount", 2)
.containsEntry("maxMessageCount", 10)
.containsEntry("minConsumerCount", 1);
}
@Test
void it_should_report_down_when_queue_is_not_found() {
Queue queue = QueueBuilder.durable("queue.name").build();
rabbitQueueCheckHealthIndicator.addQueueCheck(queue, 100, 1);
when(rabbitAdmin.getQueueProperties(queue.getName())).thenReturn(null);
var health = rabbitQueueCheckHealthIndicator.health();
assertThat(health).isNotNull();
assertThat(health.getStatus().getCode()).isEqualTo("DOWN");
assertThat(health.getDetails()).containsEntry(queue.getName(), "Queue not found");
}
@Test
void it_should_report_down_when_message_count_exceeds_max_threshold() {
Queue queue = QueueBuilder.durable("queue.name").build();
rabbitQueueCheckHealthIndicator.addQueueCheck(queue, 1000, 1);
when(rabbitAdmin.getQueueProperties(queue.getName())).thenReturn(buildQueueProperties(1001, 1));
var health = rabbitQueueCheckHealthIndicator.health();
assertThat(health).isNotNull();
assertThat(health.getStatus().getCode()).isEqualTo("DOWN");
}
@Test
void it_should_report_down_when_consumer_count_is_below_min_threshold() {
Queue queue = QueueBuilder.durable("queue.name").build();
rabbitQueueCheckHealthIndicator.addQueueCheck(queue, 10, 2);
when(rabbitAdmin.getQueueProperties(queue.getName())).thenReturn(buildQueueProperties(1, 1));
var health = rabbitQueueCheckHealthIndicator.health();
assertThat(health).isNotNull();
assertThat(health.getStatus().getCode()).isEqualTo("DOWN");
}
private Properties buildQueueProperties(int messageCount, int consumerCount) {
Properties properties = new Properties();
properties.put(RabbitAdmin.QUEUE_MESSAGE_COUNT, messageCount);
properties.put(RabbitAdmin.QUEUE_CONSUMER_COUNT, consumerCount);
return properties;
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment