Skip to content

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

Workers AI で BigQuery を使う

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

Workers AI を始めるいちばん簡単な方法は、Multi-modal PlaygroundLLM playground を試すことです。コードを Workers AI と統合する場合は、REST API エンドポイント または Worker バインディング を使えます。

では、データはどうでしょうか。これらのモデルに、Cloudflare の外に保存されているデータを取り込ませたい場合は?

このチュートリアルでは、Google BigQuery のデータを Cloudflare Worker へ取り込み、Workers AI モデルの入力として使う方法を学びます。

前提条件

次のものが必要です。

1. Cloudflare Worker を用意する

データを Cloudflare に取り込み、Workers AI へ渡すために、Cloudflare Worker を使います。まだ作成していない場合は、利用開始のチュートリアル を確認してください。

Worker の作成手順を完了すると、新しい Worker プロジェクトには次のコードがあります。

export default {
	async fetch(request, env, ctx) {
		return new Response("Hello World!");
	},
};

Worker プロジェクトが正しく作成されていれば、コンソールで npx wrangler dev を実行して、Worker をローカルで起動できます。

[wrangler:inf] Ready on http://localhost:8787

ブラウザーで http://localhost:8787/ を開いて、デプロイした Worker を確認します。ポート 8787 は環境によって異なる場合があります。

ブラウザーには Hello World! と表示されます。

Hello World!

この手順で問題が起きた場合は、Workers の利用開始ガイド を確認してください。

2. GCP サービスキーを Worker の Secrets として取り込む

Worker が正しく作成できたことを確認したら、このチュートリアルの 前提条件 で作成した Google Cloud Platform のサービスキーを参照します。

Google Cloud Platform からダウンロードしたキーの JSON ファイルは、次の形式です。

{
	"type": "service_account",
	"project_id": "<your_project_id>",
	"private_key_id": "<your_private_key_id>",
	"private_key": "<your_private_key>",
	"client_email": "<your_service_account_id>@<your_project_id>.iam.gserviceaccount.com",
	"client_id": "<your_oauth2_client_id>",
	"auth_uri": "https://accounts.google.com/o/oauth2/auth",
	"token_uri": "https://oauth2.googleapis.com/token",
	"auth_provider_x509_cert_url": "https://www.googleapis.com/oauth2/v1/certs",
	"client_x509_cert_url": "https://www.googleapis.com/robot/v1/metadata/x509/<your_service_account_id>%40<your_project_id>.iam.gserviceaccount.com",
	"universe_domain": "googleapis.com"
}

このチュートリアルで必要なのは、次のフィールドの値だけです。client_emailprivate_keyprivate_key_idproject_id

この情報を Worker に平文で置かず、Secrets を使います。暗号化されていない内容にアクセスできるのは、Worker 自身だけです。

JSON ファイルの 3 つの値を Secrets に取り込みます。まず JSON キーファイルの client_email フィールドから始め、名前は BQ_CLIENT_EMAIL とします(別の変数名でも構いません)。

npx wrangler secret put BQ_CLIENT_EMAIL

シークレット値の入力を求められます。JSON キーファイルの client_email フィールドの値を入力します。

シークレットのアップロードが成功すると、次のメッセージが表示されます。

 Success! Uploaded secret BQ_CLIENT_EMAIL

残りの 3 フィールド(private_keyprivate_key_idproject_id)も、それぞれ BQ_PRIVATE_KEYBQ_PRIVATE_KEY_IDBQ_PROJECT_ID として取り込みます。

npx wrangler secret put BQ_PRIVATE_KEY
npx wrangler secret put BQ_PRIVATE_KEY_ID
npx wrangler secret put BQ_PROJECT_ID

これで、Google Cloud Platform からダウンロードした JSON キーファイルの 3 フィールドを、Worker で使う Cloudflare Secrets へ取り込みました。

Secrets が Worker で使えるのは、デプロイ後です。開発中に使うには、.dev.vars を作成 し、認証情報をローカルに保存して環境変数として参照します。

dev.vars ファイルは次のようになります。

BQ_CLIENT_EMAIL="<your_service_account_id>@<your_project_id>.iam.gserviceaccount.com"
BQ_CLIENT_KEY="-----BEGIN PRIVATE KEY-----<content_of_your_private_key>-----END PRIVATE KEY-----\n"
BQ_PRIVATE_KEY_ID="<your_private_key_id>"
BQ_PROJECT_ID="<your_project_id>"

