この参照アーキテクチャでは、イベント駆動型の取り込み、セマンティック取得、エージェント実行向けに Performance を使用してMongoDB Atlasで耐久性のあるAIワークフローを構築する方法について説明します。
失敗、再試行、長時間実行操作が確実に継続されるように、検索拡張生成 (RAG)(RAG)と複数ステップのAIワークフローを必要とするチームをサポートします。
このアーキテクチャは 2 つの取り込みパターンをサポートしています。イベントストリーム パターンでは、一時 を呼び出す前に、ソース更新はKafkaと Atlas Stream Processing を介してフローします。直接パターンでは、ソース更新によって一時ワークフローがすぐにトリガーされます。どちらの場合も、 MongoDB Atlas は運用データ、セマンティックな知識、アプリケーションの状態を保存し、一時的な は抽出、チャンク、埋め込み、インデックス作成、検索操作を調整します。
ソリューション 1: Kafkaと Atlas Stream Processing による取り込み
このパターンは、イベントトランスポート、変更伝達、または切り離されたシステム統合にKafkaをすでに使用している環境に適しています。 Kafka は大容量またはさまざまなソース アップデート用の標準的な Ingressレイヤーを提供し、Atlas Stream Processing は一時を呼び出す前にイベントを変換、ルーティングします。これは、分離取り込みがワークフローの実行とは独立して増やす必要がある場合、またはチームがAIパイプラインに加えて複数の下流の利用者に共通のイベントバックグラウンドを必要とする場合に重要です。コンテンツの更新は、S3、API、またはデータベースやメッセージング システムなどのさまざまなプラットフォームから送信されます。
図
次の図は、このフローを示しています。

図の 1。 Kafkaと Atlas Stream Processing による取り込み
データフロー
次の手順では、このフローを説明します。
ソース システムはコンテンツを生成または公開します
コンテンツは、 Amazon S3、IoT プラットフォーム、運用データベースなどの上書きシステムから提供されています。これらのシステムは、パイプラインが処理する未加工のドキュメント、レコード、またはイベントを提供します。新しいまたは更新されたコンテンツが利用可能になると、ソース システムはKafka Sink Connector にイベントを発行し、 MongoDB
sourcesコレクションにそれを書込みます。Kafka Sink コネクタ
Kafka Sink Connector は、ストリーミングレイヤーとMongoDB間の永続的なメッセージのハンドオフです。 Kafkaトピックからのイベントを消費し、ドキュメントとしてMongoDB Atlasに書込みます。このコンポーネントは、直接 trigger ではない 2 つのことを提供します。1 つは複数のさまざまなソースを 1 つの順序付けられたストリームに送受信する方法と、逆負荷をバッファする方法です。
MongoDB (Atlas、Stream Processing、ベクトル検索)
MongoDB はこのアーキテクチャ内で二重ロールを果たし、フローの 2 つの異なるポイントに表示されます。
MongoDB Atlasは、 Kafka Sink Connector イベントのランディングゾーンです。 Sink Connector が書き込む未加工のソース ドキュメントは、外部から到着したドキュメントのステージングレコードとして Atlas に存在します。
Atlas Stream Processing は、そのランディングゾーンと一時的な間のトリガーとして機能します。変更ストリームを通じて受信ドキュメントを監視し、下流のワークフローを開始します。これにより、 MongoDB はアクティブなストアからアクティブなイベントソースに変わり、ポーリング ループが不要になり、一時イベントへの引き継ぎはスケジュールされたものではなく、一時的なイベントが実行されるようになります。
Atlas ベクトル検索 はエージェントの読み取りパスを提供し、セマンティック類似性検索をネイティブ集計パイプラインステージとして実行中、別のベクトルデータベースとクロスシステムクエリの関数アウトはありません。
MongoDB はデータフローの両方に表示されます。未加工のイベントを受け取り、最終的な埋め込みベクトルを提供します。 Atlas Stream Processing は 2 つを接続します。
一時 AI 埋め込みと投票AI埋め込み
Atlas Stream Processing は、一時的なイベントを直接受信するのではなく、一時的な を呼び出します。プロセシング コアはソリューション 2 と同じです。一時ワークフローは耐久性があり、再開可能なオーケストレーションを提供し、 AI埋め込みはこれらのワークフロー内で個別に再試行可能なアクティビティとして実行されます。
このソリューションでは、一時 はソース イベントを直接受信するのではなく、 MongoDBイベントを消費します。 Atlas Stream Processing からは、trigger が S3、IoT、またはデータベースのいずれから発生したかにかかわらず、クリーンで正規化された trigger のみを受信します。このステージは双方向的です。処理が完了すると、 は埋め込み、インデックスが作成されたチャンクをMongoDB の知識コレクションに一時的に書き込みます。
ソリューション 2: データソースを一時的な に直接取り込む
このパターンは、信頼できるオブジェクト、 API、またはアプリケーションイベントをワークフロートリガーに直接発行できるシステムに適しています。アーキテクチャ レイヤーを削減し、取り込みパスをわかりやすくすると同時に、長時間実行される抽出や埋め込みステップの再開可能性を維持します。ワークフローは未加工のコンテンツではなくソース参照でトリガーされるため、下流のパイプラインはソースに依存せず、オーケストレーションの変更を最小限に抑えながら、追加の上書きシステムをサポートできます。
図
次の図は、このフローを示しています。

