对于 AI 代理:可在 https://www.mongodb.com/zh-cn/docs/llms.txt 获取文档索引—通过在任何 URL 路径后添加 .md 可获取所有页面的 Markdown 版本。
Docs 菜单

使用时态和MongoDB Atlas构建持久的AI工作流程

此参考架构描述了如何使用时态在MongoDB Atlas上构建持久性的AI工作流程,以实现事件驱动的摄取、语义检索和代理执行。

它支持需要检索增强生成 (RAG) 和多步骤AI工作流程的团队,以便在出现故障、重试和长时间运行的操作时可靠地继续运行。

该架构支持两种摄取模式。在事件流模式中,源更新在调用时态之前流经Kafka和Atlas Stream Processing 。在直接模式中,源更新会立即触发时态工作流程。在这两种情况下, MongoDB Atlas都会存储操作数据、语义知识和应用程序状态,而时态则会协调提取、分块、嵌入、索引和检索操作。

此模式适合已使用Kafka进行事件传输、变更传播或解耦系统集成的环境。 Kafka为大容量或各种源更新提供标准入口层,而Atlas Stream Processing在调用时态之前转换和路由事件。当数据摄取必须独立于工作流程执行进行扩展时,或者当团队希望除了AI管道之外,在多个下游使用者之间有一个通用的事件主干时,这种分离就很重要。内容更新源自不同的平台,如 S3、API、数据库或消息传递系统。

下图显示了此流程。

通过Kafka和Atlas Stream Processing进行摄取

图 1。通过Kafka和Atlas Stream Processing进行摄取

点击放大

以下步骤描述了此流程:

  1. 源系统生成或公开内容

    内容源自Amazon S3、物联网(IoT)平台和操作数据库等上游系统。这些系统提供管道处理的原始文档、记录或事件。当新的或更新的内容可用时,源系统会向Kafka Sink Connector 发出事件,后者会将其写入MongoDB sources集合。

  2. Kafka Sink 连接器

    Kafka Sink Connector 是流媒体层和MongoDB之间的持久性消息切换。它使用来自Kafka主题的事件并将其作为文档写入MongoDB Atlas 。该组件提供了直接触发所没有的两个功能:它将多个不同的源扇入一个有序流,并针对背压进行缓冲。

  3. MongoDB (Atlas、流处理和向量搜索)

    MongoDB在此架构中扮演双重角色,出现在流程中的两个不同点。

    1. MongoDB Atlas是Kafka Sink Connector 事件的登陆区域。 Sink Connector 写入的原始源文档位于Atlas中,作为从外部世界到达的内容的暂存记录。

      Atlas Stream Processing充当该登陆区域和时域之间的触发。它通过变更流监视传入文档并启动下游工作流程。这将MongoDB从被动存储转变为主动事件源,无需轮询循环,并使切换到时态事件驱动而不是按计划进行。

    2. Atlas Vector Search为代理提供读取路径,将语义相似度搜索作为原生聚合管道阶段运行,没有单独的向量数据库,也没有跨系统查询扇出。

    MongoDB出现在数据流的两端:它接收原始事件并为最终的嵌入式向量提供服务。 Atlas Stream Processing将两者结合起来。

  4. 时间和 Voyage AI嵌入

    Atlas Stream Processing调用 Tempal,而不是直接接收事件的 Tempal。处理核心与解决方案 2 相同:时态工作流程提供持久性、可恢复的编排,而 Voyage AI嵌入在这些工作流程中作为可独立重试的活动运行。

    在此解决方案中,时态使用MongoDB事件,而不是直接接收源事件。它只从Atlas Stream Processing接收干净、规范化的触发,无论该触发是源自 S3、物联网(IoT)还是数据库。此阶段是双向的:处理完成后,Temporal 会将嵌入式索引数据段写回到MongoDB知识集合中。

此模式适合可将可靠对象、 API或应用程序事件直接发送到工作流触发的系统。它减少了架构层数,保持摄取路径简单明了,同时仍然保留了长时间运行的提取和嵌入步骤的可恢复性。由于工作流程是使用源引用而不是原始内容触发的,因此下游管道与源无关,并且可以通过最少的编排更改来支持其他上游系统。

下图显示了此流程。

将数据源直接引入时态

图 2。将数据源直接引入时态

点击放大

