Spring Boot 3 + MyBatis-Plusマルチデータソースアーキテクチャ完全ガイド

后端开发

なぜマルチデータソースはエンタープライズの必須要件なのか

シングルデータソースアーキテクチャはインターネット初期には十分でしたが、ビジネス規模が拡大するにつれ、1つのデータベースですべてのトラフィックと複雑さに対応できなくなります。マルチデータソースは「锦上添花」ではなく、「やらざるを得ない」選択です。

5つのコアシナリオ

シナリオ 典型的な要件 データソース関係
リードライト分離 読み多い書き少ない、読み負荷分散 1プライマリ+複数レプリカ EC商品詳細ページ
ビジネス分割 異なるビジネスドメインの物理分離 並行独立 注文DB vs ユーザーDB
シャーディング 単一テーブルのデータ量過大 水平シャーディング ログテーブル月次分割
マルチテナント 異なるテナントのデータ分離 1テナント1DB/1Schema SaaSプラットフォーム
異種データソース 異なるタイプのストレージ連携 MySQL + PG + ES 検索+トランザクション混在
┌──────────────────────────────────────────────────────────┐
│            シングルデータソースアーキテクチャ(初期)          │
│                                                          │
│    App ──────► MySQL (読み書きすべてここ)                   │
│                                                          │
│    問題:読み書き競合、単一テーブル肥大化、水平スケール不可     │
└──────────────────────────────────────────────────────────┘

┌──────────────────────────────────────────────────────────┐
│            マルチデータソースアーキテクチャ(進化)            │
│                                                          │
│    App ──┬──► MySQL-Master (書き込み)                     │
│           ├──► MySQL-Slave1 (読み込み)                    │
│           ├──► MySQL-Slave2 (読み込み)                    │
│           ├──► OrderDB (注文データベース)                   │
│           ├──► UserDB (ユーザーデータベース)               │
│           └──► Elasticsearch (検索)                       │
└──────────────────────────────────────────────────────────┘

実例: あるソーシャルプラットフォーム(DAU 500万)は、単一DBピークQPS 12万でした。リードライト分離+ビジネス分割後、プライマリDB QPSは3万に低下、全体スループットは4倍向上。


マルチデータソースアーキテクチャの3つのパターン

パターン比較概要

次元 アプリ内マルチデータソース JDBCレイヤープロキシ プロキシ層独立デプロイ
代表ソリューション dynamic-datasource ShardingSphere-JDBC ShardingSphere-Proxy / MyCat
侵入性 低(アノテーション/設定) 中(DataSource置き換え) なし(独立プロセス)
運用複雑さ
パフォーマンスオーバーヘッド 極低 中(ネットワークホップ)
機能豊富さ 基本切替 シャーディング+RW分離+暗号化 フル機能
適用規模 小中プロジェクト 中大プロジェクト 大規模/多言語プロジェクト
言語バインディング Java Java 言語非依存

パターン1:アプリ内マルチデータソース

┌─────────────────────────────────────────────┐
│              Spring Boot Application          │
│                                              │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐  │
│  │ DataSource│  │ DataSource│  │ DataSource│  │
│  │  master   │  │  slave_1  │  │  slave_2  │  │
│  └─────┬────┘  └─────┬────┘  └─────┬────┘  │
│        │             │             │        │
│  ┌─────▼─────────────▼─────────────▼────┐   │
│  │     DynamicRoutingDataSource          │   │
│  │  (ThreadLocal + @DS アノテーションルート) │   │
│  └──────────────────────────────────────┘   │
└─────────────────────────────────────────────┘

パターン2:JDBCレイヤープロキシ

┌─────────────────────────────────────────────┐
│              Spring Boot Application          │
│                                              │
│  ┌──────────────────────────────────────┐   │
│  │     ShardingSphere-JDBC (JAR)         │   │
│  │  ┌────────────────────────────────┐  │   │
│  │  │  SQL解析 → ルーティング → 書換   │  │   │
│  │  │  → 実行 → 結果マージ            │  │   │
│  │  └────────────────────────────────┘  │   │
│  └──────┬──────────┬──────────┬─────────┘   │
│         │          │          │              │
│    ds_0 │     ds_1 │     ds_2 │              │
└─────────┼──────────┼──────────┼──────────────┘
          ▼          ▼          ▼
      MySQL-0    MySQL-1    MySQL-2

パターン3:プロキシ層独立デプロイ

┌────────────┐     ┌──────────────────┐     ┌─────────┐
│  App (Java) │────►│                  │     │ MySQL-0 │
├────────────┤     │  ShardingSphere  │────►├─────────┤
│  App (Go)  │────►│  -Proxy (独立)    │────►│ MySQL-1 │
├────────────┤     │                  │     ├─────────┤
│  App (Py)  │────►│  アプリに透過      │────►│ MySQL-2 │
└────────────┘     └──────────────────┘     └─────────┘

dynamic-datasource-spring-boot-starter実践

これは現在最受欢迎のアプリ内マルチデータソースソリューションで、Baomidouチームが保守し、MyBatis-Plusと深く統合されています。

Maven依存関係

<dependencies>
    <dependency>
        <groupId>com.baomidou</groupId>
        <artifactId>dynamic-datasource-spring-boot3-starter</artifactId>
        <version>4.3.1</version>
    </dependency>
    <dependency>
        <groupId>com.baomidou</groupId>
        <artifactId>mybatis-plus-spring-boot3-starter</artifactId>
        <version>3.5.7</version>
    </dependency>
    <dependency>
        <groupId>com.mysql</groupId>
        <artifactId>mysql-connector-j</artifactId>
        <scope>runtime</scope>
    </dependency>
