Skip to content

非公式本サイトは非公式の日本語ドキュメントであり、Cloudflare 公式サイトではありません。最新情報はdevelopers.cloudflare.comをご確認ください。

REST API で D1 に一括インポートする

最終更新 Markdown で表示Agent セットアップ

このチュートリアルでは、REST API を使ってデータベースを D1 にインポートする方法を学びます。

前提条件

  1. Cloudflare アカウント に登録します。
  2. Node.js をインストールします。

Node.js のバージョンマネージャー

権限の問題を避け、Node.js のバージョンを切り替えられるよう、Voltanvm などの Node バージョンマネージャーを使います。このガイドの後半で説明する Wrangler には、Node バージョン 16.17.0 以降が必要です。

1. D1 API トークンを作成する

REST API を使うには、リクエストを認証するための API トークンが必要です。Cloudflare ダッシュボードから作成できます。

  1. Cloudflare ダッシュボードで API Tokens ページを開きます。

    Account API tokens を開く ↗
  2. API TokensCreate Token を選びます。

  3. Custom token > Create custom token までスクロールし、Get started を選びます。

  4. Token name にわかりやすいトークン名を入力します。例: Name-D1-Import-API-Token

  5. Permissions で次を選びます。

    • Account を選びます。
    • D1 を選びます。
    • Edit を選びます。
  6. Continue to summary を選びます。

  7. Create token を選びます。

  8. API トークンをコピーし、安全なファイルに保存します。

2. 取り込み先テーブルを作成する

インポートするデータのスキーマと一致する D1 テーブルが、あらかじめ必要です。

このチュートリアルでは次を使います。

  • データベース名 d1-import-tutorial
  • テーブル名 TargetD1Table
  • TargetD1Table 内の 3 つの列 idtextdate_added

テーブルを作成する手順は次のとおりです。

  1. Cloudflare ダッシュボードで、D1 ページに移動します。

    D1 SQL database を開く ↗
  2. Create database を選択します。

  3. データベースに名前を付けます。このチュートリアルでは、D1 データベース名を d1-import-tutorial にします。

  4. (任意)ロケーションヒントを指定します。ロケーションヒントは、データベースの希望する地理的な場所を示す任意のパラメーターです。詳細は ロケーションヒントを指定する を参照してください。

  5. Create を選択します。

  6. Console に移動し、次の SQL スニペットを貼り付けます。TargetD1Table というテーブルが作成されます。

    DROP TABLE IF EXISTS TargetD1Table;
    CREATE TABLE IF NOT EXISTS TargetD1Table (id INTEGER PRIMARY KEY, text TEXT, date_added TEXT);

    代わりに Wrangler CLI も使えます。

    # Create a D1 database
    npx wrangler d1 create d1-import-tutorial
    
    # Create a D1 table
    npx wrangler d1 execute d1-import-tutorial --command="DROP TABLE IF EXISTS TargetD1Table; CREATE TABLE IF NOT EXISTS TargetD1Table (id INTEGER PRIMARY KEY, text TEXT, date_added TEXT);" --remote

3. index.js ファイルを作成する

  1. 新しいディレクトリを作り、Node.js プロジェクトを初期化します。

    mkdir d1-import-tutorial
    cd d1-import-tutorial
    npm init -y
  2. このリポジトリに index.js という新しいファイルを作成します。このファイルに、REST API で D1 データベースへデータをインポートするコードを書きます。

  3. index.js で次の変数を定義します。

    • TARGET_TABLE: 取り込み先のテーブル名
    • ACCOUNT_ID: アカウント ID。Workers & PagesAccount Details を参照してください。
    • DATABASE_ID: D1 のデータベース ID。データベースを開くと確認できます。
    • D1_API_KEY: 手順 1 で作成した D1 API トークン
    index.jsjs
    const TARGET_TABLE = " "; // for the tutorial, `TargetD1Table`
    const ACCOUNT_ID = " ";
    const DATABASE_ID = " ";
    const D1_API_KEY = " ";
    const D1_URL = `https://api.cloudflare.com/client/v4/accounts/${ACCOUNT_ID}/d1/database/${DATABASE_ID}/import`;
    const filename = crypto.randomUUID(); // create a random filename
    const uploadSize = 500;
    const headers = {
    	"Content-Type": "application/json",
    	Authorization: `Bearer ${D1_API_KEY}`,
    };

4. サンプルデータを生成する(任意)

実務では、D1 にインポートしたいデータがすでに手元にあることがあります。

