跳转到内容
搜索文档

使用 REST API 批量导入 D1

最后更新 查看 MarkdownAgent 设置

在本教程中,您将学习如何使用 REST API 将数据库导入 D1。

前提条件

  1. 注册 Cloudflare 账户
  2. 安装 Node.js

Node.js 版本管理器

使用 Voltanvm 等 Node 版本管理器,以避免权限问题并切换 Node.js 版本。本指南后续将介绍的 Wrangler 需要 Node 版本 16.17.0 或更高。

1. 创建 D1 API 令牌

要使用 REST API,您需要生成一个 API 令牌来验证您的 API 请求。您可以通过 Cloudflare 仪表板执行此操作。

  1. 在 Cloudflare 仪表板中,前往 API Tokens(API 令牌) 页面。

    Go to Account API tokens ↗
  2. API Tokens(API 令牌) 下,选择 Create 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 中,包含三个名为 idtextdate_added 的列。

要创建该表,请遵循以下步骤:

  1. 在 Cloudflare 仪表板中,转到 D1 页面。

    Go to D1 SQL database ↗
  2. 选择 Create database(创建数据库)

  3. 命名您的数据库。在本教程中,将您的 D1 数据库命名为 d1-import-tutorial

  4. (可选)提供位置提示(Location hint)。位置提示是一个可选参数,您可以提供它来指示您希望数据库所在的地理位置。请参阅提供位置提示了解更多信息。

  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

    # 创建 D1 数据库
    npx wrangler d1 create d1-import-tutorial
    
    # 创建 D1 表
    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 & Pages 中的 Account Details
    • DATABASE_ID:D1 数据库 ID。转到您的数据库页面即可查看您的数据库 ID。
    • D1_API_KEY:在步骤 1中生成的 D1 API 令牌。
    index.jsjs
    const TARGET_TABLE = " "; // 本教程中为 `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(); // 创建一个随机文件名
    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 文件的开头添加以下代码。此代码将创建一个名为 data 的数组,其中包含 2500 (uploadSize) 个数组元素,每个元素包含一个带有 idtextdate_added 属性的对象。每个数组元素对应表中的一行。

    index.jsjs
    import crypto from "crypto";
    import { faker } from "@faker-js/faker";
    
    // 生成假数据
    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

导入过程包含四个步骤:

  1. 初始化上传 (Init upload):此步骤初始化上传过程。它将 SQL 命令的哈希值发送到 D1 API 并接收一个上传 URL。
  2. 上传到 R2 (Upload to R2):此步骤将 SQL 命令上传到指定的上传 URL。
  3. 开始摄取 (Start ingestion):此步骤开始数据摄取过程。
  4. 轮询 (Polling):此步骤对导入过程进行轮询,直到其完成。
  1. 创建一个名为 uploadToD1 的函数,用于执行导入过程的四个步骤。

    index.jsjs
    async function uploadToD1() {
    	// 1. 初始化上传
    	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. 上传到 R2
    		const r2Response = await fetch(uploadUrl, {
    			method: "PUT",
    			body: sqlInsert,
    		});
    
    		const r2Etag = r2Response.headers.get("ETag").replace(/"/g, "");
    
    		// 验证 etag
    		if (r2Etag !== hashStr) {
    			throw new Error("ETag mismatch");
    		}
    
    		// 3. 开始摄取
    		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. 轮询
    		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}`,
    };
    
    // 生成假数据
    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", " "),
    }));
    
    // 创建 SQL insert 语句
    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));
    	}
    }
    
    // 上传到 D1
    async function uploadToD1() {
    	// 1. 初始化上传
    	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. 上传到 R2
    		const r2Response = await fetch(uploadUrl, {
    			method: "PUT",
    			body: sqlInsert,
    		});
    
    		const r2Etag = r2Response.headers.get("ETag").replace(/"/g, "");
    
    		// 验证 etag
    		if (r2Etag !== hashStr) {
    			throw new Error("ETag mismatch");
    		}
    
    		// 3. 开始摄取
    		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. 轮询
    		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 目标表中。

这篇文档对您有帮助吗?