</dependencies>

基本設定

spring:
  datasource:
    dynamic:
      primary: master
      strict: true
      datasource:
        master:
          url: jdbc:mysql://localhost:3306/db_master?useSSL=false&serverTimezone=Asia/Shanghai
          username: root
          password: master_pwd
          driver-class-name: com.mysql.cj.jdbc.Driver
        slave_1:
          url: jdbc:mysql://localhost:3307/db_slave?useSSL=false&serverTimezone=Asia/Shanghai
          username: root
          password: slave_pwd
          driver-class-name: com.mysql.cj.jdbc.Driver
        slave_2:
          url: jdbc:mysql://localhost:3308/db_slave?useSSL=false&serverTimezone=Asia/Shanghai
          username: root
          password: slave_pwd
          driver-class-name: com.mysql.cj.jdbc.Driver
      hikari:
        minimum-idle: 5
        maximum-pool-size: 20
        idle-timeout: 30000
        max-lifetime: 1800000
        connection-timeout: 30000
        pool-name: DynamicHikariCP

アノテーション切替 @DS

@Service
public class UserServiceImpl implements UserService {

    @Autowired
    private UserMapper userMapper;

    @DS("master")
    @Override
    public void createUser(User user) {
        userMapper.insert(user);
    }

    @DS("slave_1")
    @Override
    public User getUserById(Long id) {
        return userMapper.selectById(id);
    }

    @DS("slave_2")
    @Override
    public List<User> listUsers(int pageNum, int pageSize) {
        Page<User> page = new Page<>(pageNum, pageSize);
        return userMapper.selectPage(page, null).getRecords();
    }
}

プログラマティック切替

アノテーション方式は複雑なシナリオで柔軟性に欠けます。プログラマティック切替はランタイムでデータソースを動的に決定できます:

@Service
public class OrderServiceImpl implements OrderService {

    @Autowired
    private OrderMapper orderMapper;

    @Override
    public Order getOrderWithDynamicDs(String tenantId) {
        String dsKey = "tenant_" + tenantId;
        return DynamicDataSourceContextHolder.push(dsKey, () -> {
            return orderMapper.selectById(tenantId);
        });
    }

    @Override
    public List<Order> batchQuery(List<String> tenantIds) {
        List<Order> result = new ArrayList<>();
        for (String tenantId : tenantIds) {
            DynamicDataSourceContextHolder.push("tenant_" + tenantId, () -> {
                result.addAll(orderMapper.selectList(null));
            });
        }
        return result;
    }
}

カスタム動的データソースルーティング戦略

@Component
public class CustomDsProcessor extends DsProcessor {

    @Autowired
    private TenantDataSourceManager tenantDsManager;

    @Override
    public String determineDsKey(MethodInvocation invocation, String key) {
        if (key.startsWith("tenant#")) {
            String tenantId = extractTenantId(invocation);
            return tenantDsManager.resolveDsKey(tenantId);
        }
        return key;
    }

    private String extractTenantId(MethodInvocation invocation) {
        Object[] args = invocation.getArguments();
        for (Object arg : args) {
            if (arg instanceof TenantAware tenantAware) {
                return tenantAware.getTenantId();
            }
        }
        throw new IllegalStateException("テナントIDを抽出できません");
    }
}

AOPによるリードライト分離自動ルーティング

@Aspect
@Component
@Order(Ordered.HIGHEST_PRECEDENCE)
public class ReadWriteSplitAspect {

    private static final Set<String> READ_METHOD_PREFIXES =
        Set.of("get", "query", "find", "list", "count", "page", "search");

    @Around("execution(* com.example..service.impl.*ServiceImpl.*(..))")
    public Object around(ProceedingJoinPoint point) throws Throwable {
        String methodName = point.getSignature().getName();
        boolean isRead = READ_METHOD_PREFIXES.stream()
            .anyMatch(methodName::startsWith);

        if (isRead) {
            DynamicDataSourceContextHolder.push("slave_1");
        } else {
            DynamicDataSourceContextHolder.push("master");
        }

        try {
            return point.proceed();
        } finally {
            DynamicDataSourceContextHolder.poll();
        }
    }
}

MyBatis-Plusマルチデータソース設定

クロスデータソース結合クエリ

MyBatis-PlusはクロスデータソースJOINをネイティブサポートしていません。アプリケーション層で組み立てる必要があります:

@Service
public class OrderDetailServiceImpl implements OrderDetailService {

    @Autowired
    private OrderMapper orderMapper;

    @Autowired
    private UserMapper userMapper;

    @Autowired
    private ProductMapper productMapper;

    @DS("order_db")
    public OrderDetailVO getOrderDetail(Long orderId) {
        Order order = orderMapper.selectById(orderId);

        User user = DynamicDataSourceContextHolder.push("user_db", () -> {
            return userMapper.selectById(order.getUserId());
        });

        Product product = DynamicDataSourceContextHolder.push("product_db", () -> {
            return productMapper.selectById(order.getProductId());
        });

        return OrderDetailVO.builder()
                .order(order)
                .user(user)
                .product(product)
                .build();
    }
}

クロスデータソースバッチクエリ最適化

@Service
public class BatchCrossDsService {

    @Autowired
    private OrderMapper orderMapper;

    @Autowired
    private UserMapper userMapper;

