mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-28 20:09:14 +08:00
* Schedule example add the default NamesrvAddr * Schedule example add the default NamesrvAddr * update annotation
This commit is contained in:
committed by
GitHub
parent
6f66b5d2a9
commit
febb21092b
+17
-15
@@ -17,31 +17,33 @@
|
||||
package org.apache.rocketmq.example.schedule;
|
||||
|
||||
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
|
||||
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
|
||||
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
|
||||
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
public class ScheduledMessageConsumer {
|
||||
|
||||
|
||||
public static final String CONSUMER_GROUP = "ExampleConsumer";
|
||||
public static final String DEFAULT_NAMESRVADDR = "127.0.0.1:9876";
|
||||
public static final String TOPIC = "TestTopic";
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
// Instantiate message consumer
|
||||
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ExampleConsumer");
|
||||
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(CONSUMER_GROUP);
|
||||
|
||||
// Uncomment the following line while debugging, namesrvAddr should be set to your local address
|
||||
// consumer.setNamesrvAddr(DEFAULT_NAMESRVADDR);
|
||||
|
||||
// Subscribe topics
|
||||
consumer.subscribe("TestTopic", "*");
|
||||
consumer.subscribe(TOPIC, "*");
|
||||
// Register message listener
|
||||
consumer.registerMessageListener(new MessageListenerConcurrently() {
|
||||
@Override
|
||||
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) {
|
||||
for (MessageExt message : messages) {
|
||||
// Print approximate delay time period
|
||||
System.out.printf("Receive message[msgId=%s %d ms later]\n", message.getMsgId(),
|
||||
System.currentTimeMillis() - message.getStoreTimestamp());
|
||||
}
|
||||
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
|
||||
consumer.registerMessageListener((MessageListenerConcurrently) (messages, context) -> {
|
||||
for (MessageExt message : messages) {
|
||||
// Print approximate delay time period
|
||||
System.out.printf("Receive message[msgId=%s %d ms later]\n", message.getMsgId(),
|
||||
System.currentTimeMillis() - message.getStoreTimestamp());
|
||||
}
|
||||
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
|
||||
});
|
||||
// Launch consumer
|
||||
consumer.start();
|
||||
|
||||
+16
-5
@@ -17,25 +17,36 @@
|
||||
package org.apache.rocketmq.example.schedule;
|
||||
|
||||
import org.apache.rocketmq.client.producer.DefaultMQProducer;
|
||||
import org.apache.rocketmq.client.producer.SendResult;
|
||||
import org.apache.rocketmq.common.message.Message;
|
||||
|
||||
public class ScheduledMessageProducer {
|
||||
|
||||
public static final String PRODUCER_GROUP = "ExampleProducerGroup";
|
||||
public static final String DEFAULT_NAMESRVADDR = "127.0.0.1:9876";
|
||||
public static final String TOPIC = "TestTopic";
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
// Instantiate a producer to send scheduled messages
|
||||
DefaultMQProducer producer = new DefaultMQProducer("ExampleProducerGroup");
|
||||
DefaultMQProducer producer = new DefaultMQProducer(PRODUCER_GROUP);
|
||||
|
||||
// Uncomment the following line while debugging, namesrvAddr should be set to your local address
|
||||
// producer.setNamesrvAddr(DEFAULT_NAMESRVADDR);
|
||||
|
||||
// Launch producer
|
||||
producer.start();
|
||||
int totalMessagesToSend = 100;
|
||||
for (int i = 0; i < totalMessagesToSend; i++) {
|
||||
Message message = new Message("TestTopic", ("Hello scheduled message " + i).getBytes());
|
||||
Message message = new Message(TOPIC, ("Hello scheduled message " + i).getBytes());
|
||||
// This message will be delivered to consumer 10 seconds later.
|
||||
message.setDelayTimeLevel(3);
|
||||
// Send the message
|
||||
producer.send(message);
|
||||
SendResult result = producer.send(message);
|
||||
System.out.print(result);
|
||||
}
|
||||
|
||||
|
||||
// Shutdown producer after use.
|
||||
producer.shutdown();
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user