RabbitMQ-高级-分布式事务(实战)
RabbitMQ-高级-分布式事务简介
简述
分布式事务指事务的操作位于不同的节点上,需要保证事务的 AICD 特性。 例如在下单场景下,库存和订单如果不在同一个节点上,就涉及分布式事务。
分布式事务的方式
在分布式系统中,要实现分布式事务,无外乎那几种解决方案。
两阶段提交(2PC)需要数据库产商的支持,java组件有atomikos等。
两阶段提交(Two-phase Commit,2PC),通过引入协调者(Coordinator)来协调参与者的行为,并最终决定这些参与者是否要真正执行事务。
准备阶段
协调者询问参与者事务是否执行成功,参与者发回事务执行结果。

提交阶段 如果事务在每个参与者上都执行成功,事务协调者发送通知让参与者提交事务;否则,协调者发送通知让参与者回滚事务。 需要注意的是,在准备阶段,参与者执行了事务,但是还未提交。只有在提交阶段接收到协调者发来的通知后,才进行提交或者回滚。

存在的问题
同步阻塞 所有事务参与者在等待其它参与者响应的时候都处于同步阻塞状态,无法进行其它操作。
单点问题 协调者在 2PC 中起到非常大的作用,发生故障将会造成很大影响。特别是在阶段二发生故障,所有参与者会一直等待状态,无法完成其它操作。
数据不一致 在阶段二,如果协调者只发送了部分 Commit 消息,此时网络发生异常,那么只有部分参与者接收到 Commit 消息,也就是说只有部分参与者提交了事务,使得系统数据不一致。
太过保守 任意一个节点失败就会导致整个事务失败,没有完善的容错机制。
补偿事务(TCC) 严选,阿里,蚂蚁金服。
TCC 其实就是采用的补偿机制,其核心思想是:针对每个操作,都要注册一个与其对应的确认和补偿(撤销)操作。它分为三个阶段:
- Try 阶段主要是对业务系统做检测及资源预留
- Confirm 阶段主要是对业务系统做确认提交,Try阶段执行成功并开始执行 Confirm阶段时,默认 - - - Confirm阶段是不会出错的。即:只要Try成功,Confirm一定成功。
- Cancel 阶段主要是在业务执行错误,需要回滚的状态下执行的业务取消,预留资源释放
举个例子,假入 Bob 要向 Smith 转账,思路大概是: 我们有一个本地方法,里面依次调用 1:首先在 Try 阶段,要先调用远程接口把 Smith 和 Bob 的钱给冻结起来。 2:在 Confirm 阶段,执行远程调用的转账的操作,转账成功进行解冻。 3:如果第2步执行成功,那么转账成功,如果第二步执行失败,则调用远程冻结接口对应的解冻方法 (Cancel)。
优点: 跟2PC比起来,实现以及流程相对简单了一些,但数据的一致性比2PC也要差一些 缺点: 缺点还是比较明显的,在2,3步中都有可能失败。TCC属于应用层的一种补偿方式,所以需要程序员在实现的时候多写很多补偿的代码,在些场景中,一些业务流程可能用TCC不太好定义及处理。
本地消息表(异步确保)比如:支付宝、微信支付主动查询支付状态,对账单的形式
本地消息表与业务数据表处于同一个数据库中,这样就能利用本地事务来保证在对这两个表的操作满足事务特性,并且使用了消息队列来保证最终一致性。
- 在分布式事务操作的一方完成写业务数据的操作之后向本地消息表发送一个消息,本地事务能保证这个消息一定会被写入本地消息表中。
- 之后将本地消息表中的消息转发到 Kafka 等消息队列中,如果转发成功则将消息从本地消息表中删除,否则继续重新转发。
- 在分布式事务操作的另一方从消息队列中读取一个消息,并执行消息中的操作。

