//product publicclassSyncProducer { publicstaticvoidmain(String[] args)throws Exception { //Instantiate with a producer group name. DefaultMQProducerproducer=new DefaultMQProducer("please_rename_unique_group_name");//把双引号里的改成你的rocketmq创建的group name // Specify name server addresses. producer.setNamesrvAddr("localhost:9876");//这里填你的rocketmq占用的端口 //Launch the instance. producer.start(); for (inti=0; i < 100; i++) { //Create a message instance, specifying topic, tag and message body. Messagemsg=newMessage("TopicTest"/* Topic */,//你的topic "TagA"/* Tag */, ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET) /* Message body */ ); //Call send message to deliver message to one of brokers. SendResultsendResult= producer.send(msg); System.out.printf("%s%n", sendResult); } //Shut down once the producer instance is not longer in use. producer.shutdown(); } }
// Instantiate with specified consumer group name. DefaultMQPushConsumerconsumer=newDefaultMQPushConsumer("please_rename_unique_group_name"); // Specify name server addresses.This sever runs on localhost(本地) port 9876 consumer.setNamesrvAddr("localhost:9876"); // Subscribe one more more topics to consume. consumer.subscribe("TopicTest", "*");//你要消费的topics // Register callback to execute on arrival of messages fetched from brokers. consumer.registerMessageListener(newMessageListenerConcurrently() {
@Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) { System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msgs); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } });