返回 文章 build CMS 文章

Cloudflare 如何构建统一数据平台及 AI 数据代理

Cloudflare 通过 Town Lake 统一数据平台和 Skipper AI 代理,让非技术人员也能用自然语言查询海量数据。

Cloudflare数据平台AI代理Trino
成长分 / 100 84 综合收获、行动、留存与影响

Cloudflare 如何构建统一数据平台及 AI 数据代理
为什么值得读了解超大规模公司如何解决数据蔓延、多系统查询和权限治理的痛点。

学习如何用 AI 代理(Skipper)结合多层上下文(模式、血缘、代码)生成准确 SQL 查询。

关键洞察
  1. Town Lake 基于 Trino + Iceberg + R2 构建,统一查询 Postgres、ClickHouse、R2 等异构数据源。
  2. Skipper 通过 5 层上下文(模式、人工注释、代码衍生知识、精选模型、运行时自省)减少 LLM 幻觉。
  3. 采用“代码模式”MCP 服务器,让模型用 JavaScript 编程式调用工具,减少往返次数。
转成行动

深入阅读

正文与原文对照

原文保真覆盖:全文原文字符:20579

Cloudflare 每秒处理超过十亿个事件。我们的网络覆盖 120 多个国家的 330 多个城市。在每一个 HTTP 请求、每一次 Worker 调用、每一次 R2 读取操作的背后,都有数据,而且数据量巨大。

多年来,这些数据并不容易访问。它们分布在数十个生产数据库、ClickHouse 集群、Kafka 流、Google Cloud 存储桶、BigQuery 数据集以及大量的管道中。要回答一个简单的问题,比如“今天注册的域名中有多少进入了流量前 100 名?”,Cloudflare 的分析师必须知道该询问哪个系统、使用什么凭据、编写什么查询语言,以及他们查看的数据是采样数据、新鲜数据还是七天前的陈旧数据。因此,很难从数据中获取有见地的洞察。

为了解决这个问题,我们构建了两个内部工具:Town Lake(Cloudflare 的统一数据分析平台)和 Skipper(运行在其之上的 AI 数据代理)。Town Lake 是一个单一的 SQL 接口,可以访问 Cloudflare 所知的一切,而 Skipper 让 Cloudflare 的任何人都可以用简单的英语提问,并在几秒钟内获得正确、可审计的答案。

这就是我们构建这两个工具的故事。

如果你曾在经历过超高速增长期的公司工作过,你就会知道数据蔓延是什么样子。我们的数据蔓延有一些具体症状:

太多不同的系统。 想要调查客户问题的产品工程师可能需要查询 Postgres 以获取账户元数据、ClickHouse 以获取分析事件、BigQuery 以获取使用量汇总、R2 以获取原始日志,以及 Kafka 主题以获取实时信号。每个系统都有自己的凭据、自己的语言和自己的保留策略。

采样数据。 这对于仪表板来说没问题,但对于计费等领域则不适用。我们的 分析管道 进行降采样以处理每秒超过 7 亿个事件。当你希望分析仪表板加载时,这是正确的行为,但当你试图计算某人的使用量以开具发票时,这恰恰是错误的做法。

内部数据的外部依赖。 我们之前的内部报告堆栈部分由外部供应商提供支持。除了成本之外,我们的一些关键数据还严重依赖于另一个云。

没人能找到数据。 即使你拥有所有正确的凭据,你也需要知道“按账户计费的 Workers 请求”的正确表位于特定的 ClickHouse 集群、特定的模式中,并与特定的 Postgres 维度表连接,而且该连接需要一个晦涩的客户 ID 转换。这涉及太多部落知识。

我们还面临文化挑战:数据基础设施历来被视为服务于业务的后台功能,而不是其本身的关键基础设施。

我们希望创建一个地方,让公司中任何拥有适当权限且需要了解的人都能获得关于 Cloudflare 的问题的答案:“显示上一季度收入前 100 的客户”、“列出过去 48 小时内来自特定 ASN 且评分 > 0.9 的所有 Bot Management ML 评分事件”、“查找消费超过 100 美元的客户中排名前 100 的计费支持工单”等。

