TECH JOURNAL

Shopify Bulk Operation API で 10 万件商品を高速同期する設計

ShareXB!
Shopify Bulk Operation API で 10 万件商品を高速同期する設計
目次

はじめに

Shopify で数万〜数十万件の商品を扱うストアの在庫・価格を社内システムへ同期しようとすると、通常の REST / GraphQL API のレート制限にすぐ突き当たる。

  • REST API は 1 秒 2 リクエスト(バースト 40)が上限
  • Admin GraphQL API はコスト計算式があり、Products + Variants をネストして取得するだけで 1 クエリのコストが数百になる
  • 10 万件規模では何時間もかかる、あるいはスロットリングで中断する

Shopify が提供する Bulk Operation API は、こうした大量データの読み書きをバックグラウンドジョブとして処理するための仕組みだ。ジョブを登録→完了を確認→JSONL ファイルをダウンロード、という非同期フローにより、レート制限を気にせず全件取得・全件書き込みができる。

Bulk Operation の基本

Bulk Operation には 2 種類ある。

種別ミューテーション用途
読み取りbulkOperationRunQuery商品・注文の全件エクスポート
書き込みbulkOperationRunMutation在庫・価格の一括 upsert

制約として、ストアごとに同時実行できるジョブは 1 つだけ。新しいジョブを登録する前に、既存ジョブが COMPLETED / FAILED / CANCELED になっているかを必ず確認する。

実装パターン 1: bulkOperationRunQuery で全件エクスポート

商品の全件取得は bulkOperationRunQuery にクエリを渡してジョブを起動する。

mutation {
  bulkOperationRunQuery(
    query: """
    {
      products {
        edges {
          node {
            id
            title
            updatedAt
            variants {
              edges {
                node {
                  id
                  sku
                  price
                  inventoryQuantity
                }
              }
            }
          }
        }
      }
    }
    """
  ) {
    bulkOperation {
      id
      status
    }
    userErrors {
      field
      message
    }
  }
}

