8200 万行的订单表,加个索引要 6 小时
2021 年 2 月,我们的 t_order 表数据量到了 8200 万行,物理大小 68 GB(含索引)。
mysql> SELECT table_name, table_rows,
-> ROUND(data_length/1024/1024/1024, 2) data_gb,
-> ROUND(index_length/1024/1024/1024, 2) idx_gb
-> FROM information_schema.tables WHERE table_name='t_order';
+------------+------------+---------+--------+
| table_name | table_rows | data_gb | idx_gb |
+------------+------------+---------+--------+
| t_order | 82341779 | 41.20 | 27.31 |
+------------+------------+---------+--------+
压垮我们的最后一根稻草是:想加一个索引,DBA 说要 6 小时,而且期间主库会有锁表风险。
同时慢查询日志里这张表占了 41%,最典型的:
# Query_time: 3.842011 Lock_time: 0.000121 Rows_sent: 20 Rows_examined: 1204318
SELECT * FROM t_order WHERE user_id = 882341
AND status IN (1,2,3) ORDER BY create_time DESC LIMIT 20;
扫 120 万行取 20 行。idx_user_id 是有的,但 status IN 和 ORDER BY 让优化器放弃了一部分索引。这种 SQL 在 4000 万行时不算问题,到 8200 万就崩了。
为什么选 ShardingSphere 5.0.0-beta
选型时对比了三个:
| 方案 | 形态 | 我们的顾虑 |
|---|---|---|
| ShardingSphere-JDBC 5.0.0-beta | 客户端 JAR | beta 版,但社区活跃,API 设计比 4.x 清晰 |
| ShardingSphere-JDBC 4.1.1 | 客户端 JAR | 稳定,但配置格式是老的一套 |
| MyCat / ShardingSphere-Proxy | 独立代理 | 多一跳网络、要多维护一套中间件 |
最后选了 5.0.0-beta。理由很实在:我们是 Java 单体技术栈,客户端直连少一跳网络,延迟更低;beta 这个词吓人,但我们压测了两周,核心的 CRUD、分页、聚合都正常,且社区 issue 响应很快。5.0 的 YAML 配置比 4.x 的 ShardingRuleConfiguration 编程式配置清爽太多。
<dependency>
<groupId>org.apache.shardingsphere</groupId>
<artifactId>shardingsphere-jdbc-core</artifactId>
<version>5.0.0-beta</version>
</dependency>
分片键的选择是最重要的一步
我把近三个月的慢查询日志全导出来,统计了这张表的 WHERE 条件分布:
| 查询模式 | 占比 | 举例 |
|---|---|---|
| 按 user_id 查自己的订单列表 | 68% | C 端"我的订单" |
| 按 order_id 查单条详情 | 21% | 订单详情页、支付回调 |
| 按 merchant_id 查 | 7% | 商家后台 |
| 运营多维查询(时间 + 状态 + 金额) | 4% | 运营后台 |
结论很明确:用 user_id 做分片键,覆盖 68% 的主流量。
但订单详情和支付回调都只有 order_id,没有 user_id。如果只按 user_id 分片,这两类查询会路由到全部分片,21% 的流量全部放大 32 倍,这不能接受。
基因法:让 order_id 自带分片信息
解决办法是在生成订单号时,把 user_id 的分片路由值嵌入 order_id,这样两个字段算出来的分片位置天然一致。
public class OrderIdGenerator {
private static final int SHARD_COUNT = 32; // 2 库 × 16 表
private static final int GENE_BITS = 5; // 2^5 = 32,取 user_id 末 5 位
public static long nextOrderId(long userId) {
// 末 5 位基因,与 userId % 32 一致
long gene = userId & (SHARD_COUNT - 1);
// 时间戳部分:秒级,保留 2021-01-01 起的偏移
long seconds = Instant.now().getEpochSecond() - 1609459200L;
// 序列部分:同一秒内的自增
long seq = SEQUENCE.incrementAndGet() & 0x3FF; // 10 bit, 0~1023
// 布局:41 bit 秒 + 10 bit 序列 + 5 bit 基因 = 56 bit
return (seconds << 15) | (seq << 5) | gene;
}
// 从 orderId 反推分片位置
public static int shardIndexOfOrderId(long orderId) {
return (int) (orderId & (SHARD_COUNT - 1));
}
}
验证一下这个设计:
long userId = 882341L;
long orderId = OrderIdGenerator.nextOrderId(userId);
System.out.println(userId % 32); // 21
System.out.println(OrderIdGenerator.shardIndexOfOrderId(orderId)); // 21
一致。这样两种查询都能精确定位到单个分片:
SELECT * FROM t_order WHERE user_id = 882341 ...; -- 走 user_id % 32 = 21
SELECT * FROM t_order WHERE order_id = 12345678901; -- 走基因位 = 21
基因法的代价是 order_id 不再是纯递增的,同一秒内的订单在数值上跳跃。这对我们没影响(没用 ID 排序),但如果你有依赖 ID 单调性的逻辑要提前想清楚。
分多少个片
按三年容量算:
当前 8200 万
年增量 约 3400 万(2020 年实际)
三年后 8200 + 3400 × 3 = 1.84 亿
分片 32: 1.84 亿 / 32 = 575 万/表 ✓ 单表 500 万左右是舒适区
分片 16: 1.84 亿 / 16 = 1150 万/表 ✗ 偏大,索引深度会到 4 层
最后定 2 库 × 16 表 = 32 个分片。库数取 2 是因为我们只有两台物理机做这个集群(预算有限),而且 2 的幂次方便以后翻倍扩容。
分布式主键:时钟回拨怎么办
ShardingSphere 内置了 SNOWFLAKE 分布式主键算法,但它有个我们没有选择的原因:worker.id 的配置在单机部署时靠 ServerSocket 找空闲端口生成,在容器环境下不可靠,两个 Pod 可能拿到同一个 worker.id,产生重复主键。
所以我们用了上面那套自带基因的自研方案,ID 里不含 worker 节点信息,也就没有 worker.id 冲突的问题。代价是同一秒同一台机器上的订单 ID 末 5 位固定,不随机。
但时钟回拨依然要处理。我们的做法是启动时检测 + 运行时兜底:
private static final AtomicInteger SEQUENCE = new AtomicInteger(0);
private static volatile long lastSecond = 0;
private static synchronized long nextSeconds() {
long sec = Instant.now().getEpochSecond() - 1609459200L;
if (sec < lastSecond) {
// 时钟回拨,等待追平
long back = lastSecond - sec;
if (back > 5) {
throw new IllegalStateException("时钟回拨超过 5 秒,拒绝生成订单号: " + back);
}
LockSupport.parkNanos(back * 1_000_000_000L);
sec = lastSecond;
}
if (sec > lastSecond) {
lastSecond = sec;
SEQUENCE.set(0);
}
return sec;
}
回拨 5 秒以内就等,超过 5 秒直接抛异常让请求失败。宁可失败也不能生成重复主键,这是分布式 ID 的基本原则。我们的机器配了 NTP,一年多没触发过。
配置
schemaName: order_db
dataSources:
ds0:
dataSourceClassName: com.zaxxer.hikari.HikariDataSource
props:
jdbcUrl: jdbc:mysql://10.0.1.10:3306/order_0?useSSL=false
username: order
password: ***
maximumPoolSize: 50
ds1:
dataSourceClassName: com.zaxxer.hikari.HikariDataSource
props:
jdbcUrl: jdbc:mysql://10.0.1.11:3306/order_1?useSSL=false
username: order
password: ***
maximumPoolSize: 50
rules:
- !SHARDING
tables:
t_order:
actualDataNodes: ds$->{0..1}.t_order_$->{0..15}
databaseStrategy:
standard:
shardingColumn: user_id
shardingAlgorithmName: db-inline
tableStrategy:
standard:
shardingColumn: user_id
shardingAlgorithmName: tbl-inline
t_order_item:
actualDataNodes: ds$->{0..1}.t_order_item_$->{0..15}
databaseStrategy:
standard:
shardingColumn: user_id
shardingAlgorithmName: db-inline
tableStrategy:
standard:
shardingColumn: user_id
shardingAlgorithmName: tbl-inline
bindingTables:
- t_order, t_order_item
broadcastTables:
- t_dict_region, t_dict_status
shardingAlgorithms:
db-inline:
type: INLINE
props:
algorithm-expression: ds${user_id % 32 / 16}
tbl-inline:
type: INLINE
props:
algorithm-expression: t_order_${user_id % 16}
三个配置值得单独说:
bindingTables。t_order和t_order_item用同样的分片规则,配成绑定表后,JOIN 查询不会走笛卡尔积。我们测过一个order JOIN order_item的查询,不配绑定表时 ShardingSphere 要算 32 × 32 = 1024 种组合,配上之后只有 32 次。这个差别是致命的。broadcastTables。字典表不分片,全量同步到每个库,JOIN 时本地完成。注意广播表的写操作会同步到所有库,别把频繁更新的表配成广播表。- 库路由表达式
user_id % 32 / 16。先对 32 取模再整除 16,等价于user_id % 32 >= 16 ? 1 : 0。这么写是为了和基因法里的 32 对齐。
跨库分页:绕不开的难题
C 端的查询带了 user_id,全部路由到单库单表,改完之后 P99 从 3.8 秒降到 12 ms。这部分很顺利。
真正麻烦的是运营后台。这类 SQL 没有分片键:
SELECT * FROM t_order
WHERE create_time BETWEEN '2021-02-01' AND '2021-02-20'
AND status = 2
ORDER BY create_time DESC
LIMIT 20 OFFSET 100;
ShardingSphere 的处理是:改写后发到全部 32 个分片,每个分片取 LIMIT 120(0 到 offset+limit),拿回 32 × 120 = 3840 行,在内存里做归并排序,最后取第 101 到 120 行。
看这个 OFFSET 的变化就明白问题在哪:
| SQL | 各分片取多少行 | 内存中归并 | 实测耗时 |
|---|---|---|---|
LIMIT 20 OFFSET 0 | 20 | 640 行 | 48 ms |
LIMIT 20 OFFSET 100 | 120 | 3840 行 | 92 ms |
LIMIT 20 OFFSET 10000 | 10020 | 320,640 行 | 4.2 s |
LIMIT 20 OFFSET 100000 | 100020 | 3,200,640 行 | 38 s |
深分页是灾难性的。而且这还只是改写阶段,每个分片自己扫 OFFSET + LIMIT 行也很慢。
我们的解决方案有三层:
- 运营后台的查询不走 MySQL,走 Elasticsearch。把订单数据同步到 ES(7.10),运营的所有多维查询、导出都查 ES。ES 的
search_after天然适合深分页。这一层解决了 90% 的问题。 - 必须查库的场景,禁止深分页。改成"上一页最大 ID + 下一页"的游标方式:
WHERE create_time < ? AND id < ? ORDER BY create_time DESC LIMIT 20。虽然仍要路由 32 个分片,但每个分片只取 20 行,归并量恒定在 640 行。 - 超过 3 个月的订单查归档库。我们在 ClickHouse 里存了一份全量订单,历史数据查询走那里。
跨库聚合:能算但别依赖
ShardingSphere 支持 COUNT、SUM、MAX、AVG 的归并,AVG 会被改写成 SUM 和 COUNT 分别下发再相除。GROUP BY 也能做,但 5.0.0-beta 下有些复杂场景(比如 GROUP BY 后跟 HAVING 再套子查询)支持不完整。
我们的做法:运营看板一律走预聚合。每天凌晨跑一次汇总,把"各状态订单数"、"各商家销售额"算好存到 t_order_daily_summary 表里。看板查这张表,毫秒级。
CREATE TABLE t_order_daily_summary (
stat_date DATE NOT NULL,
merchant_id BIGINT NOT NULL,
status TINYINT NOT NULL,
order_cnt INT NOT NULL DEFAULT 0,
amount_sum DECIMAL(18,2) NOT NULL DEFAULT 0,
PRIMARY KEY (stat_date, merchant_id, status)
) ENGINE=InnoDB;
这个表是单表,不分片,数据量每天几百行。
踩的坑
PageHelper 的分页插件顺序
我们用了 PageHelper 做分页。ShardingSphere 的数据源必须包在 MyBatis 的 SqlSessionFactory 里,而 PageHelper 的拦截器是按数据库方言改写 SQL 的。ShardingSphere 必须作为 DataSource 层生效,PageHelper 生成 LIMIT 后再由 ShardingSphere 改写,这个顺序是对的。但如果有人直接注入了原生 DruidDataSource,就会绕过分片,查到单库的数据。
我们加了一个启动校验:
@PostConstruct
public void checkDataSource() {
DataSource ds = applicationContext.getBean(DataSource.class);
if (!(ds instanceof ShardingSphereDataSource)) {
throw new IllegalStateException("订单数据源未走 ShardingSphere!");
}
}
不支持的 SQL
压测期间遇到两类报错:
SQLSyntaxErrorException: Can not support DML routing without sharding conditions
-- UPDATE/DELETE 没带分片键,且没有 hint 强制路由
UnsupportedOperationException: DISTINCT and GROUP BY can not be together
-- 5.0.0-beta 的归并引擎限制
第一类我们全部改成了"先按分片键查出来,再按主键更新"。第二类改写成两个查询,或者拆到 ES。
分布式事务
跨库 UPDATE 需要分布式事务。我们评估了 ShardingSphere 的 XA(Atomikos),最后没用。原因是 XA 性能损耗大(我们测的 TPS 掉 40%)且死锁排查困难。业务上我们把跨库操作改成了最终一致性:本地消息表 + 定时任务补偿。
上线效果
方案分两步:先上线分片集群(双写 + 灰度切流,这部分写在另一篇《分库分表后的数据迁移与双写方案》里),3 月底切完读流量。
| 指标 | 分片前 | 分片后 |
|---|---|---|
| C 端订单列表 P99 | 3840 ms | 12 ms |
| 订单详情(按 order_id) | 210 ms | 8 ms |
| 单表数据量 | 8200 万 | 257 万 |
| 表物理大小 | 68 GB | 2.4 GB |
| 加索引耗时 | 6 小时 | 4 分钟 |
先到这
《ShardingSphere 分库分表落地:从设计到上线》这块我前前后后踩了不止一次。今天先写这些,后面想到新的再补。