サブタイトル:需要シグナル(Data in Motion)と確定在庫(SoR)を組み合わせた Agentic AI アーキテクチャ
はじめに
本記事では、需要シグナル(需要の兆候)とSoR(System of Record)の役割の違いを整理しながら、Confluent・Apache Flink・watsonx Orchestrateを利用したイベント駆動型AIシステムを紹介します。
本文
IBM DeveloperのBuilding an event-driven agentic AI system with Apache Kafka on Confluent Cloud and watsonx Orchestrateを元に、Confluent・Flink・watsonx Orchestrate を利用したイベント駆動型 AI システムを実際に構築できます。こちらを試してみました。
実際に動かしてみると多くの学びがありましたが、その中でも特に気になったのが、「DBでも業務のイベントデータを保存できるのでは?」という疑問です。
そこで、ハンズオン内容をベースにし、私自身がより説明しやすい形にアレンジしたデモを作りました。
その中で、需要シグナル(Data in Motion)と SoR の違い、そして Agentic AI 時代における Confluent の価値についても整理してみました。
今回のデモ作成は、IBM Bobを活用しました。
デモの設計支援だけでなく、ConfluentやIBM watsonx Orchestrateへのコンソールやコマンドラインからの実行もIBM Bobから実施できました。Kafka Topic作成やFlink SQL作成の初期実装を自然言語から生成でき、調査や実装にかかる時間を大幅に短縮できます。
今回はConfluentデモ部分に集中したいので、IBM Bobで実施した流れは簡略いたします。
Qiitaでは多くのIBMBobやBobでタグ付けされた投稿があります。参考にされるとIBM Bobの凄さを感じられると思います。
このカスタマイズ・デモを作ろうとした経緯
ConfluentやKafkaを紹介していると、
「在庫データや需要シグナルならDBだけでよくないですか?」
という質問をいただくのではないか?と思いました。
確かに在庫や注文のような確定データはDBで管理できます。では、
・ECサイトの商品閲覧
・カート投入
・店舗での商品接触
・営業案件の進捗
といった需要の予兆イベントも、すべてDB中心に処理すべきなのでしょうか。
DB中心のデータ管理では、システムごとに責任範囲が決まっています。
一方で需要シグナルは、AI Agent、BI、需要予測など複数の利用者が後から現れる可能性があります。
そのため 需要シグナルは、特定システムのデータとしてではなく、企業全体で共有・再利用するデータとして考える必要があります。
今回はこの疑問に対する一つの答えとして、
・確定在庫(SoR)
・身の回りに起きている「シグナル」のような情報、需要シグナル(Data in Motion)
を分離し、Confluent + Flink + AI Agentでリアルタイム活用するデモを作成してみました。
本記事では、商品閲覧やカート投入、商談開始のような「需要を示唆するイベント」を 「需要シグナル」 と呼びます。
⭐️ この記事のポイント
ポイントは最後にまとめてみたいと思います。
デモの流れとそのための全体像を先にご覧ください。
システム・アーキテクチャ図 (比較)
全体の流れとして左が「システム毎にDBで管理」、右が「SoR + Data in Motion(Confluent)」

デモの流れ
Step 0(Kafka Topicが流れる前)
3つの Kafka Topic を作っています。
inventory.stock:在庫DBから。アーキテクチャ図の左上
demand.signals:需要シグナル。アーキテクチャ図の左下
demand.insights:在庫と需要シグナルから集計。アーキテクチャ図の右側の中央赤枠(Kafka Topic+Flink)

Kafka Topic の inventory.stock にスキーマ定義し、ConfluentのSQL Workspacesコンソールから確認してみました。
Kafka Topic にデータが流れる前なので、なにも結果は取得されません。

Step 1(在庫投入)
確定している 在庫データ(アーキテクチャ図の左上) より CDC(*)を受けてConfluent KafkaにProduceします
今回は、IBM BobからPythonで擬似的にTopicにデータを流します
流す先はinventory.stockです
👇 IBM Bobから事前に設計・実装した仕組みを実行してもらいます

*CDC:Change Data Capture
実行した結果が次の通りです
👇ConfluentのSQL Workspacesコンソール
inventory.stockから5件取得できました。

👇 BIダッシュボードからも 「在庫」 が確認できています。
顧客の商品閲覧やカート投入といった「需要シグナル」は未だ流れていない状態です。

Step 2(需要シグナル投入)
ECサイト・モバイルアプリ・店舗・営業からの 需要予兆イベント を シミュレーションします 。こちらもWebアプリケーションから投入されるデータを擬似的に Kafka Topic にデータを流します。
流す先はdemand.signalsです。
👇 BIダッシュボードから「在庫」と 「需要シグナル」 が確認できています。
BIダッシュボードは demand.insights を参照しており、需要シグナルが到着すると自動的に最新状態へ更新されます。

👇ConfluentのSQL Workspacesコンソール。
demand.signalsから5件取得できました。