    public List<OrderWithUserVO> batchQuery(List<Long> orderIds) {
        List<Order> orders = DynamicDataSourceContextHolder.push("order_db", () -> {
            return orderMapper.selectBatchIds(orderIds);
        });

        List<Long> userIds = orders.stream()
                .map(Order::getUserId)
                .distinct()
                .collect(Collectors.toList());

        Map<Long, User> userMap = DynamicDataSourceContextHolder.push("user_db", () -> {
            return userMapper.selectBatchIds(userIds).stream()
                    .collect(Collectors.toMap(User::getId, Function.identity()));
        });

        return orders.stream().map(order -> {
            OrderWithUserVO vo = new OrderWithUserVO();
            vo.setOrder(order);
            vo.setUser(userMap.get(order.getUserId()));
            return vo;
        }).collect(Collectors.toList());
    }
}

動的データソース追加・削除

SaaSシナリオでは、テナントの動的登録にランタイムデータソースプロビジョニングが必要です:

@Component
public class TenantDataSourceManager {

    @Autowired
    private DynamicRoutingDataSource dynamicRoutingDataSource;

    @Autowired
    private DataSourceProperty defaultProperty;

    private final Map<String, DataSourceProperty> tenantDsMap = new ConcurrentHashMap<>();

    public synchronized void addTenantDs(String tenantId, String url, String username, String password) {
        String dsKey = "tenant_" + tenantId;

        DataSourceProperty property = new DataSourceProperty();
        property.setUrl(url);
        property.setUsername(username);
        property.setPassword(password);
        property.setDriverClassName(defaultProperty.getDriverClassName());
        property.setLazy(true);

        HikariDataSource dataSource = property.getDataSource()
                .orElseThrow(() -> new RuntimeException("データソースの作成に失敗しました"));
        dataSource.setMaximumPoolSize(10);
        dataSource.setMinimumIdle(2);

        dynamicRoutingDataSource.addDataSource(dsKey, dataSource);
        tenantDsMap.put(dsKey, property);

        log.info("テナント[{}]データソース追加: {}", tenantId, dsKey);
    }

    public synchronized void removeTenantDs(String tenantId) {
        String dsKey = "tenant_" + tenantId;
        dynamicRoutingDataSource.removeDataSource(dsKey);
        tenantDsMap.remove(dsKey);
        log.info("テナント[{}]データソース削除: {}", tenantId, dsKey);
    }

    public String resolveDsKey(String tenantId) {
        String dsKey = "tenant_" + tenantId;
        if (!tenantDsMap.containsKey(dsKey)) {
            throw new IllegalStateException("テナントデータソースが存在しません: " + dsKey);
        }
        return dsKey;
    }

    public Set<String> getAllTenantDsKeys() {
        return Collections.unmodifiableSet(tenantDsMap.keySet());
    }
}

テナントデータソース自動登録リスナー

@Component
@Slf4j
public class TenantDataSourceInitializer implements ApplicationRunner {

    @Autowired
    private TenantDataSourceManager tenantDsManager;

    @Autowired
    private TenantConfigRepository tenantConfigRepository;

    @Override
    public void run(ApplicationArguments args) {
        List<TenantConfig> tenants = tenantConfigRepository.findAllActive();
        for (TenantConfig tenant : tenants) {
            try {
                tenantDsManager.addTenantDs(
                    tenant.getTenantId(),
                    tenant.getDbUrl(),
                    tenant.getDbUsername(),
                    tenant.getDbPassword()
                );
            } catch (Exception e) {
                log.error("テナント[{}]データソース初期化失敗", tenant.getTenantId(), e);
            }
        }
        log.info("{}テナントデータソースを初期化しました", tenants.size());
    }
}

リードライト分離:設定から実装まで

完全なリードライト分離アーキテクチャ

                    ┌─────────────────┐
                    │   Application    │
                    │  @DS("master")   │──書き込み
                    │  @DS("slave")    │──読み込み
                    └────────┬────────┘
                             │
              ┌──────────────┼──────────────┐
              │              │              │
        ┌─────▼─────┐ ┌─────▼─────┐ ┌─────▼─────┐
        │  Master    │ │  Slave-1  │ │  Slave-2  │
        │  (RW)      │ │  (RO)     │ │  (RO)     │
        └─────┬─────┘ └───────────┘ └───────────┘
              │
        ┌─────▼─────┐
        │  Binlog    │
        │  レプリケーション│
        └───────────┘

リードライト分離YAML設定

spring:
  datasource:
    dynamic:
      primary: master
      strict: true
      datasource:
        master:
          url: jdbc:mysql://mysql-master:3306/app_db?useSSL=false&serverTimezone=Asia/Shanghai&characterEncoding=utf8mb4
          username: app_writer
          password: ${MASTER_DB_PWD}
        slave_1:
          url: jdbc:mysql://mysql-slave-1:3306/app_db?useSSL=false&serverTimezone=Asia/Shanghai&characterEncoding=utf8mb4
          username: app_reader
          password: ${SLAVE1_DB_PWD}
        slave_2:
          url: jdbc:mysql://mysql-slave-2:3306/app_db?useSSL=false&serverTimezone=Asia/Shanghai&characterEncoding=utf8mb4
          username: app_reader
          password: ${SLAVE2_DB_PWD}

カスタムロードバランシング戦略

デフォルトのラウンドロビン戦略では柔軟性が不足します。本番環境ではレプリカ遅延に基づく動的ルーティングが必要です:

@Component
public class SlaveLoadBalanceStrategy implements LoadBalanceStrategy {

    @Autowired
    private MySQLReplicationMonitor replicationMonitor;