以下步骤描述了此流程:

  1. 源系统生成或公开内容

    内容源自Amazon S3、物联网(IoT)平台和操作数据库等上游系统。这些系统提供管道处理的原始文档、记录或事件。当新的或更新的内容可用时,源系统会向轻量级适配器发出事件,例如 AWS Lambda、Webhook 或Connector。适配器启动时态工作流程,仅传递源引用和所需的元数据。这会将特定于源的逻辑保留在工作流程之外,因此您可以在不更改下游管道的情况下添加新的源类型。

  2. 时间管理摄取生命周期

    工作流程启动后,时态会协调各摄取步骤的执行、重试和恢复。这使得处理路径持久性,因此长时间运行的操作在出现故障或重新启动时仍能可靠地继续运行。该工作流检索源内容并将其转换为规范化形式以供下游AI处理,从而在语义转换开始之前跨源类型创建一致的表示。

  3. Voyage AI在工作流程中生成嵌入

    Voyage AI嵌入作为工作流程的一部分运行,而不是作为外部的即发即弃步骤运行。这样可以将可观察值和可恢复值嵌入同一执行路径中,并使语义转换与摄取生命周期紧密结合。

  4. MongoDB Atlas存储内容、元数据和嵌入

    MongoDB Atlas将处理后的内容、相关元数据和嵌入向量保留在单个平台中。这将创建一个持久性知识层,支持摄取写入路径和下游检索路径。

  5. Atlas Vector Search使知识可检索

    Atlas Vector Search会对嵌入式内容编制索引,以便您可以通过语义相似度查询。由于嵌入和操作元数据保留在同一平台中,因此以后的应用程序和代理请求检索相关上下文,而无需单独的向量存储或同步层。

此流程适用于两种摄取解决方案。它使用与摄取相同的持久性原则:该架构不是将代理执行视为短期API请求,而是将检索和推理作为工作流驱动的操作来运行,您可以观察、重试和恢复这些操作。当代理执行多个检索调用、调用外部工具或需要在最终结果之前返回渐进式状态更新时,这一点很重要。

下图显示了代理的组件。

研究代理架构

图 3。研究代理架构

点击放大

以下步骤描述了此流程:

  1. 启动代理研究请求

    当代理API通过用户界面接收用户查询时,它会启动持久性时空工作流程。系统会立即返回工作流标识符,以便用户界面可以在研究代理检索回答的上下文和原因时追踪进度。

  2. 从MongoDB Atlas Vector Search 检索上下文

    MongoDB Atlas托管元数据、向量嵌入和语义索引。在用户交互期间,代理会嵌入查询并使用MongoDB Atlas Vector Search 来检索相关上下文,通常在最终合成之前应用重新排名层。有关索引和查询模式,请参阅Atlas Vector Search文档。

  3. 返回持久性的研究结果

    研究工作流程可协调工具执行、模型交互和最终综合。时态通过重试和状态持久性来维护进程,因此用户界面会提供经过验证的回答。有关持久性执行的更多信息,请参阅时态文档。

以下组件实现了此架构。

使用MongoDB Atlas将摄取管道的分阶段和嵌入式数据段、 MongoDB Atlas Vector Search索引以及代理的状态存储在单个数据库中。由于代理读取的数据与摄取管道写入的数据相同,因此无需单独的内存存储来保持同步。

使用Atlas Stream Processing在流媒体媒体源和时态工作流程之间提供可选的事件驱动型集成路径。它为高吞吐量架构处理传入事件的实时转换和路由。如果您选择直接从源到时间连接,则可以在不需要Kafka的情况下实现近乎实时的摄取。

使用 Voyage AI生成嵌入并为此架构执行结果重新排名。由于嵌入和重新排序与工作流程编排是分开的,因此团队可以独立于摄取和检索管道升级模型。

使用MongoDB Atlas Vector Search 通过语义索引检索文档上下文和代理内存。研究代理会查询用于主节点 (primary node in the replica set)存储的同一MongoDB Atlas 群集,因此检索会使用当前的操作数据。

使用时态作为持久性执行层,用于跨源系统、 MongoDB Atlas和外部AI服务的摄取和代理工作流程。它协调提取、分块、嵌入和索引等长时间运行的步骤,并提供内置重试、检查点和恢复功能。这允许工作流在出现故障或中断后从上次成功状态恢复,而不是重新启动。通过使编排持久性,时态帮助团队可靠地操作、回填和发展生产AI管道。

在采用此架构之前,请考虑直接摄取和基于 Kafka 的摄取之间的权衡。源到时态的直接连接需要操作的移动部件较少。基于 Kafka 的路径增加了一个消息代理,这会引入额外的基础架构,但如果您的组织已经通过Kafka路由源更新,则自然适合。

有关部署指导和技术文档,请参阅 mdb-temporal-pra GitHub存储库。