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テナントプラグイン
マルチデータソースは銀の弾丸ではありません。導入前に本当に必要かを評価してください。単一データベースで解決できるなら、過剰設計しないでください。
ブラウザローカルツールを無料で試す →