    @Override
    public String determineDsKey(List<String> slaveDsKeys, String masterDsKey) {
        List<String> availableSlaves = slaveDsKeys.stream()
                .filter(key -> {
                    long delay = replicationMonitor.getReplicationDelay(key);
                    return delay < 1000;
                })
                .collect(Collectors.toList());

        if (availableSlaves.isEmpty()) {
            log.warn("全レプリカの遅延が大きすぎます。プライマリからの読み取りにフォールバック");
            return masterDsKey;
        }

        int index = (int) (System.currentTimeMillis() % availableSlaves.size());
        return availableSlaves.get(index);
    }
}

レプリケーション遅延監視

@Component
@Slf4j
public class MySQLReplicationMonitor {

    private final Map<String, Long> replicationDelayMap = new ConcurrentHashMap<>();

    @Scheduled(fixedDelay = 5000)
    public void monitorReplicationDelay() {
        List<String> slaveKeys = List.of("slave_1", "slave_2");
        for (String slaveKey : slaveKeys) {
            DynamicDataSourceContextHolder.push(slaveKey, () -> {
                try (Connection conn = dataSource.getConnection();
                     Statement stmt = conn.createStatement();
                     ResultSet rs = stmt.executeQuery("SHOW SLAVE STATUS")) {
                    if (rs.next()) {
                        long secondsBehindMaster = rs.getLong("Seconds_Behind_Master");
                        replicationDelayMap.put(slaveKey, secondsBehindMaster * 1000);
                        if (secondsBehindMaster > 5) {
                            log.warn("レプリカ[{}]遅延: {}秒", slaveKey, secondsBehindMaster);
                        }
                    }
                } catch (SQLException e) {
                    replicationDelayMap.put(slaveKey, Long.MAX_VALUE);
                    log.error("レプリカ[{}]ステータスチェック失敗", slaveKey, e);
                }
            });
        }
    }

    public long getReplicationDelay(String slaveKey) {
        return replicationDelayMap.getOrDefault(slaveKey, 0L);
    }
}

プライマリ強制アノテーション

一部のシナリオ(書き込み後の即時読み取りなど)ではプライマリにアクセスする必要があります:

@Target({ElementType.METHOD, ElementType.TYPE})
@Retention(RetentionPolicy.RUNTIME)
@DS("master")
public @interface MasterOnly {
}
@Service
public class PaymentServiceImpl implements PaymentService {

    @Autowired
    private PaymentMapper paymentMapper;

    @DS("master")
    @Override
    public Payment createPayment(PaymentDTO dto) {
        Payment payment = new Payment();
        BeanUtils.copyProperties(dto, payment);
        paymentMapper.insert(payment);
        return payment;
    }

    @MasterOnly
    @Override
    public Payment getPaymentAfterCreate(Long paymentId) {
        return paymentMapper.selectById(paymentId);
    }
}

シャーディング:ShardingSphere-JDBC深層統合

シャーディングアーキテクチャ

┌──────────────────────────────────────────────────────┐
│                   Application Layer                    │
│                                                       │
│  ┌─────────────────────────────────────────────────┐ │
│  │           ShardingSphere-JDBC                    │ │
│  │                                                  │ │
│  │  SQL: SELECT * FROM t_order WHERE user_id=1001  │ │
│  │         │                                       │ │
│  │    ┌────▼────┐                                  │ │
│  │    │ SQL解析  │                                  │ │
│  │    └────┬────┘                                  │ │
│  │    ┌────▼────┐                                  │ │
│  │    │ ルーティング│──► ds_0.t_order_1             │ │
│  │    └────┬────┘                                  │ │
│  │    ┌────▼────┐                                  │ │
│  │    │ SQL書換  │──► SELECT * FROM t_order_1 ...  │ │
│  │    └────┬────┘                                  │ │
│  │    ┌────▼────┐                                  │ │
│  │    │ 結果マージ│                                  │ │
│  │    └─────────┘                                  │ │
│  └─────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────┘

Maven依存関係

<dependencies>
    <dependency>
        <groupId>org.apache.shardingsphere</groupId>
        <artifactId>shardingsphere-jdbc-core</artifactId>
        <version>5.5.1</version>
    </dependency>
    <dependency>
        <groupId>org.apache.shardingsphere</groupId>
        <artifactId>shardingsphere-jdbc-core-spring-boot-starter</artifactId>
        <version>5.5.1</version>
    </dependency>
</dependencies>

シャーディング設定

spring:
  shardingsphere:
    mode:
      type: Standalone
      repository:
        type: JDBC
    datasource:
      names: ds_0,ds_1
      ds_0:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://localhost:3306/shard_db_0?useSSL=false&serverTimezone=Asia/Shanghai
        username: root
        password: ds0_pwd
      ds_1:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://localhost:3306/shard_db_1?useSSL=false&serverTimezone=Asia/Shanghai
        username: root
        password: ds1_pwd
    rules:
      sharding:
        tables:
          t_order:
            actual-data-nodes: ds_${0..1}.t_order_${0..15}
            table-strategy:
              standard:
                sharding-column: order_id
                sharding-algorithm-name: order-table-inline
            database-strategy:
              standard:
                sharding-column: user_id
                sharding-algorithm-name: order-db-mod
            key-generate-strategy:
              column: order_id
              key-generator-name: snowflake
          t_order_item:
            actual-data-nodes: ds_${0..1}.t_order_item_${0..15}
            table-strategy:
              standard:
                sharding-column: order_id
                sharding-algorithm-name: order-item-table-inline
            database-strategy:
              standard:
                sharding-column: user_id
                sharding-algorithm-name: order-db-mod
        sharding-algorithms:
          order-table-inline:
            type: INLINE
            props:
              algorithm-expression: t_order_${order_id % 16}
          order-item-table-inline:
            type: INLINE
            props:
              algorithm-expression: t_order_item_${order_id % 16}
          order-db-mod:
            type: MOD
            props:
              sharding-count: 2
        key-generators:
          snowflake:
            type: SNOWFLAKE
            props:
              worker-id: 1
    props:
      sql-show: true

