分散システムにおけるデータ整合性確保の戦略

分散ロックによる排他制御

複数ノード間でのリソース競合を解消し、整合性を担保するための基盤技術として分散ロックが挙げられます。共有リソースへのアクセス順序を制御することで、間接的にデータの一貫性を維持します。

1. Redis を用いた簡易ロック実装

Redis の SET NX 命令(存在しない場合のみ設定)を利用することで、基本的な分散ロックを構築可能です。

public class RedisMutexProvider {
    private StringRedisTemplate redisTemplate;
    private String resourceId;
    private static final int LOCK_DURATION = 30;

    // ロック取得
    public boolean acquire() {
        String requestId = UUID.randomUUID().toString();
        return redisTemplate.opsForValue()
            .setIfAbsent(resourceId, requestId, LOCK_DURATION, TimeUnit.SECONDS);
    }

    // ロック解放(所有者確認付き)
    public void release(String requestId) {
        String luaScript = "if redis.call('get', KEYS[1]) == ARGV[1] then " +
                           "return redis.call('del', KEYS[1]) " +
                           "else return 0 end";
        redisTemplate.execute(
            new DefaultRedisScript<>(luaScript, Integer.class),
            Collections.singletonList(resourceId), 
            requestId
        );
    }
}

課題点:ロックの自動更新機構がないため、業務処理が長時間化した場合にロックが失効するリスクがあります。

2. Redisson による高度なロック管理

Redis 向け Java クライアントである Redisson は、基本的な実装の欠点を補完する機能を提供します。

主要機能:

  • 再入可能ロック:同一スレッドによる重複ロック取得を許可
  • ウォッチドッグ:ロック有効期限の自動延長
  • フェアロック:リクエスト順序に基づいた公平なロック取得
  • レッドロック:複数ノード構成における高可用性ロック

実装例:

Service
public class StockAdjustmentService {
    @Autowired
    private RedissonClient redissonClient;
    @Autowired
    private StockRepository stockRepository;

    public boolean reduceQuantity(Long itemId, int amount) {
        // 再入可能ロックの取得
        RLock mutex = redissonClient.getLock("stock:" + itemId);
        try {
            // 最大 10 秒待機し、ロック保持時間は 10 秒
            boolean isLocked = mutex.tryLock(10, 10, TimeUnit.SECONDS);
            if (!isLocked) {
                return false; // ロック取得失敗
            }

            // 在庫減算処理
            Stock stock = stockRepository.findById(itemId);
            if (stock.getQuantity() < amount) {
                return false;
            }
            stock.setQuantity(stock.getQuantity() - amount);
            stockRepository.save(stock);
            return true;
        } finally {
            // 当前スレッドが保持している場合のみ解放
            if (mutex.isHeldByCurrentThread()) {
                mutex.unlock();
            }
        }
    }
}

レッドロック(RedLock)構成:

Redis クラスタ環境では、複数の独立ノードに対してロックを取得し、過半数の成功をもってロック成立とみなすことで耐障害性を高めます。

// 複数ノードのロックインスタンス生成
RLock lockA = clientA.getLock("stock:1001");
RLock lockB = clientB.getLock("stock:1001");
RLock lockC = clientC.getLock("stock:1001");

RedissonRedLock multiLock = new RedissonRedLock(lockA, lockB, lockC);
try {
    // 30 秒以内に取得し、10 秒後に自動解放
    boolean secured = multiLock.tryLock(30, 10, TimeUnit.SECONDS);
    if (secured) {
        //  critical section
    }
} finally {
    multiLock.unlock();
}

分散トランザクションフレームワークの活用

TCC や SAGA Pattern を手実装する代わりに、成熟したフレームワークを利用することで開発負荷を軽減できます。

1. Seata AT モード

Alibaba 开源の Seata は、ローカルトランザクションとグローバルロックを組み合わせる AT モードを提供します。

処理フロー:

  1. TM(Transaction Manager):グローバルトランザクションの開始
  2. RM(Resource Manager):各マイクロサービスのローカルトランザクション登録
  3. TC(Transaction Coordinator):undo_log を利用したコミット/ロールバック調整

コード例:

// 1. グローバルトランザクション開始側
@Service
public class OrderCreationHandler {
    @Autowired
    private OrderRepository orderRepo;
    @Autowired
    private InventoryClient inventoryClient;
    @Autowired
    private PaymentClient paymentClient;

    @GlobalTransactional // Seata アノテーション
    public void placeOrder(OrderDTO order) {
        // ローカル注文登録
        orderRepo.save(order);
        // 在庫減算リモートコール
        boolean success = inventoryClient.reduce(order.getProductId(), order.getQty());
        if (!success) {
            throw new RuntimeException("Insufficient stock");
        }
        // 決済処理リモートコール
        paymentClient.process(order.getId(), order.getAmount());
    }
}

// 2. 在庫サービス(ブランチトランザクション)
@Service
public class InventoryService {
    @Autowired
    private StockRepository stockRepo;

    @Transactional
    public boolean reduce(Long productId, int qty) {
        Stock stock = stockRepo.findById(productId);
        if (stock.getQuantity() < qty) {
            return false;
        }
        stock.setQuantity(stock.getQuantity() - qty);
        stockRepo.save(stock);
        return true;
    }
}

2. Hmily による TCC パターン

Hmily は TCC(Try-Confirm-Cancel)モデルに特化しており、アノテーションによる補償処理の簡易化を実現します。

@Service
@HmilyTCC(confirmMethod = "commitReservation", cancelMethod = "rollbackReservation")
public class StockReservationComponent {
    @Autowired
    private StockRepository stockRepo;

