package com.rocketmq.controller.produce;
import com.rocketmq.util.RocketmqSampleUtils;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.springframework.stereotype.Controller;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.ResponseBody;
import java.nio.charset.StandardCharsets;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
* 通知系统生产者
*/
@Controller
@RequestMapping("/notice")
public class NoticeSystemProducer {
private static final Logger logger = LoggerFactory.getLogger(NoticeSystemProducer.class);
@RequestMapping("/produce")
@ResponseBody
public void produceNotice() {
DefaultMQProducer producer = new DefaultMQProducer("NoticeSystemProducer");
producer.setNamesrvAddr("100.85.217.111:8200;100.93.3.194:8200");
try {
producer.start();
String[] tags = new String[]{"TagA", "TagB", "TagC", "TagD", "TagE"};
for (int i = 0; i < 10; i++) {
String orderId = "order" + (i % 10);
Message msg = new Message("NoticeTopic", tags[i % tags.length], "KEY" + i, ("This is notice " + i).getBytes(StandardCharsets.UTF_8));
SendResult sendResult = RocketmqSampleUtils.sendMessage(producer, msg, orderId);
logger.info(sendResult.toString());
}
} catch (Exception e) {
logger.info(e.getMessage());
}
producer.shutdown();
}
}