Overview
Apache RocketMQ provides two consumption models:
- Push Model: Uses
DefaultMQPushConsumerwhere message are actively pushed to consumers - Pull Model: Uses
DefaultLitePullConsumerwhere 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).