千万级订单销量排行榜优化:Canal+Redis+预聚合表
前言
发现在表数据量达到500w级别时,排行榜查询会变得非常慢,慢到分钟级别.
这个数据量下即使加了索引,响应时间也是秒级的,为了优化用户和商家的体验,我打算利用redis的zset实现排行榜功能.
⏸ plan mode on
技术栈: canal,redis,MySQL,MyBatis,Springboot
初步方案: 用canal监听MySQL的binlog 变化,将MySQL的变更同步到redis的zset
方案细节:
- 在同步前,应先将全量数据同步到redis
- 处理好退款等可能导致订单量减少的问题
方案解释:
为什么要用canal而不直接双写?
- MySQL 的事务提交后才会产生 binlog,只要 binlog 到位,回放即可跟 MySQL 精确对齐.不会出现一个成功一个失败的情况
- 业务耦合低,对原代码无侵入
方案流程:
- canal会伪装成MySQL的slave,使MySQL向他推送binlog更新
- canal解析binlog内容,发送到应用层
- 应用层将数据推送到redis
⏵⏵ accept edits on
核心架构
订单完成 status→5 → MySQL binlog → Canal Server → SalesRankCanalClient → ZINCRBY Redis ZSet
关键逻辑
- Canal 客户端 (SalesRankCanalClient.java) — 独立 daemon 线程循环拉取 Canal Server 的 binlog 事件,过滤 orders 表 UPDATE 且 status 从非5→5 的行,提取orderId 后调用 SalesRankService.updateRank()。含断线重连机制。
@Override
public void run(String... args) {
listenerThread = new Thread(this::listen, "canal-sales-rank");
listenerThread.setDaemon(true);
listenerThread.start();
log.info("Canal 销量排行监听线程已启动");
}
private void listen() {
while (running) {
CanalConnector connector = CanalConnectors.newSingleConnector(
new InetSocketAddress(canalProperties.getHost(), canalProperties.getPort()),
canalProperties.getDestination(),
canalProperties.getUsername(),
canalProperties.getPassword()
);
try {
connector.connect();
connector.subscribe("sky_take_out\\.orders");
connector.rollback();
log.info("Canal 连接成功: {}:{}/{}",
canalProperties.getHost(), canalProperties.getPort(), canalProperties.getDestination());
while (running) {
Message message = connector.getWithoutAck(100, 1000L, TimeUnit.MILLISECONDS);
long batchId = message.getId();
if (batchId == -1 || message.getEntries().isEmpty()) {
continue;
}
processEntries(message.getEntries());
connector.ack(batchId);
}
} catch (Exception e) {
log.error("Canal 连接异常,10秒后重试...", e);
sleepQuietly(10_000);
} finally {
try {
connector.disconnect();
} catch (Exception ignored) {
}
}
}
}
当前代码的 while 循环是合适的,因为:
- Canal 是长连接流式监听,必须持续运行
- 双层 while 实现了消费循环 + 断线重连
- getWithoutAck 的超时机制避免了 CPU 空转
private void handleRowUpdate(CanalEntry.RowData rowData) {
String afterStatus = null;
String beforeStatus = null;
Long orderId = null;
for (CanalEntry.Column column : rowData.getAfterColumnsList()) {
if ("id".equals(column.getName())) {
orderId = Long.valueOf(column.getValue());
}
if ("status".equals(column.getName())) {
afterStatus = column.getValue();
}
}
for (CanalEntry.Column column : rowData.getBeforeColumnsList()) {
if ("status".equals(column.getName())) {
beforeStatus = column.getValue();
}
}
if ("5".equals(afterStatus) && !"5".equals(beforeStatus)) {
log.info("Canal 检测到订单完成: orderId={}", orderId);
try {
salesRankService.updateRank(orderId);
} catch (Exception e) {
log.error("更新销量排行失败: orderId={}", orderId, e);
}
}
}
- 排行更新 (SalesRankServiceImpl.updateRank()) — 根据 orderId 查 order_detail,对每条明细:
dishId != null → ZINCRBY sales:rank:dish {number} {dishId}
setmealId != null → ZINCRBY sales:rank:setmeal {number} {setmealId}
@Override
public void updateRank(Long orderId) {
List<OrderDetail> details = orderDetailMapper.getByOrderId(orderId);
if (details == null || details.isEmpty()) {
log.warn("订单 {} 无明细数据,跳过排行更新", orderId);
return;
}
for (OrderDetail detail : details) {
if (detail.getDishId() != null) {
redisTemplate.opsForZSet().incrementScore(
DISH_RANK_KEY,
String.valueOf(detail.getDishId()),
detail.getNumber().doubleValue()
);
}
if (detail.getSetmealId() != null) {
redisTemplate.opsForZSet().incrementScore(
SETMEAL_RANK_KEY,
String.valueOf(detail.getSetmealId()),
detail.getNumber().doubleValue()
);
}
}
log.info("销量排行已更新: orderId={}, items={}", orderId, details.size());
}
- 排行查询 (getDishRanking/getSetmealRanking) — ZREVRANGE … WITHSCORES 取 Top N,解析 member(ID 字符串)后批量查 DishMapper.selectBatch() /SetmealMapper.getByIds() 组装名称和图片返回。
@Override
public List<SalesRankVO> getDishRanking(int top) {
Set<ZSetOperations.TypedTuple<Object>> topSet =
redisTemplate.opsForZSet().reverseRangeWithScores(DISH_RANK_KEY, 0, top - 1);
if (topSet == null || topSet.isEmpty()) {
return Collections.emptyList();
}
List<Long> dishIds = new ArrayList<>();
List<Double> scores = new ArrayList<>();
for (ZSetOperations.TypedTuple<Object> entry : topSet) {
dishIds.add(Long.valueOf((String) entry.getValue()));
scores.add(entry.getScore());
}
List<Dish> dishes = dishMapper.selectBatch(dishIds);
Map<Long, Dish> dishMap = dishes.stream()
.collect(Collectors.toMap(Dish::getId, d -> d, (a, b) -> a));
List<SalesRankVO> result = new ArrayList<>();
for (int i = 0; i < dishIds.size(); i++) {
Long dishId = dishIds.get(i);
Dish dish = dishMap.get(dishId);
result.add(SalesRankVO.builder()
.rank(i + 1)
.id(dishId)
.name(dish != null ? dish.getName() : "未知菜品")
.sales(scores.get(i).intValue())
.image(dish != null ? dish.getImage() : null)
.build());
}
return result;
}
- 全量兜底 (syncFromDatabase) — 查询所有 status=5 的 order_detail,按 dish_id / setmeal_id 分别 SUM(number) 后 ZADD 覆盖写入 Redis。启动时若 ZSet为空自动执行,每天凌晨 3 点定时重同步。
@Override
public void syncFromDatabase() {
log.info("开始全量同步销量排行...");
List<Map<String, Object>> dishSales = orderMapper.getCompletedDishSales();
if (dishSales != null && !dishSales.isEmpty()) {
redisTemplate.delete(DISH_RANK_KEY);
for (Map<String, Object> row : dishSales) {
Object idObj = row.get("id");
Object numObj = row.get("number");
String member = String.valueOf(idObj);
double score = toDouble(numObj);
redisTemplate.opsForZSet().add(DISH_RANK_KEY, member, score);
}
log.info("菜品销量同步完成: {} 条", dishSales.size());
}
List<Map<String, Object>> setmealSales = orderMapper.getCompletedSetmealSales();
if (setmealSales != null && !setmealSales.isEmpty()) {
redisTemplate.delete(SETMEAL_RANK_KEY);
for (Map<String, Object> row : setmealSales) {
Object idObj = row.get("id");
Object numObj = row.get("number");
String member = String.valueOf(idObj);
double score = toDouble(numObj);
redisTemplate.opsForZSet().add(SETMEAL_RANK_KEY, member, score);
}
log.info("套餐销量同步完成: {} 条", setmealSales.size());
}
}
启动
超绝预热时间,查两次500w级别的库再inner join 1000w的detail表直接干到分钟级别的查询延迟,待会再加个预聚合表优化,先看这套方案行不行
Redis:

