MySQL数据不一致怎么办:电商超卖、银行转账错误的排查与完整维护方案
先说一个扎心的真实故事
2021年某大型电商平台的双11期间,出现了一起震惊业界的超卖事故。一款限量999台的手机,最终卖出了1247台,直接损失超过千万。事后排查发现,根本原因竟然是Redis缓存和MySQL数据库之间的数据延迟导致的并发竞争。
无独有偶,2022年某银行系统出现转账错误,A账户扣款后B账户没有到账,金额达数十万。最终定位到是分布式事务在两阶段提交(2PC)过程中,分支事务超时导致的数据不一致。
这些问题说到底是同一个问题:数据不一致。
今天我们就把这个问题掰开揉碎了讲清楚,从现象到根因,从排查到解决,从主从复制到分布式架构,给你一套完整的维护方案。
第一部分:数据不一致到底长什么样?
1.1 电商超卖:库存数据打架
想象一下这个场景:
库存:100件
用户A:下单,库存减1 → 99件
用户B:同时下单,库存减1 → 99件(因为A还没提交)
用户C:同时下单,库存减1 → 99件
三个用户都下单成功,但实际库存只剩97件。这就是典型的超卖。
根本原因:读-改-写(Read-Modify-Write)操作的原子性被破坏。
1.2 银行转账:扣了A没加B
转账操作的理想状态:
BEGIN;
UPDATE accounts SET balance = balance - 1000 WHERE user_id = 'A'; -- A扣1000
UPDATE accounts SET balance = balance + 1000 WHERE user_id = 'B'; -- B加1000
COMMIT;
实际可能出现问题:
- A扣了1000,系统崩溃了,B没加上
- 事务提交了,但Binlog还没同步到从库,从库数据”落后”
- 分布式场景下,A的分支事务提交了,B的分支事务超时回滚了
1.3 数据不一致的常见形态
| 类型 | 描述 | 典型场景 |
|---|---|---|
| 缓存不一致 | Redis缓存与MySQL数据不同步 | 电商库存、商品详情 |
| 主从不一致 | 主库已提交,从库未同步 | 读写分离查询 |
| 分布式不一致 | 多个节点数据状态不同 | 微服务架构、分库分表 |
| 最终不一致 | 数据最终会一致,但中间有延迟 | MQ消息队列消费 |
第二部分:如何排查数据不一致?
2.1 诊断思路:像侦探一样思考
排查数据不一致,我们遵循一个核心原则:先定位问题范围,再追踪问题根源。
第一步:确认问题存在
不要相信感觉,要用数据说话。
-- 检查MySQL主从同步状态
SHOW SLAVE STATUS\G
-- 关键指标:
-- Slave_IO_Running: Yes
-- Slave_SQL_Running: Yes
-- Seconds_Behind_Master: 0(理想状态,越大越危险)
-- Last_Error: (如果有错误信息,说明同步有问题)
第二步:检查事务状态
-- 查看当前活跃事务
SELECT * FROM information_schema.innodb_trx;
-- 关键指标解读:
-- trx_state: 事务状态(RUNNING, LOCK WAIT等)
-- trx_started: 事务开始时间
-- trx_mysql_thread_id: 对应的MySQL线程ID
-- trx_weight: 事务权重
如果某个事务长时间不提交,可能是导致锁等待和数据不一致的元凶。
第三步:检查Binlog
-- 查看当前Binlog状态
SHOW MASTER STATUS;
-- 查看最近的Binlog事件
SHOW BINLOG EVENTS IN 'mysql-bin.000001' LIMIT 10;
-- 或者使用mysqlbinlog工具解析
mysqlbinlog --base64-output=DECODE-ROWS -v mysql-bin.000001
Binlog是排查的”黑匣子”,记录了所有改变数据的操作。
2.2 使用可视化工具辅助排查
2.2.1 pt-table-checksum:主从一致性校验
这是Percona Toolkit中最强大的工具之一:
# 对指定库进行校验
pt-table-checksum \
--host=192.168.1.10 \
--user=admin \
--password=your_password \
--databases=ecommerce \
--tables=orders,inventory \
--nocheck-replication-filters \
--replicate=percona.checksums
# 查看校验结果
SELECT * FROM percona.checksums;
结果解读:
this_cnt:当前主库的行数master_cnt:从库的行数- 如果两者不一致,说明主从不一致
2.2.2 对比法快速定位
-- 方法一:直接COUNT对比
-- 主库执行
SELECT COUNT(*) FROM orders WHERE create_time > '2024-01-01';
-- 从库执行
SELECT COUNT(*) FROM orders WHERE create_time > '2024-01-01';
-- 方法二:MD5校验(适合小表)
SELECT MD5(GROUP_CONCAT(id ORDER BY id)) AS md5_hash FROM orders;
-- 方法三:使用checksum对比
SELECT CHECKSUM TABLE orders;
2.3 业务层面的数据校验
除了工具,业务层面的校验也很重要。
库存一致性校验
-- 检查库存与实际订单的关系
SELECT
i.product_id,
i.stock AS db_stock,
(
SELECT SUM(quantity)
FROM order_items
WHERE product_id = i.product_id AND status = 'PAID'
) AS sold_quantity,
i.stock - (
SELECT SUM(quantity)
FROM order_items
WHERE product_id = i.product_id AND status = 'PAID'
) AS expected_stock
FROM inventory i
WHERE i.stock != (
SELECT i2.stock - (
SELECT SUM(quantity)
FROM order_items
WHERE product_id = i2.product_id AND status = 'PAID'
)
FROM inventory i2
WHERE i2.product_id = i.product_id
);
账户余额一致性校验
-- 检查账户余额与交易流水的关系
SELECT
a.account_id,
a.balance AS db_balance,
(
SELECT SUM(amount)
FROM transactions
WHERE account_id = a.account_id
) AS calculated_balance,
a.balance - (
SELECT SUM(amount)
FROM transactions
WHERE account_id = a.account_id
) AS diff
FROM accounts a
WHERE ABS(a.balance - (
SELECT SUM(amount)
FROM transactions
WHERE account_id = a.account_id
)) > 0.01; -- 允许0.01的精度误差
第三部分:电商超卖问题的深度剖析与解决
3.1 为什么超卖这么难避免?
超卖的根因是并发控制。在分布式系统中,这个问题被进一步放大:
用户A的请求 → 网关 → 订单服务 → 库存服务(Redis)→ 数据库
用户B的请求 → 网关 → 订单服务 → 库存服务(Redis)→ 数据库
两个请求几乎同时到达,都读到了相同的库存值,都执行了减库存操作,最终导致超卖。
3.2 解决方案一:数据库乐观锁
-- 表结构
CREATE TABLE inventory (
id INT PRIMARY KEY,
product_id INT NOT NULL,
stock INT NOT NULL DEFAULT 0,
version INT NOT NULL DEFAULT 0 -- 乐观锁版本号
);
-- 扣减库存的SQL(带乐观锁)
UPDATE inventory
SET stock = stock - #{quantity}, version = version + 1
WHERE product_id = #{productId}
AND stock >= #{quantity}
AND version = #{version};
-- 检查结果
-- rows affected = 1:成功
-- rows affected = 0:库存不足或版本冲突,需要重试
Java代码示例:
@Service
public class InventoryService {
private static final int MAX_RETRY = 3;
@Transactional
public boolean deductStock(Long productId, int quantity) {
for (int i = 0; i < MAX_RETRY; i++) {
// 1. 查询当前库存和版本号
Inventory inventory = inventoryMapper.selectById(productId);
if (inventory.getStock() < quantity) {
return false; // 库存不足
}
// 2. 尝试扣减(乐观锁)
int rows = inventoryMapper.deductWithOptimisticLock(
productId, quantity, inventory.getVersion()
);
if (rows > 0) {
return true; // 扣减成功
}
// 3. 如果失败,可能是并发冲突,稍后重试
try {
Thread.sleep(50 * (i + 1)); // 指数退避
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return false;
}
}
return false;
}
}
3.3 解决方案二:数据库悲观锁
-- 使用SELECT ... FOR UPDATE进行行锁
BEGIN;
SELECT stock FROM inventory
WHERE product_id = #{productId}
FOR UPDATE; -- 加排他锁
-- 检查库存并扣减
UPDATE inventory
SET stock = stock - #{quantity}
WHERE product_id = #{productId}
AND stock >= #{quantity};
COMMIT;
注意:悲观锁会阻塞其他事务,性能较差,适合高价值商品或库存紧张的场景。
3.4 解决方案三:Redis Lua脚本原子扣减
-- Lua脚本:原子性地扣减库存
local product_id = KEYS[1]
local quantity = tonumber(ARGV[1])
-- 获取当前库存
local current_stock = redis.call('GET', 'stock:' .. product_id)
if not current_stock then
return -1 -- 商品不存在
end
current_stock = tonumber(current_stock)
if current_stock < quantity then
return 0 -- 库存不足
end
-- 扣减库存
redis.call('DECRBY', 'stock:' .. product_id, quantity)
-- 记录扣减日志(用于对账)
redis.call('LPUSH', 'stock_log:' .. product_id,
current_stock .. ':' .. quantity .. ':' .. tostring(os.time()))
return 1 -- 扣减成功
调用方式:
String script = "local current_stock = redis.call('GET', KEYS[1]) " +
"if not current_stock then return -1 end " +
"current_stock = tonumber(current_stock) " +
"if current_stock < tonumber(ARGV[1]) then return 0 end " +
"redis.call('DECRBY', KEYS[1], ARGV[1]) " +
"return 1";
Long result = (Long) redisTemplate.execute(
new DefaultRedisScript<>(script, Long.class),
Collections.singletonList("stock:" + productId),
String.valueOf(quantity)
);
3.5 解决方案四:消息队列异步扣减
@Service
public class OrderService {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
@Autowired
private InventoryService inventoryService;
public OrderResult createOrder(OrderRequest request) {
// 1. 预扣减库存(使用Redis Lua脚本)
boolean deducted = inventoryService.deductStock(
request.getProductId(), request.getQuantity()
);
if (!deducted) {
return OrderResult.fail("库存不足");
}
// 2. 创建订单(写入数据库)
Order order = new Order();
order.setProductId(request.getProductId());
order.setQuantity(request.getQuantity());
order.setStatus(OrderStatus.PENDING);
orderMapper.insert(order);
// 3. 发送消息到队列(异步)
kafkaTemplate.send("order-created",
JSON.toJSONString(order));
return OrderResult.success(order.getId());
}
}
@Service
public class OrderConsumer {
@KafkaListener(topics = "order-created")
public void handleOrder(String message) {
Order order = JSON.parseObject(message, Order.class);
try {
// 支付处理
paymentService.pay(order);
// 更新订单状态
order.setStatus(OrderStatus.PAID);
orderMapper.updateById(order);
} catch (Exception e) {
// 支付失败,回滚库存
inventoryService.refundStock(
order.getProductId(), order.getQuantity()
);
order.setStatus(OrderStatus.CANCELLED);
orderMapper.updateById(order);
}
}
}
3.6 终极方案:分布式锁 + 限流
@Service
public class SafeInventoryService {
@Autowired
private RedisTemplate<String, String> redisTemplate;
@Autowired
private InventoryMapper inventoryMapper;
private static final String LOCK_PREFIX = "inventory_lock:";
private static final long LOCK_EXPIRE = 10; // 秒
public boolean deductStockSafe(Long productId, int quantity) {
String lockKey = LOCK_PREFIX + productId;
// 1. 尝试获取分布式锁
boolean locked = redisTemplate.opsForValue()
.setIfAbsent(lockKey, "1", LOCK_EXPIRE, TimeUnit.SECONDS);
if (!locked) {
return false; // 并发太高,直接拒绝
}
try {
// 2. 在锁的保护下操作数据库
Inventory inventory = inventoryMapper.selectById(productId);
if (inventory.getStock() < quantity) {
return false;
}
inventoryMapper.deductStock(productId, quantity);
return true;
} finally {
// 3. 释放锁
redisTemplate.delete(lockKey);
}
}
}
第四部分:银行转账错误的数据一致性保障
4.1 转账场景的特殊性
银行转账对一致性要求极高,不能容忍任何数据丢失或错误。这里我们讨论的是金融级的解决方案。
4.2 核心原则:ACID + 分布式事务
本地事务方案(单库场景)
-- 转账存储过程
DELIMITER $$
CREATE PROCEDURE transfer_money(
IN from_account VARCHAR(32),
IN to_account VARCHAR(32),
IN amount DECIMAL(15,2)
)
BEGIN
DECLARE exit_handler EXCEPTION FOR SQLEXCEPTION;
-- 出错时回滚
DECLARE EXIT HANDLER FOR SQLEXCEPTION
BEGIN
ROLLBACK;
RESIGNAL;
END;
START TRANSACTION;
-- 1. 锁住转出账户(防止并发)
SELECT balance FROM accounts
WHERE account_id = from_account
FOR UPDATE;
-- 2. 检查余额
IF (SELECT balance FROM accounts WHERE account_id = from_account) < amount THEN
SIGNAL SQLSTATE '45000'
SET MESSAGE_TEXT = '余额不足';
END IF;
-- 3. 扣减转出账户
UPDATE accounts
SET balance = balance - amount,
version = version + 1,
updated_at = NOW()
WHERE account_id = from_account;
-- 4. 增加转入账户
UPDATE accounts
SET balance = balance + amount,
version = version + 1,
updated_at = NOW()
WHERE account_id = to_account;
-- 5. 记录交易流水
INSERT INTO transactions (
id, from_account, to_account, amount, status, created_at
) VALUES (
UUID(), from_account, to_account, amount, 'SUCCESS', NOW()
);
COMMIT;
END$$
DELIMITER ;
分布式事务方案(多库场景)
方案A:Seata AT模式
# seata配置
seata:
enabled: true
tx-service-group: my_test_group
service:
vgroup-mapping:
my_test_group: default
grouplist:
default: 127.0.0.1:8091
registry:
type: nacos
nacos:
server-addr: 127.0.0.1:8848
namespace: ""
group: SEATA_GROUP
config:
type: nacos
nacos:
server-addr: 127.0.0.1:8848
namespace: ""
group: SEATA_GROUP
// 转账服务
@Service
public class TransferService {
@Autowired
private AccountMapper accountMapper;
@GlobalTransactional // Seata全局事务注解
public void transfer(String from, String to, BigDecimal amount) {
// 1. 扣减转出账户
accountMapper.deduct(from, amount);
// 2. 增加转入账户
accountMapper.add(to, amount);
// 3. 记录流水
transactionMapper.insert(new Transaction(from, to, amount));
}
}
方案B:TCC模式(更可靠)
// Try阶段:预留资源
public class AccountTccServiceImpl implements AccountTccService {
@Override
@TransactionMode(TransactionMode.TCC)
public boolean tryTransfer(String from, String to, BigDecimal amount, String xid) {
// Try:冻结转出账户的余额
int rows = accountMapper.freeze(from, amount, xid);
if (rows == 0) {
return false; // 余额不足
}
// 预创建转入记录
frozenAccountMapper.insert(new FrozenAccount(to, amount, xid));
return true;
}
// Confirm阶段:确认提交
@Override
public boolean confirmTransfer(String from, String to, BigDecimal amount, String xid) {
// 扣除冻结余额
accountMapper.deductFrozen(from, amount, xid);
// 增加实际余额
accountMapper.add(to, amount);
return true;
}
// Cancel阶段:回滚
@Override
public boolean cancelTransfer(String from, String to, BigDecimal amount, String xid) {
// 解冻转出账户
accountMapper.unfreeze(from, amount, xid);
// 删除预创建的转入记录
frozenAccountMapper.deleteByXid(xid);
return true;
}
}
4.3 对账系统:最后的防线
无论用什么方案,对账都是必须的。
-- 对账表结构
CREATE TABLE reconciliation (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
batch_no VARCHAR(64) NOT NULL COMMENT '对账批次号',
source_system VARCHAR(32) NOT NULL COMMENT '源系统',
business_no VARCHAR(64) NOT NULL COMMENT '业务流水号',
amount DECIMAL(15,2) NOT NULL COMMENT '金额',
status VARCHAR(32) NOT NULL COMMENT '状态',
created_at DATETIME NOT NULL,
matched_at DATETIME,
diff_reason VARCHAR(256) COMMENT '差异原因',
INDEX idx_batch_no (batch_no),
INDEX idx_business_no (business_no)
);
-- 生成对账差异报告
SELECT
a.business_no,
a.amount AS system_a_amount,
b.amount AS system_b_amount,
a.amount - b.amount AS diff_amount,
CASE
WHEN a.amount IS NULL THEN '仅系统B有记录'
WHEN b.amount IS NULL THEN '仅系统A有记录'
WHEN ABS(a.amount - b.amount) > 0.01 THEN '金额不一致'
ELSE '一致'
END AS diff_type
FROM (
SELECT business_no, amount FROM system_a_transactions
) a
FULL OUTER JOIN (
SELECT business_no, amount FROM system_b_transactions
) b ON a.business_no = b.business_no
WHERE a.amount IS NULL OR b.amount IS NULL
OR ABS(a.amount - b.amount) > 0.01;
第五部分:主从复制场景的数据一致性维护
5.1 主从不一致的常见原因
1. 网络中断导致Binlog传输延迟
2. 从库执行SQL出错(如数据类型不匹配)
3. 主库手动执行了未记录到Binlog的操作
4. 从库手动修改了数据
5. 大事务导致从库追不上主库
6. 锁等待导致从库SQL线程阻塞
5.2 实时监控系统
监控主从延迟
-- 主库执行
SHOW MASTER STATUS;
-- 从库执行
SHOW SLAVE STATUS\G
-- 关键字段解读:
-- Master_Host: 主库地址
-- Master_Log_File: 正在读取的Binlog文件
-- Read_Master_Log_Pos: 读取位置
-- Relay_Master_Log_File: 正在执行的Binlog文件
-- Exec_Master_Log_Pos: 执行位置
-- Seconds_Behind_Master: 延迟秒数(NULL表示从库未运行)
使用Prometheus + Grafana实时监控
# prometheus.yml
scrape_configs:
- job_name: 'mysql_replication'
static_configs:
- targets: ['mysql-exporter:9104']
metrics_path: /metrics
# 自定义SQL采集脚本
# mysql_replication.sql
SELECT
VARIABLE_NAME,
VARIABLE_VALUE
FROM information_schema.GLOBAL_STATUS
WHERE VARIABLE_NAME IN (
'Seconds_Behind_Master',
'Slave_IO_Running',
'Slave_SQL_Running',
'Slave_IO_Running_state',
'Slave_SQL_Running_state'
);
// Go语言的监控代码示例
package main
import (
"database/sql"
"log"
"github.com/prometheus/client_golang/prometheus"
_ "github.com/go-sql-driver/mysql"
)
type MySQLReplicationCollector struct {
db *sql.DB
secondsBehindMaster *prometheus.GaugeVec
ioRunning *prometheus.GaugeVec
sqlRunning *prometheus.GaugeVec
}
func NewMySQLReplicationCollector(db *sql.DB) *MySQLReplicationCollector {
return &MySQLReplicationCollector{
db: db,
secondsBehindMaster: prometheus.NewGaugeVec(
prometheus.GaugeOpts{
Name: "mysql_replication_seconds_behind_master",
Help: "Seconds behind master",
},
[]string{"instance"},
),
ioRunning: prometheus.NewGaugeVec(
prometheus.GaugeOpts{
Name: "mysql_replication_io_running",
Help: "IO thread running status",
},
[]string{"instance"},
),
sqlRunning: prometheus.NewGaugeVec(
prometheus.GaugeOpts{
Name: "mysql_replication_sql_running",
Help: "SQL thread running status",
},
[]string{"instance"},
),
}
}
func (c *MySQLReplicationCollector) Collect(ch chan<- prometheus.Metric) {
rows, err := c.db.Query("SHOW SLAVE STATUS")
if err != nil {
log.Printf("Error querying slave status: %v", err)
return
}
defer rows.Close()
for rows.Next() {
var secondsBehind interface{}
var ioRunning, sqlRunning interface{}
err := rows.Scan(&secondsBehind, &ioRunning, &sqlRunning)
if err != nil {
log.Printf("Error scanning row: %v", err)
continue
}
c.secondsBehindMaster.WithLabelValues("slave1").Set(
secondsBehind.(float64),
)
c.ioRunning.WithLabelValues("slave1").Set(
boolToInt(ioRunning.(string) == "Yes"),
)
c.sqlRunning.WithLabelValues("slave1").Set(
boolToInt(sqlRunning.(string) == "Yes"),
)
}
}
5.3 主从不一致时的修复策略
策略一:跳过错误继续同步
-- 临时跳过错误
STOP SLAVE;
SET GLOBAL SQL_SLAVE_SKIP_COUNTER = 1;
START SLAVE;
-- 或者设置跳过特定错误
STOP SLAVE;
SET GLOBAL slave_skip_errors = 1062; -- 跳过重复键错误
START SLAVE;
策略二:使用pt-table-sync修复
# 生成修复SQL
pt-table-sync \
--execute \
--databases=ecommerce \
--tables=orders \
mysql://admin:password@primary/db \
mysql://admin:password@replica/db
# 只打印不执行(先预览)
pt-table-sync \
--print \
--databases=ecommerce \
--tables=orders \
mysql://admin:password@primary/db \
mysql://admin:password@replica/db
策略三:重建从库(最彻底)
# 1. 主库备份
mysqldump --single-transaction --master-data=2 \
--routines --triggers \
-u root -p ecommerce > ecommerce_backup.sql
# 2. 传输到从库
scp ecommerce_backup.sql replica:/tmp/
# 3. 从库恢复
mysql -u root -p ecommerce < /tmp/ecommerce_backup.sql
# 4. 配置主从复制
CHANGE MASTER TO
MASTER_HOST='primary',
MASTER_USER='repl',
MASTER_PASSWORD='password',
MASTER_LOG_FILE='mysql-bin.000001',
MASTER_LOG_POS=154;
START SLAVE;
5.4 GTID复制的优势
-- 启用GTID复制
-- my.cnf配置
[mysqld]
gtid_mode=ON
enforce_gtid_consistency=ON
log_slave_updates=ON
GTID(全局事务标识符)让主从同步更加可靠:
- 每个事务有唯一ID,不会重复执行
- 自动追踪同步位置,不需要手动指定Binlog文件和位置
- 故障切换时更简单
第六部分:分布式场景下的数据一致性维护
6.1 分布式场景的挑战
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ 订单服务 │────▶│ 库存服务 │────▶│ 支付服务 │
│ (数据库A) │ │ (数据库B) │ │ (数据库C) │
└─────────────┘ └─────────────┘ └─────────────┘
│ │
└─────────────┬────────────────────────┘
│
┌─────────────┐
│ 对账服务 │
│ (数据库D) │
└─────────────┘
在分布式系统中,单个事务可能涉及多个数据库、多个服务。传统的ACID已经不够用了。
6.2 分布式事务解决方案对比
| 方案 | 一致性级别 | 性能 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 2PC/XA | 强一致 | 低 | 高 | 金融核心 |
| Seata AT | 最终一致 | 中 | 中 | 电商订单 |
| Seata TCC | 最终一致 | 高 | 高 | 高并发交易 |
| 本地消息表 | 最终一致 | 高 | 中 | 订单通知 |
| RocketMQ事务消息 | 最终一致 | 高 | 低 | 异步解耦 |
| Saga模式 | 最终一致 | 高 | 中 | 长流程业务 |
6.3 本地消息表方案详解
这是最经典、最可靠的最终一致性方案。
核心思路
业务操作和消息记录在同一个本地事务中提交。
然后通过定时任务将消息发送到MQ。
代码实现
// 1. 消息表结构
CREATE TABLE local_message (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
business_id VARCHAR(64) NOT NULL COMMENT '业务ID',
message_type VARCHAR(32) NOT NULL COMMENT '消息类型',
content TEXT NOT NULL COMMENT '消息内容',
status TINYINT NOT NULL DEFAULT 0 COMMENT '0:待发送 1:已发送 2:发送失败',
retry_count INT NOT NULL DEFAULT 0,
next_retry_time DATETIME COMMENT '下次重试时间',
created_at DATETIME NOT NULL,
updated_at DATETIME NOT NULL,
INDEX idx_status (status, next_retry_time)
);
// 2. 业务操作 + 消息记录在同一事务
@Service
public class OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private LocalMessageMapper messageMapper;
@Transactional
public void createOrder(OrderRequest request) {
// 业务操作
Order order = new Order();
order.setProductId(request.getProductId());
order.setQuantity(request.getQuantity());
order.setStatus(OrderStatus.CREATED);
orderMapper.insert(order);
// 本地消息记录(与业务操作在同一事务)
LocalMessage message = new LocalMessage();
message.setBusinessId(String.valueOf(order.getId()));
message.setMessageType("ORDER_CREATED");
message.setContent(JSON.toJSONString(order));
message.setCreatedAt(new Date());
messageMapper.insert(message);
}
}
// 3. 消息发送定时任务
@Service
public class MessageSendJob {
@Autowired
private LocalMessageMapper messageMapper;
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Scheduled(fixedRate = 5000) // 每5秒执行
public void sendMessages() {
// 查询待发送的消息
List<LocalMessage> messages = messageMapper.selectToSend();
for (LocalMessage message : messages) {
try {
// 发送消息
rocketMQTemplate.syncSend(
"order-topic",
MessageBuilder.withPayload(message.getContent()).build()
);
// 更新状态
messageMapper.updateStatus(message.getId(), 1);
} catch (Exception e) {
// 更新重试次数
messageMapper.incrementRetry(message.getId());
// 设置下次重试时间(指数退避)
long nextRetry = System.currentTimeMillis() +
Math.min(1000 * Math.pow(2, message.getRetryCount()), 60000);
messageMapper.updateNextRetryTime(message.getId(), new Date(nextRetry));
}
}
}
}
-- Mapper SQL
-- 查询待发送的消息
<select id="selectToSend" resultType="LocalMessage">
SELECT * FROM local_message
WHERE status = 0
AND (next_retry_time IS NULL OR next_retry_time <= NOW())
ORDER BY created_at ASC
LIMIT 100
</select>
-- 更新状态
<update id="updateStatus">
UPDATE local_message
SET status = #{status}, updated_at = NOW()
WHERE id = #{id}
</update>
6.4 RocketMQ事务消息
@Component
public class TransactionalMessageListener implements TransactionalMessageListener {
@Autowired
private OrderService orderService;
@Override
public LocalTransactionExecuted checkLocalTransaction(MessageExt msg) {
String businessId = msg.getKeys();
// 查询业务状态
Order order = orderService.getByBusinessId(businessId);
if (order != null && order.getStatus() == OrderStatus.CREATED) {
return LocalTransactionExecuted.COMMIT;
}
return LocalTransactionExecuted.ROLLBACK;
}
}
// 发送事务消息
public void sendMessage(String businessId, String content) {
Message msg = new Message("order-topic", content.getBytes());
msg.setKeys(businessId);
// 发送事务消息
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
msg,
null
);
if (result.getLocalTransactionState() == LocalTransactionState.COMMIT_MESSAGE) {
log.info("事务消息提交成功: {}", businessId);
} else {
log.warn("事务消息回滚: {}", businessId);
}
}
6.5 分布式ID生成(雪花算法)
@Component
public class SnowflakeIdGenerator {
private long workerId;
private long datacenterId;
private long sequence = 0L;
private long lastTimestamp = -1L;
public SnowflakeIdGenerator(long workerId, long datacenterId) {
if (workerId > 31 || workerId < 0) {
throw new IllegalArgumentException("worker id can't be greater than 31 or less than 0");
}
if (datacenterId > 31 || datacenterId < 0) {
throw new IllegalArgumentException("datacenter id can't be greater than 31 or less than 0");
}
this.workerId = workerId;
this.datacenterId = datacenterId;
}
public synchronized long nextId() {
long timestamp = System.currentTimeMillis();
if (timestamp < lastTimestamp) {
throw new RuntimeException("Clock moved backwards");
}
if (timestamp == lastTimestamp) {
sequence = (sequence + 1) & 4095L; // 12位序列号
if (sequence == 0) {
timestamp = waitNextMillis(lastTimestamp);
}
} else {
sequence = 0L;
}
lastTimestamp = timestamp;
// 64位ID结构:
// 1位符号位 + 41位时间戳 + 5位datacenter + 5位worker + 12位序列号
return ((timestamp - 1577836800000L) << 22)
| (datacenterId << 17)
| (workerId << 12)
| sequence;
}
private long waitNextMillis(long lastTimestamp) {
long timestamp = System.currentTimeMillis();
while (timestamp <= lastTimestamp) {
timestamp = System.currentTimeMillis();
}
return timestamp;
}
}
第七部分:完整的维护方案与最佳实践
7.1 预防优于修复:架构设计阶段就要考虑一致性
7.1.1 数据分层策略
┌─────────────────────────────────────────────────────────────┐
│ 数据一致性保障层级 │
├─────────────────────────────────────────────────────────────┤
│ L1: 业务层校验(参数校验、业务规则) │
│ L2: 数据库约束(唯一索引、外键、CHECK约束) │
│ L3: 事务管理(本地事务、分布式事务) │
│ L4: 对账系统(定时对账、差异告警) │
│ L5: 人工复核(大额交易、异常预警) │
└─────────────────────────────────────────────────────────────┘
7.1.2 数据库设计规范
-- 1. 所有表必须包含:
-- - 主键(推荐使用雪花算法生成的BIGINT)
-- - 创建时间(created_at)
-- - 更新时间(updated_at)
-- - 逻辑删除标记(deleted)
-- - 版本号(version,用于乐观锁)
-- 2. 关键字段必须加索引
ALTER TABLE orders ADD INDEX idx_user_id (user_id);
ALTER TABLE orders ADD INDEX idx_status_create_time (status, create_time);
ALTER TABLE inventory ADD UNIQUE KEY uk_product_id (product_id);
-- 3. 使用合适的字段类型
-- - 金额字段使用DECIMAL(15,2),禁止使用FLOAT/DOUBLE
-- - 状态字段使用TINYINT,配合枚举类
-- - 时间字段使用DATETIME或TIMESTAMP,统一时区
-- 4. 添加外键约束(或至少添加索引)
ALTER TABLE order_items
ADD CONSTRAINT fk_order_id
FOREIGN KEY (order_id) REFERENCES orders(id);
-- 5. 设置合理的默认值
ALTER TABLE orders
ALTER COLUMN status SET DEFAULT 0;
-- 6. 使用CHECK约束(MySQL 8.0+)
ALTER TABLE orders
ADD CONSTRAINT chk_order_amount
CHECK (amount > 0);
7.2 监控告警体系
7.2.1 关键监控指标
# 监控指标定义
monitoring:
metrics:
# 数据库层面
- name: mysql_replication_lag_seconds
description: "主从复制延迟(秒)"
threshold_critical: 60
threshold_warning: 10
- name: mysql_transactions_per_second
description: "每秒事务数"
threshold_critical: 1000
- name: mysql_lock_wait_count
description: "锁等待次数"
threshold_critical: 100
# 业务层面
- name: order_create_success_rate
description: "订单创建成功率"
threshold_critical: 0.95
- name: inventory_deduct_failure_count
description: "库存扣减失败次数"
threshold_warning: 10
- name: transfer_balance_diff_count
description: "转账余额差异数"
threshold_critical: 1
7.2.2 Prometheus + Grafana配置
# prometheus.yml
global:
scrape_interval: 15s
evaluation_interval: 15s
scrape_configs:
- job_name: 'mysql'
static_configs:
- targets: ['mysql-exporter:9104']
- job_name: 'order-service'
static_configs:
- targets: ['order-service:8080']
# 告警规则
# alert_rules.yml
groups:
- name: mysql_alerts
rules:
- alert: MySQLReplicationLag
expr: mysql_replication_seconds_behind_master > 60
for: 5m
labels:
severity: critical
annotations:
summary: "MySQL主从延迟超过60秒"
description: "当前延迟: {{ $value }}秒"
- alert: MySQLTransactionsPerSecond
expr: mysql_transactions_per_second > 1000
for: 2m
labels:
severity: warning
annotations:
summary: "MySQL事务数过高"
7.3 定期巡检清单
┌─────────────────────────────────────────────────────────────┐
│ 每日巡检清单 │
├─────────────────────────────────────────────────────────────┤
│ □ 检查主从同步状态(Seconds_Behind_Master < 10) │
│ □ 检查Binlog是否按时清理 │
│ □ 检查慢查询日志(slow_query_log) │
│ □ 检查错误日志(error log) │
│ □ 检查磁盘空间使用率 │
│ □ 检查表空间碎片率 │
├─────────────────────────────────────────────────────────────┤
│ 每周巡检清单 │
├─────────────────────────────────────────────────────────────┤
│ □ 运行pt-table-checksum校验主从一致性 │
│ □ 检查死锁情况(SHOW ENGINE INNODB STATUS) │
│ □ 检查大表增长情况 │
│ □ 检查索引使用情况 │
│ □ 备份恢复演练 │
├─────────────────────────────────────────────────────────────┤
│ 每月巡检清单 │
├─────────────────────────────────────────────────────────────┤
│ □ 全量数据一致性校验 │
│ □ 容量规划与性能评估 │
│ □ 安全审计(权限、密码) │
│ □ 备份文件完整性验证 │
└─────────────────────────────────────────────────────────────┘
7.4 应急响应流程
发现数据不一致
│
▼
立即止损
(暂停写入/只读模式)
│
▼
定位问题范围
(哪些表/哪些数据)
│
▼
分析根因
(日志分析/代码审查)
│
▼
制定修复方案
(回填/回滚/跳过)
│
▼
执行修复
(先在测试环境验证)
│
▼
验证修复结果
(对账/业务验证)
│
▼
复盘总结
(预防措施/流程优化)
7.5 代码级别的防御
/**
* 通用幂等性检查工具类
*/
@Component
public class IdempotentChecker {
@Autowired
private RedisTemplate<String, String> redisTemplate;
private static final String IDEMPOTENT_KEY_PREFIX = "idempotent:";
private static final long EXPIRE_SECONDS = 24 * 60 * 60; // 24小时
/**
* 检查操作是否幂等
* @param bizKey 业务唯一标识
* @return true表示可以执行,false表示重复操作
*/
public boolean tryAcquire(String bizKey) {
String key = IDEMPOTENT_KEY_PREFIX + bizKey;
Boolean isSet = redisTemplate.opsForValue()
.setIfAbsent(key, "1", EXPIRE_SECONDS, TimeUnit.SECONDS);
return Boolean.TRUE.equals(isSet);
}
/**
* 带重试的幂等检查
*/
public boolean tryAcquireWithRetry(String bizKey, int retries) {
for (int i = 0; i < retries; i++) {
if (tryAcquire(bizKey)) {
return true;
}
try {
Thread.sleep(100 * (i + 1));
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return false;
}
}
return false;
}
}
/**
* 库存扣减的完整防御性代码
*/
@Service
public class DefensiveInventoryService {
@Autowired
private InventoryMapper inventoryMapper;
@Autowired
private IdempotentChecker idempotentChecker;
@Autowired
private RedisTemplate<String, String> redisTemplate;
/**
* 扣减库存(多层防御)
*/
@Transactional(rollbackFor = Exception.class)
public boolean deductStock(String productId, int quantity, String requestId) {
// 1. 幂等性检查
if (!idempotentChecker.tryAcquire(requestId)) {
log.warn("重复请求,跳过: requestId={}", requestId);
return true; // 幂等返回成功
}
// 2. 参数校验
if (quantity <= 0) {
throw new IllegalArgumentException("库存扣减量必须大于0");
}
// 3. 分布式锁(防止并发超卖)
String lockKey = "inventory_lock:" + productId;
Boolean locked = redisTemplate.opsForValue()
.setIfAbsent(lockKey, "1", 10, TimeUnit.SECONDS);
if (!locked) {
log.warn("获取锁失败,请稍后重试: productId={}", productId);
throw new BusinessException("系统繁忙,请稍后重试");
}
try {
// 4. 查询当前库存
Inventory inventory = inventoryMapper.selectById(productId);
if (inventory == null) {
throw new BusinessException("商品不存在");
}
if (inventory.getStock() < quantity) {
throw new BusinessException("库存不足");
}
// 5. 乐观锁扣减
int rows = inventoryMapper.deductWithVersion(
productId, quantity, inventory.getVersion()
);
if (rows == 0) {
throw new BusinessException("库存扣减失败,请重试");
}
// 6. 记录操作日志
inventoryLogMapper.insert(new InventoryLog(
productId, quantity, "DEDUCT", System.currentTimeMillis()
));
return true;
} finally {
// 7. 释放锁
redisTemplate.delete(lockKey);
}
}
}
第八部分:总结与核心要点
8.1 数据不一致的根因总结
┌────────────────────────────────────────────────────────────────┐
│ 数据不一致的根因图谱 │
├────────────────────────────────────────────────────────────────┤
│ │
│ ┌──────────────┐ │
│ │ 并发竞争 │──▶ 超卖、重复扣款、重复转账 │
│ └──────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────┐ │
│ │ 网络故障 │──▶ 主从不一致、消息丢失 │
│ └──────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────┐ │
│ │ 系统故障 │──▶ 事务中断、数据丢失 │
│ └──────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────┐ │
│ │ 设计缺陷 │──▶ 缺少约束、缺少校验 │
│ └──────────────┘ │
│ │
└────────────────────────────────────────────────────────────────┘
8.2 核心解决方案速查表
| 场景 | 首选方案 | 备选方案 | 兜底方案 |
|---|---|---|---|
| 电商超卖 | 乐观锁 + Redis Lua | 分布式锁 + 数据库行锁 | 对账 + 人工补偿 |
| 银行转账 | 本地事务 / XA | Seata TCC | 本地消息表 + 对账 |
| 主从同步 | GTID复制 | Binlog位点复制 | pt-table-sync修复 |
| 分布式事务 | Seata AT | 本地消息表 | 定期对账 |
| 消息可靠性 | RocketMQ事务消息 | 本地消息表 | 手动补偿 |
8.3 记住这五点
永远不要相信单一系统:主从会有延迟,缓存会有不一致,网络会故障,永远要有兜底方案。
对账是最后的防线:无论架构多完美,都应该有定期对账机制。差异不可怕,可怕的是不知道差异存在。
幂等性是分布式系统的基石:任何可能重复执行的操作,都必须做好幂等处理。
监控告警要提前建设:不要等问题发生了再去搭建监控,要在架构设计阶段就规划好监控体系。
预案比技术更重要:技术问题总有解决方案,但如果没有预案,问题发生时慌乱应对只会让情况更糟。
写在最后
数据不一致是分布式系统永远要面对的难题,没有银弹,只有层层防御。
电商超卖、银行转账错误,这些问题的背后是架构设计的得失。好的架构能在问题发生前就预防,能在问题发生时快速定位,能在问题发生后及时恢复。
希望这篇文章能帮你建立起完整的数据一致性保障体系。记住:防御性编程不是悲观,而是专业。
如果你有具体的场景或问题,欢迎进一步交流!