図の 2。データソースを一時停止に直接取り込む
データフロー
次の手順では、このフローを説明します。
ソース システムはコンテンツを生成または公開します
コンテンツは、 Amazon S3、IoT プラットフォーム、運用データベースなどの上書きシステムから提供されています。これらのシステムは、パイプラインが処理する未加工のドキュメント、レコード、またはイベントを提供します。新しいまたは更新されたコンテンツが利用可能になると、ソースシステムはAWS Lambda、 Webhook、またはコネクタなどの軽量アダプターにイベントを発行します。アダプターは一時ワークフローを開始し、ソース参照と必要なメタデータのみを渡します。これにより、ワークフローの外部でソース固有のロジックが保持されるため、 下流のパイプライン を変更せずに新しいソースタイプを追加できます。
取り込みライフサイクルの一時的な管理
ワークフローが開始されると、一時的な は取り込みステップ全体で実行、再試行、リカバリを調整します。これにより、プロセシング パスが 耐久性 が確保されるため、長時間実行される操作は障害や再起動後も確実に継続されます。ワークフローはソース コンテンツを取得し、それを下流のAI処理用の正規化された形式に変換し、セマンティック変換が開始される前にソースタイプ間で一貫した表現を作成します。
投票AI がワークフロー内に埋め込みを生成
投票AI埋め込みは、外部の Fire-and-foreget ステップとしてではなく、ワークフローの一部として実行されます。これにより、同じ実行パス内で埋め込みが観察可能であり、回復可能であり、セマンティック変換が取り込みライフサイクルに密に結合された状態が維持されます。
MongoDB Atlas はコンテンツ、メタデータ、 埋め込み を保存します。
MongoDB Atlas は、処理されたコンテンツ、関連するメタデータ、埋め込みベクトルを単一のプラットフォームに永続化します。これにより、取り込み書込み (write) パスと下流への検索パスの両方をサポートする耐久性のある知識レイヤーが作成されます。
Atlas ベクトル検索 は知識を検索可能にします
Atlas ベクトル検索 は埋め込みコンテンツにインデックスを付けると、セマンティック類似性でクエリが可能になります。埋め込みと運用メタデータは同じプラットフォームに残るため、後でのアプリケーションやエージェントのリクエストでは、別のベクトルストアや同期レイヤーを必要とせずに関連コンテキストを取得します。
エージェント リクエスト フロー
このフローは、両方の取り込みソリューションに適用されます。取り込みと同じ 耐久性 の原則を使用します。エージェントの実行を短時間のAPIリクエストとして扱う代わりに、アーキテクチャは取得と理由付けをワークフロー駆動型操作として実行し、監視、再試行、再開が可能です。これは、エージェントが複数の検索呼び出しを実行したり、外部ツールを呼び出したり、最終結果の前にプログレス ステータスの更新を返す必要がある場合に重要です。
図
次の図は、エージェントのコンポーネントを示しています。

