public void deleteMessageByMessageId(String messageId) throws MessageBizException;

/**
* 重发某个消息队列中的全部已死亡的消息.
*/
public void reSendAllDeadMessageByQueueName(String queueName, int batchSize) throws MessageBizException;

/**
* 获取分页数据
*/
PageBean listPage(PageParam pageParam, Map<String, Object> paramMap) throws MessageBizException;

}
@Service(“rpTransactionMessageService”)
public class RpTransactionMessageServiceImpl implements RpTransactionMessageService {

private static final Log log = LogFactory.getLog(RpTransactionMessageServiceImpl.class);

@Autowired
private RpTransactionMessageDao rpTransactionMessageDao;

@Autowired
private JmsTemplate notifyJmsTemplate;

public int saveMessageWaitingConfirm(RpTransactionMessage message) {
if (message == null) {
throw new MessageBizException(MessageBizException.SAVA_MESSAGE_IS_NULL, “保存的消息为空”);
}
if (StringUtil.isEmpty(message.getConsumerQueue())) {
throw new MessageBizException(MessageBizException.MESSAGE_CONSUMER_QUEUE_IS_NULL, "消息的消费队列不能为空 ");
}
message.setEditTime(new Date());
message.setStatus(MessageStatusEnum.WAITING_CONFIRM.name());
message.setAreadlyDead(PublicEnum.NO.name());
message.setMessageSendTimes(0);
return rpTransactionMessageDao.insert(message);
}

public void confirmAndSendMessage(String messageId) {
final RpTransactionMessage message = getMessageByMessageId(messageId);
if (message == null) {
throw new MessageBizException(MessageBizException.SAVA_MESSAGE_IS_NULL, “根据消息id查找的消息为空”);
}
message.setStatus(MessageStatusEnum.SENDING.name());
message.setEditTime(new Date());
rpTransactionMessageDao.update(message);
notifyJmsTemplate.setDefaultDestinationName(message.getConsumerQueue());
notifyJmsTemplate.send(new MessageCreator() {
public Message createMessage(Session session) throws JMSException {
return session.createTextMessage(message.getMessageBody());
}
});
}

public int saveAndSendMessage(final RpTransactionMessage message) {
if (message == null) {
throw new MessageBizException(MessageBizException.SAVA_MESSAGE_IS_NULL, “保存的消息为空”);
}
if (StringUtil.isEmpty(message.getConsumerQueue())) {
throw new MessageBizException(MessageBizException.MESSAGE_CONSUMER_QUEUE_IS_NULL, "消息的消费队列不能为空 ");
}
message.setStatus(MessageStatusEnum.SENDING.name());
message.setAreadlyDead(PublicEnum.NO.name());
message.setMessageSendTimes(0);
message.setEditTime(new Date());
int result = rpTransactionMessageDao.insert(message);
notifyJmsTemplate.setDefaultDestinationName(message.getConsumerQueue());
notifyJmsTemplate.send(new MessageCreator() {
public Message createMessage(Session session) throws JMSException {
return session.createTextMessage(message.getMessageBody());
}
});
return result;
}

public void directSendMessage(final RpTransactionMessage message) {
if (message == null) {
throw new MessageBizException(MessageBizException.SAVA_MESSAGE_IS_NULL, “保存的消息为空”);
}
if (StringUtil.isEmpty(message.getConsumerQueue())) {
throw new MessageBizException(MessageBizException.MESSAGE_CONSUMER_QUEUE_IS_NULL, "消息的消费队列不能为空 ");
}
notifyJmsTemplate.setDefaultDestinationName(message.getConsumerQueue());
notifyJmsTemplate.send(new MessageCreator() {
public Message createMessage(Session session) throws JMSException {
return session.createTextMessage(message.getMessageBody());
}
});
}

public void reSendMessage(final RpTransactionMessage message) {
if (message == null) {
throw new MessageBizException(MessageBizException.SAVA_MESSAGE_IS_NULL, “保存的消息为空”);
}
if (StringUtil.isEmpty(message.getConsumerQueue())) {
throw new MessageBizException(MessageBizException.MESSAGE_CONSUMER_QUEUE_IS_NULL, "消息的消费队列不能为空 ");
}
message.addSendTimes();
message.setEditTime(new Date());
rpTransactionMessageDao.update(message);
notifyJmsTemplate.setDefaultDestinationName(message.getConsumerQueue());
notifyJmsTemplate.send(new MessageCreator() {
public Message createMessage(Session session) throws JMSException {
return session.createTextMessage(message.getMessageBody());
}
});
}

public void reSendMessageByMessageId(String messageId) {
final RpTransactionMessage message = getMessageByMessageId(messageId);
if (message == null) {
throw new MessageBizException(MessageBizException.SAVA_MESSAGE_IS_NULL, “根据消息id查找的消息为空”);
}
int maxTimes = Integer.valueOf(PublicConfigUtil.readConfig(“message.max.send.times”));
if (message.getMessageSendTimes() >= maxTimes) {
message.setAreadlyDead(PublicEnum.YES.name());
}
message.setEditTime(new Date());
message.setMessageSendTimes(message.getMessageSendTimes() + 1);
rpTransactionMessageDao.update(message);
notifyJmsTemplate.setDefaultDestinationName(message.getConsumerQueue());
notifyJmsTemplate.send(new MessageCreator() {
public Message createMessage(Session session) throws JMSException {
return session.createTextMessage(message.getMessageBody());
}
});
}

public void setMessageToAreadlyDead(String messageId) {
RpTransactionMessage message = getMessageByMessageId(messageId);
if (message == null) {
throw new MessageBizException(MessageBizException.SAVA_MESSAGE_IS_NULL, “根据消息id查找的消息为空”);
}
message.setAreadlyDead(PublicEnum.YES.name());
message.setEditTime(new Date());
rpTransactionMessageDao.update(message);
}

public RpTransactionMessage getMessageByMessageId(String messageId) {
Map<String, Object> paramMap = new HashMap<String, Object>();
paramMap.put(“messageId”, messageId);
return rpTransactionMessageDao.getBy(paramMap);
}

public void deleteMessageByMessageId(String messageId) {
Map<String, Object> paramMap = new HashMap<String, Object>();
paramMap.put(“messageId”, messageId);
rpTransactionMessageDao.delete(paramMap);
}

@SuppressWarnings(“unchecked”)
public void reSendAllDeadMessageByQueueName(String queueName, int batchSize) {
log.info(“>reSendAllDeadMessageByQueueName");
int numPerPage = 1000;
if (batchSize > 0 && batchSize < 100) {
numPerPage = 100;
} else if (batchSize > 100 && batchSize < 5000) {
numPerPage = batchSize;
} else if (batchSize > 5000) {
numPerPage = 5000;
} else {
numPerPage = 1000;
}
int pageNum = 1;
Map<String, Object> paramMap = new HashMap<String, Object>();
paramMap.put(“consumerQueue”, queueName);
paramMap.put(“areadlyDead”, PublicEnum.YES.name());
paramMap.put(“listPageSortType”, “ASC”);
Map<String, RpTransactionMessage> messageMap = new HashMap<String, RpTransactionMessage>();
List recordList = new ArrayList();
int pageCount = 1;
PageBean pageBean = rpTransactionMessageDao.listPage(new PageParam(pageNum, numPerPage), paramMap);
recordList = pageBean.getRecordList();
if (recordList == null || recordList.isEmpty()) {
log.info("
>recordList is empty”);
return;
}
pageCount = pageBean.getTotalPage();
for (final Object obj : recordList) {
final RpTransactionMessage message = (RpTransactionMessage) obj;
messageMap.put(message.getMessageId(), message);
}
for (pageNum = 2; pageNum <= pageCount; pageNum++) {
pageBean = rpTransactionMessageDao.listPage(new PageParam(pageNum, numPerPage), paramMap);
recordList = pageBean.getRecordList();
if (recordList == null || recordList.isEmpty()) {
break;
}
for (final Object obj : recordList) {
final RpTransactionMessage message = (RpTransactionMessage) obj;
messageMap.put(message.getMessageId(), message);
}
}
recordList = null;
pageBean = null;
for (Map.Entry<String, RpTransactionMessage> entry : messageMap.entrySet()) {
final RpTransactionMessage message = entry.getValue();
message.setEditTime(new Date());
message.setMessageSendTimes(message.getMessageSendTimes() + 1);
rpTransactionMessageDao.update(message);
notifyJmsTemplate.setDefaultDestinationName(message.getConsumerQueue());
notifyJmsTemplate.send(new MessageCreator() {
public Message createMessage(Session session) throws JMSException {
return session.createTextMessage(message.getMessageBody());
}
});
}
}

@SuppressWarnings(“unchecked”)
public PageBean listPage(PageParam pageParam, Map<String, Object> paramMap) {
return rpTransactionMessageDao.listPage(pageParam, paramMap);
}

}
@Component(“messageBiz”)
public class MessageBiz {

private static final Log log = LogFactory.getLog(MessageBiz.class);

@Autowired
private RpTradePaymentQueryService rpTradePaymentQueryService;

@Autowired
private RpTransactionMessageService rpTransactionMessageService;

/**
* 处理[waiting_confirm]状态的消息
* @param messages
*/
public void handleWaitingConfirmTimeOutMessages(Map<String, RpTransactionMessage> messageMap) {
log.debug(“开始处理[waiting_confirm]状态的消息,总条数[” + messageMap.size() + “]”);
// 单条消息处理(目前该状态的消息,消费队列全部是accounting,如果后期有业务扩充,需做队列判断,做对应的业务处理。)
for (Map.Entry<String, RpTransactionMessage> entry : messageMap.entrySet()) {
RpTransactionMessage message = entry.getValue();
try {
log.debug(“开始处理[waiting_confirm]消息ID为[” + message.getMessageId() + “]的消息”);
String bankOrderNo = message.getField1();
RpTradePaymentRecord record = rpTradePaymentQueryService.getRecordByBankOrderNo(bankOrderNo);
// 如果订单成功,把消息改为待处理,并发送消息
if (TradeStatusEnum.SUCCESS.name().equals(record.getStatus())) {
// 确认并发送消息
rpTransactionMessageService.confirmAndSendMessage(message.getMessageId());
} else if (TradeStatusEnum.WAITING_PAYMENT.name().equals(record.getStatus())) {
// 订单状态是等到支付,可以直接删除数据
log.debug(“订单没有支付成功,删除[waiting_confirm]消息id[” + message.getMessageId() + “]的消息”);
rpTransactionMessageService.deleteMessageByMessageId(message.getMessageId());
}
log.debug(“结束处理[waiting_confirm]消息ID为[” + message.getMessageId() + “]的消息”);
} catch (Exception e) {
log.error(“处理[waiting_confirm]消息ID为[” + message.getMessageId() + “]的消息异常:”, e);
}
}
}

/**
* 处理[SENDING]状态的消息
* @param messages
*/
public void handleSendingTimeOutMessage(Map<String, RpTransactionMessage> messageMap) {
SimpleDateFormat sdf = new SimpleDateFormat(“yyyy-MM-dd HH:mm:ss”);
log.debug(“开始处理[SENDING]状态的消息,总条数[” + messageMap.size() + “]”);
// 根据配置获取通知间隔时间
Map<Integer, Integer> notifyParam = getSendTime();
// 单条消息处理
for (Map.Entry<String, RpTransactionMessage> entry : messageMap.entrySet()) {
RpTransactionMessage message = entry.getValue();
try {
log.debug(“开始处理[SENDING]消息ID为[” + message.getMessageId() + “]的消息”);
// 判断发送次数
int maxTimes = Integer.valueOf(PublicConfigUtil.readConfig(“message.max.send.times”));
log.debug(“[SENDING]消息ID为[” + message.getMessageId() + “]的消息,已经重新发送的次数[”
+ message.getMessageSendTimes() + “]”);
// 如果超过最大发送次数直接退出
if (maxTimes < message.getMessageSendTimes()) {
// 标记为死亡
rpTransactionMessageService.setMessageToAreadlyDead(message.getMessageId());
continue;
}
// 判断是否达到发送消息的时间间隔条件
int reSendTimes = message.getMessageSendTimes();
int times = notifyParam.get(reSendTimes == 0 ? 1 : reSendTimes);
long currentTimeInMillis = Calendar.getInstance().getTimeInMillis();
long needTime = currentTimeInMillis - times * 60 * 1000;
long hasTime = message.getEditTime().getTime();
// 判断是否达到了可以再次发送的时间条件
if (hasTime > needTime) {
log.debug(“currentTime[” + sdf.format(new Date()) + “],[SENDING]消息上次发送时间[”
+ sdf.format(message.getEditTime()) + “],必须过了[” + times + “]分钟才可以再发送。”);
continue;
}
// 重新发送消息
rpTransactionMessageService.reSendMessage(message);
log.debug(“结束处理[SENDING]消息ID为[” + message.getMessageId() + “]的消息”);
} catch (Exception e) {
log.error(“处理[SENDING]消息ID为[” + message.getMessageId() + “]的消息异常:”, e);
}
}
}

/**
* 根据配置获取通知间隔时间
* @return
*/
private Map<Integer, Integer> getSendTime() {
Map<Integer, Integer> notifyParam = new HashMap<Integer, Integer>();
notifyParam.put(1, Integer.valueOf(PublicConfigUtil.readConfig(“message.send.1.time”)));
notifyParam.put(2, Integer.valueOf(PublicConfigUtil.readConfig(“message.send.2.time”)));
notifyParam.put(3, Integer.valueOf(PublicConfigUtil.readConfig(“message.send.3.time”)));
notifyParam.put(4, Integer.valueOf(PublicConfigUtil.readConfig(“message.send.4.time”)));
notifyParam.put(5, Integer.valueOf(PublicConfigUtil.readConfig(“message.send.5.time”)));
return notifyParam;
}

}
public class AccountingMessageListener implements SessionAwareMessageListener {

private static final Log LOG = LogFactory.getLog(AccountingMessageListener.class);

/**
* 会计队列模板(由Spring创建并注入进来)
*/
@Autowired
private JmsTemplate notifyJmsTemplate;

@Autowired
private RpAccountingVoucherService rpAccountingVoucherService;

@Autowired
private RpTransactionMessageService rpTransactionMessageService;

public synchronized void onMessage(Message message, Session session) {
RpAccountingVoucher param = null;
String strMessage = null;
try {
ActiveMQTextMessage objectMessage = (ActiveMQTextMessage) message;
strMessage = objectMessage.getText();
LOG.info(“strMessage1 accounting:” + strMessage);
param = JSONObject.parseObject(strMessage, RpAccountingVoucher.class);
// 这里转换成相应的对象还有问题
if (param == null) {
LOG.info(“param参数为空”);
return;
}
int entryType = param.getEntryType();
double payerChangeAmount = param.getPayerChangeAmount();
String voucherNo = param.getVoucherNo();
String payerAccountNo = param.getPayerAccountNo();
int fromSystem = param.getFromSystem();
int payerAccountType = 0;
if (param.getPayerAccountType() != null && !param.getPayerAccountType().equals(“”)) {
payerAccountType = param.getPayerAccountType();
}
double payerFee = param.getPayerFee();
String requestNo = param.getRequestNo();
double bankChangeAmount = param.getBankChangeAmount();
double receiverChangeAmount = param.getReceiverChangeAmount();
String receiverAccountNo = param.getReceiverAccountNo();
String bankAccount = param.getBankAccount();
String bankChannelCode = param.getBankChannelCode();
double profit = param.getProfit();
double income = param.getIncome();
double cost = param.getCost();
String bankOrderNo = param.getBankOrderNo();
int receiverAccountType = 0;
double payAmount = param.getPayAmount();
if (param.getReceiverAccountType() != null && !param.getReceiverAccountType().equals(“”)) {
receiverAccountType = param.getReceiverAccountType();
}
double receiverFee = param.getReceiverFee();
String remark = param.getRemark();
rpAccountingVoucherService.createAccountingVoucher(entryType, voucherNo, payerAccountNo, receiverAccountNo,
payerChangeAmount, receiverChangeAmount, income, cost, profit, bankChangeAmount, requestNo,
bankChannelCode, bankAccount, fromSystem, remark, bankOrderNo, payerAccountType, payAmount,
receiverAccountType, payerFee, receiverFee);
//删除消息
rpTransactionMessageService.deleteMessageByMessageId(param.getMessageId());
} catch (BizException e) {
// 业务异常,不再写会队列
LOG.error(“>BizException", e);
} catch (Exception e) {
// 不明异常不再写会队列
LOG.error("
>Exception”, e);
}
}

public JmsTemplate getNotifyJmsTemplate() {
return notifyJmsTemplate;
}

public void setNotifyJmsTemplate(JmsTemplate notifyJmsTemplate) {
this.notifyJmsTemplate = notifyJmsTemplate;
}

public RpAccountingVoucherService getRpAccountingVoucherService() {
return rpAccountingVoucherService;
}

public void setRpAccountingVoucherService(RpAccountingVoucherService rpAccountingVoucherService) {
this.rpAccountingVoucherService = rpAccountingVoucherService;
}

}

与常规MQ的ACK机制对比

常规MQ确认机制:

