Administrator
发布于 2021-02-28 / 2388 阅读
54

ShardingSphere 分库分表落地:从设计到上线

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 INORDER BY 让优化器放弃了一部分索引。这种 SQL 在 4000 万行时不算问题,到 8200 万就崩了。

为什么选 ShardingSphere 5.0.0-beta

选型时对比了三个:

方案形态我们的顾虑
ShardingSphere-JDBC 5.0.0-beta客户端 JARbeta 版,但社区活跃,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}

三个配置值得单独说:

  • bindingTablest_ordert_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 020640 行48 ms
LIMIT 20 OFFSET 1001203840 行92 ms
LIMIT 20 OFFSET 1000010020320,640 行4.2 s
LIMIT 20 OFFSET 1000001000203,200,640 行38 s

深分页是灾难性的。而且这还只是改写阶段,每个分片自己扫 OFFSET + LIMIT 行也很慢。

我们的解决方案有三层:

  1. 运营后台的查询不走 MySQL,走 Elasticsearch。把订单数据同步到 ES(7.10),运营的所有多维查询、导出都查 ES。ES 的 search_after 天然适合深分页。这一层解决了 90% 的问题。
  2. 必须查库的场景,禁止深分页。改成"上一页最大 ID + 下一页"的游标方式:WHERE create_time < ? AND id < ? ORDER BY create_time DESC LIMIT 20。虽然仍要路由 32 个分片,但每个分片只取 20 行,归并量恒定在 640 行。
  3. 超过 3 个月的订单查归档库。我们在 ClickHouse 里存了一份全量订单,历史数据查询走那里。

跨库聚合:能算但别依赖

ShardingSphere 支持 COUNTSUMMAXAVG 的归并,AVG 会被改写成 SUMCOUNT 分别下发再相除。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 端订单列表 P993840 ms12 ms
订单详情(按 order_id)210 ms8 ms
单表数据量8200 万257 万
表物理大小68 GB2.4 GB
加索引耗时6 小时4 分钟

先到这

《ShardingSphere 分库分表落地:从设计到上线》这块我前前后后踩了不止一次。今天先写这些,后面想到新的再补。

参考