AWS:实时流数据为Agentic AI提供统一流式骨干
DataHot 速览
AWS介绍面向Agentic AI时代的三种实时流数据架构模式:流式特征工程+实时推理+行动、事件驱动Agent调用、CDC与流管道保持Agent记忆实时同步。文章指出,Agentic AI系统已从研究原型进入生产,需要流式基础设施持续为自主Agent提供上下文,并保持实时湖仓新鲜,同时支撑生成式BI等多类消费。本文基于AWS 2024年博客进一步展开。
为什么值得关注:数据从业者可借鉴AWS提出的统一流式骨干架构,用于构建实时数据驱动的Agent系统和智能数据平台。
译文
AI 逐段翻译两年前,关于流式数据和生成式人工智能的讨论集中在一个直截了当的问题上:如何将实时上下文输入大型语言模型(LLM),使其能够使用最新数据回答问题?我们在2024年的博客文章中探讨了这个问题,“探索生成式人工智能应用的实时流式处理”,其中介绍了将流式管道连接到基础模型的模式。
如今情况已经发生变化。今天的生成式人工智能系统不仅回答问题,它们还能观察、推理和行动。智能体人工智能应用已从研究原型进入生产现实。由智能体人工智能驱动的数据管道现在可以监控实时遥测、检测异常、决定补救策略并无人工干预地执行操作。它们跨会话保持记忆,按需查询实时数据源,并与其他智能体协调解决复杂问题。
这种转变要求流式基础设施与人工智能之间建立根本不同的关系。仅仅将上下文注入提示已经不够了。您需要这样的架构:流式数据持续驱动智能体自主行动,并保持实时湖仓新鲜,以便用于训练和检索。这些数据还流入多种消费模式,例如面向人类的生成式商业智能(BI)、用于智能体查询的标准化协议,以及为低延迟智能体上下文进行主动内存补充。
本文介绍了三种架构模式,它们共同构成智能体人工智能时代的统一流式骨干:
- 流式特征工程 → 实时推理 → 行动:连续数据流构建特征、调用人工智能模型并在单个管道中采取行动。
- 事件驱动的智能体调用:流式管道检测数百万事件中的模式,并在上下文已完全组装好的情况下触发智能体工作流。
- 实时上下文同步:变更数据捕获(CDC)和流式管道保持智能体记忆的最新状态,使其能够立即响应,而无需进行昂贵的外部调用。
以下各节将深入探讨每种模式。
模式1:流式特征工程 → 实时推理 → 行动
您正在观看一场现场足球比赛。当一名前锋在禁区内接球时,屏幕上出现人工智能生成的评论:“这是史密斯在过去3分钟内第三次在禁区触球。本赛季他在该区域的转化率为34%。”这一见解是根据流式事件数据计算得出的,经过特征管道处理,并输入生成式人工智能模型。所有这些都在前锋转身射门的时间内完成。
此模式结合了通常分开处理的两种能力:使用实时数据持续改进人工智能模型,以及使用实时数据调用这些模型以立即采取行动。流式管道两者兼做:它构建用于训练模型的特征,以及用于驱动推理的特征。
流式事件(用户交互、传感器读数、游戏事件和交易记录)流入Amazon Managed Streaming for Apache Kafka(Amazon MSK)或Amazon Kinesis Data Streams。Amazon Managed Service for Apache Flink通过窗口聚合(滚动窗口、滑动窗口或会话窗口)处理这些事件以生成特征:滚动平均值、计数、比率、行为序列或与您的用例相关的其他派生信号。
这些特征同时服务于两条路径:
推理路径:在每个窗口结束时(或根据延迟要求在每个事件时),特征传递给生成式人工智能或机器学习(ML)推理端点:Amazon Bedrock用于生成式输出,或Amazon SageMaker用于自定义模型。模型产生结果(评论、建议、个性化决策或风险评分),管道采取行动:向用户发布内容、更新推荐源、发送通知或写入下游系统。
训练路径:相同的流式特征被持续写入实时数据仓库或湖仓,例如Apache Iceberg表在Amazon S3 Tables上,这是Amazon Simple Storage Service(Amazon S3)的一项能力,保持训练数据集的新鲜度。Amazon SageMaker lakehouse架构为训练作业和微调管道提供统一访问。随着新数据的流入,您的模型可以在几分钟前的数据上进行重新训练或微调,而不是几天前的数据。这对于模式快速变化的领域(如欺诈检测、个性化和行业动态)非常重要。
Amazon S3 Tables自动处理Iceberg表管理,包括压缩、快照管理和元数据优化。您的团队专注于特征逻辑,而不是存储操作。AWS Glue Data Catalog使这些表在训练作业、推理管道和分析消费者之间可被发现。Glue Data Catalog支持业务上下文和语义搜索。这种上下文帮助模型发现并选择适合任何给定任务的正确数据资产。
场景
实时体育评论:流式游戏事件(传球、射门、球员位置)通过Managed Service for Apache Flink上的Apache Flink处理,该服务计算滚动特征(控球率、分区射门频率、球员热图)。这些特征通过Amazon Bedrock输入生成式人工智能模型,该模型实时生成自然语言评论和统计见解。同时,这些特征被写入S3 Tables,以改善模型随时间对比赛模式的理解。
流式个性化:用户点击流数据流经Apache Flink托管服务,该服务计算行为特征(会话时长、品类偏好得分、基于近期行为的购买历史)。这些特征调用个性化模型,通过重新排序产品推荐、调整内容流或触发定向优惠实时更新用户体验。相同的特征也馈送到湖仓,用于每夜重新训练个性化模型。