优点: 一种非常经典的实现,避免了分布式事务,实现了最终一致性。 缺点: 消息表会耦合到业务系统中,如果没有封装好的解决方案,会有很多杂活需要处理。
MQ 事务消息 异步场景,通用性较强,拓展性较高
有一些第三方的MQ是支持事务消息的,比如RocketMQ,他们支持事务消息的方式也是类似于采用的二阶段提交,但是市面上一些主流的MQ都是不支持事务消息的,比如 Kafka 不支持。 以阿里的 RabbitMQ 中间件为例,其思路大致为:
-
第一阶段Prepared消息,会拿到消息的地址。 第二阶段执行本地事务,第三阶段通过第一阶段拿到的地址去访问消息,并修改状态。
-
也就是说在业务方法内要想消息队列提交两次请求,一次发送消息和一次确认消息。如果确认消息发送失败了RabbitMQ会定期扫描消息集群中的事务消息,这时候发现了Prepared消息,它会向消息发送者确认,所以生产方需要实现一个check接口,RabbitMQ会根据发送端设置的策略来决定是回滚还是继续发送确认消息。这样就保证了消息发送与本地事务同时成功或同时失败。

优点: 实现了最终一致性,不需要依赖本地数据库事务。 缺点: 实现难度大,主流MQ不支持,RocketMQ事务消息部分代码也未开源。
具体实现
分布式事务的完整架构图

美团外卖架构:

系统与系统之间的分布式事务问题

基于MQ的分布式事务整体设计思路

如果这个时候MQ服务器出现了异常和故障,那么消息是无法获取到回执信息。怎么解决呢?
基于MQ的分布式事务消息的可靠生产问题-定时重发

基于MQ的分布式事务消息的可靠消费

基于MQ的分布式事务消息的消息重发

基于MQ的分布式事务消息的消息重发