カスタムシャーディングアルゴリズム

標準の剰余アルゴリズムはスケールアウト時にデータ移行が必要です。コンシステントハッシングは移行量を最小化します:

@Component
public class ConsistentHashShardingAlgorithm implements StandardShardingAlgorithm<Long> {

    private ConsistentHash<String> consistentHash;

    @Override
    public void init(Properties props) {
        int virtualNodes = Integer.parseInt(props.getProperty("virtual-nodes", "160"));
        List<String> actualDataNodes = Arrays.stream(props.getProperty("actual-data-nodes").split(","))
                .collect(Collectors.toList());

        HashFunction hashFunction = new Murmur3HashFunction();
        consistentHash = new ConsistentHash<>(hashFunction, virtualNodes, actualDataNodes);
    }

    @Override
    public String doSharding(Collection<String> availableTargetNames, PreciseShardingValue<Long> shardingValue) {
        return consistentHash.get(shardingValue.getValue());
    }

    @Override
    public Collection<String> doSharding(Collection<String> availableTargetNames, RangeShardingValue<Long> shardingValue) {
        return availableTargetNames;
    }

    @Override
    public String getType() {
        return "CONSISTENT_HASH";
    }
}

MyBatis-Plus + ShardingSphere統合の注意点

@Configuration
@MapperScan("com.example.mapper")
public class MybatisPlusConfig {

    @Bean
    public MybatisPlusInterceptor mybatisPlusInterceptor() {
        MybatisPlusInterceptor interceptor = new MybatisPlusInterceptor();
        interceptor.addInnerInterceptor(new PaginationInnerInterceptor(DbType.MYSQL));
        return interceptor;
    }

    @Bean
    public IdentifierGenerator identifierGenerator() {
        return new CustomIdGenerator();
    }
}

public class CustomIdGenerator implements IdentifierGenerator {
    @Override
    public Number nextId(Object entity) {
        return ShardingSphereIdGenerator.nextId();
    }
}

マルチテナントデータ分離:SchemaレベルとTableレベル

分離モード比較

次元 Schemaレベル分離 Tableレベル分離
分離レベル 高い 低い
リソースオーバーヘッド テナントごとに独立Schema 共有Schema、テーブル名で区別
運用複雑さ
データ漏洩リスク 中(厳格なフィルタリング必要)
適用シナリオ 大口顧客/金融 小規模顧客/SaaS
拡張性 Schema数に制限 テーブル数に制限

Schemaレベル分離の実装

@Component
public class SchemaTenantInterceptor implements InnerInterceptor {

    @Override
    public void beforeQuery(Executor executor, MappedStatement ms, Object parameter, RowBounds rowBounds, ResultHandler resultHandler, BoundSql boundSql) {
        String tenantId = TenantContextHolder.getTenantId();
        if (tenantId != null) {
            DynamicDataSourceContextHolder.push("tenant_" + tenantId);
        }
    }
}

Tableレベル分離(MyBatis-Plusテナントプラグイン)

@Configuration
public class TenantConfig {

    @Bean
    public MybatisPlusInterceptor mybatisPlusInterceptor(TenantLineInnerInterceptor tenantInterceptor) {
        MybatisPlusInterceptor interceptor = new MybatisPlusInterceptor();
        interceptor.addInnerInterceptor(tenantInterceptor);
        interceptor.addInnerInterceptor(new PaginationInnerInterceptor(DbType.MYSQL));
        return interceptor;
    }

    @Bean
    public TenantLineInnerInterceptor tenantLineInnerInterceptor() {
        return new TenantLineInnerInterceptor(new TenantLineHandler() {
            @Override
            public Expression getTenantId() {
                Long tenantId = TenantContextHolder.getTenantId();
                if (tenantId == null) {
                    throw new IllegalStateException("テナントIDはnullにできません");
                }
                return new LongValue(tenantId);
            }

            @Override
            public boolean ignoreTable(String tableName) {
                return Set.of("sys_config", "sys_dict", "sys_menu")
                        .contains(tableName);
            }

            @Override
            public String getTenantIdColumn() {
                return "tenant_id";
            }
        });
    }
}

テナントコンテキスト管理

public class TenantContextHolder {

    private static final ThreadLocal<Long> TENANT_ID = new ThreadLocal<>();

    public static void setTenantId(Long tenantId) {
        TENANT_ID.set(tenantId);
    }

    public static Long getTenantId() {
        return TENANT_ID.get();
    }

    public static void clear() {
        TENANT_ID.remove();
    }
}

Web層テナントID注入

@Component
@WebFilter(urlPatterns = "/*")
public class TenantFilter implements Filter {

    @Override
    public void doFilter(ServletRequest request, ServletResponse response, FilterChain chain)
            throws IOException, ServletException {
        HttpServletRequest httpRequest = (HttpServletRequest) request;
        String tenantHeader = httpRequest.getHeader("X-Tenant-Id");

        if (tenantHeader != null) {
            try {
                TenantContextHolder.setTenantId(Long.parseLong(tenantHeader));
            } catch (NumberFormatException e) {
                throw new ServletException("無効なテナントID: " + tenantHeader);
            }
        }

        try {
            chain.doFilter(request, response);
        } finally {
            TenantContextHolder.clear();
        }
    }
}