Kafka Topic に流れてきた在庫DBと各需要シグナルのイベントをFlinkで集計し、今回は需要シグナルごとに重み付けを行い、需要スコアを算出しています。
👇 ConfluentのFlink Statement (INSERT INTO demand.insights)
Kafka Topic の需要シグナルと在庫データを集計し、需要スコアと在庫状況を組み合わせた(Join)結果を demand.insights へ Insert しています。

例えば、
・PRODUCT_VIEW(閲覧) = 1
・WISHLIST_ADD(お気に入り) = 2
・CART_ADD(カート投入) = 3
・DEAL_STARTED(商談開始) = 5
といった重みを与え、Flink SQLでリアルタイム集計しています。

実際の業務では機械学習やAIモデルによるスコアリングも考えられますが、今回は仕組みを分かりやすくするためルールベースで実装しています。
Confluent Cloud for Apache Flinkでは、提供される機能・契約条件に応じて、Flink SQLからAI/ML関数を利用することもできます。

👇ConfluentのSQL Workspacesコンソール。
demand.insightsへの集計データはこちら。

AIエージェント(カスタマーサポート・エージェント)からも確認してみました。
今回はAIエージェントとしてIBM watsonx Orchestrateを利用しています。
Kafka Topic に流れている需要シグナルは、watsonx Orchestrate に登録した MCP Server を通して demand.insights を参照(※)しています。

※)watsonx Orchestrate に登録した MCP Server


Step 3(さらに需要シグナルを Kafka Topic に流します。需要スコアが急増します。)
さらに、いくつかのAIエージェントに聞いてみました
参考までにIBM watsonx Orchestrateで用意した3つのエージェントは次のとおりです

「カスタマーサポート・エージェント」に聞いてみた
👉「iPhone 18 の在庫ありますか?(Shibuya 店)」

「調達エージェント」に聞いてみた
👉「今すぐ緊急発注が必要な商品はありますか?(Shibuya 店)」

「マーケティング・エージェント」に聞いてみた
👉「iPhone 18 の販促施策を提案してください(Shibuya 店)」

いかがでしょうか。
もちろん需要シグナルも含めて、すべてのデータを DB で処理するという考え方もあります。
しかし、需要シグナルは複数のシステムから継続的に発生し、AI Agent、BI、分析基盤など複数の利用者が同時に活用したいデータです。
今回のようにストリームとして扱うことで、
・データが特定システムに閉じない
・リアルタイムに 共有できる
・新しい利用先を追加しやすい
⭐️ この記事のポイント
1. 「Kafka は DB の代替ではない」
Kafka(Confluent)が担うのは、今この瞬間に組織中で起きているイベントを全員が同時に受け取れる共有基盤を作ることです。
在庫の正確な数はDBが持ちます。Kafka はその「変更イベント」をリアルタイムに流す通路です。
2. 「SoR(*) は DB、Kafka は Data in Motion の共有基盤」
注文確定・入出荷・在庫調整などの確定データは DB(SoR)が管理します。
一方で「カートに入れた」「商談が始まった」「店舗で商品を手に取った」といった需要の予兆は Kafka(Data in Motion の共有基盤)が束ねます。この2つを混同しないことが設計の鍵です。
(*)System of Record
3. 「Kafka は従来型MQとは異なる — 1:N のファンアウト」
従来のメッセージキュー(MQ)は 1 件のメッセージを 1 Consumer が読んだら消えます。
Confluent Kafka は Consumer が好きなタイミングで何度でも参照できます。同じイベントを AI Agent・BI ダッシュボード・需要予測・マーケティング分析が同時に利用できる — これが 「ファンアウト」 です。
4. 「Flink が異なるシステムをリアルタイムに結合する」
DB の確定在庫(CDC 経由)と、EC サイト・モバイル・店舗・営業からのリアルタイム需要シグナルは、そのままでは別世界のデータです。
Confluent の Flink SQL が 2 つのストリームを JOIN し、「需要スコア × 確定在庫」の突き合わせ結果を派生トピック demand.insights に出力します。
5. 「IBM Bob で AI 駆動開発 — 自然言語でインフラを設計する」
Kafka トピックの作成、Flink SQL の記述、メッセージの Produce まで、IBM Bob に日本語で指示するだけでコードが生成され、実行できます。
CLI や GUI ツールの知識は依然として重要です。
一方で IBM Bob を活用することで、Kafka Topic 作成や Flink SQL 作成の初期実装を自然言語から生成できるため、今回調査や実装にかかる時間を大幅に短縮できました。
伝えたいこと、まとめ
今回作成してみて改めて伝えたいと感じたことは、
・Kafka は DB の代替ではなく、企業内で発生するイベント(需要シグナル)を共有するための基盤であること
・Flink が「SoR」と「Data in Motion」の世界を結び付ける役割を果たしていること
です。
需要シグナル のような 「今起きている需要を示唆するイベント」 を 鮮度の高い状態で扱うことで、AI Agent は単なる問い合わせ対応ではなく、需要の変化を踏まえた意思決定支援へ近づいていきます。
そのためには、 システムごとの責任範囲に閉じず、需要シグナルを企業全体で共有・再利用できる仕組みの価値が高くなる と感じます。
記事:Harukichi