  • Producer生成消息并发送给MQ(同步、异步);
  • MQ接收消息并将消息数据持久化到消息存储(持久化操作为可选配置);
  • MQ向Producer返回消息的接收结果(返回值、异常);
  • Consumer监听并消费MQ中的消息;
  • Consumer获取到消息后执行业务处理;
  • Consumer对已成功消费的消息向MQ进行ACK确认(确认后的消息将从MQ中删除);

常规MQ队列消息的处理流程无法实现 消息发送一致性, 因此直接使用现成的MQ中间件产品无法实现可靠消息最终一致性的分布式事务解决方案

消息发送一致性:是指产生消息的业务动作与消息发送的一致。也就是说,如果业务操作成功,那么由这个业务操作所产生的消息一定要成功投递出去(一般是发送到kafka、rocketmq、rabbitmq等消息中间件中),否则就丢消息。

下面用伪代码进行演示消息发送和投递的不可靠性:

先进行数据库操作,再发送消息:

public void test1(){
//1 数据库操作
//2 发送MQ消息
}

这种情况下无法保证数据库操作与发送消息的一致性,因为可能数据库操作成功,发送消息失败。

先发送消息,再操作数据库:

public void test1(){
//1 发送MQ消息
//2 数据库操作
}

这种情况下无法保证数据库操作与发送消息的一致性,因为可能发送消息成功,数据库操作失败。

在数据库事务中,先发送消息,后操作数据库:

@Transactional
public void test1(){

最后

一次偶然,从朋友那里得到一份“java高分面试指南”,里面涵盖了25个分类的面试题以及详细的解析:JavaOOP、Java集合/泛型、Java中的IO与NIO、Java反射、Java序列化、Java注解、多线程&并发、JVM、Mysql、Redis、Memcached、MongoDB、Spring、Spring Boot、Spring Cloud、RabbitMQ、Dubbo 、MyBatis 、ZooKeeper 、数据结构、算法、Elasticsearch 、Kafka 、微服务、Linux。

这不,马上就要到招聘季了,很多朋友又开始准备“金三银四”的春招啦,那我想这份“java高分面试指南”应该起到不小的作用,所以今天想给大家分享一下。

image

请注意:关于这份“java高分面试指南”,每一个方向专题(25个)的题目这里几乎都会列举,在不看答案的情况下,大家可以自行测试一下水平 且由于篇幅原因,这边无法展示所有完整的答案解析

那里得到一份“java高分面试指南”,里面涵盖了25个分类的面试题以及详细的解析:JavaOOP、Java集合/泛型、Java中的IO与NIO、Java反射、Java序列化、Java注解、多线程&并发、JVM、Mysql、Redis、Memcached、MongoDB、Spring、Spring Boot、Spring Cloud、RabbitMQ、Dubbo 、MyBatis 、ZooKeeper 、数据结构、算法、Elasticsearch 、Kafka 、微服务、Linux。

这不,马上就要到招聘季了,很多朋友又开始准备“金三银四”的春招啦,那我想这份“java高分面试指南”应该起到不小的作用,所以今天想给大家分享一下。

[外链图片转存中…(img-og4gzee6-1714469027032)]

请注意:关于这份“java高分面试指南”,每一个方向专题(25个)的题目这里几乎都会列举,在不看答案的情况下,大家可以自行测试一下水平 且由于篇幅原因,这边无法展示所有完整的答案解析

本文已被CODING开源项目:【一线大厂Java面试题解析+核心总结学习笔记+最新讲解视频+实战项目源码】收录

Logo

魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。

更多推荐