我们希望这个地方能为需要它的查询(如计费或安全调查)提供新鲜、准确、未采样的数据,并为不需要的查询(如仪表板或探索)提供快速、降采样的数据。

我们希望安全性和治理能力内置于平台中,能够自动检测个人身份信息(PII),并且默认锁定敏感表。所有访问都应可审计,并具有有时间限制的权限授予,以便用户仅在积极处理需要数据的任务时才能访问数据。

我们希望它建立在 Cloudflare 自己的平台上:R2 用于存储,Workers 用于计算,Cloudflare Access 用于身份验证,Workflows 用于编排。如果我们要对数据基础设施进行重大投资,那么它将建立在与我们销售给客户相同的产品之上。

最终,我们希望有一个不需要了解任何 SQL 的界面。目标是让公司中任何拥有适当权限且需要了解的人都能查看流经我们网络的数据流,而不仅仅是分析师。

最后一个需求成为了 Skipper。

核心上,我们的数据平台架构是一个 数据湖仓:一个从对象存储读取数据的查询引擎,带有一个元数据层,使存储表现得像数据库。我们称之为 Town Lake,以其在德克萨斯州奥斯汀的同名湖泊命名。

其最重要的组件包括:

查询引擎。 我们选择了 Apache Trino:一个 SQL 查询可以连接 Postgres 表、ClickHouse 表和 R2 上的 Iceberg 表,而无需将中间结果物化到另一个系统。一个查询“本周按 Workers 请求数排名前 100 的付费客户是什么”会被编译成一个计划,将过滤器推送到 ClickHouse,与 Postgres 中的账户维度连接,并针对 R2 中的计费汇总进行排名,一气呵成。

R2 数据目录,我们的托管 Apache Iceberg 服务,是冷数据和温数据所在之处。Iceberg 提供了模式演化、时间旅行、分区演化以及随着数据老化进行压缩的能力。上周的每分钟使用量变成每小时,上季度的每小时变成每天,等等。存储成本随着数据的新鲜度降低而降低,同时数据仍然可查询。与将相同数据保存在 OLAP 数据库中相比,R2 中的 Parquet 文件要便宜得多。

DataHub 是我们的元数据目录。每个表、列、所有者、血缘边和术语表条目都在其中。当用户问“townlake.dim.accounts 里有什么”时,DataHub 会提供答案,包括表描述、列描述、所属团队、提供数据的上游表以及消费数据的下游表。

Lifeguard 是我们的访问控制服务:它在 D1 中存储访问规则,从我们的内部访问管理系统动态拉取用户和组成员身份,并生成一个组合的 JSON 策略,Trino 通过 HTTP 读取该策略。Lifeguard 还将基本访问信息提供给 Skipper 和 Gateway,因此用户在查询之前就会被阻止在前门。

Skimmer 是一个 PII 检测扫描器。它持续运行,从每个表的每一列中采样行,并使用 Workers AI 对每列是否包含 PII 进行分类。它分两遍进行:首先,一个快速的逐列分类器;然后,如果有任何标记,第二遍是代理性的,获取完整的表上下文,并可以直接查询 Trino 进行验证。发现结果流入 DataHub 和 Lifeguard 的白名单,以便进行人工审核。

Transformer 是我们基于 Workflows 构建的 ELT(提取、加载、转换)引擎。用户通过 YAML 前置元数据(目标表、物化模式、依赖关系、调度)定义 SQL 转换的有向无环图(DAG)。Transformer 编译该图并在 Trino 上运行,状态由 Durable Objects 管理,定义存储在 R2 中,运行历史记录在 D1 中。

Ingestion 是从运营系统到数据湖的桥梁。一个编排器作为长期运行的 Kubernetes 部署运行,读取管道配置,并生成短期工作任务,从 Postgres 或 ClickHouse 提取数据,转换为 Parquet,然后作为 Iceberg 表加载到 R2 中。每个管道以全量替换或增量追加模式运行。