トランザクション管理:ローカル/Seata分散/MQ結果整合性

トランザクション方式比較

方式 一貫性 パフォーマンス 複雑さ 適用シナリオ
ローカル(@DSTransactional) 強一貫性 単一DSまたは短い不一致を許容
Seata ATモード 強一貫性 クロスDB強一貫性要件
Seata TCCモード 強一貫性 中低 非常に高い 資金/在庫などコア業務
MQ結果整合性 結果整合性 非同期シナリオ/遅延許容
Sagaパターン 結果整合性 長時間プロセスオーケストレーション

ローカルトランザクション:@DSTransactional

dynamic-datasourceは@DSTransactionalを提供し、データソース切替時にトランザクションを維持します:

@Service
public class OrderServiceImpl implements OrderService {

    @Autowired
    private OrderMapper orderMapper;

    @Autowired
    private OrderItemMapper orderItemMapper;

    @DSTransactional
    @DS("master")
    @Override
    public void createOrder(OrderCreateDTO dto) {
        Order order = new Order();
        order.setOrderNo(generateOrderNo());
        order.setUserId(dto.getUserId());
        order.setStatus(OrderStatus.CREATED);
        orderMapper.insert(order);

        for (OrderItemDTO itemDTO : dto.getItems()) {
            OrderItem item = new OrderItem();
            item.setOrderId(order.getId());
            item.setProductId(itemDTO.getProductId());
            item.setQuantity(itemDTO.getQuantity());
            orderItemMapper.insert(item);
        }
    }
}

Seata ATモード統合

┌──────────┐    ┌──────────┐    ┌──────────┐
│  Service A│───►│  Service B│───►│  Service C│
│  (注文DB) │    │  (在庫DB) │    │  (口座DB) │
└─────┬────┘    └─────┬────┘    └─────┬────┘
      │               │               │
      └───────────────┼───────────────┘
                      │
              ┌───────▼───────┐
              │  Seata Server  │
              │  (TCコーディネータ)│
              └───────────────┘
<dependency>
    <groupId>io.seata</groupId>
    <artifactId>seata-spring-boot-starter</artifactId>
    <version>2.2.0</version>
</dependency>
<dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-seata</artifactId>
    <version>2023.0.3.2</version>
</dependency>
seata:
  enabled: true
  application-id: order-service
  tx-service-group: my_test_tx_group
  service:
    vgroup-mapping:
      my_test_tx_group: default
  registry:
    type: nacos
    nacos:
      server-addr: localhost:8848
      namespace: seata
      group: SEATA_GROUP
  config:
    type: nacos
    nacos:
      server-addr: localhost:8848
      namespace: seata
      group: SEATA_GROUP
@Service
public class BusinessServiceImpl implements BusinessService {

    @Autowired
    private OrderService orderService;

    @Autowired
    private StockService stockService;

    @Autowired
    private AccountService accountService;

    @GlobalTransactional(name = "create-order-tx", rollbackFor = Exception.class)
    @Override
    public void purchase(PurchaseDTO dto) {
        orderService.createOrder(dto);
        stockService.deductStock(dto.getProductId(), dto.getQuantity());
        accountService.deductBalance(dto.getUserId(), dto.getAmount());
    }
}

MQ結果整合性

@Service
@Slf4j
public class OrderMQService {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Autowired
    private OrderMapper orderMapper;

    @Autowired
    private TransactionLogMapper transactionLogMapper;

    @Transactional
    public void createOrderWithMQ(OrderCreateDTO dto) {
        Order order = new Order();
        BeanUtils.copyProperties(dto, order);
        order.setStatus(OrderStatus.CREATED);
        orderMapper.insert(order);

        TransactionLog log = new TransactionLog();
        log.setTransactionId(UUID.randomUUID().toString());
        log.setBizType("CREATE_ORDER");
        log.setBizId(order.getId());
        log.setPayload(JSON.toJSONString(dto));
        transactionLogMapper.insert(log);
    }

    @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
    public void onOrderCreated(OrderCreatedEvent event) {
        rocketMQTemplate.asyncSend(
            "order-create-topic",
            MessageBuilder.withPayload(event)
                .setHeader("KEYS", event.getOrderNo())
                .build(),
            new SendCallback() {
                @Override
                public void onSuccess(SendResult sendResult) {
                    log.info("注文メッセージ送信成功: {}", event.getOrderNo());
                }

                @Override
                public void onException(Throwable e) {
                    log.error("注文メッセージ送信失敗: {}", event.getOrderNo(), e);
                }
            }
        );
    }
}

メッセージ補償スケジュールタスク

@Component
@Slf4j
public class TransactionLogCompensator {

    @Autowired
    private TransactionLogMapper transactionLogMapper;

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Scheduled(fixedDelay = 30000)
    public void compensate() {
        List<TransactionLog> pendingLogs = transactionLogMapper.selectList(
            new LambdaQueryWrapper<TransactionLog>()
                .eq(TransactionLog::getStatus, 0)
                .lt(TransactionLog::getRetryCount, 5)
                .le(TransactionLog::getCreateTime, LocalDateTime.now().minusMinutes(1))
        );

        for (TransactionLog log : pendingLogs) {
            try {
                rocketMQTemplate.syncSend("order-create-topic", log.getPayload());
                log.setStatus(1);
                transactionLogMapper.updateById(log);
            } catch (Exception e) {
                log.setRetryCount(log.getRetryCount() + 1);
                transactionLogMapper.updateById(log);
                log.warn("トランザクションログ[{}]補償失敗、リトライ回数: {}", log.getTransactionId(), log.getRetryCount());
            }
        }
    }
}

