目次
はじめに
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 フィールドで親子解決 |
| 書き込みの一括 upsert | Staged Upload + bulkOperationRunMutation |
| 同期の欠損・鮮度 | Webhook 差分 + Bulk フル同期の併用 |
| 再実行時の重複 | SKU 単位の upsert で冪等保証 |