图1:流式特征工程馈送实时推理路径和持续训练路径
模式2:事件驱动的代理调用
凌晨2点47分,生产线上一个压力传感器开始漂移。几秒内,流式管道检测到异常,汇总完整上下文(设备历史、维护计划、相关传感器读数),并调用一个代理,该代理打开维护工单、调整设备采样率,并通知值班工程师。这一切都发生在人类看到警报之前。
模式1对每个窗口或事件调用推理。它持续运行。模式2在此基础上增加:流式管道持续分析数据,并在特定条件满足或检测到模式时调用代理工作流。管道是传感器。代理是响应者。动态规则是它们之间的桥梁。
关键区别在于事件和触发条件是动态的。它们由编程到流式管道中的规则或用于预测或检测的传统ML模型定义。管道决定何时和如何触发代理,使系统流畅且自适应。您无需重新部署代理即可更新检测逻辑。您无需更改响应逻辑即可添加新的异常模式。
流式遥测流入Amazon MSK或Amazon Kinesis数据流。Apache Flink托管服务运行连续异常检测逻辑,例如统计模型、窗口聚合、基于阈值的规则或基于ML的评分。关键是,当Flink检测到异常时,它不仅仅发布原始警报。它会组装一个上下文包:异常详情、相关历史数据、来自其他流的相关信号,以及代理立即行动所需的元数据。
该上下文包发布到下游主题,并由一个Amazon Bedrock AgentCore代理消费。由于管道已组装完整上下文,代理不会浪费时间收集信息。它可以立即推理和行动。AgentCore Runtime托管代理,AgentCore Observability提供跟踪和日志记录,AgentCore Memory维护跨调用状态(因此代理知道,例如,这是本周该设备的第三次异常)。
这种模式相对于基于轮询或计划的方法有两个好处:
- 延迟:代理在异常发生几秒内被调用,而不是在下一个轮询间隔。
- 上下文丰富性:管道已完成相关信号关联和上下文组装的工作。基于轮询的代理需要多次查询来重建管道已知的信息。
触发调用的规则是一个强大的抽象。它们可以是简单阈值(“温度超过95°C”)、统计(“值与滚动均值的偏差超过3σ”)或基于ML的(“嵌入模型的异常得分超过0.85”)。您可以通过添加新的检测模式、调整灵敏度或将不同异常类型路由到不同代理来动态更新这些规则。

