eBay Inventory Mapping API②:タスク完了をポーリングしてプレビュー結果を安全に取得する

前回の記事はこちら

【連載#19】eBay Inventory Mapping API②:タスク完了をポーリングしてプレビュー結果を安全に取得する

はじめに

本記事は、全42回にわたる「eBay API 実践ガイド」の第19回です。

前回(#18)では、GraphQL mutation の startListingPreviewsCreation を呼び出し、商品データを投入してタスク ID を取得するところまでを実装しました。しかし、ここで多くの開発者が立ち止まります。「タスク ID をもらったはいいが、結果はどうやって受け取るのか? いつ完了するのか?」——本記事はまさにその問いに答えます。

Inventory Mapping API の処理は AI による推論を伴うため、同期的に結果が返ってきません。listingPreviewsCreationTaskById という Query を一定間隔で呼び出し(ポーリング)、タスクが完了するまで待つ設計が必要です。ただし「ただ繰り返し呼べばいい」というものでもなく、レート制限・タイムアウト・エラー判別という3つの壁が待ち構えています。

この記事で得られること:

  • listingPreviewsCreationTaskById の result フィールドを使った「null 判定パターン」の完全理解——IN_PROGRESS のような enum は存在しない、という eBay API 特有の設計思想を把握できます。
  • COMPLETED / COMPLETED_WITH_ERROR / FAILED という 3 種類の completionStatus と、それに付随する invalidProducts・unmappedProducts・unprocessedProducts の違いを実務レベルで理解できます。
  • 指数バックオフ付きポーリング・タイムアウト制御・エラー別ハンドリングを備えた、本番稼働に耐えるプロダクションレベルの Python 実装を習得できます。

背景・なぜこれが重要か (Motivation)

「GraphQL の mutation を呼んだら、その場で結果も返ってくるんじゃないの?」

Inventory Mapping API を初めて触る開発者のほぼ全員が、最初にこの疑問を抱きます。確かに、多くの GraphQL API はリクエストに対してそのまま結果を同期的に返します。しかし eBay の Inventory Mapping API は根本的に異なる設計を採用しています。

なぜ非同期なのでしょうか。Inventory Mapping API が内部で行っていることを考えると納得できます。あなたが投入した商品データ(タイトル・説明・外部商品 ID など)に対して、eBay の AI エンジンが「この商品は eBay カタログの何に対応するか」「最適なカテゴリはどこか」「Item Specifics は何を付けるべきか」という推論処理を走らせています。数十件であれば数秒で終わることもありますが、数百件のバッチや複雑な商品では数分から最長 10 分程度かかることもあります。これを同期的に処理しようとすると、HTTP 接続がタイムアウトしてしまいます。

そのため、eBay は「タスクを受け付けた」という証明としてタスク ID(前回取得)を即時返却し、実際の処理は非同期で進める設計を採用しています。クライアント側はそのタスク ID を使って「もう終わったか?」と定期的に問い合わせる——これがポーリングパターンです。

補足: eBay の「null 判定パターン」について

処理中かどうかを判定するために、eBay は IN_PROGRESS のような明示的な enum 値を使いません。その代わり、result フィールドそのものが null かどうかで状態を表現します。result が null → まだ処理中。result に値が入っている → 処理完了(成否は completionStatus で判断)。このパターンを知らずにコードを書くと、null チェックを忘れてデータをパースしようとして AttributeError や TypeError に悩まされることになります。

基本的な使い方(ベースライン):listingPreviewsCreationTaskById で結果を取得する

まず最小限のポーリングループを実装します。GraphQL の Query を requests で送り、result フィールドが null かどうかをチェックして、値が返ってくるまでループします。

以下が最小実装のコードです。

# polling_baseline.py import time import requests

GRAPHQL_URL = "https://api.ebay.com/sell/listing_preview/graphql" QUERY = """
query GetTask($id: ID!) {
  listingPreviewsCreationTaskById(input: { id: $id }) {
    requestedId
    listingPreviewsCreationTask {
      id
      result {
        completionStatus
        listingPreviews {
          sku
          title
          mappingReferenceId
          category { id }
        }
        invalidProducts {
          externalProductId
          errors { message }
        }
        unmappedProducts {
          externalProductId
        }
        unprocessedProducts {
          externalProductId
        }
      }
    }
  }
}
""" def poll_task(task_id: str, access_token: str, interval: int = 10) -> dict: """
    タスクが完了するまでポーリングし、result を返す(最小実装)。
    interval: ポーリング間隔(秒)
    """ headers = { "Authorization": f"Bearer {access_token}", "Content-Type": "application/json",
    }
    variables = {"id": task_id} while True:
        resp = requests.post(
            GRAPHQL_URL,
            json={"query": QUERY, "variables": variables},
            headers=headers,
            timeout=30,
        )
        resp.raise_for_status()
        data = resp.json()

        task = (
            data.get("data", {})
            .get("listingPreviewsCreationTaskById", {})
            .get("listingPreviewsCreationTask")
        ) if task is None: raise ValueError(f"タスクが見つかりません: {task_id}") # null 判定パターン: result が None なら処理中 result = task.get("result") if result is not None: return result # 完了(status は呼び出し側で判定) print(f"[polling] タスク {task_id} は処理中です。{interval}秒後に再確認します...")
        time.sleep(interval) # 使用例 if __name__ == "__main__":
    TASK_ID = "xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx" # 第18回で取得したタスクID TOKEN   = "v^1.1.xxxx..." result = poll_task(TASK_ID, TOKEN)
    print(f"completionStatus: {result['completionStatus']}")
    print(f"listingPreviews 件数: {len(result['listingPreviews'])}")
補足: completionStatus の 3 種類の値

result が返ってきたら、その中の completionStatus を確認します。取り得る値は 3 種類です。
【COMPLETED】すべての商品が正常にマッピングされ、listingPreviews にプレビューが格納されています。
【COMPLETED_WITH_ERROR】一部の商品は正常にマッピングされましたが、一部の商品に問題がありました。listingPreviews と問題商品リスト(invalidProducts / unmappedProducts / unprocessedProducts)が混在します。
【FAILED】タスク全体が失敗しました。listingPreviews は空です。

この上記の最小実装には、実務で致命的になる欠陥がいくつかあります。次のセクションで詳しく解説します。

実務で躓く場面・深いポイント (Core)

ベースライン実装を本番環境に持ち込むと、必ずいくつかの壁にぶつかります。ここでは、現場で頻出する3つの落とし穴とその回避策を解説します。

1. 固定間隔の無限ループはレート制限を引き起こす

上記の最小実装では、interval=10(秒)で固定間隔のポーリングを行っています。一見問題なさそうですが、次のシナリオを考えてみてください。「10件の商品を並行して処理する5本のポーリングループを、1サーバーで同時実行している」——この場合、1分間に 5本 × 6回 = 30 回の API コールが発生します。商品数が増えてバッチサイズが大きくなると、あっという間に eBay の API レート制限に抵触します。

さらに深刻な問題は、「タスクが高速で完了しそうな場合」です。商品数が少なければ数秒で完了することもありますが、固定間隔だと最大 interval 秒のロスが生じます。逆に「処理が長引いているとき」に同じ短い間隔でポーリングし続けるのは明らかな無駄です。

解決策は「指数バックオフ(Exponential Backoff)」です。最初は短い間隔(例: 5秒)でポーリングを開始し、完了していなければ間隔を徐々に延ばしていきます(10秒 → 20秒 → 最大 60秒)。これにより、高速完了のタスクへの応答性を保ちながら、長時間かかるタスクへの無駄なコールを削減できます。

2. COMPLETED_WITH_ERROR を見落として不完全なデータで後続処理を進める罠

Inventory Mapping API 実装で最も多いバグは「completionStatus の確認漏れ」です。具体的には、result が返ってきた瞬間に「完了 = 成功」と判断して listingPreviews をそのまま処理してしまうパターンです。

例えば 100件の商品を投入し、80件は正常にマッピングされたが 20件は問題があった場合、completionStatus は COMPLETED_WITH_ERROR になります。この状態で listingPreviews だけを見ると 80件は問題なく見えます。しかし残りの 20件がどうなったのかを確認しないまま Inventory API に渡すと、「投入した 100件のうち 80件しか出品されていない」というサイレントバグが生まれます。このようなバグは検知が難しく、実務では非常に発見が遅れます。

必ず completionStatus を明示的に確認し、COMPLETED_WITH_ERROR の場合は問題商品を別途ログ・再試行キューに振り分ける処理を実装してください。

3. invalidProducts / unmappedProducts / unprocessedProducts の区別を誤ると原因調査が長引く

COMPLETED_WITH_ERROR のとき、問題のあった商品は 3 種類のリストに振り分けられます。それぞれの意味を正確に理解しないと、「なぜこの商品がマッピングされなかったのか」の原因調査に無駄な時間がかかります。

【invalidProducts】:投入したデータ自体に不備がある商品です。例として、externalProductId が空、必須フィールドが欠落している、フォーマットが不正などが該当します。この場合の対処は「入力データの修正」であり、同じデータを再送しても同じ結果になります。errors フィールドに具体的なエラーメッセージが入っているため、必ず参照してください。

【unmappedProducts】:入力データ自体は有効だが、eBay カタログに対応する商品が見つからなかった商品です。つまり「データは問題ないが、AI が eBay の商品カタログと照合できなかった」状態です。対処としては、タイトルや説明を補強して再試行する、または手動でカテゴリ・Item Specifics を指定する方法があります。

【unprocessedProducts】:システム的な理由(タイムアウト、内部エラーなど)で処理が完了しなかった商品です。データに問題があるわけではないため、そのまま再試行(リトライ)することで成功する可能性が高いです。これを invalidProducts と混同して「データを修正」しようとすると、無駄な作業が発生します。

注意: タイムアウト設計について

AI による推論処理は、商品数・複雑さによって処理時間が大きく変動します。eBay の公式ドキュメントでは最長 10 分程度かかる可能性が示唆されています。そのため、ポーリングには必ず上限時間(タイムアウト)を設定してください。推奨は 15〜20 分(余裕を持って設定)。タイムアウトした場合は、タスク ID を記録した上でエラーとして扱い、アラートを発報して運用チームが確認できるようにしてください。タイムアウト後のタスクが実は完了していた場合でも、タスク ID は有効なため後から結果を取得できます。

頻出エラーコード早見表

ポーリング中に発生しうる HTTP / GraphQL エラーと対処法をまとめます。

HTTP ステータス / エラー内容 原因と対処法
HTTP 401 Unauthorized アクセストークンの有効期限切れ。OAuth 2.0 のトークンリフレッシュを行い、新しいトークンで再試行してください。
HTTP 429 Too Many Requests API レート制限に抵触。ポーリング間隔を延ばす(指数バックオフ)か、並行ポーリング数を削減してください。
HTTP 500 / 503 eBay 側の一時的なサーバーエラー。最大3回まで指数バックオフでリトライ。継続する場合は eBay Developer Support に連絡。
GraphQL: "Task not found" タスク ID が無効または別の OAuth スコープのトークンを使用している。スコープを確認し、正しいトークンを使用してください。
GraphQL: "Access denied" 必要な OAuth スコープ(sell.listing_preview)がトークンに含まれていない。

堅牢な実装:指数バックオフとエラー別ハンドリングを備えたポーリングエンジン

上記の3つの落とし穴をすべてクリアした、本番稼働に耐えるポーリング実装を示します。型アノテーション・docstring・例外処理・入力バリデーション・指数バックオフ・タイムアウト・completionStatus 別の結果ハンドリングを完備しています。

# inventory_mapping_poller.py """
Inventory Mapping API タスクポーリングエンジン。
指数バックオフ・タイムアウト・completionStatus 別ハンドリングを実装。

Usage:
    poller = InventoryMappingPoller(access_token=os.environ["EBAY_ACCESS_TOKEN"])
    outcome = poller.poll(task_id="xxxxxx-xxxx-xxxx")
    if outcome.is_success:
        for preview in outcome.listing_previews:
            print(preview["sku"], preview["title"])
""" from __future__ import annotations import logging import time from dataclasses import dataclass, field from typing import Any import requests

logger = logging.getLogger(__name__)

GRAPHQL_URL = "https://api.ebay.com/sell/listing_preview/graphql" QUERY = """
query GetTaskById($id: ID!) {
  listingPreviewsCreationTaskById(input: { id: $id }) {
    requestedId
    listingPreviewsCreationTask {
      id
      result {
        completionStatus
        listingPreviews {
          sku
          title
          mappingReferenceId
          description
          category { id }
          images { value }
          aspects {
            name
            values
          }
        }
        invalidProducts {
          externalProductId
          errors { message errorId }
        }
        unmappedProducts {
          externalProductId
          title
        }
        unprocessedProducts {
          externalProductId
          title
        }
      }
    }
  }
}
""" @dataclass class PollOutcome: """ポーリング結果を格納するデータクラス。""" task_id: str
    completion_status: str # COMPLETED / COMPLETED_WITH_ERROR / FAILED listing_previews: list[dict[str, Any]] # 正常にマッピングされた商品 invalid_products: list[dict[str, Any]] # 入力データ不備 unmapped_products: list[dict[str, Any]] # カタログ照合失敗 unprocessed_products: list[dict[str, Any]] # システムエラーで未処理 timed_out: bool = False error_message: str = "" @property def is_success(self) -> bool: """COMPLETED のみ True(COMPLETED_WITH_ERROR は False)。""" return self.completion_status == "COMPLETED" @property def has_partial_results(self) -> bool: """一部成功・一部失敗の場合 True。""" return self.completion_status == "COMPLETED_WITH_ERROR" @property def problem_count(self) -> int: """問題のあった商品の合計件数。""" return (
            len(self.invalid_products)
            + len(self.unmapped_products)
            + len(self.unprocessed_products)
        )


@dataclass class PollerConfig: """ポーリング設定。""" initial_interval: float = 5.0 # 初回待機時間(秒) backoff_factor: float = 1.8 # 間隔の乗数 max_interval: float = 60.0 # 最大ポーリング間隔(秒) timeout_seconds: float = 900.0 # タイムアウト(15分) max_server_retries: int = 3 # HTTP 5xx 連続エラーの最大リトライ数 class InventoryMappingPoller: """
    listingPreviewsCreationTaskById を指数バックオフでポーリングし、
    タスク完了後に PollOutcome を返すポーリングエンジン。
    """ def __init__(self, access_token: str, config: PollerConfig | None = None) -> None: if not access_token: raise ValueError("access_token が空です。OAuth 2.0 トークンを設定してください。")
        self._token = access_token
        self._config = config or PollerConfig() def poll(self, task_id: str) -> PollOutcome: """
        タスクが完了するまでポーリングを実行し、PollOutcome を返す。

        Args:
            task_id: 第18回で取得した startListingPreviewsCreation のタスク ID。

        Returns:
            PollOutcome: completionStatus 別の結果データ。

        Raises:
            ValueError: task_id が空の場合。
            RuntimeError: 最大リトライ数を超えたサーバーエラーが発生した場合。
        """ if not task_id or not task_id.strip(): raise ValueError("task_id が空です。")

        cfg = self._config
        interval = cfg.initial_interval
        start_time = time.monotonic()
        server_error_count = 0 logger.info("ポーリング開始: task_id=%s, timeout=%ss", task_id, cfg.timeout_seconds) while True: # タイムアウト判定 elapsed = time.monotonic() - start_time if elapsed >= cfg.timeout_seconds:
                logger.warning("タイムアウト: task_id=%s (%.0fs経過)", task_id, elapsed) return PollOutcome(
                    task_id=task_id,
                    completion_status="",
                    listing_previews=[],
                    invalid_products=[],
                    unmapped_products=[],
                    unprocessed_products=[],
                    timed_out=True,
                    error_message=f"タイムアウト({cfg.timeout_seconds}秒経過)。タスクIDを記録して後から再確認してください。",
                ) try:
                result = self._call_api(task_id)
                server_error_count = 0 # 成功したらリセット except requests.HTTPError as exc:
                status_code = exc.response.status_code if exc.response else 0 if status_code == 401: # トークン切れは再試行しても意味がないのですぐ終了 raise RuntimeError( "HTTP 401: アクセストークンが無効または期限切れです。" "トークンをリフレッシュして再実行してください。" ) from exc if status_code == 429: # レート制限: 最大 interval の 2 倍待つ wait = min(interval * 2, cfg.max_interval)
                    logger.warning("HTTP 429 レート制限。%.0f秒待機します...", wait)
                    time.sleep(wait) continue if status_code >= 500:
                    server_error_count += 1 if server_error_count > cfg.max_server_retries: raise RuntimeError( f"HTTP {status_code} が {cfg.max_server_retries} 回連続しました。" "eBay サーバー側の問題の可能性があります。" ) from exc
                    logger.warning( "HTTP %s サーバーエラー(%d/%d回目)。リトライします...",
                        status_code, server_error_count, cfg.max_server_retries,
                    )
                    time.sleep(interval) continue raise # その他の HTTP エラーは再スロー # GraphQL エラーチェック if "errors" in result:
                messages = [e.get("message", "") for e in result["errors"]] raise RuntimeError(f"GraphQL エラー: {'; '.join(messages)}")

            task_data = (
                result.get("data", {})
                .get("listingPreviewsCreationTaskById", {})
                .get("listingPreviewsCreationTask")
            ) if task_data is None: raise RuntimeError( f"タスクが見つかりません(task_id={task_id})。" "IDが正しいか、同じ OAuth スコープのトークンを使っているか確認してください。" ) # null 判定パターン: result が None なら処理中 task_result = task_data.get("result") if task_result is None:
                logger.info( "処理中: task_id=%s (経過 %.0fs) → %.0f秒後に再確認",
                    task_id, elapsed, interval,
                )
                time.sleep(interval) # 指数バックオフで interval を更新 interval = min(interval * cfg.backoff_factor, cfg.max_interval) continue # --- タスク完了 --- status = task_result.get("completionStatus", "FAILED")
            previews    = task_result.get("listingPreviews", []) or []
            invalid     = task_result.get("invalidProducts", []) or []
            unmapped    = task_result.get("unmappedProducts", []) or []
            unprocessed = task_result.get("unprocessedProducts", []) or []

            outcome = PollOutcome(
                task_id=task_id,
                completion_status=status,
                listing_previews=previews,
                invalid_products=invalid,
                unmapped_products=unmapped,
                unprocessed_products=unprocessed,
            )

            self._handle_outcome(outcome) return outcome def _call_api(self, task_id: str) -> dict[str, Any]: """GraphQL API を 1 回呼び出して生のレスポンス dict を返す。""" headers = { "Authorization": f"Bearer {self._token}", "Content-Type": "application/json",
        }
        response = requests.post(
            GRAPHQL_URL,
            json={"query": QUERY, "variables": {"id": task_id}},
            headers=headers,
            timeout=30,
        )
        response.raise_for_status() return response.json() def _handle_outcome(self, outcome: PollOutcome) -> None: """completionStatus に応じてログ出力・問題商品の分類を行う。""" status = outcome.completion_status if status == "COMPLETED":
            logger.info( "COMPLETED: task_id=%s | listingPreviews=%d件",
                outcome.task_id, len(outcome.listing_previews),
            ) elif status == "COMPLETED_WITH_ERROR":
            logger.warning( "COMPLETED_WITH_ERROR: task_id=%s | 正常=%d件 / 問題=%d件",
                outcome.task_id,
                len(outcome.listing_previews),
                outcome.problem_count,
            ) # invalidProducts: 入力データ不備 → 修正が必要 for item in outcome.invalid_products:
                errors = [e.get("message", "") for e in (item.get("errors") or [])]
                logger.error( "  [INVALID] externalProductId=%s | errors=%s",
                    item.get("externalProductId"), errors,
                ) # unmappedProducts: カタログ照合失敗 → タイトル補強 or 手動設定 for item in outcome.unmapped_products:
                logger.warning( "  [UNMAPPED] externalProductId=%s | title=%s",
                    item.get("externalProductId"), item.get("title", ""),
                ) # unprocessedProducts: システムエラーで未処理 → そのまま再試行 for item in outcome.unprocessed_products:
                logger.warning( "  [UNPROCESSED] externalProductId=%s → リトライキューへ追加",
                    item.get("externalProductId"),
                ) elif status == "FAILED":
            logger.error( "FAILED: task_id=%s | listingPreviews は空です。" "投入データを確認して再実行してください。",
                outcome.task_id,
            ) else:
            logger.error("不明な completionStatus: %s (task_id=%s)", status, outcome.task_id) # ============================================================
# 実行例
# ============================================================ if __name__ == "__main__": import os import json

    logging.basicConfig(
        level=logging.INFO,
        format="%(asctime)s [%(levelname)s] %(message)s",
    )

    TOKEN   = os.environ["EBAY_ACCESS_TOKEN"]
    TASK_ID = os.environ["TASK_ID"] # 第18回で取得したタスクID poller  = InventoryMappingPoller(access_token=TOKEN)
    outcome = poller.poll(task_id=TASK_ID) if outcome.timed_out:
        print(f"タイムアウト: {outcome.error_message}") elif outcome.is_success:
        print(f"完全成功: {len(outcome.listing_previews)} 件のプレビューを取得しました。") # 第20回: Inventory API へ渡す処理をここに追加 elif outcome.has_partial_results:
        print(f"部分成功: {len(outcome.listing_previews)} 件成功 / {outcome.problem_count} 件に問題あり") # unprocessed を再試行キューへ retry_ids = [p.get("externalProductId") for p in outcome.unprocessed_products]
        print(f"再試行対象: {json.dumps(retry_ids, ensure_ascii=False)}") else:
        print("タスク全体が失敗しました。ログを確認してください。")

このポーリングエンジンの設計上の重要ポイントをまとめます。

【1】指数バックオフの実装: initial_interval=5秒から開始し、backoff_factor=1.8 で間隔を拡大(5 → 9 → 16 → 29 → 52 → 60秒でキャップ)。処理が速く終わるタスクへの応答性を保ちつつ、長時間処理時の無駄なコールを抑制します。

【2】タイムアウトのハードコードを避ける: PollerConfig の timeout_seconds をデフォルト 900秒(15分)に設定しつつ、呼び出し元が上書きできる設計にしています。処理規模に応じて柔軟に調整できます。

【3】completionStatus 別の分岐処理: _handle_outcome メソッドが各種問題商品を種別ごとに記録します。invalidProducts は修正対象、unprocessedProducts は即時リトライ対象、unmappedProducts は人手レビュー対象——それぞれの対処方針が異なることをコード上で明示しています。

【4】PollOutcome データクラスの活用: タスク ID・ステータス・プレビューリスト・問題商品リストをひとつのオブジェクトにまとめることで、後続処理(第20回の Inventory API への受け渡し)がシンプルに書けるようになります。

パフォーマンス・スケーリング視点 (深度)

1日数十件の商品を単発で処理するフェーズでは、上記のシンプルな同期ポーリングで十分です。しかし、商品カタログが数千件規模になり、複数のタスクを並行処理する必要が出てくると、設計を見直す必要があります。

複数タスクの並行ポーリング設計

例えば、1000件の商品を 100件ずつ 10個のタスクに分割して startListingPreviewsCreation に投入するケースを考えます。各タスクを順番にポーリングすると、最後のタスクの結果を受け取るまでに、最初のタスクの処理待ち時間 × タスク数 だけの余分な時間がかかってしまいます。

解決策は concurrent.futures.ThreadPoolExecutor を使った並行ポーリングです。以下のコードで、複数タスクを同時にポーリングできます。

# concurrent_polling.py import concurrent.futures import os import logging from inventory_mapping_poller import InventoryMappingPoller, PollerConfig, PollOutcome

logger = logging.getLogger(__name__) # 同時ポーリング数は3〜5が実務的な上限。
#    多すぎると HTTP 429 が頻発し、かえって遅くなります。 MAX_WORKERS = 4 def poll_all_tasks(
    task_ids: list[str],
    access_token: str,
) -> dict[str, PollOutcome]: """
    複数のタスクを並行してポーリングし、task_id -> PollOutcome の dict を返す。

    Args:
        task_ids:     第18回で取得したタスク ID のリスト。
        access_token: eBay OAuth 2.0 アクセストークン。

    Returns:
        dict[str, PollOutcome]: タスクIDをキーとした結果の辞書。
    """ poller = InventoryMappingPoller(
        access_token=access_token,
        config=PollerConfig(
            initial_interval=8.0, # 並行時は間隔を少し広げる backoff_factor=2.0,
            max_interval=60.0,
            timeout_seconds=1200.0, # 並行タスクが多い場合はタイムアウトを延ばす ),
    )
    results: dict[str, PollOutcome] = {} with concurrent.futures.ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor:
        future_to_id = {
            executor.submit(poller.poll, task_id): task_id for task_id in task_ids
        } for future in concurrent.futures.as_completed(future_to_id):
            task_id = future_to_id[future] try:
                outcome = future.result()
                results[task_id] = outcome
                logger.info("完了: %s → %s", task_id, outcome.completion_status) except Exception as exc:
                logger.error("エラー: %s → %s", task_id, exc) return results if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO,
                        format="%(asctime)s [%(levelname)s] %(message)s")

    TASK_IDS = [ "task-id-001", "task-id-002", "task-id-003",
    ]
    TOKEN = os.environ["EBAY_ACCESS_TOKEN"]

    all_results = poll_all_tasks(TASK_IDS, TOKEN)
    total_previews = sum(len(r.listing_previews) for r in all_results.values())
    print(f"合計 {len(all_results)} タスク完了 | 取得プレビュー数: {total_previews}")
注意: MAX_WORKERS を大きくしすぎないでください。

4〜5 を超えると HTTP 429(Too Many Requests)が頻発し、かえってスループットが低下します。実務では MAX_WORKERS=3〜4 から始めて、エラーレートを監視しながら調整してください。

補足: Notification API 連携による「プッシュ型」への移行展望

現状のポーリング設計はシンプルですが、大規模になると「無駄なコール」が増えます。将来的には eBay の Notification API(SetNotificationPreferences / 第15回参照)と組み合わせることで、タスク完了時に eBay 側からコールバックを受け取る「プッシュ型」アーキテクチャに移行できます。プッシュ型では、ポーリングコール自体が不要になるため、API コール数を大幅に削減できます。大規模オペレーション(1日 1,000 件超)を目指す場合は、この移行を検討してください。

まとめ

本記事では、Inventory Mapping API の非同期タスク処理において、最も重要な「結果の取得」を実装しました。

  • ベースライン: listingPreviewsCreationTaskById の result フィールドを使った null 判定パターンを理解しました。result が null なら処理中、値が入ったら completionStatus で成否を判断——この設計思想が eBay 非同期 API の基本です。
  • 深いポイント: 固定間隔ポーリングによるレート制限の罠、COMPLETED_WITH_ERROR の見落とし、そして invalidProducts / unmappedProducts / unprocessedProducts という3種類の問題商品リストの区別——これらを正確に理解することで、原因不明のサイレントバグを予防できます。
  • スケーリング: 指数バックオフ付きのポーリングエンジンに加え、ThreadPoolExecutor を使った複数タスクの並行ポーリング設計を習得しました。将来的には Notification API 連携によるプッシュ型への移行で、さらなる効率化が可能です。

これで、Inventory Mapping API からプレビュー結果を安全かつ効率的に取得する基盤が整いました。次はいよいよ、取得したプレビューデータを活用して、実際に eBay に商品を出品するステップです。

次のステップ

今回取得した listingPreviews には、AI が推奨するカテゴリ・Item Specifics・タイトルが格納されています。次回(#20)は、このプレビューデータを Inventory API(createOrReplaceInventoryItem / publishOffer)に渡し、高品質な出品を自動作成する完全な自動出品パイプラインを構築します。Inventory Mapping API → ポーリング → Inventory API という一連のフローが完成する、この連載のひとつの到達点となります。お楽しみに!

トップに戻る