    // Try: 在庫の仮確保
    public boolean prepareStock(Long productId, int qty) {
        Stock stock = stockRepo.findById(productId);
        if (stock.getAvailable() < qty) {
            throw new BusinessException("Stock shortage");
        }
        // 凍結枠の増加と可用在庫の減少
        stock.setReserved(stock.getReserved() + qty);
        stock.setAvailable(stock.getAvailable() - qty);
        return stockRepo.update(stock) > 0;
    }

    // Confirm: 確定処理(凍結枠の解除)
    public boolean commitReservation(Long productId, int qty) {
        Stock stock = stockRepo.findById(productId);
        stock.setReserved(stock.getReserved() - qty);
        return stockRepo.update(stock) > 0;
    }

    // Cancel: 取消処理(可用在庫の回復)
    public boolean rollbackReservation(Long productId, int qty) {
        Stock stock = stockRepo.findById(productId);
        stock.setReserved(stock.getReserved() - qty);
        stock.setAvailable(stock.getAvailable() + qty);
        return stockRepo.update(stock) > 0;
    }
}

データベース層での整合性対策

1. DB 行ロックを用いたmutex

データベースのユニーク制約を利用し、低負荷環境向けの分散ロックを実現します。

-- ロック管理テーブル
CREATE TABLE mutex_table (
    resource_key VARCHAR(64) PRIMARY KEY,
    owner_id VARCHAR(64) NOT NULL,
    expires_at TIMESTAMP NOT NULL
);

-- ロック取得(INSERT 成功がロック取得)
INSERT INTO mutex_table (resource_key, owner_id, expires_at)
VALUES ('stock:1001', 'uuid-abc', NOW() + INTERVAL 30 SECOND)
ON DUPLICATE KEY UPDATE resource_key = resource_key;

-- ロック解放(所有者検証)
DELETE FROM mutex_table
WHERE resource_key = 'stock:1001' AND owner_id = 'uuid-abc';

懸念点:パフォーマンスボトルネックになりやすく、DB 障害時にロックが解放されなくなる可能性があります。

2. シャーディング環境でのトランザクション

ShardingSphere を利用し、XA プロトコルまたは BASE トランザクションにより分庫分表時の整合性を保証します。

# ShardingSphere 設定例
spring:
  shardingsphere:
    rules:
      transaction:
        default-type: XA # または BASE
        provider-type: Atomikos # XA トランザクションマネージャ

メッセージキューを活用した非同期整合性

1. RocketMQ トランザクションメッセージ

ローカルトランザクションとメッセージ送信の原子性を保証する仕組みです。

@Service
public class OrderMessageProducer {
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    @Autowired
    private OrderRepository orderRepo;

    public void placeOrder(OrderDTO order) {
        // 半メッセージ送信
        rocketMQTemplate.sendMessageInTransaction(
            "order-tx-group",
            "order-topic",
            MessageBuilder.withPayload(order).build(),
            order
        );
    }

    // ローカルトランザクション状態確認
    @RocketMQTransactionListener(txProducerGroup = "order-tx-group")
    public class OrderTxListener implements RocketMQLocalTransactionListener {
        @Override
        public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
            OrderDTO order = (OrderDTO) arg;
            try {
                orderRepo.save(convert(order));
                return RocketMQLocalTransactionState.COMMIT;
            } catch (Exception e) {
                return RocketMQLocalTransactionState.ROLLBACK;
            }
        }

        // 状態不明時のチェックバック
        @Override
        public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
            String oid = msg.getHeaders().get("orderId", String.class);
            Order order = orderRepo.findById(oid);
            return order != null ? 
                RocketMQLocalTransactionState.COMMIT : 
                RocketMQLocalTransactionState.ROLLBACK;
        }
    }
}

2. Kafka トランザクション機能

Kafka 0.11 以降では、トランザクション ID を使用してパーティション跨ぎの書き込み原子性を実現できます。

@Service
public class KafkaEventCoordinator {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    public void processEvent() {
        // トランザクション開始
        kafkaTemplate.executeInTransaction(template -> {
            template.send("topic-alpha", "data-1");
            template.send("topic-beta", "data-2");
            updateLocalDB();
            return true; // コミット
        });
    }
}

アーキテクチャ選定基準

各手法の特性を理解し、ビジネス要件に合わせて最適なパターンを選択する必要があります。

手法 代表ツール 整合性レベル パフォーマンス 推奨ユースケース
分散ロック Redisson 最終整合性 在庫減算などリソース競合
TCC フレームワーク Hmily 最終整合性 中〜高 決済、注文など核心業務
グローバルトランザクション Seata AT 最終整合性 マイクロサービス間調整
トランザクションメッセージ RocketMQ 最終整合性 非同期通知、データ同期
DB ロック ユニーク索引 強整合性 低Concurrency、単純场景
レッドロック Redisson RedLock 強整合性 高可用性が求められるロック

選定における主要観点:

  1. ビジネス許容度:金融取引などは強整合性が必要だが、非核心業務は最終整合性で十分
  2. 性能要件:高Concurrency 環境では Redis や MQ ベースの方案を優先
  3. 実装コスト:Seata や Redisson などの成熟フレームワークの活用を推奨

設計上の注意点として、過度な設計を避け、最終整合性で要件を満たせる場合はそちらを採用します。また、すべての分散方案には補償機制(定时任务によるデータ検証など)を用意し、防御的な設計を心がけます。読み込みが多い场景ではキャッシュによる負荷分散を行い、書き込みが多い场景では TCC やトランザクションメッセージによりデータ精度を確保します。さらに、分散トランザクションの監視体系を構築し、不整合発生時に即座に検知できる体制を整備することが重要です。

タグ: distributed-lock Redisson Seata tcc-pattern RocketMQ

8月14日 10:59 投稿