レシートの写真1枚を構造化データに変換するパイプラインを作った。APIのレスポンスを15秒から100msに短縮し、GPTのプロンプトトークンを53%削減し、スキーマ違反率を0%にするまでの過程。
問題の背景
コア機能はシンプルだ: レシートの写真をアップロードすると、構造化データが返ってくる。
しかしこのフローの中に、厄介な要件が潜んでいた:
- OCRに3〜5秒、GPTの解析に5〜10秒 → APIで同期処理すると1リクエストあたり最大15秒のブロッキング
- OCRが失敗しても、解析は別途リトライできなければならない
- ユーザーに処理の進行状況をリアルタイムで表示する必要がある
- AIの結果を人間がレビューして修正できなければならない
全体アーキテクチャ
ステートマシンでレシートのライフサイクルを管理する
レシートの状態を明示的に定義し、許可された遷移のみを可能にした。
なぜステートマシンなのか? if-elseで状態を管理するとコードが散在する。「解析失敗の状態からOCRを再実行できますか?」 — ステートマシンなら遷移テーブル1つで答えが出る。
# app/domain/state_machine.py — 外部依存ゼロ、ユニットテストが容易
_TRANSITIONS = {
UPLOADED: [OCR_PROCESSING],
OCR_PROCESSING: [OCR_COMPLETED, OCR_FAILED],
OCR_COMPLETED: [PARSING],
OCR_FAILED: [OCR_PROCESSING], # リトライ
PARSING: [PARSE_COMPLETED, PARSE_FAILED],
PARSE_COMPLETED: [REVIEW_NEEDED],
PARSE_FAILED: [PARSING], # リトライ
REVIEW_NEEDED: [EXPORTED],
EXPORTED: [], # 最終状態
}
@staticmethod
def transition(current, target):
if target not in _TRANSITIONS.get(current, []):
raise InvalidTransitionError(f"{current} → {target} は不可")
return target
9つの状態を持つのは過剰に見えるかもしれない。しかし REVIEW_NEEDED なしで即座に出力すればAIのエラーを検出できず、OCR_FAILED と PARSE_FAILED を分離しなければ「どこで失敗したか」がわからない。各状態には存在する理由がある。
KafkaでAPIレスポンスを15秒 → 100msへ
APIサーバーが直接OCR/解析を行うと、1リクエストにつき15秒かかる。Kafkaを導入してアップロード(API)と処理(Worker)を分離したことで、APIはメッセージをキューに入れて即座に202 Acceptedを返すようになった。
Workerが落ちても、Kafkaがメッセージを保持する。再起動すれば未処理のメッセージから再消費される。
Kafkaはこの規模には過剰であることは認める。Redis Streamで十分だったかもしれない。しかしConsumer Groupによる水平スケールを見据えて選択した。Workerを複数台に増やせば、スループットは線形に伸びる。
2フェーズコミット: 「何がクリティカルか」を区別する
Workerで最も難しいのは部分的な失敗への対処だ。DBの状態は更新したのに、Redisへの通知が失敗したら?
解決策: すべての処理を**Phase 1(クリティカル)とPhase 2(ベストエフォート)**に分ける。
# app/workers/receipt_worker.py
async def _process_ocr(self, receipt, db):
# Phase 1: DB状態遷移 — 必ずコミット
receipt.start_ocr()
await db.commit()
file_data = await minio_service.download_file(receipt.file_path)
ocr_text = await ocr_service.extract_text_from_bytes(file_data)
receipt.complete_ocr(ocr_text)
await db.commit()
# Phase 2: 通知 — 失敗してもデータ整合性に影響なし
try:
await event_log_service.log(receipt.id, "OCR_COMPLETED")
await redis_service.update_progress(receipt.id, progress=50)
except Exception:
pass # RedisがダウンしてもDBの状態はすでに正常
Phase 1が失敗すると例外が上がり、Kafkaがメッセージを再配信する。Phase 2が失敗してもデータ整合性には問題ない。Taskレベルの例外が発生した場合、緊急セーフガードが新規DBセッションでレシートを強制的に失敗状態に遷移させる。
GPT Structured Output: トークン53%削減、スキーマ違反0%
最初は response_format={"type": "json_object"}(JSONモード)を使っていた。
問題:
- GPTがスキーマに違反するケースがあった (
total_priceを文字列で返すなど) - プロンプトにJSONスキーマ + 例3つをテキストで書く必要があり、〜1,800トークンを消費
Structured Outputに切り替えた:
# app/services/gpt_service.py
class ParsedReceiptSchema(BaseModel):
store_name: Optional[str] = None
store_address: Optional[str] = None
store_phone: Optional[str] = None
date: Optional[str] = None
time: Optional[str] = None
items: List[ReceiptItemSchema] # name, quantity, unit_price, total_price
subtotal: Optional[float] = None
tax: Optional[float] = None
total: float
payment_method: Optional[str] = None
card_number: Optional[str] = None
category: str
response = client.beta.chat.completions.parse(
model="gpt-4o-mini",
messages=[...],
temperature=0.1,
response_format=ParsedReceiptSchema, # PydanticモデルがそのままスキーマになるPydantic
)
parsed: ParsedReceiptSchema = response.choices[0].message.parsed
結果:
| 指標 | JSONモード | Structured Output |
|---|---|---|
| プロンプトトークン | 〜1,800 | 〜850 (53%削減) |
| スキーマ違反率 | 〜5% | 0% |
| エラー処理コード |
json.loads() + try/except
|
不要 |
Pydanticモデルがそのままスキーマになるため、プロンプトからスキーマの説明と例3つをすべて省ける。ルール8行だけを残した。json.loads() のパースエラーハンドリングも消えた。
ただし、GPTの出力をそのまま信頼するわけではない。Structured Outputが保証するのは形式であって、値の正確さではないからだ。ドメインのValue Objectでもう一段階検証する:
raw_parsed = await gpt_service.parse_receipt(ocr_text, categories) # Structured Output
raw_parsed = await gpt_service.validate_parsed_data(raw_parsed) # VO正規化
parsed_vo = ParsedReceiptVO.from_dict(raw_parsed) # ドメイン不変条件の検証
# → unit_price × quantity = total_price のクロスチェック
# → アイテム合計 ≈ total (±1円の誤差を許容)
リアルタイム進捗: Redis Pub/Sub → WebSocket
処理が非同期のため、ユーザーに進行状況を伝える必要がある。
@router.websocket("/ws/{receipt_id}")
async def websocket_receipt_updates(websocket: WebSocket, receipt_id: str):
await manager.connect(websocket, receipt_id)
# 接続直後に現在の状態を送信 — 再接続時の空白を補完
current_status = await redis_service.get_status(receipt_id)
if current_status:
await websocket.send_json(current_status)
# Redis Pub/Sub を購読 → クライアントにフォワード
pubsub_task = asyncio.create_task(
_listen_pubsub_and_forward(websocket, redis_service, receipt_id)
)
フロントエンドでは指数バックオフで再接続(1s → 2s → 4s → 8s → 16s)し、切断されても自動復旧する。再接続時にサーバーが現在の状態を即座に送信するため、切断中に見逃した更新も補完される。
結果
| 指標 | Before | After |
|---|---|---|
| APIレスポンス時間 | 最大15秒 (同期処理) | < 100ms (Kafka非同期) |
| プロンプトトークン | 〜1,800 | 〜850 (53%削減) |
| GPTスキーマ違反率 | 〜5% | 0% |
| レシート処理時間 | - | 8〜15秒 (OCR 3〜5秒 + 解析 5〜10秒) |
| 障害復旧 | 手動 | Worker再起動時に自動再処理 |
学んだこと
- ステートマシンは過剰ではない。 非同期パイプラインで「今どの状態にあるか」を明確に把握することが、デバッグの半分を占める。遷移テーブル1つですべてのルールが整理され、ドメイン層に置けば外部依存なしでユニットテストができる。
- 2フェーズコミットの核心は分類だ。 DBの状態遷移はクリティカル(失敗時はロールバック)、Redisへの通知はベストエフォート(失敗は無視)。この分類さえ正しくできていれば、Redisがダウンしてもデータは安全だ。
- Structured OutputはJSONモードの完全な上位互換だ。 スキーマ説明に使っていたトークンが53%削減され、パースエラーのハンドリングコードが消える。ただし「形式の保証 ≠ 値の正確さの保証」なので、ドメイン検証は依然として必要だ。
- プログレスバー1本がUXを変える。 「処理中です…」の1行 vs リアルタイムのプログレスバーでは体感がまったく異なる。Redis Pub/Sub + WebSocketのコードを追加するコストは、十分に正当化される。





