RabbitMQ DLXを活用した決済タイムアウト処理
受注システムの状態遷移において、30分以内に決済が完遂しない場合に自動取消しを行う機構を実装しました。メッセージの有効期限(TTL)切れまたは消費者側での明示的拒否を検知し、デッドレターエクスチェンジ(DLX)経由で別キューへ迂回させる構成としています。
@Component
public class OrderStateTransitionConsumer {
@RabbitListener(queues = "order.timeout.handler.queue")
public void processExpiredOrder(TimedOutOrderDto dto) {
if (dto.getPaymentStatus() == PaymentStatus.WAITING_AUTH) {
executeCancellation(dto.getOrderId());
}
}
private void executeCancellation(String targetOrderId) {
log.info("期限超過に基づく注文取消処理を開始: {}", targetOrderId);
// orderRepository.transitionToCancelled(targetOrderId);
}
}
キューおよび交換機設定
Spring設定クラスにて、プライマリキューにDLX向けパラメータを付与します。これにより、正常な順序立てたルーティングとは別に、期限切れメッセージを確実に回収するパスを確保できます。
@Configuration
public class MessagingTopologyConfig {
@Bean
DirectExchange primaryOrderExchange() {
return new DirectExchange("order.routing.exchange");
}
@Bean
Queue primaryOrderQueue() {
return QueueBuilder.durable("order.transaction.queue")
.withArgument("x-dead-letter-exchange", "order.deadletter.exchange")
.withArgument("x-dead-letter-routing-key", "dl.order.timeout.key")
.build();
}
@Bean
DirectExchange deadLetterExchange() {
return new DirectExchange("order.deadletter.exchange");
}
@Bean
Queue deadLetterQueue() {
return new Queue("order.timeout.handler.queue");
}
@Bean
Binding routePrimary() {
return BindingBuilder.bind(primaryOrderQueue())
.to(primaryOrderExchange())
.with("primary.order.key");
}
@Bean
Binding routeDeadLetter() {
return BindingBuilder.bind(deadLetterQueue())
.to(deadLetterExchange())
.with("dl.order.timeout.key");
}
}
RedisとLuaを用いた高負荷在庫管理の実装
同時アクセスが集中する販売イベント対策として、Redis上のカウンター値を対象としたアトミックな減算ロジックを構築しました。初期段階ではLuaスクリプト単体で実装していましたが、後の運用ではRedissonクライアントを用いた分布式排他制御を追加し、ネットワーク分断時やスクリプト実行途中のエラーによるロックリークを防止する設計に変更しています。
@Service
public class StockReservationManager {
private final RedissonClient redissonClient;
private final String atomicDeductionScript;
public StockReservationManager(RedissonClient client) {
this.redissonClient = client;
this.atomicDeductionScript = """
local stockVal = tonumber(redis.call('GET', KEYS[1]))
local demand = tonumber(ARGV[1])
if stockVal >= demand then
redis.call('DECRBY', KEYS[1], demand)
return 1
else
return 0
end
""";
}
public boolean reserveInventory(String productId, int amount) {
RLock mutex = redissonClient.getLock("res:lock:" + productId);
try {
mutex.lock();
Object evalResult = redissonClient.getScript().eval(
ScriptMode.READ_WRITE,
atomicDeductionScript,
ReturnType.INTEGER,
Collections.singletonList(productId),
String.valueOf(amount)
);
return Integer.valueOf(1).equals(evalResult);
} finally {
if (mutex.isLocked() && mutex.isHeldByCurrentThread()) {
mutex.unlock();
}
}
}
}
Elasticsearch導入によるシステム監視基盤
アプリケーションの稼働ログ、APIレイテンシー、および予期せぬ例外発生時のスタックトレースについて、ファイルストレージからElasticsearchへの一元管理へ移行しました。JSON形式で構造化されたログをインデックス化することで、特定エンドポイントの性能ボトルネック特定や、定期アラート発報のための複雑なクエリ実行を容易にしています。
Redisデータ構造を活用した複合ドメイン機能
1つのインメモリキャッシュ層に対し、ユースケースに応じて最適化されたデータ型を採用し、メモリ効率とクエリパフォーマンスの両立を図りました。
- コンテンツ評価(いいね): `Sorted Set` に投稿IDをメンバー、評価時刻や加权スコアをバリューとする方式で、リアルタイムなトレンド表示に対応。
- 地理空間検索: `GeoHash` アルゴリズムに基づく `GEO` コマンド群を使用し、ユーザ当前位置からの半径指定による店舗一覧抽出を実装。
- 継続アクション記録: `Bitmap` で日次フラグを管理し、月間連続ログイン日数や離脱リスク顧客の判定をO(1)計算量で処理。
- 抽象ユーザ数計測: `HyperLogLog` により重複を除いたユニークアカウント数を低メモリコストで概算しつつ、アクティブセッション数は原子性を担保したカウンターキーで並列安全に追跡。
メタデータの分離とセグメント単位アップロード
帯域の低い環境やモバイル通信を考慮し、大容量バイナリデータを固定長チャンクへ分割して転送するプロトコルを採用しました。サーバ側では各フラグメントのハッシュ値およびアップロード済みオフセットをリレーショナルデータベースまたはRedisのSet構造に記録します。クライアントは切断復旧時に該当するchunk_idの一覧を照会し、未送信区画のみを再送することで、全体やり直しのオーバーヘッドをゼロに抑えています。