バージョン管理を使う場合に認証情報がリポジトリへ上がらないよう、プロジェクトの .gitignore.dev.vars を追加してください。

src/index.js で値をコンソール出力し、シークレットが正しく読み込まれているか確認します。

export default {
	async fetch(request, env, ctx) {
		console.log("BQ_CLIENT_EMAIL: ", env.BQ_CLIENT_EMAIL);
		console.log("BQ_PRIVATE_KEY: ", env.BQ_PRIVATE_KEY);
		console.log("BQ_PRIVATE_KEY_ID: ", env.BQ_PRIVATE_KEY_ID);
		console.log("BQ_PROJECT_ID: ", env.BQ_PROJECT_ID);
		return new Response("Hello World!");
	},
};

Worker を再起動し、npx wrangler dev を実行します。サーバーが、追加した変数に言及しているはずです。

Using vars defined in .dev.vars
Your worker has access to the following bindings:
- Vars:
  - BQ_CLIENT_EMAIL: "(hidden)"
  - BQ_PRIVATE_KEY: "(hidden)"
  - BQ_PRIVATE_KEY_ID: "(hidden)"
  - BQ_PROJECT_ID: "(hidden)"
[wrangler:inf] Ready on http://localhost:8787

ブラウザーで http://localhost:8787 を開くと、npx wrangler dev を実行しているコンソールに変数の値が表示されます。ブラウザーウィンドウには、これまでどおり Hello World! だけが表示されます。

これで Worker から GCP の認証情報にアクセスできます。次は、GCP の API と連携するために必要な JSON Web Token の作成を助けるライブラリをインストールします。

3. JWT 操作用のライブラリをインストールする

BigQuery の REST API と連携するには、前の手順で Worker secrets に読み込んだ認証情報を使い、リクエストを認証する JSON Web Token を生成する必要があります。

このチュートリアルでは、JWT 関連の操作に jose ライブラリを使います。コンソールで次のコマンドを実行してインストールします。

npm i jose

インストールが成功したかは、インストール済みパッケージを一覧する npm list を実行し、jose 依存関係が追加されているかで確認できます。

<project_name>@0.0.0
/<path_to_your_project>/<project_name>
├── @cloudflare/vitest-pool-workers@0.4.29
├── jose@5.9.2
├── vitest@1.5.0
└── wrangler@3.75.0

4. JSON Web Token を生成する

jose ライブラリをインストールしたので、インポートし、署名付き JSON Web Token(JWT)を生成する関数をコードへ追加します。

import * as jose from 'jose';
...
const generateBQJWT = async (aCryptoKey, env) => {
const algorithm = "RS256";
const audience = "https://bigquery.googleapis.com/";
const expiryAt = (new Date().valueOf() / 1000);
	const privateKey = await jose.importPKCS8(env.BQ_PRIVATE_KEY, algorithm);

	// Generate signed JSON Web Token (JWT)
	return new jose.SignJWT()
    	.setProtectedHeader({
        	typ: 'JWT',
        	alg: algorithm,
        	kid: env.BQ_PRIVATE_KEY_ID
    	})
    	.setIssuer(env.BQ_CLIENT_EMAIL)
    	.setSubject(env.BQ_CLIENT_EMAIL)
    	.setAudience(audience)
    	.setExpirationTime(expiryAt)
    	.setIssuedAt()
    	.sign(privateKey)
}

export default {
	async fetch(request, env, ctx) {
       ...
// Create JWT to authenticate the BigQuery API call
    	let bqJWT;
    	try {
        	bqJWT = await generateBQJWT(env);
    	} catch (e) {
        	return new Response('An error has occurred while generating the JWT', { status: 500 })
    	}
	},
       ...
};

JWT を作成したので、次は BigQuery へ API 呼び出しを行い、データを取得します。

5. Google BigQuery へ認証済みリクエストを送る

前の手順で作成した JWT トークンを使い、BigQuery の API へリクエストして、テーブルからデータを取得します。

このチュートリアルの前半で BigQuery に作成したテーブルをクエリします。この例では、MIT ライセンスの Hacker News Corpus のサンプルを BigQuery にアップロードして使っています。