監視と運用:データソースヘルスチェック+Micrometerメトリクス収集

データソースヘルスチェック

@Component
@Slf4j
public class DataSourceHealthChecker {

    @Autowired
    private DynamicRoutingDataSource dynamicRoutingDataSource;

    @Scheduled(fixedDelay = 10000)
    public void checkAllDataSources() {
        Map<String, DataSource> dataSources = dynamicRoutingDataSource.getCurrentDataSources();

        for (Map.Entry<String, DataSource> entry : dataSources.entrySet()) {
            String dsKey = entry.getKey();
            try (Connection conn = entry.getValue().getConnection();
                 Statement stmt = conn.createStatement();
                 ResultSet rs = stmt.executeQuery("SELECT 1")) {

                if (rs.next() && rs.getInt(1) == 1) {
                    log.debug("データソース[{}]ヘルスチェック通過", dsKey);
                }
            } catch (SQLException e) {
                log.error("データソース[{}]ヘルスチェック失敗", dsKey, e);
                alertService.sendAlert("データソース異常: " + dsKey, e.getMessage());
            }
        }
    }
}

Micrometerメトリクス収集

@Configuration
public class DataSourceMetricsConfig {

    @Bean
    public MeterDataSourceAspect meterDataSourceAspect(MeterRegistry registry) {
        return new MeterDataSourceAspect(registry);
    }

    public static class MeterDataSourceAspect {

        private final MeterRegistry registry;
        private final Counter dsSwitchCounter;
        private final Timer dsQueryTimer;

        public MeterDataSourceAspect(MeterRegistry registry) {
            this.registry = registry;
            this.dsSwitchCounter = Counter.builder("datasource.switch.count")
                    .description("データソース切替回数")
                    .register(registry);
            this.dsQueryTimer = Timer.builder("datasource.query.duration")
                    .description("データソースクエリ所要時間")
                    .tag("type", "dynamic")
                    .register(registry);
        }

        public void recordSwitch(String dsKey) {
            dsSwitchCounter.increment();
            registry.counter("datasource.switch.count", "ds", dsKey).increment();
        }

        public Timer.Sample startQuery() {
            return Timer.start(registry);
        }

        public void endQuery(Timer.Sample sample, String dsKey) {
            sample.stop(Timer.builder("datasource.query.duration")
                    .tag("ds", dsKey)
                    .register(registry));
        }
    }
}

HikariCPコネクションプール監視

spring:
  datasource:
    dynamic:
      hikari:
        register-mbeans: true
        metrics-tracker: true
@Component
public class HikariMetricsExporter {

    @Autowired
    private DynamicRoutingDataSource dynamicRoutingDataSource;

    @Autowired
    private MeterRegistry meterRegistry;

    @Scheduled(fixedDelay = 15000)
    public void exportHikariMetrics() {
        Map<String, DataSource> dataSources = dynamicRoutingDataSource.getCurrentDataSources();

        for (Map.Entry<String, DataSource> entry : dataSources.entrySet()) {
            if (entry.getValue() instanceof HikariDataSource hikari) {
                HikariPoolMXBean pool = hikari.getHikariPoolMXBean();

                Gauge.builder("hikari.active.connections", pool, HikariPoolMXBean::getActiveConnections)
                        .tag("pool", entry.getKey())
                        .description("アクティブ接続数")
                        .register(meterRegistry);

                Gauge.builder("hikari.idle.connections", pool, HikariPoolMXBean::getIdleConnections)
                        .tag("pool", entry.getKey())
                        .description("アイドル接続数")
                        .register(meterRegistry);

                Gauge.builder("hikari.threads.awaiting", pool, HikariPoolMXBean::getThreadsAwaitingConnection)
                        .tag("pool", entry.getKey())
                        .description("接続待ちスレッド数")
                        .register(meterRegistry);
            }
        }
    }
}

Prometheus + Grafanaダッシュボード主要メトリクス

メトリクス 意味 アラート閾値
hikari.active.connections 現在のアクティブ接続数 > maxPool * 0.8
hikari.threads.awaiting 接続待ちスレッド数 > 5
datasource.switch.count データソース切替頻度 異常な急増
datasource.query.duration クエリ所要時間 P99 > 1s
hikari.connection.creation 接続作成レート 持続的高頻度作成

本番級アーキテクチャとベストプラクティス

完全な本番アーキテクチャ

┌─────────────────────────────────────────────────────────────┐
│                       Nginx / Gateway                        │
└──────────────────────────┬──────────────────────────────────┘
                           │
┌──────────────────────────▼──────────────────────────────────┐
│                    Spring Boot Application                    │
│                                                              │
│  ┌──────────────────────────────────────────────────────┐   │
│  │                 Controller Layer                      │   │
│  │        TenantFilter → @DS → @GlobalTransactional     │   │
│  └──────────────────────────┬───────────────────────────┘   │
│                             │                                │
│  ┌──────────────────────────▼───────────────────────────┐   │
│  │                 Service Layer                         │   │
│  │     @DSTransactional / @GlobalTransactional           │   │
│  └──────────────────────────┬───────────────────────────┘   │
│                             │                                │
│  ┌──────────────────────────▼───────────────────────────┐   │
│  │          DynamicRoutingDataSource                     │   │
│  │    ┌─────────┬──────────┬──────────┬──────────┐      │   │
│  │    │ master  │ slave_1  │ order_db │ tenant_X │      │   │
│  │    └────┬────┴────┬─────┴────┬─────┴────┬─────┘      │   │
│  └─────────┼─────────┼──────────┼──────────┼────────────┘   │
└────────────┼─────────┼──────────┼──────────┼────────────────┘
             │         │          │          │
        ┌────▼───┐ ┌──▼───┐ ┌───▼──┐ ┌────▼────┐
        │Master  │ │Slave │ │Order │ │Tenant   │
        │MySQL   │ │MySQL │ │MySQL │ │MySQL    │
        └────────┘ └──────┘ └──────┘ └─────────┘

