Files
agent-tcp/src/main/java/com/tongran/agenttcp/rocketmq/AgentProducer.java
T
2025-08-20 16:50:00 +08:00

56 lines
1.7 KiB
Java

package com.tongran.agenttcp.rocketmq;
import com.tongran.agenttcp.rocketmq.core.RocketMqService;
import com.tongran.agenttcp.rocketmq.core.model.MqMsg;
import com.tongran.agenttcp.utils.AssertLog;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Component
public class AgentProducer implements RocketMqService {
@Resource
private RocketMQTemplate rocketMQTemplate;
@Override
public void send(MqMsg msg) {
AssertLog.info("send发送消息:==>{}", msg);
rocketMQTemplate.send(msg.getTopic(), MessageBuilder.withPayload(msg.getContent()).build());
}
@Override
public void asyncSend(MqMsg msg) {
AssertLog.info("asyncSend发送消息:==>{}", msg);
rocketMQTemplate.asyncSend(msg.getTopic(), msg.getContent(), new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
// AssertLog.info("asyncSend发送消息成功:==>{}", sendResult.getSendStatus());
}
@Override
public void onException(Throwable throwable) {
AssertLog.error("asyncSend发送消息失败:==>{}", throwable.getMessage());
}
});
}
@Override
public void syncSendOrderly(MqMsg msg) {
AssertLog.info("syncSendOrderly发送消息:==>{}", msg);
rocketMQTemplate.sendOneWay(msg.getTopic(), msg.getContent());
}
@Override
public void delayedSendOrderly(MqMsg msg) {
AssertLog.info("delayedSendOrderly发送消息:==>{}", msg);
rocketMQTemplate.syncSend(msg.getTopic(), MessageBuilder.withPayload(msg.getContent()).build(), 30000, msg.getLevel());
}
}