Configuring RocketMQ Push Consumer in Spring Boot for Message Flow Control

Overview

Apache RocketMQ provides two consumption models:

  • Push Model: Uses DefaultMQPushConsumer where message are actively pushed to consumers
  • Pull Model: Uses DefaultLitePullConsumer where consumers actively pull messages

While the pull model offers more control over message retrieval timing, it lacks built-in flow control mechanisms. For scenarios requiring message rate limiting and batch size configuration, the push model with DefaultMQPushConsumer is preferred.

Maven Dependency

Replace the auto-configuration starter with the core client library:

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>${rocketmq.version}</version>
</dependency>

Consumer Configuration Class

Create a configuration class to initialize the push consumer with flow control parameters:

@Configuration
@Slf4j
public class RocketMqConsumerConfig {
    
    @Value("${rocketmq.namesrv.address}")
    private String serverAddress;
    
    @Value("${consumer.poll.delay.ms}")
    private long pollingDelayMs;
    
    @Value("${consumer.batch.size}")
    private int batchSize;
    
    @Autowired
    private BusinessMessageHandler messageHandler;
    
    @PostConstruct
    public void initializeConsumer() {
        configureConsumer("USER_NOTIFICATION_GROUP", "NOTIFICATION_TOPIC", messageHandler);
    }
    
    private void configureConsumer(String consumerGroup, String topicName, MessageListener listener) {
        try {
            log.info("Initializing consumer for group: {}", consumerGroup);
            
            DefaultMQPushConsumer mqConsumer = new DefaultMQPushConsumer(consumerGroup);
            mqConsumer.setNamesrvAddr(serverAddress);
            
            // Configure message polling interval in milliseconds
            mqConsumer.setPullInterval(pollingDelayMs);
            
            // Set maximum messages per poll operation
            mqConsumer.setPullBatchSize(batchSize);
            
            // Subscribe to specified topic
            mqConsumer.subscribe(topicName, "*");
            
            // Register appropriate listener type
            if (listener instanceof MessageListenerConcurrently) {
                mqConsumer.registerMessageListener((MessageListenerConcurrently) listener);
            } else if (listener instanceof MessageListenerOrderly) {
                mqConsumer.registerMessageListener((MessageListenerOrderly) listener);
            } else {
                log.error("Unsupported listener type for group: {}", consumerGroup);
                return;
            }
            
            mqConsumer.start();
            log.info("Consumer successfully started for group: {}", consumerGroup);
        } catch (Exception exception) {
            log.error("Failed to start consumer for group: {}", consumerGroup, exception);
        }
    }
}

Message Handler Implementation

Implement business logic in a separate component:

@Component
@Slf4j
public class BusinessMessageHandler implements MessageListenerOrderly {
    
    @Autowired
    private NotificationRepository notificationRepo;
    
    @Override
    public ConsumeOrderlyStatus consumeMessage(List<MessageExt> messages, ConsumeOrderlyContext context) {
        try {
            MessageExt receivedMsg = messages.get(0);
            NotificationData data = deserializePayload(receivedMsg.getBody());
            
            log.info("Processing notification: {}", data.toString());
            
            // Process business logic here
            return ConsumeOrderlyStatus.SUCCESS;
        } catch (Exception error) {
            log.error("Message processing failed: {}", error.getMessage(), error);
            return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT;
        }
    }
    
    private NotificationData deserializePayload(byte[] payload) {
        try {
            ByteArrayInputStream inputStream = new ByteArrayInputStream(payload);
            ObjectInputStream objectStream = new ObjectInputStream(inputStream);
            
            NotificationData result = (NotificationData) objectStream.readObject();
            
            objectStream.close();
            inputStream.close();
            
            return result;
        } catch (Exception ex) {
            log.error("Deserialization error: {}", ex.getMessage(), ex);
            throw new ProcessingException("Failed to deserialize message");
        }
    }
}

Configuration Properties

Add these properties to your application.yml:

rocketmq:
  namesrv:
    address: localhost:9876

consumer:
  poll:
    delay:
      ms: 1000
  batch:
    size: 16

The poll.delay.ms controls how frequently the consumer polls for new messages (default 0 means immediate polling). The batch.size determines how many messages are retrieved in each polling cycle (default 32).

Tags: spring-boot rocketmq messaging flow-control message-consumer

Posted on Thu, 08 Oct 2026 16:47:52 +0000 by Delqath