INTEGRITY Документация

Массовый импорт в D1 через REST API

В этом руководстве вы узнаете, как импортировать базу данных в D1 с помощью REST API.

Предварительные требования

  1. Зарегистрируйтесь для получения Аккаунт Cloudflare.
  2. Установка Node.js.

менеджер версий Node.js

Используйте менеджер версий Node, например Volta или nvm чтобы избежать проблем с правами доступа и переключать версии Node.js. Wrangler, о котором пойдёт речь далее в этом руководстве, требует версию Node 16.17.0 или более поздней версии.

1. Создайте API-токен D1

Чтобы использовать REST API, вам нужно создать токен API для аутентификации запросов API. Сделать это можно через панель управления Cloudflare.

  1. На панели управления Cloudflare перейдите к разделу API-токены страницу.

    Перейдите в Токены API аккаунта ↗
  2. В разделе API-токены, выберите Create Token.

  3. Прокрутите до Custom token > Create custom token, затем выберите Начало работы.

  4. В разделе Имя токена, введите понятное название токена. Например, Name-D1-Import-API-Token.

  5. В разделе Разрешения:

    • Выберите Аккаунт.
    • Выберите D1.
    • Выберите Изменить.
  6. Выберите Continue to summary.

  7. Выберите Создать токен.

  8. Скопируйте API-токен и сохраните его в защищенном файле.

2. Создайте целевую таблицу

У вас должна быть существующая таблица D1, соответствующая схеме импортируемых данных.

В этом руководстве используется следующее:

Чтобы создать таблицу, выполните следующие шаги:

  1. На панели управления Cloudflare перейдите к разделу D1 страницу.

    Перейдите в SQL-база данных D1 ↗
  2. Выберите Create database.

  3. Присвойте имя своей базе данных. Для этого руководства назовите базу данных D1 d1-import-tutorial.

  4. (Необязательно) Укажите подсказку расположения. Подсказка расположения представляет собой необязательный параметр, с помощью которого можно указать желаемое географическое расположение вашей базы данных. См. Укажите подсказку расположения, где это описано подробнее.

  5. Выберите Создание.

  6. Перейдите в Консоль, затем вставьте следующий фрагмент 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: Идентификатор аккаунта. См. Сведения об аккаунте в Workers & Pages.
    • DATABASE_ID: Идентификатор базы данных D1. Его можно найти на странице вашей базы данных.
    • D1_API_KEY: Токен D1 API, созданный в шаг 1
    index.js
    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 файл. Этот код создаёт массив с именем data с 2500 (uploadSize) элементов массива, где каждый элемент массива содержит объект с id, text, а также date_added. Каждый элемент массива соответствует строке таблицы.

    index.js
    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.js
    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: Этот шаг загружает SQL-команду по URL-адресу для загрузки.
  3. Запуск приёма данных: Этот шаг запускает процесс приёма данных.
  4. Опрос: Этот шаг периодически проверяет процесс импорта до его завершения.
  1. Создайте функцию с именем uploadToD1 который выполняет четыре этапа процесса импорта.

    index.js
    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";
    	}
    }

    В приведенном выше коде:

    • Одна md5 хеш SQL команды генерируется.
    • initResponse инициализирует процесс загрузки и получает URL для загрузки.
    • r2Response загружает SQL-команду на URL для загрузки.
    • Перед началом приёма данных выполняется проверка ETag.
    • ingestResponse запускает процесс приёма данных.
    • pollImport опрашивает процесс импорта до его завершения.
  2. Добавьте pollImport функцию в index.js файл.

    index.js
    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));
    	}
    }

    Приведённый выше код выполняет следующее:

    • Отправляет poll действие в D1 API.
    • Опрашивает процесс импорта до его завершения.
  3. Наконец, добавьте runImport функцию в index.js файл, чтобы запустить процесс импорта.

    index.js
    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. Импортированы примеры данных в целевую таблицу D1 с помощью REST API.