このチュートリアルでは、インポート手順を示すためにサンプルデータを生成します。

  1. @faker-js/faker モジュールをインストールします。

    npm i @faker-js/faker
  2. index.js の先頭に次のコードを追加します。このコードは、要素数 2500(uploadSize)の data 配列を作ります。各要素は idtextdate_added を持つオブジェクトです。各要素がテーブルの 1 行に対応します。

    index.jsjs
    import crypto from "crypto";
    import { faker } from "@faker-js/faker";
    
    // Generate Fake data
    const data = Array.from({ length: uploadSize }, () => ({
    	id: Math.floor(Math.random() * 1000000),
    	text: faker.lorem.paragraph(),
    	date_added: new Date().toISOString().slice(0, 19).replace("T", " "),
    }));

5. SQL コマンドを生成する

  1. 取り込み先テーブルへデータを挿入する SQL コマンドを生成する関数を作成します。この関数は、前の手順で作った data 配列を使います。

    index.jsjs
    function makeSqlInsert(data, tableName, skipCols = []) {
    	const columns = Object.keys(data[0]).join(",");
    	const values = data
    		.map((row) => {
    			return (
    				"(" +
    				Object.values(row)
    					.map((val) => {
    						if (skipCols.includes(val) || val === null || val === "") {
    							return "NULL";
    						}
    						return `'${String(val).replace(/'/g, "").replace(/"/g, "'")}'`;
    					})
    					.join(",") +
    				")"
    			);
    		})
    		.join(",");
    
    	return `INSERT INTO ${tableName} (${columns}) VALUES ${values};`;
    }

6. データを D1 にインポートする