默认关闭:通过构造实现治理

构建统一数据平台时的一个实际问题是,你刚刚构建了一个大型敏感数据面。对此的传统答案是:默认开放,例外限制。允许访问所有内容,然后在有人注意到时审计并锁定敏感表。

Town Lake 采取相反的方法。表在审核之前无法查询。当新数据库连接到 Trino 或创建新表时,Skimmer 会扫描它,对其列进行分类,并将其注册到中央允许列表中,状态为待定。在审核者批准该表及其特定列之前,用户无法查询它。这听起来很痛苦,而且确实如此,但有两件事例外。

首先,它是自动化的。Skimmer 的分类器相当不错:它能捕获明显的 PII(电子邮件、IP、姓名、电话号码)以及长尾的非明显敏感数据(匹配特定前缀的 API 令牌、可追溯到用户的不透明 ID)。审核者会看到检测到的内容,并批准、覆盖或拒绝。大多数审核只需几秒钟。

其次,工作流程是自助式的。如果你查询一个没有权限的表,错误消息不是“权限被拒绝”,而是“此表需要审核,请点击此处请求审核”。AI 代理 Skipper 甚至会建议正确的 RBAC 组并直接链接到它。

我们将模式发现与数据访问分开。用户可以看到存在哪些表,但未审核的列会从 DESCRIBESHOW COLUMNS 以及 SELECT * 中隐藏。这种微妙的区别很重要:这意味着新的未审核列不会破坏基于已批准表其余部分构建的现有仪表板。

PII 是按会话选择加入的。默认情况下,Trino 会在敏感列到达你的屏幕之前对其进行编辑。如果你有正当理由需要原始 PII(例如欺诈调查),你可以翻转会话中的位,检查你的权限,然后解除编辑。每次翻转和每次查询都会被记录。

Skipper:AI 数据代理

如今,仅靠查询引擎是不够的。SQL 仍然是一个障碍,同样,知道要查询成千上万个表中的哪一个也是一个障碍——你需要知道规范模式。

Skipper 是我们对对话式 AI 代理的尝试,它从自然语言问题到经过验证的答案,基于公司的实际数据、代码和机构知识。我们将其构建在 Town Lake 之上,并基于我们的开发者平台:Workers、Workers AI、Durable Objects、D1、R2、Workflows、KV。

界面是一个聊天框。提问:

显示过去 30 天内按 R2 存储成本排名前 10 的客户,以及与之前 30 天相比的变化。

Skipper 找到正确的表(DataHub 搜索),提取其模式和血缘,编写 SQL,提交给 Trino,轮询结果,并向你展示表格或图表。后续操作:

现在按区域细分,并忽略内部 Cloudflare 账户。

它携带上下文,优化查询,并重新运行。如果出现错误,例如连接产生零行或过滤器排除了预期内容,Skipper 会在闭环推理中进行调查、调整并重试。难点在于拥有正确的上下文。

Skipper 还可以将图表打包成仪表板,这些仪表板可以在内部共享并嵌入到其他内部应用程序中。它还有通过 Transformer 构建转换图以及通过 Lifeguard 检查访问和权限的工具。

Skipper 随时随地满足用户需求。所有这些工具都可以通过一个由 Workers AI 驱动的内置代理框架支持的 Worker 获得。另一方面,我们的许多内部用户通过本地代理流程工作,Skipper 的工具还可以通过 MCP 服务器获得。

LLM 在给定 SQL 提示和表名列表时,可能会幻觉连接、误用列,并自信地产生完全错误的数字。我们在早期实验中通过艰难的方式学到了这一点。解决方案是模型在检索时可以提取的多层接地上下文。

第 1 层:模式和使用元数据。DataHub 知道每个表的每一列、每个类型、每个主键、每个外键。它还知道哪些表基于历史查询模式经常连接在一起。Skipper 的 search_datasets

get_entity_details

工具直接提供这些信息。