よくある落とし穴まとめ

落とし穴 現象 解決策
@DSアノテーションが効かない AOPプロキシが切替を妨害 publicメソッドにアノテーション、同クラス内部呼出を避ける
トランザクション内データソース切替 切替後も元DSトランザクションに留まる @Transactionalの代わりに@DSTransactionalを使用
ThreadLocalリーク 非同期スレッドでDS識別子が残留 finallyでDynamicDataSourceContextHolderをクリーンアップ
コネクションプール枯渇 高並行で接続待ちタイムアウト データソースQPSに応じてプールサイズを適切に設定
レプリカから古いデータを読む 書き込み後すぐレプリカから読めない 書き込み後読み取りはプライマリへ、または@MasterOnlyを使用
シャードキーなしフルテーブルスキャン シャードキーなしクエリが全シャードにルーティング シャードキー必須またはESワイドテーブルを使用
テナントID未注入 SQLにtenant_id条件が欠落 グローバルテナントインターセプター + ホワイトリストテーブル
分散トランザクションタイムアウト Seataグローバルトランザクションがタイムアウトロールバック 適切なtimeoutを設定、長トランザクションを回避
動的データソース未クローズ テナント登録解除後も接続プールが解放されない removeDataSource時に接続プールをクローズ
HikariCP設定不適切 minimumIdle=maximumPoolSizeでリソース浪費 コアとピーク設定を区別

コネクションプール設定リファレンス

spring:
  datasource:
    dynamic:
      hikari:
        minimum-idle: 5
        maximum-pool-size: 20
        idle-timeout: 300000
        max-lifetime: 1800000
        connection-timeout: 30000
        connection-test-query: SELECT 1
        pool-name: DynamicHikariCP
      datasource:
        master:
          hikari:
            minimum-idle: 10
            maximum-pool-size: 50
        slave_1:
          hikari:
            minimum-idle: 10
            maximum-pool-size: 30
        slave_2:
          hikari:
            minimum-idle: 10
            maximum-pool-size: 30

設定暗号化

spring:
  datasource:
    dynamic:
      datasource:
        master:
          url: jdbc:mysql://mysql-master:3306/app_db
          username: ENC(g2U3x8vKpQ==)
          password: ENC(aB3dE7fG9hJ==)
@Configuration
public class DataSourceEncryptConfig {

    @Bean
    public DataSourcePropertySourceProcessor dataSourcePropertySourceProcessor() {
        return new DataSourcePropertySourceProcessor() {
            @Override
            public String decrypt(String cipherText) {
                if (cipherText.startsWith("ENC(")) {
                    String encrypted = cipherText.substring(4, cipherText.length() - 1);
                    return AesUtil.decrypt(encrypted, getSecretKey());
                }
                return cipherText;
            }
        };
    }
}

グレースフルシャットダウン

@Configuration
public class GracefulShutdownConfig {

    @Autowired
    private DynamicRoutingDataSource dynamicRoutingDataSource;

    @PreDestroy
    public void gracefulShutdown() {
        log.info("データソースのグレースフルシャットダウンを開始...");

        Map<String, DataSource> dataSources = dynamicRoutingDataSource.getCurrentDataSources();
        for (Map.Entry<String, DataSource> entry : dataSources.entrySet()) {
            if (entry.getValue() instanceof HikariDataSource hikari) {
                log.info("データソース[{}]をクローズ、アクティブ接続: {}", entry.getKey(), hikari.getHikariPoolMXBean().getActiveConnections());
                hikari.close();
            }
        }

        log.info("全データソースをクローズしました");
    }
}

マルチ環境設定管理

# application-dev.yml
spring:
  datasource:
    dynamic:
      strict: false
      datasource:
        master:
          url: jdbc:mysql://localhost:3306/dev_db

# application-prod.yml
spring:
  datasource:
    dynamic:
      strict: true
      hikari:
        minimum-idle: 10
        maximum-pool-size: 50
      datasource:
        master:
          url: jdbc:mysql://mysql-master.internal:3306/prod_db
          username: ${DB_MASTER_USER}
          password: ${DB_MASTER_PWD}

まとめ

Spring Boot 3 + MyBatis-Plusマルチデータソースアーキテクチャは、エンタープライズJavaアプリケーションの標準ソリューションです。主要な選定推奨:

  • 小中プロジェクトdynamic-datasource-spring-boot-starter + アノテーション切替 — シンプルで効率的
  • 中大プロジェクトShardingSphere-JDBC + シャーディング — 機能豊富
  • 多言語プロジェクトShardingSphere-Proxy スタンドアロンプロキシ層
  • トランザクション要件:単一DBは@DSTransactional、クロスDBはSeata AT、非同期はMQ結果整合性
  • マルチテナント:大口顧客はSchema分離、小規模顧客はTable分離 + MyBatis-Plusテナントプラグイン

マルチデータソースは銀の弾丸ではありません。導入前に本当に必要かを評価してください。単一データベースで解決できるなら、過剰設計しないでください。

ブラウザローカルツールを無料で試す →

#MyBatis-Plus#多数据源#分库分表#读写分离#ShardingSphere#Spring Boot