ジョブ登録後に受け取る bulkOperation.id(例: gid://shopify/BulkOperation/123456)を保存しておく。完了確認のポーリングに使う。

実装パターン 2: currentBulkOperation ポーリングと指数バックオフ

ジョブ完了は currentBulkOperation クエリで確認する。完了したら url フィールドに JSONL ファイルの一時ダウンロード URL が入る。

async function pollBulkOperation(jobId: string): Promise<string> {
  const query = `
    query {
      currentBulkOperation {
        id
        status
        errorCode
        objectCount
        url
      }
    }
  `
 
  let delay = 5_000 // 最初は 5 秒待つ
  const MAX_DELAY = 120_000 // 最大 2 分
 
  while (true) {
    await sleep(delay)
 
    const { data } = await shopify.graphql(query)
    const op = data.currentBulkOperation
 
    // 別のジョブが割り込んでいる場合は id が変わる
    if (op.id !== jobId) {
      throw new Error(`Unexpected job in progress: ${op.id}`)
    }
 
    if (op.status === "COMPLETED") {
      if (!op.url) throw new Error("COMPLETED but url is null")
      return op.url
    }
 
    if (op.status === "FAILED" || op.status === "CANCELED") {
      throw new Error(`Bulk operation ${op.status}: ${op.errorCode}`)
    }
 
    // RUNNING / CANCELING → 待機継続
    delay = Math.min(delay * 1.5, MAX_DELAY) // 指数バックオフ
  }
}
 
function sleep(ms: number) {
  return new Promise((resolve) => setTimeout(resolve, ms))
}

ステータスは CREATED → RUNNING → COMPLETED / FAILED / CANCELED の順に遷移する。FAILED 時の errorCode は ACCESS_DENIED / TIMEOUT / INTERNAL_SERVER_ERROR などがある。TIMEOUT はクエリが複雑すぎる場合に発生するので、クエリを分割して再試行する。

実装パターン 3: JSONL のストリーム処理

完了後の url は AWS S3 の署名付き URL で、有効期限は約 1 時間。JSONL(1 行 = 1 オブジェクト)を行単位でストリーム処理することで、数百 MB になるファイルもメモリに乗せずに処理できる。

JSONL の構造には注意が必要。ProductVariant は対応する Product の 後の行に続く形で出力され、__parentId フィールドで親を参照している。

import { createInterface } from "readline"
import { Readable } from "stream"
 
interface ProductNode {
  id: string
  title: string
  updatedAt: string
}
 
interface VariantNode {
  id: string
  sku: string
  price: string
  inventoryQuantity: number
  __parentId: string // 親 Product の id
}
 
async function processBulkResult(url: string): Promise<void> {
  const response = await fetch(url)
  if (!response.ok || !response.body) {
    throw new Error(`Failed to fetch JSONL: ${response.status}`)
  }
 
  const productMap = new Map<string, ProductNode>()
  const rl = createInterface({
    input: Readable.fromWeb(response.body as any),
    crlfDelay: Infinity,
  })
 
  for await (const line of rl) {
    if (!line.trim()) continue
    const obj = JSON.parse(line)
 
    if (obj.__parentId) {
      // Variant 行 → 親 Product と紐づけて処理
      const variant = obj as VariantNode
      const product = productMap.get(variant.__parentId)
      if (product) {
        await upsertVariant(product, variant)
      }
    } else {
      // Product 行 → Map に保持(後続 Variant のために)
      const product = obj as ProductNode
      productMap.set(product.id, product)
    }
  }
}

10 万件超の場合は productMap のメモリ使用量が増える。Product 行を読んだ直後にバッファに書き出し、処理済み Product を Map から削除する実装にすると安全だ。

実装パターン 4: bulkOperationRunMutation で在庫・価格を一括書き込み

書き込み方向(社内システム → Shopify)は bulkOperationRunMutation を使う。事前に JSONL ファイルを Staged Upload に送信し、そのファイルを参照してジョブを起動する流れになる。

// Step 1: Staged Upload を作成(POST 形式で受け取る)
const stagedUploadMutation = `
  mutation {
    stagedUploadsCreate(input: {
      resource: BULK_MUTATION_VARIABLES
      filename: "variants.jsonl"
      mimeType: "text/jsonl"
      httpMethod: POST
    }) {
      stagedTargets {
        url
        parameters { name value }
      }
      userErrors { field message }
    }
  }
`
const { data: uploadData } = await shopify.graphql(stagedUploadMutation)
const target = uploadData.stagedUploadsCreate.stagedTargets[0]
const params = target.parameters as { name: string; value: string }[]
 
// Step 2: 返ってきた parameters とファイルを multipart/form-data で POST
const jsonl = variantsToUpsert
  .map((v) =>
    JSON.stringify({
      input: {
        id: v.shopifyVariantId,
        price: v.price,
        inventoryQuantities: {
          availableQuantity: v.quantity,
          locationId: LOCATION_GID,
        },
      },
    })
  )
  .join("\n")
 
const form = new FormData()
// parameters を先に append し、file は必ず最後に append する
for (const param of params) {
  form.append(param.name, param.value)
}
form.append("file", new Blob([jsonl], { type: "text/jsonl" }), "variants.jsonl")
await fetch(target.url, { method: "POST", body: form })
 
// Step 3: ジョブ登録
// stagedUploadPath には resourceUrl ではなく、parameters の "key"
// (アップロード済みオブジェクトのパス)を渡す
const stagedUploadPath = params.find((p) => p.name === "key")!.value
 
const bulkMutation = `
  mutation bulkRunMutation($stagedUploadPath: String!) {
    bulkOperationRunMutation(
      mutation: "mutation updateVariant($input: ProductVariantInput!) { productVariantUpdate(input: $input) { productVariant { id price } userErrors { field message } } }"
      stagedUploadPath: $stagedUploadPath
    ) {
      bulkOperation { id status }
      userErrors { field message }
    }
  }
`
await shopify.graphql(bulkMutation, { stagedUploadPath })

書き込みジョブの完了後も currentBulkOperation でポーリングし、url に返ってくる JSONL を確認する。userErrors が入っている行があれば、個別にリトライまたはアラートを飛ばす。

実装パターン 5: SKU での冪等 upsert と Webhook との組み合わせ

Bulk Operation でのフル同期は時間がかかる(規模によっては 10〜30 分)。日常的な在庫・価格の鮮度を維持するには Webhook での差分更新 を主系に据え、Bulk Operation は週次・日次のバックフィルとして使うアーキテクチャが現実的だ。

Webhook (products/update) ─── 差分同期 ──→ 社内在庫DB (鮮度優先)
                                                   ↑
Bulk Operation (夜間バッチ) ─ フル同期 ──────────┘ (欠損補完)

冪等性の担保は SKU を一意キーとした upsert が基本。Shopify の productVariantId は削除・再作成で変わることがあるが、SKU は業務側でコントロールできる。

async function upsertVariant(product: ProductNode, variant: VariantNode) {
  await db
    .insertInto("variants")
    .values({
      sku: variant.sku,
      shopify_variant_id: variant.id,
      shopify_product_id: product.id,
      price: variant.price,
      inventory_quantity: variant.inventoryQuantity,
      synced_at: new Date().toISOString(),
    })
    .onConflict((oc) =>
      oc.column("sku").doUpdateSet({
        shopify_variant_id: (eb) => eb.ref("excluded.shopify_variant_id"),
        price: (eb) => eb.ref("excluded.price"),
        inventory_quantity: (eb) => eb.ref("excluded.inventory_quantity"),
        synced_at: (eb) => eb.ref("excluded.synced_at"),
      })
    )
    .execute()
}

ON CONFLICT (sku) DO UPDATE により、同じ SKU を何度 upsert しても最終的な状態は同じになる。Bulk Operation のリラン時も安全に再実行できる。

まとめ

課題対策
REST / GraphQL のレート制限Bulk Operation API で非同期ジョブ化
ポーリングの過負荷指数バックオフ(5秒→2分上限)
巨大 JSONL のメモリ枯渇readline ストリームで行単位処理
Variant と Product の紐づけ__parentId フィールドで親子解決
書き込みの一括 upsertStaged Upload + bulkOperationRunMutation
同期の欠損・鮮度Webhook 差分 + Bulk フル同期の併用
再実行時の重複SKU 単位の upsert で冪等保証
ShareXB!

この記事を書いた人

渡部 誠也

執行役員 / CTO

独立系 SIer で Web・組み込み・基幹システムの開発を経験し、2017 年に illustrious へ。CTO としてシステム開発事業を立ち上げ、要件定義からコーディングまで一貫して担う。EC に特化した Web アプリケーションを数多く手がける。

ECの業務やシステムについて、
ご相談ください。

いまの運用で困っていること、実現したいことから、一緒に整理します。