第 2 层:人工注释。当拥有 dim.accounts

的团队编写描述如 "账户级实体。每个 account_id 一行。每个账户恰好属于一个客户(通过 customer_id 外键)," 该描述存在于 DataHub 中并最终进入 Skipper 的上下文。像 curated

这样的标签标记了 Skipper 应优先于临时空间的已验证表。

第 3 层:代码衍生知识。一些最有价值的上下文不在任何目录中:它存在于生成表的 SQL 中。Transformer 管道在每次成功运行时向 DataHub 发出每个节点的 .meta.json

文档。因此,当 Skipper 查看 fct.billings_allocated

时,它不仅看到模式;它看到这是一个从 dim.accounts

dim.customers

seed.product_classification

构建的预连接事实表,其 alloc_amount

列计算为 billed_amount / 12 for annual; billed_amount for monthly

。这种细微差别区分了正确答案和自信的错误答案。

第 4 层:精选数据模型。我们维护一小部分“数据模型”页面:简短的人工编写文档,描述如何思考计费、客户、账户和区域。*"优先选择标记为 'curated' 的表。避免 *scratch_r2

  • 和标记为 'internal' 的表。使用数据模型术语(例如 'billing product revenue')搜索,而不是自然语言。"* 这些作为 MCP 资源提供,当问题匹配时代理可以提取。

第 5 层:运行时自省。当其他所有方法都失败时,Skipper 可以向 Trino 发出实时查询:DESCRIBE table, SELECT DISTINCT col LIMIT 20, SELECT COUNT(*)

。它谨慎使用这些,因为运行时上下文成本高昂,但它是使系统其余部分健壮的安全网。

Skipper 作为 MCP:代码模式

有一个具体的实现细节值得单独拿出来讲,因为它完全是 Cloudflare 风格的解决方案。

当你用工具构建 AI 代理时,标准模式是在提示词中定义工具,让模型逐个调用它们,解析响应,执行并返回结果。这没问题,但很啰嗦:一个五工具的工作流需要五次模型往返,每次都要重新建立上下文。

对于我们的 MCP 服务器,我们使用 代码模式。我们不定义 30 个单独的工具,而是暴露两个:search

execute

。模型编写一个 JavaScript 片段,以编程方式调用我们的整个工具集:

const datasets = await skipper.search_datasets({ query: "billing product revenue" })
const queryId = await skipper.start_query({ sql: "SELECT ..." })
const results = await skipper.fetch_results({ queryId, mode: "inject" })
return skipper.create_chart({ chartType: "bar", data: results.rows, ... })

该 JavaScript 通过 WorkerLoader 在沙盒化的 Dynamic Worker 隔离环境中运行。模型能够用其已非常熟悉的语言,在单次往返中表达复杂的多步骤工作流。这更快、更便宜,且生成的工作流可作为代码进行审计。

安全模型即数据模型

Skipper 的所有操作均以调用用户的身份执行。如果你无权访问某个表,Skipper 就无法为你查询它。如果你请求 PII,系统会检查你的权限。如果你保存的查询与队友共享,他们的访问权限会在查看时而非保存时进行检查,因为组成员身份会发生变化。

共享仪表板有其独特之处。它们可以通过单个占位符 div 和 script 标签嵌入到任何内部 Cloudflare 工具中:

<div data-skipper-dashboard="dash-123"></div>
<script src="https://skipper.cloudflare.com/embed.js" async></script>

iframe 会自动调整大小以适应内容。内容安全策略(CSP)的 frame-ancestors 指令会阻止来自企业域外任何地方的嵌入。Cloudflare Access 仍然对 iframe 内容进行访问控制,因此未认证的查看者会在 iframe 中看到 Access 登录页面,而不是数据。非所有者查看者会根据底层表进行检查:如果他们没有访问权限,系统会引导他们向正确的群组请求权限。

它的作用:极速回答