図の 3。検索エージェントのアーキテクチャ
データフロー
次の手順では、このフローを説明します。
エージェント検索リクエストの開始
エージェントAPI がUIを介してユーザー クエリを受信すると、永続的な一時ワークフローを開始します。システムはワークフロー識別子をすぐに返すため、 UI は、検索エージェントが応答に関するコンテキストと理由を検索する際に進行状況を追跡できます。
MongoDB Atlas Vector Searchからコンテキストを取得
MongoDB Atlas は、メタデータ、ベクトル埋め込み、セマンティック インデックスをホストします。ユーザーとの対話中に、エージェントはクエリを埋め込み、 MongoDB Atlas Vector Searchを使用して関連するコンテキストを検索し、多くの場合、最終統合の前にリランク付けレイヤーを適用します。インデックス作成とクエリ パターンについては、Atlas ベクトル検索 のドキュメントを参照してください。
永続的な調査結果を返す
このワークフローは、ツールの実行、モデル インタラクション、最終合成を調整します。一時的な は再試行と状態の永続性を通じてプロセスを維持するため、 UI は検証された応答を提供します。永続的な実行の詳細については、 一時ドキュメント を参照してください。
コンポーネント
次のコンポーネントは、このアーキテクチャを実装します。
MongoDB Atlas
MongoDB Atlasを使用して、取り込みパイプラインのステージと埋め込みチャンク、 MongoDB Atlas Vector Searchインデックス、およびエージェントの状態を 単一のデータベースに保存します。エージェントは、取り込みパイプラインが書き込むのと同じデータを読み取るため、同期するためのメモリストアが別途必要になることはありません。
Atlas Stream Processing
Atlas Stream Processing を使用して、ストリーミングソースと 一時ワークフロー間のイベント駆動型の任意の統合パスを提供します。高スループット アーキテクチャの受信イベントのリアルタイム変換とルーティングを処理します。これにより、代わりにソースから一時停止への直接接続を選択する場合には、 Kafka を必要とせずにほぼリアルタイムの取り込みを実装できます。
Voyage AI
このAIの場合、埋め込みを生成し、結果の再ランク付けを実行します。埋め込みと再ランク付けはワークフローのオーケストレーションとは別であるため、チームは取り込みパイプラインと取得パイプラインとは無関係にモデルをアップグレードできます。
Atlas Vector Search
MongoDB Atlas Vector Search を使用して、セマンティックインデックス作成を通じてドキュメントコンテキストとエージェントメモリを検索します。検索エージェントは、 プライマリストレージに使用されるのと同じMongoDB Atlasクラスターをクエリするため、検索には現在の運用データが使用されます。
一時的な
ソース システム、 MongoDB Atlas、および外部AIサービス全体にわたる取り込みおよびエージェント的なワークフローの耐久性がある実行レイヤーとして Performance を使用します。抽出、チャンク、埋め込み、インデックス作成などの実行時間が長いステップを調整し、組み込みの再試行、チェックポイント、リカバリを提供します。これにより、ワークフローは再起動ではなく、障害または中断が発生した後に、最後の成功した状態から再開できます。一時的な は、オーケストレーションを永続的にすることで、チームが本番環境のAIパイプラインを確実に運用し、バックフィルして開発するのに役立ちます。
例外、注意事項、トレードオフ
このアーキテクチャを採用する前に、直接取り込みと Kafka ベースの取り込みのトレードオフを考慮してください。ソースと一時接続の直接接続では、操作する移動部分は少なくなります。 Kafka ベースの パスでは メッセージ プロバイダーが追加され、追加のインフラストラクチャが導入されますが、組織がすでにKafkaを通じてソース アップデートをルーティングしている場合には、自然に適しています。