const queryBQ = async (bqJWT, path) => {
	const bqEndpoint = `https://bigquery.googleapis.com${path}`
	// In this example, text is a field in the BigQuery table that is being queried (hn.news_sampled)
	const query = 'SELECT text FROM hn.news_sampled LIMIT 3';
	const response = await fetch(bqEndpoint, {
    	method: "POST",
    	body: JSON.stringify({
        	"query": query
    	}),
    	headers: {
        	Authorization: `Bearer ${bqJWT}`
    	}
	})
	return response.json()
}
...
export default {
	async fetch(request, env, ctx) {
		...
    		let ticketInfo;
    		try {
    		ticketInfo = await queryBQ(bqJWT);
    	} catch (e) {
        	return new Response('An error has occurred while querying BQ', { status: 500 });
    	}
	...
	},
};

BigQuery から生の行データを取得できたので、次は JSON に近い形式へ整えます。

6. クエリ結果を整形する

BigQuery からデータを取得したあとの API レスポンスは、次のようになります。

{
	...
	"schema": {
    	"fields": [
        	{
            	"name": "title",
            	"type": "STRING",
            	"mode": "NULLABLE"
        	},
        	{
            	"name": "text",
            	"type": "STRING",
            	"mode": "NULLABLE"
        	}
    	]
	},
	...
	"rows": [
    	{
        	"f": [
            	{
                	"v": "<some_value>"
            	},
            	{
                	"v": "<some_value>"
            	}
        	]
    	},
    	{
        	"f": [
            	{
                	"v": "<some_value>"
            	},
            	{
                	"v": "<some_value>"
            	}
        	]
    	},
    	{
        	"f": [
            	{
                	"v": "<some_value>"
            	},
            	{
                	"v": "<some_value>"
            	}
        	]
    	}
	],
	...
}

この形式は読みにくく、結果を繰り返し処理するときにも扱いにくいです。そこで、スキーマを各値へ対応付ける関数を実装します。出力は次のように読みやすくなります。各行は、配列内のオブジェクトになります。

[
	{
		title: "<some_value>",
		text: "<some_value>",
	},
	{
		title: "<some_value>",
		text: "<some_value>",
	},
	{
		title: "<some_value>",
		text: "<some_value>",
	},
];

BigQuery レスポンス本文の行とフィールドを受け取り、名前付きフィールドを持つオブジェクトの配列を返す formatRows 関数を作成します。

const formatRows = (rowsWithoutFieldNames, fields) => {
	// Index to fieldName
	const fieldsByIndex = new Map();

	// Load all fields by name and have their index in the array result as their key
	fields.forEach((field, index) => {
    	fieldsByIndex.set(index, field.name)
	})

	// Iterate through rows
	const rowsWithFieldNames = rowsWithoutFieldNames.map(row => {
    	// Per each row represented by an array f, iterate through the unnamed values and find their field names by searching them in the fieldsByIndex.
    	let newRow = {}
    	row.f.forEach((field, index) => {
        	const fieldName = fieldsByIndex.get(index);
        	if (fieldName) {
		// For every field in a row, add them to newRow
            	newRow = ({ ...newRow, [fieldName]: field.v });
        	}
    	})
    	return newRow
	})

	return rowsWithFieldNames
}

export default {
	async fetch(request, env, ctx) {
		...
    	// Transform output format into array of objects with named fields
    	let formattedResults;

    	if ('rows' in ticketInfo) {
        	formattedResults = formatRows(ticketInfo.rows, ticketInfo.schema.fields);
        	console.log(formattedResults)
    	} else if ('error' in ticketInfo) {
        	return new Response(ticketInfo.error.message, { status: 500 })
    	}
	...
	},
};

7. データを Workers AI へ渡す

BigQuery API のレスポンスを結果の配列へ変換したので、Workers AI 経由の LLM でタグを生成し、感情スコアを付けます。

const generateTags = (data, env) => {
	return env.AI.run("@cf/meta/llama-3.1-8b-instruct", {
    	prompt: `Create three one-word tags for the following text. return only these three tags separated by a comma. don't return text that is not a category.Lowercase only. ${JSON.stringify(data)}`,
	});
}

const generateSentimentScore = (data, env) => {
	return env.AI.run("@cf/meta/llama-3.1-8b-instruct", {
    	prompt: `return a float number between 0 and 1 measuring the sentiment of the following text. 0 being negative and 1 positive. return only the number, no text. ${JSON.stringify(data)}`,
	});
}