计费。 这是最初的用例。我们的可计费使用量仪表盘(面向客户的仪表盘,向按量付费用户精确显示其欠费金额)由计量管道驱动,其数据源是 R2 中的一组 Iceberg 表,通过 Trino 查询。该仪表盘的 API 拉取与计费系统相同的紧凑行 (date, account_id, metric_name, usage),因此仪表盘上的数字与账单上的数字一致。

与计费相关的查询占 Town Lake 所有查询的 53%:在最近的测量周期中,来自 324 名不同 Cloudflare 员工的 91,760 次查询。过去用于按客户计算收入汇总的 200–300 行传统 SQL 查询现在只需五行。

商业智能。 现在在 Skipper 中,“按收入排名前 100 的客户”这个问题大约需要三秒钟。“今天注册的域名中有多少在前 100 名”也是如此。我们过去需要提交 Jira 工单的大多数数据相关问题也是如此。

安全分析。 我们的机器人管理团队使用 Town Lake 查询过去 48 小时内 ML 评分事件中得分 > 0.9 且按 ASN 和地理位置过滤的数据。威胁研究人员在其上构建了自己的查询工具包。信任与安全团队拉取信号以协助监管滥用行为。

客户支持。 “查找消费超过 100 美元的客户中排名前 100 的计费支持工单”过去是一个需要多天的项目。现在只需一个 Skipper 查询。

有几件事让我们感到惊讶。

少提示反而更好。早期版本的 Skipper 有详细且指令性的系统提示:“首先,使用 search_datasets。然后,使用 get_entity_details。然后,如果需要,使用 list_schema_fields...” 质量反而下降了。模型擅长推理分析工作流;不需要微观管理。我们用高层指导替换了指令性提示,让模型自行选择路径。结果变得更好。

工具重叠是毒药。我们最初暴露了每个工具的每个变体:三个不同的“获取结果”工具,两个“搜索”工具,几个“列出”工具。模型感到困惑并调用了错误的工具。我们进行了整合。现在 fetch_results 有一个 mode 参数 (inject / display / both),而不是三个独立的工具。每个工具只有一个存在的理由。

代码而非元数据捕获含义。最大的准确性提升来自于我们开始摄入生成表的实际 SQL,而不仅仅是其模式。一个值为 contractpaygofreecustomer_type 列在任何上下文中看起来都一样,但 SQL 告诉你当 Salesforce 数据缺失时 customer_type 默认为 paygo。这种上下文永远不会存在于列描述中。

记忆比我们预期的更重要。存在一长串修正,例如“你必须像这样过滤 X”或“忽略标记为 Y 的表”。没有记忆层,代理会在每次对话中重新发现并重新学习这些内容。有了记忆层,它会在团队实际提出的重复问题上单调地变得更好。

枯燥的基础设施才是难点。Trino + Iceberg 并非新技术。真正困难的工作在于那些枯燥的部分:按行访问控制、默认关闭的表白名单、查询审计、有时限的凭证、PII 检测、幂等摄取、模式演化。这些才是让数据平台安全可用的关键。

我们正在扩展智能体(agent)的覆盖范围。Skipper 已经作为 MCP 服务器集成到任何支持它的 IDE 中。下一步是更深入地集成我们自己的内部聊天和工单系统,这样“询问数据”将成为任何人调试事件、规划项目或验证假设时的自然首选。

我们正在大力投资 Transformer 管道。目标是让 Cloudflare 的任何团队都能通过几个 SQL 文件和一个 .meta.json 描述来构建精选数据集,将其部署为工作流,自动进行调度和监控,并使其出现在 DataHub 和 Skipper 中,无需任何额外工作。其理念是自助式数据工程,与自助式软件工程具有相同的形态。

R2 SQL,Cloudflare 的无服务器分布式分析查询引擎,正日益强大。随着其功能集的扩展,我们计划将 Town Lake 工作流的许多部分迁移到它上面。

我们所下的赌注——下一个突破性产品来自某个查看数据并看到别人看不到的东西的人——仍然是我们正在押注的。Town Lake 就是我们确保他们能找到它的方式。