正常,查询看看耗时
优化成功,最后压测看看性能,JMeter(3000 3 2)
够用了
预聚合
预聚合是指提前对原始数据进行计算和汇总,将聚合结果(如求和、计数、平均值等)存储起来,以便后续查询时直接读取这些预先计算好的结果,而不需要每次都重新扫描全部原始数据进行计算。
非常适合我的情况:千万级表查询
索引优化
改造前先看看能不能加索引优化:
SELECT od.dish_id, 1, SUM(od.number)
FROM order_detail od
INNER JOIN orders o ON od.order_id = o.id
WHERE o.status = 5 AND od.dish_id IS NOT NULL
GROUP BY od.dish_id
SELECT od.setmeal_id, 2, SUM(od.number)
FROM order_detail od
INNER JOIN orders o ON od.order_id = o.id
WHERE o.status = 5 AND od.setmeal_id IS NOT NULL
GROUP BY od.setmeal_id
对应的表
create table orders
(
id bigint auto_increment comment '主键'
primary key,
number varchar(50) null comment '订单号',
status int default 1 not null comment '订单状态 1待付款 2待接单 3已接单 4派送中 5已完成 6已取消 7退款',
user_id bigint not null comment '下单用户',
address_book_id bigint not null comment '地址id',
order_time datetime not null comment '下单时间',
checkout_time datetime null comment '结账时间',
pay_method int default 1 not null comment '支付方式 1微信,2支付宝',
pay_status tinyint default 0 not null comment '支付状态 0未支付 1已支付 2退款',
amount decimal(10, 2) not null comment '实收金额',
remark varchar(100) null comment '备注',
phone varchar(11) null comment '手机号',
address varchar(255) null comment '地址',
user_name varchar(32) null comment '用户名称',
consignee varchar(32) null comment '收货人',
cancel_reason varchar(255) null comment '订单取消原因',
rejection_reason varchar(255) null comment '订单拒绝原因',
cancel_time datetime null comment '订单取消时间',
estimated_delivery_time datetime null comment '预计送达时间',
delivery_status tinyint(1) default 1 not null comment '配送状态 1立即送出 0选择具体时间',
delivery_time datetime null comment '送达时间',
pack_amount int null comment '打包费',
tableware_number int null comment '餐具数量',
tableware_status tinyint(1) default 1 not null comment '餐具数量状态 1按餐量提供 0选择具体数量'
)
comment '订单表' collate = utf8mb3_bin;
create index idx_cover_query
on orders (user_id, status, order_time, pay_status, id);
create table order_detail
(
id bigint auto_increment comment '主键'
primary key,
name varchar(32) null comment '名字',
image varchar(255) null comment '图片',
order_id bigint not null comment '订单id',
dish_id bigint null comment '菜品id',
setmeal_id bigint null comment '套餐id',
dish_flavor varchar(50) null comment '口味',
number int default 1 not null comment '数量',
amount decimal(10, 2) not null comment '金额'
)
comment '订单明细表' collate = utf8mb3_bin;
create index idx_order_dish_number
on order_detail (order_id, dish_id, number);
这个覆盖索引是之前那篇千万级表单分页查询优化时加的.
- 对于order表,我们先考虑where后面的字段status和join后面的id.
- 对于order_detail表,我们考虑join后的order_id和第一条sql的group后的dish_id以及第二条sql的group后的setmeal_id,同时可以将SUM中的number也考虑进去,建覆盖索引
注意,在这里order_detail是被驱动表,被驱动表建索引时先考虑join后的字段再考虑其他列
建:
ALTER TABLE orders ADD INDEX idx_status_id (status, id);
ALTER TABLE order_detail ADD INDEX idx_order_dish_number (order_id, dish_id, number);
ALTER TABLE order_detail ADD INDEX idx_order_setmeal_number(order_id,setmeal_id,number);
来对比下建索引前后的查询销量:
2026-05-07 19:50:39.112 INFO 37072 — [nio-8080-exec-2] c.sky.service.impl.SalesRankServiceImpl : 开始从大表重建汇总表…
2026-05-07 19:56:06.897 INFO 37072 — [nio-8080-exec-2] c.sky.service.impl.SalesRankServiceImpl : 汇总表重建完成, 耗时 327785ms
建orders表的索引后:
2026-05-07 20:20:42.618 INFO 37072 — [nio-8080-exec-8] c.sky.service.impl.SalesRankServiceImpl : 开始从大表重建汇总表…
2026-05-07 20:23:53.464 INFO 37072 — [nio-8080-exec-8] c.sky.service.impl.SalesRankServiceImpl : 汇总表重建完成, 耗时 190846ms
建order_detail的索引后:
2026-05-07 21:43:49.882 INFO 37072 — [nio-8080-exec-2] c.sky.service.impl.SalesRankServiceImpl : 开始从大表重建汇总表…
2026-05-07 21:44:34.733 INFO 37072 — [nio-8080-exec-2] c.sky.service.impl.SalesRankServiceImpl : 汇总表重建完成, 耗时 44851ms
327s->190s->44s,优化的幅度还行,考虑到重建的频率较低,将重建方法放在凌晨执行即可
建表与SQL
create table sky_take_out.sales_summary
(
id bigint auto_increment primary key,
item_id bigint not null comment 'dish_id 或 setmeal_id',
item_type tinyint not null comment '1=菜品 2=套餐',
total_sales int default 0 not null comment '累计销量',
updated_at datetime default CURRENT_TIMESTAMP null on update CURRENT_TIMESTAMP,
constraint uk_item
unique (item_type, item_id)
)
comment '销量汇总表';
unique (item_type, item_id)保证菜品和套餐不冲突
对应全量同步SQL:
INSERT INTO sales_summary (item_id, item_type, total_sales)
SELECT od.dish_id, 1, SUM(od.number)
FROM order_detail od
INNER JOIN orders o ON od.order_id = o.id
WHERE o.status = 5 AND od.dish_id IS NOT NULL
GROUP BY od.dish_id
ON DUPLICATE KEY UPDATE total_sales = VALUES(total_sales)
INSERT INTO sales_summary (item_id, item_type, total_sales)
SELECT od.setmeal_id, 2, SUM(od.number)
FROM order_detail od
INNER JOIN orders o ON od.order_id = o.id
WHERE o.status = 5 AND od.setmeal_id IS NOT NULL
GROUP BY od.setmeal_id
ON DUPLICATE KEY UPDATE total_sales = VALUES(total_sales)
ON DUPLICATE KEY UPDATE就是有则更新,无则插入,也叫 Upsert(Update + Insert),和表的unique (item_type, item_id)相配合,VALUES(total_sales)是引用INSERT那边的新值
业务层
改造
更新更新排行榜的方法,这里是提速的核心
@Override
public void updateRank(Long orderId) {
List<OrderDetail> details = orderDetailMapper.getByOrderId(orderId);
if (details == null || details.isEmpty()) {
log.warn("订单 {} 无明细数据,跳过排行更新", orderId);
return;
}
for (OrderDetail detail : details) {
if (detail.getDishId() != null) {
int number = detail.getNumber();
// 双写:MySQL 汇总表
salesSummaryMapper.insertOrIncrement(SalesSummary.builder()
.itemId(detail.getDishId())
.itemType(SalesSummary.ITEM_TYPE_DISH)
.totalSales(number)
.build());
// 双写:Redis ZSet
redisTemplate.opsForZSet().incrementScore(
DISH_RANK_KEY,
String.valueOf(detail.getDishId()),
number
);
}
if (detail.getSetmealId() != null) {
int number = detail.getNumber();
salesSummaryMapper.insertOrIncrement(SalesSummary.builder()
.itemId(detail.getSetmealId())
.itemType(SalesSummary.ITEM_TYPE_SETMEAL)
.totalSales(number)
.build());
redisTemplate.opsForZSet().incrementScore(
SETMEAL_RANK_KEY,
String.valueOf(detail.getSetmealId()),
number
);
}
}
log.info("销量排行已更新: orderId={}, items={}", orderId, details.size());
}
同步到redis的方法更新为从汇总表查询
@Override
public void syncFromDatabase() {
log.info("开始从汇总表全量同步销量排行...");
long start = System.currentTimeMillis();
List<SalesSummary> summaries = salesSummaryMapper.selectAll();
if (summaries == null || summaries.isEmpty()) {
log.info("汇总表无数据,跳过同步");
return;
}
redisTemplate.delete(DISH_RANK_KEY);
redisTemplate.delete(SETMEAL_RANK_KEY);
int dishCount = 0;
int setmealCount = 0;
for (SalesSummary s : summaries) {
String key = SalesSummary.ITEM_TYPE_DISH.equals(s.getItemType()) ? DISH_RANK_KEY : SETMEAL_RANK_KEY;
String member = String.valueOf(s.getItemId());
redisTemplate.opsForZSet().add(key, member, s.getTotalSales().doubleValue());
if (SalesSummary.ITEM_TYPE_DISH.equals(s.getItemType())) {
dishCount++;
} else {
setmealCount++;
}
}
long elapsed = System.currentTimeMillis() - start;
log.info("销量排行同步完成: 菜品 {} 项, 套餐 {} 项, 耗时 {}ms", dishCount, setmealCount, elapsed);
}
新增
重建汇总表方法
@Override
public void syncSummaryTable() {
log.info("开始从大表重建汇总表...");
long start = System.currentTimeMillis();
salesSummaryMapper.rebuildDishSummary();
salesSummaryMapper.rebuildSetmealSummary();
long elapsed = System.currentTimeMillis() - start;
log.info("汇总表重建完成, 耗时 {}ms", elapsed);
}
改造完成,测试下同步速度
响应时间很短
再新增订单,查看逻辑是否正常
2026-05-07 23:08:43.180 INFO 37072 — [nio-8080-exec-2] com.sky.controller.user.OrderController : 用户下单:OrdersSubmitDTO(addressBookId=2, payMethod=1, remark=, estimatedDeliveryTime=2026-05-28T21:24:58, deliveryStatus=1, tablewareNumber=1, tablewareStatus=0, packAmount=1, amount=45)
2026-05-07 23:10:13.146 INFO 37072 — [nio-8080-exec-5] com.sky.controller.user.OrderController : 订单支付:OrdersPaymentDTO(orderNumber=1778166523182, payMethod=1)
2026-05-07 23:10:18.526 INFO 37072 — [nio-8080-exec-5] com.sky.controller.user.OrderController : 生成预支付交易单:OrderPaymentVO(nonceStr=null, paySign=null, timeStamp=null, signType=null, packageStr=null)
2026-05-07 23:11:37.207 INFO 37072 — [anal-sales-rank] com.sky.canal.SalesRankCanalClient : Canal 检测到订单完成: orderId=5220352
2026-05-07 23:11:37.222 INFO 37072 — [anal-sales-rank] c.sky.service.impl.SalesRankServiceImpl : 销量排行已更新: orderId=5220352, items=1