インポート処理は次の 4 ステップです。

  1. Init upload: アップロードを初期化します。SQL コマンドのハッシュを D1 API に送り、アップロード URL を受け取ります。
  2. Upload to R2: SQL コマンドをアップロード URL へ送ります。
  3. Start ingestion: 取り込み処理を開始します。
  4. Polling: インポートが完了するまでポーリングします。
  1. インポートの 4 ステップを実行する uploadToD1 関数を作成します。

    index.jsjs
    async function uploadToD1() {
    	// 1. Init upload
    	const hashStr = crypto.createHash("md5").update(sqlInsert).digest("hex");
    
    	try {
    		const initResponse = await fetch(D1_URL, {
    			method: "POST",
    			headers,
    			body: JSON.stringify({
    				action: "init",
    				etag: hashStr,
    			}),
    		});
    
    		const uploadData = await initResponse.json();
    		const uploadUrl = uploadData.result.upload_url;
    		const filename = uploadData.result.filename;
    
    		// 2. Upload to R2
    		const r2Response = await fetch(uploadUrl, {
    			method: "PUT",
    			body: sqlInsert,
    		});
    
    		const r2Etag = r2Response.headers.get("ETag").replace(/"/g, "");
    
    		// Verify etag
    		if (r2Etag !== hashStr) {
    			throw new Error("ETag mismatch");
    		}
    
    		// 3. Start ingestion
    		const ingestResponse = await fetch(D1_URL, {
    			method: "POST",
    			headers,
    			body: JSON.stringify({
    				action: "ingest",
    				etag: hashStr,
    				filename,
    			}),
    		});
    
    		const ingestData = await ingestResponse.json();
    		console.log("Ingestion Response:", ingestData);
    
    		// 4. Polling
    		await pollImport(ingestData.result.at_bookmark);
    
    		return "Import completed successfully";
    	} catch (e) {
    		console.error("Error:", e);
    		return "Import failed";
    	}
    }

    上記のコードでは次を行います。

    • SQL コマンドの md5 ハッシュを生成します。
    • initResponse でアップロードを初期化し、アップロード URL を受け取ります。
    • r2Response で SQL コマンドをアップロード URL へ送ります。
    • 取り込み開始前に ETag を検証します。
    • ingestResponse で取り込み処理を開始します。
    • pollImport でインポートが完了するまでポーリングします。
  2. pollImport 関数を index.js に追加します。

    index.jsjs
    async function pollImport(bookmark) {
    	const payload = {
    		action: "poll",
    		current_bookmark: bookmark,
    	};
    
    	while (true) {
    		const pollResponse = await fetch(D1_URL, {
    			method: "POST",
    			headers,
    			body: JSON.stringify(payload),
    		});
    
    		const result = await pollResponse.json();
    		console.log("Poll Response:", result.result);
    
    		const { success, error } = result.result;
    
    		if (
    			success ||
    			(!success && error === "Not currently importing anything.")
    		) {
    			break;
    		}
    
    		await new Promise((resolve) => setTimeout(resolve, 1000));
    	}
    }

    上記のコードでは次を行います。

    • D1 API に poll アクションを送ります。
    • インポートが完了するまでポーリングします。
  3. 最後に、インポートを実行する runImport 関数を index.js に追加します。

    index.jsjs
    async function runImport() {
    	const result = await uploadToD1();
    	console.log(result);
    }
    
    runImport();

7. 最終コードを書く

これまでの手順で、D1 へのデータインポートに必要な各処理の関数を作成しました。最終コードは、これらの関数を実行してサンプルデータを取り込み先の D1 テーブルへインポートします。

  1. 次に示す index.js の最終コードをコピーし、変数をコード先頭で定義します。

    import crypto from "crypto";
    import { faker } from "@faker-js/faker";
    
    const TARGET_TABLE = "";
    const ACCOUNT_ID = "";
    const DATABASE_ID = "";
    const D1_API_KEY = "";
    const D1_URL = `https://api.cloudflare.com/client/v4/accounts/${ACCOUNT_ID}/d1/database/${DATABASE_ID}/import`;
    const uploadSize = 500;
    const headers = {
    	"Content-Type": "application/json",
    	Authorization: `Bearer ${D1_API_KEY}`,
    };
    
    // Generate Fake data
    const data = Array.from({ length: uploadSize }, () => ({
    	id: Math.floor(Math.random() * 1000000),
    	text: faker.lorem.paragraph(),
    	date_added: new Date().toISOString().slice(0, 19).replace("T", " "),
    }));
    
    // Make SQL insert statements
    function makeSqlInsert(data, tableName, skipCols = []) {
    	const columns = Object.keys(data[0]).join(",");
    	const values = data
    		.map((row) => {
    			return (
    				"(" +
    				Object.values(row)
    					.map((val) => {
    						if (skipCols.includes(val) || val === null || val === "") {
    							return "NULL";
    						}
    						return `'${String(val).replace(/'/g, "").replace(/"/g, "'")}'`;
    					})
    					.join(",") +
    				")"
    			);
    		})
    		.join(",");
    
    	return `INSERT INTO ${tableName} (${columns}) VALUES ${values};`;
    }
    
    const sqlInsert = makeSqlInsert(data, TARGET_TABLE);
    
    async function pollImport(bookmark) {
    	const payload = {
    		action: "poll",
    		current_bookmark: bookmark,
    	};
    
    	while (true) {
    		const pollResponse = await fetch(D1_URL, {
    			method: "POST",
    			headers,
    			body: JSON.stringify(payload),
    		});
    
    		const result = await pollResponse.json();
    		console.log("Poll Response:", result.result);
    
    		const { success, error } = result.result;
    
    		if (
    			success ||
    			(!success && error === "Not currently importing anything.")
    		) {
    			break;
    		}
    
    		await new Promise((resolve) => setTimeout(resolve, 1000));
    	}
    }
    
    // Upload to D1
    async function uploadToD1() {
    	// 1. Init upload
    	const hashStr = crypto.createHash("md5").update(sqlInsert).digest("hex");
    
    	try {
    		const initResponse = await fetch(D1_URL, {
    			method: "POST",
    			headers,
    			body: JSON.stringify({
    				action: "init",
    				etag: hashStr,
    			}),
    		});
    
    		const uploadData = await initResponse.json();
    		const uploadUrl = uploadData.result.upload_url;
    		const filename = uploadData.result.filename;
    
    		// 2. Upload to R2
    		const r2Response = await fetch(uploadUrl, {
    			method: "PUT",
    			body: sqlInsert,
    		});
    
    		const r2Etag = r2Response.headers.get("ETag").replace(/"/g, "");
    
    		// Verify etag
    		if (r2Etag !== hashStr) {
    			throw new Error("ETag mismatch");
    		}
    
    		// 3. Start ingestion
    		const ingestResponse = await fetch(D1_URL, {
    			method: "POST",
    			headers,
    			body: JSON.stringify({
    				action: "ingest",
    				etag: hashStr,
    				filename,
    			}),
    		});
    
    		const ingestData = await ingestResponse.json();
    		console.log("Ingestion Response:", ingestData);
    
    		// 4. Polling
    		await pollImport(ingestData.result.at_bookmark);
    
    		return "Import completed successfully";
    	} catch (e) {
    		console.error("Error:", e);
    		return "Import failed";
    	}
    }
    
    async function runImport() {
    	const result = await uploadToD1();
    	console.log(result);
    }
    
    runImport();

8. コードを実行する

  1. コードを実行します。

    node index.js

取り込み先の D1 テーブルに、サンプルデータが入ります。

まとめ

このチュートリアルを完了すると、次のことを行えています。

  1. API トークンを作成しました。
  2. 取り込み先のデータベースとテーブルを作成しました。
  3. サンプルデータを生成しました。
  4. サンプルデータ用の SQL コマンドを作成しました。
  5. REST API でサンプルデータを D1 の取り込み先テーブルへインポートしました。

役に立ちましたか?