// Iterates through values, sends them to an AI handler and encapsulates all responses into a single Promise
const getAIGeneratedContent = (data, env, aiHandler) => {
	let results = data?.map(dataPoint => {
    	return aiHandler(dataPoint, env)
	})
	return Promise.all(results)
}
...
export default {
	async fetch(request, env, ctx) {
		...
let summaries, sentimentScores;
    	try {
        	summaries = await getAIGeneratedContent(formattedResults, env, generateTags);
        	sentimentScores = await getAIGeneratedContent(formattedResults, env, generateSentimentScore)
    	} catch {
        	return new Response('There was an error while generating the text summaries or sentiment scores')
    	}
},

formattedResults = formattedResults?.map((formattedResult, i) => {
        	if (sentimentScores[i].response && summaries[i].response) {
            	return {
                	...formattedResult,
                	'sentiment': parseFloat(sentimentScores[i].response).toFixed(2),
                	'tags': summaries[i].response.split(',').map((result) => result.trim())
            	}
        	}
    	}
};

プロジェクトの Wrangler ファイルで、次の行のコメントを外します。

{
	"ai": {
		"binding": "AI"
	}
}
[ai]
binding = "AI"

ローカルで動いている Worker を再起動し、アプリケーションのエンドポイントへアクセスします。

curl http://localhost:8787

Workers AI を使うとき、Cloudflare アカウントへのログインと、Wrangler(Cloudflare CLI)への一時アクセス許可を求められることがあります。

http://localhost:8787 にアクセスすると、次のような出力が表示されます。

{
  "data": [
	{
  	"text": "You can see a clear spike in submissions right around US Thanksgiving.",
  	"sentiment": "0.61",
  	"tags": [
    	"trends",
    	"submissions",
    	"thanksgiving"
  	]
	},
	{
  	"text": "I didn't test the changes before I published them.  I basically did development on the running server. In fact for about 30 seconds the comments page was broken due to a bug.",
  	"sentiment": "0.35",
  	"tags": [
    	"software",
    	"deployment",
    	"error"
  	]
	},
	{
  	"text": "I second that. As I recall, it's a very enjoyable 700-page brain dump by someone who's really into his subject. The writing has a personal voice; there are lots of asides, dry wit, and typos that suggest restrained editing. The discussion is intelligent and often theoretical (and Bartle is not scared to use mathematical metaphors), but the tone is not academic.",
  	"sentiment": "0.86",
  	"tags": [
    	"review",
    	"game",
    	"design"
  	]
	}
  ]
}

実際の値とフィールドは、手順 5 のクエリと、それを渡した LLM に大きく依存します。

最終結果

各手順のコードをまとめると、src/index.js は次のようになります。

import * as jose from "jose";

const generateBQJWT = async (env) => {
	const algorithm = "RS256";
	const audience = "https://bigquery.googleapis.com/";
	const expiryAt = new Date().valueOf() / 1000;
	const privateKey = await jose.importPKCS8(env.BQ_PRIVATE_KEY, algorithm);

	// Generate signed JSON Web Token (JWT)
	return new jose.SignJWT()
		.setProtectedHeader({
			typ: "JWT",
			alg: algorithm,
			kid: env.BQ_PRIVATE_KEY_ID,
		})
		.setIssuer(env.BQ_CLIENT_EMAIL)
		.setSubject(env.BQ_CLIENT_EMAIL)
		.setAudience(audience)
		.setExpirationTime(expiryAt)
		.setIssuedAt()
		.sign(privateKey);
};

const queryBQ = async (bgJWT, path) => {
	const bqEndpoint = `https://bigquery.googleapis.com${path}`;
	const query = "SELECT text FROM hn.news_sampled LIMIT 3";
	const response = await fetch(bqEndpoint, {
		method: "POST",
		body: JSON.stringify({
			query: query,
		}),
		headers: {
			Authorization: `Bearer ${bgJWT}`,
		},
	});
	return response.json();
};

const formatRows = (rowsWithoutFieldNames, fields) => {
	// Index to fieldName
	const fieldsByIndex = new Map();

	fields.forEach((field, index) => {
		fieldsByIndex.set(index, field.name);
	});

	const rowsWithFieldNames = rowsWithoutFieldNames.map((row) => {
		// Map rows into an array of objects with field names
		let newRow = {};
		row.f.forEach((field, index) => {
			const fieldName = fieldsByIndex.get(index);
			if (fieldName) {
				newRow = { ...newRow, [fieldName]: field.v };
			}
		});
		return newRow;
	});

	return rowsWithFieldNames;
};

const generateTags = (data, env) => {
	return env.AI.run("@cf/meta/llama-3.1-8b-instruct", {
		prompt: `Create three one-word tags for the following text. return only these three tags separated by a comma. don't return text that is not a category.Lowercase only. ${JSON.stringify(data)}`,
	});
};

const generateSentimentScore = (data, env) => {
	return env.AI.run("@cf/meta/llama-3.1-8b-instruct", {
		prompt: `return a float number between 0 and 1 measuring the sentiment of the following text. 0 being negative and 1 positive. return only the number, no text. ${JSON.stringify(data)}`,
	});
};

const getAIGeneratedContent = (data, env, aiHandler) => {
	let results = data?.map((dataPoint) => {
		return aiHandler(dataPoint, env);
	});
	return Promise.all(results);
};

export default {
	async fetch(request, env, ctx) {
		// Create JWT to authenticate the BigQuery API call
		let bqJWT;
		try {
			bqJWT = await generateBQJWT(env);
		} catch (error) {
			console.log(error);
			return new Response("An error has occurred while generating the JWT", {
				status: 500,
			});
		}

		// Fetch results from BigQuery
		let ticketInfo;
		try {
			ticketInfo = await queryBQ(
				bqJWT,
				`/bigquery/v2/projects/${env.BQ_PROJECT_ID}/queries`,
			);
		} catch (error) {
			console.log(error);
			return new Response("An error has occurred while querying BQ", {
				status: 500,
			});
		}

		// Transform output format into array of objects with named fields
		let formattedResults;
		if ("rows" in ticketInfo) {
			formattedResults = formatRows(ticketInfo.rows, ticketInfo.schema.fields);
		} else if ("error" in ticketInfo) {
			return new Response(ticketInfo.error.message, { status: 500 });
		}

		// Generate AI summaries and sentiment scores
		let summaries, sentimentScores;
		try {
			summaries = await getAIGeneratedContent(
				formattedResults,
				env,
				generateTags,
			);
			sentimentScores = await getAIGeneratedContent(
				formattedResults,
				env,
				generateSentimentScore,
			);
		} catch {
			return new Response(
				"There was an error while generating the text summaries or sentiment scores",
			);
		}

		// Add AI summaries and sentiment scores to previous results
		formattedResults = formattedResults?.map((formattedResult, i) => {
			if (sentimentScores[i].response && summaries[i].response) {
				return {
					...formattedResult,
					sentiment: parseFloat(sentimentScores[i].response).toFixed(2),
					tags: summaries[i].response.split(",").map((result) => result.trim()),
				};
			}
		});

		const response = { data: formattedResults };

		return new Response(JSON.stringify(response), {
			headers: { "Content-Type": "application/json" },
		});
	},
};

この Worker をデプロイするには、npx wrangler deploy を実行します。

Total Upload: <size_of_your_worker> KiB / gzip: <compressed_size_of_your_worker> KiB
Uploaded <name_of_your_worker> (x sec)
Deployed <name_of_your_worker> triggers (x sec)
  https://<your_public_worker_endpoint>
Current Version ID: <worker_script_version_id>

これで、世界中から Worker へアクセスできる公開エンドポイントが作成されます。本番データを扱うときはこの点に注意し、追加のアクセス制御を必ず入れてください。

まとめ

このチュートリアルでは、GCP サービスアカウントキーを作成し、その一部を Worker secrets として保存して、Google BigQuery と Cloudflare Workers を連携しました。コードでシークレットを読み込み、jose npm ライブラリで JSON Web Token を作成し、BigQuery への API クエリを認証しました。

結果を取得したあと、Workers AI 経由の生成 AI モデルへ渡せる形式へ整え、抽出データからタグを生成し、感情分析を行いました。

次のステップ

AI モデルへ取り込んだ結果をブラウザーに表示するのではなく、一定間隔でデータを取得して保存する(R2D1 など)ワークフローなら、この Worker に scheduled handler を追加することを検討してください。Cron Trigger で、決まった間隔で Worker を起動できます。BigQuery データを Workers AI へ取り込む のリファレンスアーキテクチャ図も確認してください。

このチュートリアルのように他ソースからデータを取り込む用途のひとつが、RAG システムの作成です。関心がある場合は、Retrieval Augmented Generation (RAG) AI を構築するチュートリアル を参照してください。

Cloudflare で使える他の AI モデルについては、ドキュメントの Workers AI セクションを参照してください。

役に立ちましたか?