-レガシー大量送信情報の受信区間を同期処理から保存後終了の構造へ移行した経験-
1. はじめに — どのような問題に直面したのか
私が担当したアダプターサービスは、病院のレガシー医療情報システム(以下、レガシー)が送信するアンケート配信情報を受信し、内部マイクロサービス(アカウントサービス、アンケートサービス)に登録する窓口の役割を担っています。一見すると、1件の電文を受け取って登録する単純な連携ですが、運用を開始すると3つの問題が同時に明らかになりました。
-
処理時間がそのままレガシーのトランザクション占有時間になっていました。 1件を処理するために、アダプターはレガシーへの照会を2回、アカウントサービスへの呼び出しを4〜5回、アンケートサービスへの呼び出しを10回余り、合計約20回のREST往復を同期的に実行します。レガシーは大量のデータを自身のトランザクション内で連続して呼び出すため、レガシーがアダプターの処理時間全体を保持したまま待機する構造になっていました。トレースを確認すると、1件あたりのHTTP応答は1.4〜2.2秒で、数百件が集中する時間帯にはレガシー側のタイムアウトが常に発生するリスクがありました。
-
部分的な失敗によって孤児リソースが残りました。 アダプターは独自のDBを持たない(DB-less)設計だったため、入口に付いていた@Transactionalは実質的に何の効果もありませんでした。途中のREST呼び出しが失敗すると、前の段階ですでに作成されたゲストアカウントや受信者情報が、ロールバックされずそのまま残っていました。
-
1つの同期応答に多すぎる責任が集中していました。電文の解析結果から業務検証(書式の存在有無、グループの存在有無、書式の配布有無)までをすべて1つの応答の成功/失敗ヘッダーで返していたため、どの段階でなぜ失敗したのかが応答に含まれず、再処理の手段もありませんでした。
この記事では、この受信区間を「受信後すぐに保存して終了する(Inboxパターン)」構造へ変更し、実際の登録ロジックをKafkaイベントベースの非同期処理として分離した過程を整理します。設計判断の根拠や途中で覆した決定、実測で確認した結果まで併せて記録します。
2. 技術選定の背景 — 選択肢を絞り込んだ過程
2.1 非同期化方式の選択
「先に応答を返し、後で処理する」という方針はすぐに決まりましたが、その方法には複数の候補がありました。それぞれを除外した理由は次のとおりです。
-
アダプター内部のスレッドプールによる非同期処理 — コード変更は最も少ないものの、アダプターはDBを持たないため、受け付けたデータをどこにも記録できません。Podが再起動すると、受け付けたと応答したデータが痕跡なく消えてしまいます。レガシーは成功応答を受け取っているため再送しないので、これはそのまま配信の欠落につながります。
-
アダプターにDBを追加 — 受付記録は残せますが、「アダプターはデータを所有しない」というサービスのアイデンティティを損ないます。データを持った瞬間に、元のサービスとの同期問題が新たに発生します。
-
メッセージブローカーに原文を直接投入 — 欠落は防げますが、個々のデータの状態(再試行回数、失敗理由、終了有無)を照会したり、運用担当者が目視で確認したりする手段がありません。大量失敗の状況で「現在何件が滞留しているのか」に答えられない構造は、運用で使うのは難しいと判断しました。
最終的な選択はInboxパターン + イベントトリガーでした。原文を、すでにDBを所有しているアンケートサービスに保存します。ただし、アンケートサービスはそれを解釈しない原文(不透明なペイロード、opaque payload)と処理メタデータとしてのみ扱います。電文の意味を理解しているのはあくまでアダプターだけという境界を維持しながら、保存・状態管理というインフラの責任のみをDB所有サービスに委譲した形です。
2.2 処理トリガー — Feignによる逆方向呼び出しではなくイベント
保存後、実際の処理を誰が起動するのかも論点になりました。アンケートサービスがアダプターを直接呼び出す(Feign)と配線は単純ですが、これまで一方向だったアダプター→アンケートの依存関係に逆方向の同期依存が新たに生じます。通信マトリクスが循環した瞬間に、障害の伝播経路も同時に生まれます。
そこでアンケートサービスは、保存のコミット後にドメインイベントのみを発行し、アダプターがコンシューマーとしてこれを受け取って処理する方式にしました。結果として、サービス間の契約には「アダプターがアンケートイベントを購読する」という1行だけが追加され、逆方向の同期呼び出しは作りませんでした。パーティションキーをアンケート作成番号に指定し、同じ対象のイベントが同じパーティションで発行順に消費されるようにしたことも、この選択によって追加コストなしで得られた利点です。
2.3 適用した技術スタック
-
Apache Kafka — 保存完了イベントの発行と購読。パーティションキーによって順序保証のインフラ上の根拠を確保しました。
-
PostgreSQL & JPA — Inboxの永続化。原文はJSONテキストとして、状態・再試行回数・失敗理由は正規化したカラムで管理します。
-
ShedLock — fallbackバッチが複数インスタンスで重複実行されないよう、分散ロックを適用しました。
-
Redis — 外部システム照会のキャッシュ(§4.4)および配信メタデータの保存層として使用しました。
-
Spring Cloud OpenFeign — Inboxの登録・占有・状態更新APIの呼び出し。契約は独立したclientモジュールに分離し、社内リポジトリにデプロイしました。
3. 適用過程
3.1 受信区間 — 何を残し、何を取り除いたのか
移行の最初の段階は、受信時に行う処理を最小限に減らすことでした。ただし、検証をすべてなくしたわけではありません。電文自体が壊れているデータは保存しても恒久的な失敗になるため保存せず、直ちに同期失敗として応答し、それ以外の業務検証はすべて非同期経路に引き渡しました。
// 전환 전: 진입점이 등록 전체를 동기로 수행 (@Transactional은 DB가 없어 무효)
@Transactional
public void registerSurveySendInfo(SurveySendInfo info) {
// 레거시 조회 2회 + 계정 4~5회 + 설문 10여 회 ... 총 20여 회 REST 왕복
// 중간 실패 시 앞 단계 산출물(게스트, 수신자정보 등)이 롤백 없이 잔존
}
// 전환 후: 저장-후-종료
public void acceptSurveySendInfo(SrvySendInfoDto ipd) {
validateRequired(ipd); // srvyDrawupNo, patno, formCd 존재만 확인
smsMsgInfoAction.saveMsgInfo(ipd); // 발송 메타데이터 보존(Redis)
surveySendInfoInboxProxy.register( // 원문 그대로 저장 + 커밋 후 이벤트 발행
SurveySendInfoInboxCdo.from(ipd));
}
[コード1] 受信入口の移行前後。レガシーとの対面契約(URL・電文仕様・応答ヘッダー)は変更していません。
ここで実際に直面したのが 必須フィールド一覧の確定でした。最初は発送日時も必須としていましたが、ステージングでの実測により、発送日時が空の電文が本番環境に実在することを確認しました。既存のロジックがすでにパース失敗時に受信時刻へフォールバックしていたため、検証対象から除外し、空のまま保存したうえで、処理時のフォールバックに任せるよう修正しました。検証を強化しようとして、問題なく処理されていたトラフィックを遮断しかけた事例です。
3.2 状態機械と単一リトライ台帳
保存されたレコードは5つの状態を遷移します。受信(RECEIVED) → 処理中(PROCESSING) → 完了(COMPLETED)となり、失敗時はFAILEDとして記録され、リトライ上限を超えるとDEADへ遷移して手動介入の対象になります。
設計で最も重視した原則は、リトライ台帳を一つに維持することでした。メッセージング層のリトライ(デフォルト10回後にDLT)とアプリケーションのリトライが重なると、実際の試行回数を誰も説明できなくなります。そのため、コンシューマーは例外を外へ投げず、すべて捕捉してInboxに記録するだけにし、リトライの判断はInboxの試行回数とfallbackバッチが専任で行います。
3.3 原子的な占有(claim) — 重複処理を防ぐ唯一の関門
非同期化で最も危険なシナリオは、イベントの重複消費とバッチによる再発行が競合し、同じレコードが2回登録されることです。これを条件付きの単一UPDATE文で防ぎました。遷移に成功したインスタンスだけが原文を受け取って処理を進め、残りは静かに終了します。
@Modifying
@Query("UPDATE SurveySendInfoInboxJpo i " +
" SET i.status = 'PROCESSING', i.modifiedOn = :now " +
" WHERE i.id = :id AND i.status IN ('RECEIVED', 'FAILED')")
int claim(@Param("id") Long id, @Param("now") LocalDateTime now);
// 갱신 행 수가 1이면 점유 성공 → 원문 반환, 0이면 이미 점유/종결된 건 → null
public SurveySendInfoInboxRdo claimAndLoad(Long id) {
if (repository.claim(id, LocalDateTime.now()) == 0) {
return null;
}
return SurveySendInfoInboxRdo.from(repository.getById(id));
}
[コード2] 占有は取得後に更新するのではなく、条件付きの単一UPDATEでなければなりません。取得と更新を分けると、その間が競合区間になります。
設計初期には占有APIと原文取得APIを別々にしていましたが、実装時に、占有に成功した際に原文も併せて返すよう統合しました。2回の呼び出しに分けると、呼び出し回数が増えるだけでなく、占有には成功したものの取得に失敗する中途半端な状態が発生するためです。
3.4 コンシューマー — 例外を外へ投げない
@KafkaListener(topics = SurveyTopic.SURVEY_SEND_INFO_RECEIVED, groupId = "...")
public void handle(SurveySendInfoReceivedDomainEvent event) {
SurveySendInfoInboxRdo rdo = inboxProxy.claim(event.getInboxId());
if (rdo == null) {
return; // 이미 점유·종결됐거나 선행 미종결 건이 있어 보류된 경우
}
try {
surveyFlow.registerSurveySendInfo(rdo.toSurveySendInfo());
inboxProxy.complete(rdo.getId());
} catch (PromsBizException e) { // 업무 예외 → 메시지를 그대로 사유로
handleFailure(rdo, e.getMessage());
} catch (Exception e) { // 시스템 예외 → 고정 문구
log.error("inbox 처리 실패. id={}", rdo.getId(), e);
handleFailure(rdo, "시스템 오류");
}
}
private void handleFailure(SurveySendInfoInboxRdo rdo, String reason) {
String status = InboxProxy.fail(rdo.getId(), reason); // FAILED 또는 DEAD 반환
if ("DEAD".equals(status)) {
amcFlowProxy.sendResult(rdo.getSrvyDrawupNo(), "N", reason); // 레거시 회신
}
}
[コード3] 失敗記録APIがDEADへの遷移可否を判定して返し、返信トリガーはアダプターが担います。
業務例外とシステム例外をリトライポリシーでは、区別しないことにしました。書式未登録のような業務エラーも、担当者が書式を登録すれば解消される性質であり、リトライ間隔が1分なのでコストもほとんどかからないためです。両者の違いは、レガシーに返す理由メッセージにのみ現れます。
3.5 Fallbackバッチ — 直接処理せず、イベントだけを再発行
イベントの紛失やコンシューマーの停止を補正するバッチを設けましたが、バッチが登録ロジックを直接呼び出さないようにしました。イベントだけを再び流せば、正常経路と復旧経路が完全に同じコードへ収束するためです。復旧専用のコードパスは通常実行されないため、本当に必要な瞬間に動作を信頼しにくいという点も考慮しました。
@Scheduled(cron = "0 * * * * *") // 1분 주기
@SchedulerLock(name = "surveySendInfoInboxRedeliver") // 다중 인스턴스 중복 방지
public void redeliver() {
// ① 이벤트 유실 의심: RECEIVED 이면서 생성 후 2분 경과
// ② 재시도 대상: FAILED 이면서 retryCount < 5
// ③ 처리중 고착: PROCESSING 이면서 15분간 갱신 없음 → RECEIVED 복귀 후 재발행
for (SurveySendInfoInbox target : InboxStore.findRedeliverTargets()) {
try {
eventPublisher.publish(SurveySendInfoReceivedDomainEvent.of(target));
} catch (Exception e) { // 한 건 실패가 라운드 전체를 막지 않도록
log.warn("재발행 실패. id={}", target.getId(), e);
}
}
}
[コード4] 3種類の滞留状況を1つのバッチで吸収します。レコード単位のtry-catchは、バッチを作成する際の習慣として付けておくほうが安全です。
4. 問題解決の経験
4.1 リトライの安全性 — ガードを分散させるのではなく原子化する
非同期リトライを導入するには、登録ロジックが複数回実行されても安全でなければなりません。段階ごとに確認したところ、受信者情報・スケジュール・割り当て・作業登録の4段階が、いずれも無条件に新規作成されるため、重複のリスクがありました。特にスケジュールと割り当てが重複作成されると、作業のない有効な発送レコードが残り、後続バッチが永久に失敗する事故パターンが、すでに実際に発生した前例もありました。
最初はアダプターに「既存レコードの有無を逆方向に取得するAPI」を追加し、各段階にガードを付ける方向を検討しました。しかしこの方式では、ガードが4か所に分散し、新しい段階が追加されるたびに同じミスを繰り返す余地が残ります。
そこで、アンケートサービス側に、4つの段階を1つのトランザクションにまとめるバンドルAPIを作成して呼び出すよう整理しました。バンドルが失敗すればアンケート側のデータはまったく残らないため、リトライが自動的に安全になります。入口ですでに完了済みのレコードを再送分岐から除外していたため、この原子化だけでアダプター側の冪等性ガードは、すべて不要になりました。原子性の責任をデータを所有するサービスへ移したことが核心でした。
4.2 順序保証と残存リスク
レガシーはバッチと単件を区別して通知しないため、到着順をそのまま業務上の順序とみなすことを明示的な決定として残しました。パーティションキーがこれをインフラレベルで支えますが、リトライが入り込むと前提が揺らぐ区間が1つ残ります。最初の発送電文が失敗してリトライを待っている間に、同じ対象の提出電文が先に到着すると、順序が逆転します。
この区間は、占有直前に、同じ対象について、より早い未完了レコードがあれば現在のレコードを処理せず終了するガードで防ぎました。保留されたレコードはバッチが順番どおりに再発行するため、最終的には到着順に処理されます。インデックス1つで解決できる低コストなガードであり、後に「同一電文の再受信を再送依頼として解釈する」というポリシーを追加した際にも、このガードがそのまま安全装置として機能しました。
4.3 同期契約縮小の代償 — 失敗を通知する経路を作る
応答の意味が「登録完了」から「受付完了」へ変わったことで、業務検証の失敗をレガシーに通知する方法がなくなりました。この空白を埋めるため、ちょうど別途設計中だった一方向の結果通知インターフェースを、失敗返信チャネルとして併用するよう拡張しました。
-
最終失敗(DEAD)レコード — 失敗理由を含め、発送が成立しなかったことを通知します。
-
発送が省略されたレコード — 問診のみを作成し、通知メッセージは送信しない依頼であることを理由とともに通知します。
-
すでに提出済みで再発送できないレコード — 再発送できない理由を通知します。
判定はデータを所有するアンケートサービスが行って応答で通知し、返信自体はレガシー電文を把握しているアダプターが担当するよう責任を分担しました。ただし、返信に失敗した場合の再返信手段がない点は、現在も残っている課題です。
4.4 ボトルネックは推測せず計測する — 外部取得のキャッシュ
移行後も処理量が期待どおりに向上しなかったため、モニタリングツールで複数のトレースを突き合わせて確認しました。その結果、1件の登録ごとに発生するレガシー照会2件のうち同意書照会の一つが一貫して760~984ミリ秒を要する一方、もう一つの照会は50~150ミリ秒にとどまることを確認しました。2つの呼び出しはコード上ほぼ同じ形に見えたため、計測していなければ同程度に遅いと推測し、見当違いの箇所を修正していた可能性が高いです。
レガシー側の改善が不可能であることを確認した後、アダプターで短いTTLのキャッシュに吸収しました。同じアンケートを対象とする複数のメッセージが同時に到着する業務特性により、総照会回数が「メッセージ数」から「異なるアンケート数」へと減少します。
public List<AmcSurveyConsent> findAmcSurveyConsent(String srvyNo) {
String key = "proms:amc:survey-consents:" + srvyNo;
try {
String cached = redisTemplate.opsForValue().get(key);
if (cached != null) {
return objectMapper.readValue(cached, new TypeReference<>() {});
}
} catch (Exception e) {
log.warn("캐시 조회 실패 — 원 호출로 폴백. srvyNo={}", srvyNo, e);
}
List<AmcSurveyConsent> result = amcProxy.findSurveyConsent(srvyNo);
// 동의서는 빈 목록도 캐시한다 — 동의서 없는 설문이 일반적이라 여기가 절감의 핵심
cachePut(key, result, Duration.ofMinutes(5));
return result;
}
[コード5] キャッシュ障害が機能障害にならないよう、キャッシュ層の例外は警告ログを出した後、元の呼び出しにフォールバックします。
キャッシュポリシーはデータの性質に応じて差をつけました。設定が存在しないという結果はキャッシュしません。担当者が設定を修正した直後も最大5分間古い結果が維持されると、修正したのになぜ反映されないのかという問い合わせにつながるためです。逆に、同意書の空リストはキャッシュします。同意書のないアンケートのほうがむしろ一般的で、このケースが削減効果の大部分を占めます。トレードオフとして、レガシー設定の変更反映が最大5分遅延することは運用側に周知しました。
4.5 例外とトランザクション境界 — noRollbackForを使わないことにした理由
「処理は失敗したが、失敗の記録は残さなければならない」という要件が生じたとき、最も簡単な解決策はロールバック除外設定を付けることです。しかしこの方法では、例外の種類が増えるたびに一覧を管理する必要があり、どの状態が残り、どの状態が消えるのかをコード全体を読まなければ把握できません。
そこで、トランザクション境界の分離だけで解決するという原則を立てました。占有・完了・失敗記録がそれぞれ独立したREST呼び出しであり、独立したトランザクションであることが、この原則の実装です。本処理が失敗しても、失敗記録は別トランザクションであるためロールバックされず、試行回数と理由が完全に残ります。
失敗記録まで失敗する最悪のケースは、ロールバックではなく、時間による復旧に委ねます。状態が処理中のまま残った場合は15分後にバッチが戻し、占有の原子性によって二重処理を防ぐため安全です。ただし、この経路では試行回数が増加しないことを文書に明記しました。
全件検証の過程で、共通ライブラリの空白も一つ見つかりました。社内業務例外クラスがコンシューマーの再試行除外一覧から漏れていたため、即座にDLTへ送るべき業務例外が10回再試行された後にようやくDLTへ送られていました。Inbox経路ではコンシューマーが例外を投げないため影響はありませんでしたが、他サービスのコンシューマーにはそのまま影響する問題だったため、共通ライブラリの改善項目として登録しました。
4.6 巻き戻しの方法まで設計に含める
デプロイ後に問題が発生した場合の巻き戻し手順をあらかじめ整理する中で、今回の切り替えには、一般的なイメージロールバックが通用しない制約があることが分かりました。レガシーはすでに受付応答を返しているため再送しませんが、旧バージョンのアダプターにはコンシューマーがなく、残っている未完了の案件が永久に未処理のまま残るためです。つまり、未完了案件の消化(drain)確認が、どのロールバック経路でも先行して必要になります。
そこで、第一の手段をイメージロールバックではなく設定トグルとしました。受信経路だけを旧来の同期方式に戻し、コンシューマーとバッチは稼働させておけば、新規流入は同期処理されながら残存案件は引き続き消化され、自然に消化が完了します。
// 롤백 토글: false면 Inbox를 거치지 않고 구 동기 경로로 즉시 처리한다.
// off 상태에서도 컨슈머·fallback은 유지되어 잔여 inbox 건은 계속 소화된다.
@Value("${proms.survey-send-info.inbox-enabled:true}")
private boolean inboxEnabled;
if (inboxEnabled) {
surveySendInfoInboxFlow.acceptSurveySendInfo(ipd); // 저장-후-종료
} else {
surveyFlow.registerSurveySendInfo(new SurveySendInfo(ipd)); // 구 동기 경로
}
[コード6] イメージを再デプロイせず、設定変更と再起動だけで受信経路を戻せるようにしました。
トグルを設計する際、見落としかけた点もありました。送信メタデータの照会をInbox優先に変更していた状態でトグルをオフにすると、再受信された案件ではRedisだけが更新されるため、古いInbox行の値で応答する逆転が発生します。照会箇所でも同じトグルを共有し、オフの状態ではRedisだけを見るよう補完しました。トグルを一つ追加したら、それを参照するすべての箇所を併せて確認する必要があることを学びました。
5. 検証結果
開発環境で4回、ステージングで1回の実トラフィック検証を実施し、ログ・トレースツールで数値を実測しました。
-
受信応答時間 — 従来の1.4~2.2秒から中央値33ミリ秒となり、約50倍短縮されました。レガシートランザクションの占有問題が根本的に解消されました。
-
全件完走 — 111件を38秒、273件を94秒で、失敗0件・fallback介入0件のまま完走しました。処理率は毎秒約3件であり、増強が必要になった場合はパーティションとコンシューマーの同時実行性がレバーになることも併せて確認しました。
-
外部照会の削減 — キャッシュ導入後、レガシー照会がメッセージ数比で約98%減少しました(111件処理時のキャッシュミス照会は約5回)。
-
失敗通知経路の実証 — 書式未登録29件が5回の再試行後にDEADへ遷移し、レガシーへ理由とともに応答される一連の過程を、実トラフィックで確認しました。
-
障害自動復旧の実証 — ステージング検証が偶然にもレガシーDBの瞬断と重なりました。273件がすべて正常に受け付けられた後、1ラウンド目で全件が失敗しましたが、DB復旧後にバッチを再発行することで、3.5分以内に273件すべてが自動復旧(最終失敗0件)しました。意図した設計が計画されたシナリオではなく実際の障害で動作したため、個人的には今回の切り替えで最も意義のある結果でした。
6. 振り返り — 次回も生かすこと
-
応答の意味が変わったら、その波及を最後まで追跡する必要があります。「登録完了」が「受付完了」に変わったのはコード数行の変更ですが、失敗通知経路の新設と、ロールバック時の消化を先行させるという2つの大きな要件がここから派生しました。契約の文言ではなく、意味が変わったかどうかをまず問う習慣が身につきました。
-
冪等性はガードを増やすのではなく、境界を移すことで得るほうが望ましいです。 4か所にガードを付ける代わりに、データ所有サービスでアトミック化したことで、ガードがすべてなくなり、新しいステップが追加されても安全性が維持される構造になりました。
-
復旧経路は通常経路と同じコードに収束させてこそ、信頼できるものになります。 バッチが直接処理せず、イベントを再発行するだけにした判断のおかげで、実際の障害で復旧経路が初めて実行されたときも、別途検証することなく動作しました。
-
ボトルネックは計測してから手を入れるべきです。 形が同じように見える2つのクエリのうち、片方だけが15倍遅いという事実は、トレースを見なければ分かりませんでした。
-
元に戻す手順を設計に含めると、設計の盲点が明らかになります。 ロールバック手順を書いている途中で、「受付応答をすでに返した案件は、旧バージョンでは処理する手段がない」という制約を発見しました。そのおかげで、トグルという、より優れた第一の手段を用意できました。
残っている課題もあります。返信に失敗した場合に再返信する手段がなく、原文に平文の個人情報が含まれるため、完了した案件の保存期間とクリーンアップバッチの実行周期を確定する必要があります。どちらも後続作業として管理しています。
conley