数据库准备
dispatcher
SET FOREIGN_KEY_CHECKS=0;
-- ----------------------------
-- Table structure for ksd_dispather_order
-- ----------------------------
DROP TABLE IF EXISTS `ksd_dispather_order`;
CREATE TABLE `ksd_dispather_order` (
`dispatch_id` varchar(100) DEFAULT NULL,
`order_id` varchar(100) NOT NULL DEFAULT '',
`status` varchar(255) DEFAULT NULL,
`order_content` varchar(255) DEFAULT NULL,
`create_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP,
`user_id` varchar(32) DEFAULT NULL,
PRIMARY KEY (`order_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8;
dispatcher-order
SET FOREIGN_KEY_CHECKS=0;
-- ----------------------------
-- Table structure for ksd_order
-- ----------------------------
DROP TABLE IF EXISTS `ksd_order`;
CREATE TABLE `ksd_order` (
`order_id` varchar(32) DEFAULT NULL,
`user_id` varchar(32) DEFAULT NULL,
`order_content` varchar(255) DEFAULT NULL,
`create_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8;
-- ----------------------------
-- Table structure for ksd_order_message
-- ----------------------------
DROP TABLE IF EXISTS `ksd_order_message`;
CREATE TABLE `ksd_order_message` (
`order_id` varchar(100) NOT NULL,
`status` int(1) DEFAULT NULL,
`order_content` varchar(255) DEFAULT NULL,
`unique_id` varchar(100) DEFAULT NULL
) ENGINE=InnoDB DEFAULT CHARSET=utf8;
pom依赖
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-devtools</artifactId>
<scope>runtime</scope>
<optional>true</optional>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>8.0.28</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.dataformat</groupId>
<artifactId>jackson-dataformat-avro</artifactId>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>3.6</version>
</dependency>
</dependencies>
RabbitMQ实战(服务端)

application.yml
server:
port: 9000
spring:
datasource:
url: jdbc:mysql://localhost:3306/dispatcher?characterEncoding=utf8&useSSL=false&serverTimezone=UTC&allowPublicKeyRetrieval=true
username: root
password: root
driver-class-name: com.mysql.cj.jdbc.Driver
rabbitmq:
#port: 5672
#host: 47.104.141.27
username: admin
password: admin
virtual-host: /
listener:
simple:
acknowledge-mode: manual # 这里是开启手动ack,让程序去控制MQ的消息的重发和删除和转移
retry:
enabled: true # 开启重试
max-attempts: 3 #最大重试次数
initial-interval: 2000ms #重试间隔时间
addresses: 120.25.213.140:5672
logging:
level:
root: debug
mq包下的DeadMqConsumer
package com.dispatcherservice.mq;
import com.dispatcherservice.pojo.Order;
import com.dispatcherservice.service.DispatchService;
import com.dispatcherservice.util.JsonUtil;
import com.rabbitmq.client.Channel;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Service;
/**
* @description: DeadMqConsumer
*/
@Service
public class DeadMqConsumer {
@Autowired
private DispatchService dispatchService;
// 解决消息重试的集中方案:
// 1: 控制重发的次数 + 死信队列
// 2: try+catch+手动ack
// 3: try+catch+手动ack + 死信队列处理 + 人工干预
@RabbitListener(queues = {"dead.order.queue"})
public void messageconsumer(String ordermsg, Channel channel,
CorrelationData correlationData,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) throws Exception {
try {
// 1:获取消息队列的消息
System.out.println("收到MQ的消息是: " + ordermsg );
// 2: 获取订单服务的信息
Order order = JsonUtil.string2Obj(ordermsg, Order.class);
// 3: 获取订单id
String orderId = order.getOrderId();
// 幂等性问题
//int count = countOrderById(orderId);
// 4:保存运单
//if(count==0)dispatchService.dispatch(orderId);
//if(count>0)dispatchService.updateDispatch(orderId);
dispatchService.dispatch(orderId);
// 3:手动ack告诉mq消息已经正常消费
channel.basicAck(tag, false);
} catch (Exception ex) {
System.out.println("人工干预");
System.out.println("发短信预警");
System.out.println("同时把消息转移别的存储DB");
channel.basicNack(tag, false,false);
}
}
}
mq包下的OrderMqConsumer
package com.dispatcherservice.mq;
import com.dispatcherservice.pojo.Order;
import com.dispatcherservice.service.DispatchService;
import com.dispatcherservice.util.JsonUtil;
import com.rabbitmq.client.Channel;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Service;
/**
* @description: OrderMqConsumer
*/
@Service
public class OrderMqConsumer {
@Autowired
private DispatchService dispatchService;
private int count = 1;
// 解决消息重试的集中方案:
// 1: 控制重发的次数 + 死信队列
// 2: try+catch+手动ack
// 3: try+catch+手动ack + 死信队列处理 + 人工干预
@RabbitListener(queues = {"order.queue"})
public void messageconsumer(String ordermsg, Channel channel,
CorrelationData correlationData,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) throws Exception {
try {
// 1:获取消息队列的消息
System.out.println("收到MQ的消息是: " + ordermsg + ",count = " + count++);
// 2: 获取订单服务的信息
Order order = JsonUtil.string2Obj(ordermsg, Order.class);
// 3: 获取订单id
String orderId = order.getOrderId();
// 4:保存运单
dispatchService.dispatch(orderId);
// 3:手动ack告诉mq消息已经正常消费
System.out.println(1 / 0); //出现异常
channel.basicAck(tag, false);
} catch (Exception ex) {
//如果出现异常的情况下,根据实际的情况去进行重发
//重发一次后,丢失,还是日记,存库根据自己的业务场景去决定
//参数1:消息的tag 参数2:false 多条处理 参数3:requeue 重发
// false 不会重发,会把消息打入到死信队列
// true 的会会死循环的重发,建议如果使用true的话,不加try/catch否则就会造成死循环
channel.basicNack(tag, false, false);// 死信队列
}
}
}
实体类Order
package com.dispatcherservice.pojo;
import java.util.Date;
/**
* springboot+jdbctemplate/mybatis
*/
public class Order implements java.io.Serializable {
public String orderId;
public Integer userId;
public String orderContent;
public Date createTime;
public String getOrderId() {
return orderId;
}
public void setOrderId(String orderId) {
this.orderId = orderId;
}
public Integer getUserId() {
return userId;
}
public void setUserId(Integer userId) {
this.userId = userId;
}
public String getOrderContent() {
return orderContent;
}
public void setOrderContent(String orderContent) {
this.orderContent = orderContent;
}
public Date getCreateTime() {
return createTime;
}
public void setCreateTime(Date createTime) {
this.createTime = createTime;
}
@Override
public String toString() {
return "Order [orderId=" + orderId + ", userId=" + userId + ", orderContent=" + orderContent + ", createTime="
+ createTime + "]";
}
}
service下的DispatchService
package com.dispatcherservice.service;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.util.UUID;
@Service
@Transactional(rollbackFor = Exception.class)
public class DispatchService {
@Autowired
private JdbcTemplate jdbcTemplate;
/**
* @Author xuke
* @Description 运单的接收
* @Date 15:23 2021/3/7
* @Param [orderId]
* @return void
**/
public void dispatch(String orderId) throws Exception {
// 定义保存sql
String sqlString = "insert into ksd_dispather_order(order_id,dispatch_id,status,order_content,user_id)values(?,?,?,?,?)";
// 添加运动记录
int count = jdbcTemplate.update(sqlString, orderId, UUID.randomUUID().toString(), 0, "木子鱼买了一个泡面", "1");
if (count != 1) {
throw new Exception("订单创建失败,原因[数据库操作失败]");
}
}
}
util下的JsonUtil
package com.dispatcherservice.util;
import org.codehaus.jackson.map.DeserializationConfig;
import org.codehaus.jackson.map.ObjectMapper;
import org.codehaus.jackson.map.SerializationConfig;
import org.codehaus.jackson.map.annotate.JsonSerialize.Inclusion;
import org.codehaus.jackson.type.JavaType;
import org.codehaus.jackson.type.TypeReference;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.util.StringUtils;
import java.text.SimpleDateFormat;
public class JsonUtil {
private static ObjectMapper objectMapper = new ObjectMapper();
private static Logger log = LoggerFactory.getLogger(JsonUtil.class);
static {
// 对象的所有字段全部列入
objectMapper.setSerializationInclusion(Inclusion.ALWAYS);
// 取消默认转换timestamps形式
objectMapper.configure(SerializationConfig.Feature.WRITE_DATES_AS_TIMESTAMPS, false);
// 忽略空Bean转json的错误
objectMapper.configure(SerializationConfig.Feature.FAIL_ON_EMPTY_BEANS, false);
// 所有的日期格式都统一为以下的样式,即yyyy-MM-dd HH:mm:ss
objectMapper.setDateFormat(new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"));
// 忽略 在json字符串中存在,但是在java对象中不存在对应属性的情况。防止错误
objectMapper.configure(DeserializationConfig.Feature.FAIL_ON_UNKNOWN_PROPERTIES, false);
// 精度的转换问题
objectMapper.configure(DeserializationConfig.Feature.USE_BIG_DECIMAL_FOR_FLOATS, true);
objectMapper.configure(DeserializationConfig.Feature.ACCEPT_SINGLE_VALUE_AS_ARRAY, true);
}
public static <T> String obj2String(T obj) {
if (obj == null) {
return null;
}
try {
return obj instanceof String ? (String) obj : objectMapper.writeValueAsString(obj);
} catch (Exception e) {
log.warn("Parse Object to String error", e);
return null;
}
}
public static <T> String obj2StringPretty(T obj) {
if (obj == null) {
return null;
}
try {
return obj instanceof String ? (String) obj
: objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(obj);
} catch (Exception e) {
log.warn("Parse Object to String error", e);
return null;
}
}
public static <T> T string2Obj(String str, Class<T> clazz) {
if (StringUtils.isEmpty(str) || clazz == null) {
return null;
}
try {
return clazz.equals(String.class) ? (T) str : objectMapper.readValue(str, clazz);
} catch (Exception e) {
log.warn("Parse String to Object error", e);
return null;
}
}
public static <T> T string2Obj(String str, TypeReference<T> typeReference) {
if (StringUtils.isEmpty(str) || typeReference == null) {
return null;
}
try {
return (T) (typeReference.getType().equals(String.class) ? str
: objectMapper.readValue(str, typeReference));
} catch (Exception e) {
log.warn("Parse String to Object error", e);
return null;
}
}
public static <T> T string2Obj(String str, Class<?> collectionClass, Class<?>... elementClasses) {
JavaType javaType = objectMapper.getTypeFactory().constructParametricType(collectionClass, elementClasses);
try {
return objectMapper.readValue(str, javaType);
} catch (Exception e) {
log.warn("Parse String to Object error", e);
return null;
}
}
public static void main(String[] args) {
//String json = "{\"name\":\"Geely\",\"color\":\"blue\",\"id\":1274670369972846594}";
/*
* User user = new User(); user.setId(2);
* user.setAccount("geely@happymmall.com"); user.setCreateTime(new Date());
* String userJsonPretty = JsonUtil.obj2StringPretty(user);
* log.info("userJson:{}",userJsonPretty);
*
*
* User user2 = JsonUtil.string2Obj(userJsonPretty, User.class);
* System.out.println(user2);
*/
}
}
web下的DispatchController
package com.dispatcherservice.web;
import com.dispatcherservice.service.DispatchService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
@RequestMapping("/dispatch")
public class DispatchController {
@Autowired
public DispatchService dispathService;
// 添加订单后,添加调度信息
@GetMapping("/order")
public String lock(String orderId) throws Exception {
if(orderId.equals("1000001")) {
Thread.sleep(3000L); // 模拟业务耗时,接口调用者会认为超时
}
dispathService.dispatch(orderId); // 将外卖订单分配给小哥
return "success";
}
}
RabbitMQ实战(客户端)

application.yml
server:
port: 8081
spring:
datasource:
url: jdbc:mysql://localhost:3306/dispatcher-order?characterEncoding=utf8&useSSL=false&serverTimezone=UTC&allowPublicKeyRetrieval=true
username: root
password: root
driver-class-name: com.mysql.cj.jdbc.Driver
rabbitmq:
#port: 5672
#host: 47.104.141.27
username: admin
password: admin
virtual-host: /
addresses: 120.25.213.140:5672
publisher-confirm-type: correlated
#springboot.rabbitmq.publisher-confirm 新版本已被弃用,现在使用 spring.rabbitmq.publisher-confirm-type = correlated 实现相同效果
#NONE值是禁用发布确认模式,是默认值
#CORRELATED值是发布消息成功到交换器后会触发回调方法,如1示例
#SIMPLE值经测试有两种效果,其一效果和CORRELATED值一样会触发回调方法,其二在发布消息成功后使用rabbitTemplate调用waitForConfirms或waitForConfirmsOrDie方法等待broker节点返回发送结果,
#根据返回结果来判定下一步的逻辑,要注意的点是waitForConfirmsOrDie方法如果返回false则会关闭channel,则接下来无法发送消息到broker;
logging:
level:
root: debug
config包下的RabbitMQConfiguration
package com.orderservice.config;
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
/**
* @description: RabbitMQConfiguration
*/
@Configuration
public class RabbitMQConfiguration {
@Bean
public FanoutExchange deadExchange() {
return new FanoutExchange("dead_order_fanout_exchange", true, false);
}
@Bean
public Queue deadOrderQueue() {
return new Queue("dead.order.queue", true);
}
@Bean
public Binding bindDeadOrder() {
return BindingBuilder.bind(deadOrderQueue()).to(deadExchange());
}
@Bean
public FanoutExchange fanoutExchange() {
return new FanoutExchange("order_fanout_exchange", true, false);
}
@Bean
public Queue orderQueue() {
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", "dead_order_fanout_exchange");
return new Queue("order.queue", true, false, false, args);
}
@Bean
public Binding bindorder() {
return BindingBuilder.bind(orderQueue()).to(fanoutExchange());
}
}
controller包下的OrderController
package com.orderservice.controller;
import com.orderservice.pojo.Order;
import com.orderservice.service.MQOrderService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
/**
* @description: OrderController
*/
@RestController
public class OrderController {
@Autowired
private MQOrderService mqOrderService;
@GetMapping("/test/order")
public String testOrder() throws Exception {
//订单生成
String orderId = "1000001";
Order orderInfo = new Order();
orderInfo.setOrderId(orderId);
orderInfo.setUserId(1);
orderInfo.setOrderContent("买了一个方便面");
mqOrderService.createOrder(orderInfo);
System.out.println("订单创建成功.......");
return "success";
}
}
dao包下的OrderDataBaseService
package com.orderservice.dao;
import com.orderservice.pojo.Order;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@Service
@Transactional(rollbackFor = Exception.class)
public class OrderDataBaseService {
@Autowired
private JdbcTemplate jdbcTemplate;
/**
* 保存订单记录
*/
public void saveOrder(Order order) throws Exception{
// 定义保存sql
String sqlString = "insert into ksd_order(order_id,user_id,order_content)values(?,?,?)";
// 1:添加订单记录
int count = jdbcTemplate.update(sqlString,order.getOrderId(),order.getUserId(),order.getOrderContent());
if(count!=1) {
throw new Exception("订单创建失败,原因[数据库操作失败]");
}
//因为在下单可能会会rabbit会出现宕机,就引发消息是没有放入MQ.为来消息可靠生产,对消息做一次冗余
saveLocalMessage(order);
}
/**
* 保存信息到本地
* @param order
*/
public void saveLocalMessage(Order order) throws Exception{
// 定义保存sql
String sqlString = "insert into ksd_order_message(order_id,order_content,status,unique_id)values(?,?,?,?)";
// 添加运动记录
int count = jdbcTemplate.update(sqlString,order.getOrderId(),order.getOrderContent(),0,1);
if(count!=1) {
throw new Exception("出现异常,原因[数据库操作失败]");
}
}
}
mq包下的OrderMQService
package com.orderservice.mq;
import com.orderservice.pojo.Order;
import com.orderservice.util.JsonUtil;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import javax.annotation.PostConstruct;
/**
* @description: MQService
*/
@Service
public class OrderMQService {
@Autowired
private RabbitTemplate rabbitTemplate;
@Autowired
private JdbcTemplate jdbcTemplate;
//@PostConstruct注解好多人以为是Spring提供的。其实是Java自己的注解。
//Java中该注解的说明:@PostConstruct该注解被用来修饰一个非静态的void()方法。被@PostConstruct修饰的方法会在服务器加载Servlet的时候运行,
// 并且只会被服务器执行一次。PostConstruct在构造函数之后执行,init()方法之前执行。
@PostConstruct
public void regCallback() {
// 消息发送成功以后,给予生产者的消息回执,来确保生产者的可靠性
rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
@Override
public void confirm(CorrelationData correlationData, boolean ack, String cause) {
System.out.println("cause:"+cause);
// 如果ack为true代表消息已经收到
String orderId = correlationData.getId();
if (!ack) {
// 这里可能要进行其他的方式进行存储
System.out.println("MQ队列应答失败,orderId是:" + orderId);
return;
}
try {
String updatesql = "update ksd_order_message set status = 1 where order_id = ?";
int count = jdbcTemplate.update(updatesql, orderId);
if (count == 1) {
System.out.println("本地消息状态修改成功,消息成功投递到消息队列中...");
}
} catch (Exception ex) {
System.out.println("本地消息状态修改失败,出现异常:" + ex.getMessage());
}
}
});
}
public void sendMessage(Order order) {
// 通过MQ发送消息
rabbitTemplate.convertAndSend("order_fanout_exchange", "", JsonUtil.obj2String(order),
new CorrelationData(order.getOrderId()));
}
}
实体类Order
package com.orderservice.pojo;
import java.util.Date;
/**
* springboot+jdbctemplate/mybatis
*/
public class Order implements java.io.Serializable {
public String orderId;
public Integer userId;
public String orderContent;
public Date createTime;
public String getOrderId() {
return orderId;
}
public void setOrderId(String orderId) {
this.orderId = orderId;
}
public Integer getUserId() {
return userId;
}
public void setUserId(Integer userId) {
this.userId = userId;
}
public String getOrderContent() {
return orderContent;
}
public void setOrderContent(String orderContent) {
this.orderContent = orderContent;
}
public Date getCreateTime() {
return createTime;
}
public void setCreateTime(Date createTime) {
this.createTime = createTime;
}
@Override
public String toString() {
return "Order [orderId=" + orderId + ", userId=" + userId + ", orderContent=" + orderContent + ", createTime="
+ createTime + "]";
}
}
service包下的MQOrderService
package com.orderservice.service;
import com.orderservice.dao.OrderDataBaseService;
import com.orderservice.mq.OrderMQService;
import com.orderservice.pojo.Order;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@Service
public class MQOrderService {
@Autowired
private OrderDataBaseService orderDataBaseService;
@Autowired
private OrderMQService orderMQService;
// 创建订单
public void createOrder(Order orderInfo) throws Exception {
// 1: 订单信息--插入丁订单系统,订单数据库事务
orderDataBaseService.saveOrder(orderInfo);
// 2:通過Http接口发送订单信息到运单系统
orderMQService.sendMessage(orderInfo);
}
}
service包下的OrderService
package com.orderservice.service;
import com.orderservice.dao.OrderDataBaseService;
import com.orderservice.pojo.Order;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.client.SimpleClientHttpRequestFactory;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.web.client.RestTemplate;
@Service
public class OrderService {
@Autowired
private OrderDataBaseService orderDataBaseService;
// 创建订单
@Transactional(rollbackFor = Exception.class) // 订单创建整个方法添加事务 acid
public void createOrder(Order orderInfo) throws Exception {
// 1: 订单信息--插入丁订单系统,订单数据库事务
orderDataBaseService.saveOrder(orderInfo);
// 2:通過Http接口发送订单信息到运单系统
String result = dispatchHttpApi(orderInfo.getOrderId());
if(!"success".equals(result)) {
throw new Exception("订单创建失败,原因是运单接口调用失败!");
}
}
/**
* 模拟http请求接口发送,运单系统,将订单号传过去 springcloud
* @return
*/
private String dispatchHttpApi(String orderId) {
SimpleClientHttpRequestFactory factory = new SimpleClientHttpRequestFactory();
// 链接超时 > 3秒
factory.setConnectTimeout(3000);
// 处理超时 > 2秒
factory.setReadTimeout(2000);
// 发送http请求
String url = "http://localhost:9000/dispatch/order?orderId="+orderId;
RestTemplate restTemplate = new RestTemplate(factory);//异常
String result = restTemplate.getForObject(url, String.class);
return result;
}
}
stsk包下的TaskService
package com.orderservice.task;
//
//import com.xuexiangban.rabbitmq.pojo.Order;
//import org.springframework.amqp.rabbit.core.RabbitTemplate;
//import org.springframework.beans.factory.annotation.Autowired;
//import org.springframework.scheduling.annotation.EnableScheduling;
//import org.springframework.scheduling.annotation.Scheduled;
//
//import java.util.List;
//
///**
// * @description:
// * @author: xuke
// * @time: 2021/3/7 16:07
// */
//@EnableScheduling
//public class TaskService {
//
// @Autowired
// private RabbitTemplate rabbitTemplate;
//
// @Scheduled(cron = "0 0 0/2 ?")
// public void sendMessage(){
// // 把消息为0的状态消息重新查询出来,投递到MQ中。
// List<Order> orderList = orderService.selectOrderMessage(0);
// for (Order order : orderList) {
// rabbitTemplate.convertAndSend("order-fanout_exchange","",order);
// }
// }
//}
util下的JsonUtil
package com.orderservice.util;
import org.codehaus.jackson.map.DeserializationConfig;
import org.codehaus.jackson.map.ObjectMapper;
import org.codehaus.jackson.map.SerializationConfig;
import org.codehaus.jackson.map.annotate.JsonSerialize.Inclusion;
import org.codehaus.jackson.type.JavaType;
import org.codehaus.jackson.type.TypeReference;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.util.StringUtils;
import java.text.SimpleDateFormat;
/**
<dependency>
<groupId>com.fasterxml.jackson.dataformat</groupId>
<artifactId>jackson-dataformat-avro</artifactId>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>3.6</version>
</dependency>
*/
public class JsonUtil {
private static ObjectMapper objectMapper = new ObjectMapper();
private static Logger log = LoggerFactory.getLogger(JsonUtil.class);
static {
// 对象的所有字段全部列入
objectMapper.setSerializationInclusion(Inclusion.ALWAYS);
// 取消默认转换timestamps形式
objectMapper.configure(SerializationConfig.Feature.WRITE_DATES_AS_TIMESTAMPS, false);
// 忽略空Bean转json的错误
objectMapper.configure(SerializationConfig.Feature.FAIL_ON_EMPTY_BEANS, false);
// 所有的日期格式都统一为以下的样式,即yyyy-MM-dd HH:mm:ss
objectMapper.setDateFormat(new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"));
// 忽略 在json字符串中存在,但是在java对象中不存在对应属性的情况。防止错误
objectMapper.configure(DeserializationConfig.Feature.FAIL_ON_UNKNOWN_PROPERTIES, false);
// 精度的转换问题
objectMapper.configure(DeserializationConfig.Feature.USE_BIG_DECIMAL_FOR_FLOATS, true);
objectMapper.configure(DeserializationConfig.Feature.ACCEPT_SINGLE_VALUE_AS_ARRAY, true);
}
public static <T> String obj2String(T obj) {
if (obj == null) {
return null;
}
try {
return obj instanceof String ? (String) obj : objectMapper.writeValueAsString(obj);
} catch (Exception e) {
log.warn("Parse Object to String error", e);
return null;
}
}
public static <T> String obj2StringPretty(T obj) {
if (obj == null) {
return null;
}
try {
return obj instanceof String ? (String) obj
: objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(obj);
} catch (Exception e) {
log.warn("Parse Object to String error", e);
return null;
}
}
public static <T> T string2Obj(String str, Class<T> clazz) {
if (StringUtils.isEmpty(str) || clazz == null) {
return null;
}
try {
return clazz.equals(String.class) ? (T) str : objectMapper.readValue(str, clazz);
} catch (Exception e) {
log.warn("Parse String to Object error", e);
return null;
}
}
public static <T> T string2Obj(String str, TypeReference<T> typeReference) {
if (StringUtils.isEmpty(str) || typeReference == null) {
return null;
}
try {
return (T) (typeReference.getType().equals(String.class) ? str
: objectMapper.readValue(str, typeReference));
} catch (Exception e) {
log.warn("Parse String to Object error", e);
return null;
}
}
public static <T> T string2Obj(String str, Class<?> collectionClass, Class<?>... elementClasses) {
JavaType javaType = objectMapper.getTypeFactory().constructParametricType(collectionClass, elementClasses);
try {
return objectMapper.readValue(str, javaType);
} catch (Exception e) {
log.warn("Parse String to Object error", e);
return null;
}
}
public static void main(String[] args) {
//String json = "{\"name\":\"Geely\",\"color\":\"blue\",\"id\":1274670369972846594}";
/*
* User user = new User(); user.setId(2);
* user.setAccount("geely@happymmall.com"); user.setCreateTime(new Date());
* String userJsonPretty = JsonUtil.obj2StringPretty(user);
* log.info("userJson:{}",userJsonPretty);
*
*
* User user2 = JsonUtil.string2Obj(userJsonPretty, User.class);
* System.out.println(user2);
*/
}
}
测试类
package com.orderservice;
import com.orderservice.pojo.Order;
import com.orderservice.service.MQOrderService;
import com.orderservice.service.OrderService;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
@SpringBootTest
class OrderServiceApplicationTests {
@Autowired
public OrderService orderService;
@Autowired
public MQOrderService mqOrderService;
@Test
public void orderCreated() throws Exception {
//订单生成
String orderId = "1000001";
Order orderInfo = new Order();
orderInfo.setOrderId(orderId);
orderInfo.setUserId(1);
orderInfo.setOrderContent("买了一个方便面");
orderService.createOrder(orderInfo);
System.out.println("订单创建成功.......");
}
@Test
public void orderCreatedMQ() throws Exception {
//订单生成
String orderId = "1000001";
Order orderInfo = new Order();
orderInfo.setOrderId(orderId);
orderInfo.setUserId(1);
orderInfo.setOrderContent("买了一个方便面");
mqOrderService.createOrder(orderInfo);
System.out.println("订单创建成功.......");
//Thread.sleep(2000);
}
}