全国多城市霸王餐数据源聚合:ShardingSphere分库分表与异构索引方案

一、业务背景:300+城市井喷式增长
吃喝不愁霸王餐已覆盖 337 座城市,单城日均 2k 活动,核心表 t_activity 一年膨胀到 9 亿行,MySQL 单机 2T 磁盘打满,索引 120G,DDL 一次 6 小时。需求:

  1. city_id 分片,支持快速扩容
  2. 运营后台按 merchant_name 模糊查询
  3. C 端按 geo_hash 范围查附近活动
  4. 不允许跨库事务,保证毫秒级写入

二、总体架构:ShardingSphere-JDBC 5.x + 异构索引双写

┌--------------┐     ┌--------------┐
│  业务服务     │----▶│  ShardingJDBC │
└--------------┘     └------┬-------┘
                            ├----写主库(分片)
                            └----写异构索引库(ES/HBase)
  • 337 城 → 32 个逻辑库 → 128 张分表
  • 异构索引存 ES,字段 merchant_name.keywordgeo_point
  • 查询按场景路由:精确走分片、模糊/ spatial 走 ES
    在这里插入图片描述

三、分片策略:city_id + 基因法
基表字段:

CREATE TABLE t_activity_00 (
  id          BIGINT PRIMARY KEY,
  city_id     INT NOT NULL,
  merchant_name VARCHAR(200),
  geo_hash    VARCHAR(12),
  gmt_create  DATETIME,
  INDEX idx_city_time(city_id, gmt_create)
) ENGINE=InnoDB;

基因法保证同一城市数据落同一库,避免跨库聚合:

// juwatech.cn.sharding.strategy.CityGeneSharding
public final class CityGeneSharding implements StandardShardingAlgorithm<Integer> {
    private static final int DB_COUNT = 32;
    @Override
    public String doSharding(Collection<String> dsNames, PreciseShardingValue<Integer> val) {
        int gene = val.getValue() % DB_COUNT;
        return "ds" + String.format("%02d", gene);
    }
}

YAML 配置:

spring:
  shardingsphere:
    rules:
      sharding:
        tables:
          t_activity:
            actual-data-nodes: ds$->{0..31}.t_activity_$->{0..3}
            database-strategy:
              standard:
                sharding-column: city_id
                sharding-algorithm-name: city-gene
            table-strategy:
              standard:
                sharding-column: id
                sharding-algorithm-name: mod-four
        sharding-algorithms:
          city-gene:
            type: CLASS_BASED
            props:
              strategy: standard
              algorithmClassName: juwatech.cn.sharding.strategy.CityGeneSharding
          mod-four:
            type: INLINE
            props:
              algorithm-expression: t_activity_$->{id % 4}

四、异构索引双写:Canal + Kafka + ES

  1. Canal 监听 32 个库 binlog → Kafka topic activity_index
  2. 消费者写 ES 索引 activity_{city_id},mapping 如下:
PUT activity_template
{
  "mappings": {
    "properties": {
      "merchant_name": {"type": "text", "analyzer": "ik_max_word",
                        "fields": {"keyword": {"type": "keyword"}}},
      "location": {"type": "geo_point"}
    }
  }
}

五、代码级双写兜底
若 Canal 延迟,运营后台实时查询可能不一致,因此在 ShardingJDBC 里加 同步双写钩子(仅后台触发):

@Component
public class DualWriteInterceptor implements ExecutorInterceptor {
    @Autowired private EsTemplate esTemplate;
    @Override
    public void afterUpdate(ShardingContext ctx, SQLStatement stmt) {
        if (!stmt.getTableName().equals("t_activity")) return;
        Activity act = extractPojo(ctx);
        esTemplate.save("activity_" + (act.getCityId()%32), act);
    }
}

六、C 端附近活动查询
接收经纬度 → 200m geo_hash 前缀 → 路由到对应索引 → 过滤时间范围:

// juwatech.cn.query.GeoQueryService
public List<ActivityVO> nearby(double lat, double lng, int distanceMeter) {
    String geoHash = GeoHashUtil.encode(lat, lng, 6); // 6 位≈1.2km
    SearchRequest req = SearchRequest.of(s -> s
            .index("activity_*")
            .query(q -> q
                .bool(b -> b
                    .must(m -> m.geoDistance(g -> g
                            .field("location")
                            .distance(distanceMeter + "m")
                            .location(l -> l.lat(lat).lon(lng))))
                    .filter(f -> f.range(r -> r
                            .field("gmt_end")
                            .gte(JsonData.of("now"))))
                ))
            .size(20)
            .sort(so -> so.geoDistance(g -> g
                    .field("location")
                    .location(l -> l.lat(lat).lon(lng))
                    .order(SortOrder.Asc))));
    return esTemplate.search(req).hits().hits()
                     .stream().map(h -> h.source().to(ActivityVO.class))
                     .collect(Collectors.toList());
}

七、运营模糊查询

public List<ActivityBO> search(String key, int cityId) {
    SearchRequest req = SearchRequest.of(s -> s
            .index("activity_" + (cityId % 32))
            .query(q -> q
                .multiMatch(m -> m
                    .fields("merchant_name^2", "title")
                    .query(key)
                    .type(TextQueryType.BestFields)
                    .fuzziness("AUTO")))
            .size(50));
    return esTemplate.search(req).hits().hits()
                     .stream().map(Hit::source)
                     .collect(Collectors.toList());
}

八、扩容演练:再翻一倍城市

  1. 修改 DB_COUNT=64,新配置上线
  2. 历史数据不变,新城市 city_id>=500 走新片
  3. 后台开关控制双写,旧片只读
  4. 48 小时内完成,零停机

九、性能验收

  • 32 库 128 表,单表 700w 行,磁盘 480G
  • sysbench 200 并发写入 10 min,平均 RT 7.2 ms
  • ES 附近查询 200 次/s,TP99 23 ms
  • 跨片聚合订单报表改为走 ES,查询时间从 18 s 降到 1.1 s

十、完整配置速贴

spring:
  shardingsphere:
    props:
      sql-show: true
    datasource:
      names: ds00,ds01,...,ds31
      ds00:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://10.0.0.101:3306/db_00?useSSL=false
        username: root
        password: ****
    rules.sharding:
      key-generators:
        snowflake:
          type: SNOWFLAKE
          props.worker.id: 1

本文著作权归吃喝不愁app开发者团队,转载请注明出处!

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