图2:由流式管道中的异常检测触发的事件驱动代理调用
模式3:实时代理上下文
一位客户向他们的银行发送消息:“机场那笔847美元的收费合法吗?”代理在两秒内响应,提供完整上下文(客户最近的旅行模式、商户的欺诈风险评分和交易详情),因为所有这些都已通过流式CDC加载到代理的上下文层。没有这种同步的响应式代理需要跨三个系统进行五次单独的API调用,耗时8-12秒,并有超时失败的风险。
这种模式解决了一个基本问题:您的代理在收集上下文方面应该多主动?
主动式代理拥有完整上下文,与世界状态持续同步。当用户提问时,代理已经从上下文中拥有相关知识。它从记忆响应,而不是进行昂贵的外部调用。响应式代理从冷启动开始。在查询信息之前它一无所知,跨安全边界进行多次调用,处理认证,并从不同来源拼接数据。对于延迟敏感的用例,用户发送提示并期望快速响应,这种差异至关重要。
实时上下文同步使用CDC和流式管道保持代理记忆最新。代理的知识图谱成为它需要推理的分布式系统的同步副本。
没有代理是纯粹的主动或纯粹的反应式。设计决策是:哪些数据应预先加载,哪些应按需获取?这是一个谱系,您所处的位置取决于三个因素:
- 延迟敏感性:如果用户期望快速、上下文相关的响应,请预加载代理最频繁需要的数据。
- 数据量:同步一切是不切实际的。高效、快速且仍能产生准确结果的搜索比详尽预加载更重要。要有选择性地推送内容。
- 数据新鲜度要求:有些数据每秒都在变化(股票价格、会话状态)。其他数据很少变化(客户偏好、账户配置)。加载那些变化频繁且立即需要的数据。
流式管道(从Amazon MSK、Kinesis Data Streams或操作数据库的CDC流读取的托管Flink)持续处理事件,并将聚合结果写入代理的知识图谱或上下文层。这些存储可以根据访问模式采取多种形式:
- AWS Context自动将现有数据中的关系映射到知识图谱,并支持代理搜索,使AI代理能够在运行时访问受管控的数据关系、业务规则和领域知识。数据管理员通过直观的控制台管理图谱,审查推断出的关系,将其提升到生产环境,并附加特定领域的知识,如业务定义和使用规则。
- Amazon Bedrock AgentCore Memory用于跨会话持久化的结构化代理上下文。
- Amazon DynamoDB用于低延迟键值查找(客户资料、账户状态)。
- Amazon OpenSearch Serverless用于对非结构化上下文(历史对话、文档)进行语义搜索。
- Amazon Neptune用于关系丰富的数据(知识图谱)。
- Amazon S3 TablesAmazon S3中完全托管的Apache Iceberg表,用于多个查询引擎之间的互操作性。
对于未预加载的数据,代理回退到按需检索。这适用于数据太大、变化太频繁而不足以证明流式处理的合理性,或仅在边缘情况下需要的情况。模型上下文协议(MCP)为此提供了标准化接口。MCP服务器通过统一协议暴露异构数据源。当代理需要其同步内存中不存在的上下文时,代理会查询MCP。
这种相同的实时上下文同步模式服务不同的消费者:
AI代理通过实时知识图谱或上下文层以及MCP服务器(拉取层)访问新鲜上下文,如前述章节所述。
人类分析师和高管访问相同的上下文层,该层可以通过其直接查询模式直接查询S3 Tables上的Apache Iceberg表。Amazon Quick chat提供对实时湖仓数据的自然语言访问。不需要中间仓库。这是同一底层模式的生成式BI表达:流式数据保持湖仓最新,Amazon Quick为人类提供对话访问。
训练和微调管道通过Amazon SageMaker Lakehouse访问同步的湖仓,保持模型新鲜(如模式1所述)。
跨消费者的基本原则是相同的:流式管道将分布式数据同步到可访问的存储中,每个消费者通过适合其需求的接口访问这些存储。

图3:实时上下文同步从共享存储服务代理、分析师和训练管道
整合在一起
本文中的三种模式形成了一个统一架构,构建在单一流式骨干之上:
模式1使用流式管道构建特征,同时驱动实时推理并保持训练数据新鲜。您的模型在实时提供预测的同时不断改进。
模式2使用流式管道作为智能传感器,检测异常并调用已经组装好完整上下文的代理。这将检测逻辑与响应逻辑分开,以实现最大的灵活性。
模式3使用流式管道将分布式系统状态同步到代理的上下文层,使代理更加主动,并从相同的预加载数据服务多个消费者(代理、人类和训练任务)。
您构建的流式基础设施(Amazon MSK、Amazon Kinesis Data Streams、Amazon Managed Service for Apache Flink和Amazon S3 Tables)同时服务所有三种模式。Flink应用程序可以计算推理特征(模式1),检测触发代理的异常(模式2),并将状态同步到代理内存(模式3)。
要亲身体验本文描述的模式,请参阅基于代理的AI驱动的异常检测:实时发现异常。
您不必一次实现所有三种模式。从解决您最紧迫需求的那个开始。但设计您的流式基础设施时要意识到它将服务多种模式。在代理AI时代,每个流都是代理、模型和人类决策者的潜在输入。
这篇内容对你有用吗?
反馈只用于改善内容筛选,不等同于收藏