0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

レシートOCRパイプラインを設計して学んだこと

0
Posted at

レシートの写真1枚を構造化データに変換するパイプラインを作った。APIのレスポンスを15秒から100msに短縮し、GPTのプロンプトトークンを53%削減し、スキーマ違反率を0%にするまでの過程。

問題の背景

コア機能はシンプルだ: レシートの写真をアップロードすると、構造化データが返ってくる。

しかしこのフローの中に、厄介な要件が潜んでいた:

  1. OCRに3〜5秒、GPTの解析に5〜10秒 → APIで同期処理すると1リクエストあたり最大15秒のブロッキング
  2. OCRが失敗しても、解析は別途リトライできなければならない
  3. ユーザーに処理の進行状況をリアルタイムで表示する必要がある
  4. AIの結果を人間がレビューして修正できなければならない

全体アーキテクチャ

image.png

ステートマシンでレシートのライフサイクルを管理する

レシートの状態を明示的に定義し、許可された遷移のみを可能にした。

image.png

なぜステートマシンなのか? 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を返すようになった。

image.png

Workerが落ちても、Kafkaがメッセージを保持する。再起動すれば未処理のメッセージから再消費される。

Kafkaはこの規模には過剰であることは認める。Redis Streamで十分だったかもしれない。しかしConsumer Groupによる水平スケールを見据えて選択した。Workerを複数台に増やせば、スループットは線形に伸びる。

2フェーズコミット: 「何がクリティカルか」を区別する

Workerで最も難しいのは部分的な失敗への対処だ。DBの状態は更新したのに、Redisへの通知が失敗したら?

解決策: すべての処理を**Phase 1(クリティカル)とPhase 2(ベストエフォート)**に分ける。

image.png

# 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に切り替えた:

image.png

# 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

処理が非同期のため、ユーザーに進行状況を伝える必要がある。

image.png

@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. ステートマシンは過剰ではない。 非同期パイプラインで「今どの状態にあるか」を明確に把握することが、デバッグの半分を占める。遷移テーブル1つですべてのルールが整理され、ドメイン層に置けば外部依存なしでユニットテストができる。
  2. 2フェーズコミットの核心は分類だ。 DBの状態遷移はクリティカル(失敗時はロールバック)、Redisへの通知はベストエフォート(失敗は無視)。この分類さえ正しくできていれば、Redisがダウンしてもデータは安全だ。
  3. Structured OutputはJSONモードの完全な上位互換だ。 スキーマ説明に使っていたトークンが53%削減され、パースエラーのハンドリングコードが消える。ただし「形式の保証 ≠ 値の正確さの保証」なので、ドメイン検証は依然として必要だ。
  4. プログレスバー1本がUXを変える。 「処理中です…」の1行 vs リアルタイムのプログレスバーでは体感がまったく異なる。Redis Pub/Sub + WebSocketのコードを追加するコストは、十分に正当化される。
0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?