数据库更新时间能和日志对上,更新正常
redis和数据库数量对的上,一致性正常
退款
若是给用户看的排行榜,可以不考虑退款,因为是强调热度.但在商家端需要考核真实业绩,必须考虑退款情况.
改造
实现类不能单纯增加,把zadd改为zincreby
@Override
public void updateRank(Long orderId) {
adjustRank(orderId, true);
}
@Override
public void decreaseRank(Long orderId) {
adjustRank(orderId, false);
}
private void adjustRank(Long orderId, boolean isIncrement) {
List<OrderDetail> details = orderDetailMapper.getByOrderId(orderId);
if (details == null || details.isEmpty()) {
log.warn("订单 {} 无明细数据,跳过排行调整", orderId);
return;
}
String action = isIncrement ? "增加" : "减少";
for (OrderDetail detail : details) {
int delta = isIncrement ? detail.getNumber() : -detail.getNumber();
if (detail.getDishId() != null) {
salesSummaryMapper.insertOrIncrement(SalesSummary.builder()
.itemId(detail.getDishId())
.itemType(SalesSummary.ITEM_TYPE_DISH)
.totalSales(delta)
.build());
redisTemplate.opsForZSet().incrementScore(
DISH_RANK_KEY,
String.valueOf(detail.getDishId()),
delta
);
}
if (detail.getSetmealId() != null) {
salesSummaryMapper.insertOrIncrement(SalesSummary.builder()
.itemId(detail.getSetmealId())
.itemType(SalesSummary.ITEM_TYPE_SETMEAL)
.totalSales(delta)
.build());
redisTemplate.opsForZSet().incrementScore(
SETMEAL_RANK_KEY,
String.valueOf(detail.getSetmealId()),
delta
);
}
}
log.info("销量排行已{}: orderId={}, items={}", action, orderId, details.size());
}
调用改为
private void handleRowUpdate(CanalEntry.RowData rowData) {
Long orderId = null;
String afterStatus = null;
String beforeStatus = null;
String afterPayStatus = null;
String beforePayStatus = null;
for (CanalEntry.Column column : rowData.getAfterColumnsList()) {
if ("id".equals(column.getName())) {
orderId = Long.valueOf(column.getValue());
} else if ("status".equals(column.getName())) {
afterStatus = column.getValue();
} else if ("pay_status".equals(column.getName())) {
afterPayStatus = column.getValue();
}
}
for (CanalEntry.Column column : rowData.getBeforeColumnsList()) {
if ("status".equals(column.getName())) {
beforeStatus = column.getValue();
} else if ("pay_status".equals(column.getName())) {
beforePayStatus = column.getValue();
}
}
// 订单完成: status 变为 5 → 加销量
if ("5".equals(afterStatus) && !"5".equals(beforeStatus)) {
log.info("Canal 检测到订单完成: orderId={}", orderId);
try {
salesRankService.updateRank(orderId);
} catch (Exception e) {
log.error("增加销量排行失败: orderId={}", orderId, e);
}
return;
}
// 订单取消/退款(通过status): status 从 5 变为其他 → 减销量
if (!"5".equals(afterStatus) && "5".equals(beforeStatus)) {
log.info("Canal 检测到订单取消(状态变更): orderId={}, {}→{}", orderId, beforeStatus, afterStatus);
try {
salesRankService.decreaseRank(orderId);
} catch (Exception e) {
log.error("减少销量排行失败: orderId={}", orderId, e);
}
return;
}
// 支付级退款: status 仍是 5, 但 pay_status 变为 2(退款) → 减销量
if ("5".equals(afterStatus) && !"2".equals(beforePayStatus) && "2".equals(afterPayStatus)) {
log.info("Canal 检测到支付退款: orderId={}, payStatus {}→{}", orderId, beforePayStatus, afterPayStatus);
try {
salesRankService.decreaseRank(orderId);
} catch (Exception e) {
log.error("减少销量排行失败: orderId={}", orderId, e);
}
}
}
来测试看下,改前数据

发起退款
2026-05-08 10:02:09.317 INFO 26788 — [nio-8080-exec-7] c.sky.controller.admin.OrderController : 退款: OrdersRefundDTO(id=5220352)
2026-05-08 10:02:09.438 INFO 26788 — [anal-sales-rank] com.sky.canal.SalesRankCanalClient : Canal 检测到支付退款: orderId=5220352, payStatus 1→2
2026-05-08 10:02:09.451 INFO 26788 — [anal-sales-rank] c.sky.service.impl.SalesRankServiceImpl : 销量排行已减少: orderId=5220352, items=1


成功双减
我给你最直接最不绕弯最直白最不废话的结论
面对百万、千万级大表的排行榜查询,单纯依赖 MySQL 索引已经难以满足低延迟需求。通过Canal 异步同步 + Redis 做热点排行 + 预聚合表兜底的组合方案,既解决了慢查询痛点,又保证了实时性与数据准确性,同时优雅处理了退款逆向业务,架构解耦、维护简单,是高并发榜单场景非常实用的实践。
还有点小问题
Canal 无幂等、位点不持久化,重复消费导致销量错乱
MySQL+Redis 双写无事务无补偿,数据可能会不一致
这些问题先用手动重建兜底
更多推荐




所有评论(0)