はじめに
この記事では IBM Bob に Confluent MCP サーバー を接続することで、日本語のチャットだけで Kafka トピックの作成・データ投入・Flink SQL 実行 まで一気通貫でやってみた記録をまとめます。
コマンドはセットアップ時にほんの少し打つだけ。あとは Bob に話しかけるだけで Confluent Cloud を操作できるようになります。
この記事でやること
- MCP サーバーのセットアップ
- Bob からトピック作成・Connector 作成・メッセージ確認
- Flink SQL をチャットで実行してストリーミング処理
- スキーマ管理もチャットで完結
前提・環境
| 必要なもの | 備考 |
|---|---|
| IBM Bob | インストール済み、一度は起動済みであること |
| Confluent Cloud アカウント | クラスターを1つ作成済みであること |
| Node.js 22 以上 |
node -v で確認。なければ nvm install 22
|
| Confluent CLI | この記事の Step 2 でインストール |
まず Node.js のバージョンを確認しておきます。ここだけは必須チェックです。
node -v
# v22.x.x 以上であればOK
MCP サーバーとは?
MCP(Model Context Protocol) は、AIエージェントが外部ツール・サービスを呼び出すための標準プロトコルです。
Confluent が公式に公開している @confluentinc/mcp-confluent を IBM Bob に登録すると、Bob が自動的に Confluent Cloud の API を呼び出してくれます。
対応しているツールは Kafka・Flink SQL・Schema Registry・Connectors など50種類以上。
Step 1. 作業フォルダを作って config のひな形を生成する
mkdir ~/mcp-confluent-setup
cd ~/mcp-confluent-setup
npx @confluentinc/mcp-confluent --init-config
実行すると config.yaml が生成されます。これに接続情報を書き込んでいきます。
Step 2. Confluent CLI をインストールしてログインする
brew install confluentinc/tap/confluent
confluent version
# Confluent Cloud にログイン(初回のみ)
confluent login --save
ログイン後、使いたい Environment と Kafka クラスターを選択します。
confluent environment list
confluent environment use <env-id> # 例: env-pxdn3y
confluent kafka cluster list
confluent kafka cluster use <lkc-xxxxxxx> # 例: lkc-nvwq0jz
Step 3. 接続情報を控える
config.yaml に書き込む URL を CLI で確認します。
# Kafka の bootstrap server と REST endpoint
confluent kafka cluster describe <lkc-xxxxxxx>
# → Endpoint: pkc-xxxx.ap-northeast-1.aws.confluent.cloud:9092
# Schema Registry の endpoint
confluent schema-registry cluster describe
# → Endpoint URL: https://psrc-xxxx.us-east-2.aws.confluent.cloud
# Flink compute pool の確認
confluent flink compute-pool list --environment <env-id>
# Organization ID(Flink で必要)
confluent organization list
Step 4. API Key を4種類発行する
Confluent Cloud では サービスごとにスコープが異なる API Key を使います。4種類、別々に作る必要があります。
# ① Kafka クラスター用
confluent api-key create --resource <lkc-xxxxxxx>
# ② Schema Registry 用
confluent api-key create --resource <lsrc-xxxxxxx>
# ③ Confluent Cloud コントロールプレーン用(billing / environment 管理など)
confluent api-key create --resource cloud
# ④ Flink 用(compute pool ではなく region に紐づく)
confluent api-key create --resource flink \
--cloud aws --region ap-northeast-1 \
--environment <env-id>
| Key 種別 | 用途 |
|---|---|
| Kafka API Key | トピック操作・メッセージ送受信 |
| Schema Registry API Key | スキーマの登録・参照 |
| Cloud API Key | 環境・クラスター・billing の参照 |
| Flink API Key | Flink SQL ステートメントの実行 |
⚠️ Secret はその場でしか表示されません。 必ずコピーしてから次に進んでください。
Step 5. config.yaml に書き込む
Step 3・4 で集めた情報を config.yaml に書き込みます。
server:
transports: [stdio]
connections:
default:
type: direct
description: "My Confluent Cloud"
kafka:
bootstrap_servers: "pkc-xxxx.ap-northeast-1.aws.confluent.cloud:9092"
auth:
type: api_key
key: "<kafka key>"
secret: "<kafka secret>"
rest_endpoint: "https://pkc-xxxx.ap-northeast-1.aws.confluent.cloud:443"
cluster_id: "lkc-xxxxxxx"
env_id: "<env-id>"
schema_registry:
endpoint: "https://psrc-xxxx.us-east-2.aws.confluent.cloud"
auth:
type: api_key
key: "<sr key>"
secret: "<sr secret>"
confluent_cloud:
endpoint: "https://api.confluent.cloud"
auth:
type: api_key
key: "<cloud key>"
secret: "<cloud secret>"
flink:
endpoint: "https://flink.ap-northeast-1.aws.confluent.cloud"
auth:
type: api_key
key: "<flink key>"
secret: "<flink secret>"
organization_id: "<org-id>"
environment_id: "<env-id>"
compute_pool_id: "<lfcp-xxxxxxx>"
Step 6. MCP サーバーが起動するか単体確認する
Bob に組み込む前に、まずコマンドラインで動作確認します。
npx @confluentinc/mcp-confluent --config ~/mcp-confluent-setup/config.yaml --list-tools
kafka / schema-registry / flink / confluent-cloud カテゴリのツール一覧が表示されれば OK です。
Step 7. IBM Bob に MCP サーバーを登録する
Bob の設定ファイルに追記します。
~/.bob/settings/mcp_settings.json
{
"mcpServers": {
"confluent": {
"command": "npx",
"args": [
"-y",
"@confluentinc/mcp-confluent",
"--config",
"/Users/{user_name}/mcp-confluent-setup/config.yaml"
]
}
}
}
{user_name}はwhoamiコマンドで確認できます。絶対パス で書くのがポイントです(相対パスは認識されません)。
設定を保存したら IBM Bob を完全に終了して再起動 します。
Bob の Settings → MCP を開くと、登録した MCP サーバーの一覧が表示されます。ここに confluent が表示され、ステータスが緑色の 「接続済み」 になっていれば OK です。
Step 8. 疎通確認(チャットで試す)
Bob を開いて、以下のプロンプトを投げてみます。
この Kafka クラスターのトピック一覧を教えて
list-topics ツールが呼ばれ、トピック名が返ってくれば接続成功です 🎉
実際に触ってみた
セットアップ完了後、実際に Bob と対話しながら Confluent Cloud を操作してみました。以下はそのハイライトです。
トピック作成
confluentMCP経由で新しいtopicを作成します。
test_orders_raw。他の設定はデフォルトのままで
→ Bob が create-topics を呼び出し、1秒ほどで作成完了。
Confluent Cloud の画面で確認できます。
Datagen Connector でテストデータを流す
Datagen Source Connectorを作成する。
quickstartテンプレートは「Orders」を選択し、
出力先topicはtest_orders_rawに設定する
→ Orders テンプレートを使った Connector が自動作成され、サンプルの注文データが流れ始めます。
メッセージを覗いてみる
test_orders_raw の中身を少し見せてください
→ Bob が consume-messages を呼び出し、最新5件を表形式で見せてくれました。
トピック test_orders_raw にデータが継続的に流入しているのが確認できます。
スキーマを登録する
Connector の出力フォーマットが JSON(Plain JSON)だったため、Schema Registry には自動登録されていませんでした。
スキーマを登録して
→ Bob がトピックからサンプルメッセージを取得し、JSON Schema(draft-07)を自動で作って登録してくれました。コードを一行も書いていません。
Flink SQL をチャットで実行
以下のFlink SQLクエリを実行してください:
SELECT * FROM test_orders_raw WHERE orderunits > 5
→ Bob が create-flink-statement を呼び出し、ストリーミングクエリが実行されました。250件以上の結果が返ってきました。
クエリが正常に実行されました。結果を見やすいレポート形式で表示します。
CTAS でフィルタ結果を新トピックへ書き込む
上記クエリの結果をCREATE TABLE ASで
新しいtopic test_orders_filtered に書き込んでください
→ CTAS ステートメントが実行され、test_orders_filtered トピックが自動作成され、orderunits > 5 のデータがリアルタイムで書き込まれ続けました。
Connector を一時停止
サンプルデータの生成を一時停止できますか
→ pause-connector が呼ばれ、Datagen Connector が PAUSED 状態に。
他にもいろいろ試してみました。
- 「test_orders_raw の中身を少し見せてください」
- 「商品(itemid)ごとに注文数の合計を集計して、ストリーミングで確認したい」
- 「今 Kafka クラスターにどんなトピックがありますか?」
つまずきやすいポイント
| 問題 | 解決策 |
|---|---|
node: command not found / バージョンが古い |
nvm install 22 && nvm use 22 |
| Bob に登録後ツールが表示されない |
mcp_settings.json の config.yaml パスが 絶対パス か確認。Bob を 再起動 する |
| Kafka 用と SR 用の Key を混同した | スコープが別なので必ず4種類別々に作成する |
confluent api-key create がエラー |
confluent environment use / confluent kafka cluster use でクラスターを選択してから実行 |
--list-tools でツールが出ない |
config.yaml のインデントや引用符の YAML 構文エラーを確認 |
まとめ
IBM Bob に Confluent MCP サーバーを接続することで、Kafka の知識がなくても自然言語で Confluent Cloud を操作 できるようになりました。
今回やったこと:
- ✅ トピック作成
- ✅ Datagen Connector でテストデータ投入
- ✅ メッセージの確認
- ✅ Schema Registry へのスキーマ登録
- ✅ Flink SQL の実行(SELECT・CTAS)
- ✅ Connector の一時停止
MCP サーバーさえ設定してしまえば、あとは Bob に日本語で話しかけるだけ。Confluent Cloud の入門として非常におすすめです。





