最新发布:使用 Arroyo 和 Pipelines 在 Cloudflare 上进行流式数据摄取

今天,我们推出流式摄取产品 Pipelines 的 公测版 。Pipelines 让您能够获取大量结构化、实时数据,并将其加载到我们的对象存储服务 R2中。您无需管理任何底层基础设施,不必担心扩展分片或元数据服务,而且只需要为处理的数据量付费(而不是按小时)。任何Workers付费计划的任何人都可以开始使用它来摄取和批处理数据(每秒数万个请求)并直接传输到 R2 中。
但这只是冰山一角:您通常需要转换正在获取的数据,从其他来源实时整合后,并将其写入开放的表格格式(如 Apacheiceberg),以便一旦数据进入对象存储,你就可以高效地查询这些数据
好消息是,我们也已经考虑到了这一点,并且我们很高兴地宣布,我们收购了 Arroyo —— 一个云原生的分布式流处理引擎——来实现这个目标。
通过 Arroyo 和我们刚刚推出的 R2 Data Catalog,我们越来越认真地致力于建立一个数据平台,让您能够在全球范围内获取数据,大规模存储数据,并对其进行计算。
首先,您可以深入研究管道开发人员文档,或直接运行以下Wrangler命令来创建您的第一个管道:
$ npx wrangler@latest pipelines create my-clickstream-pipeline --r2-bucket my-bucket
...
✅ Successfully created Pipeline my-clickstream-pipeline with ID 0e00c5ff09b34d018152af98d06f5a1xv…… 然后撰写您的第一条记录:
$ curl -d '[{"payload": [],"id":"abc-def"}]'
"https://0e00c5ff09b34d018152af98d06f5a1xvc.pipelines.cloudflarestorage.com/"然而,真正的力量来自于在摄取数据和将数据写入 R2 等接收器之间对数据流的处理。能够编写 SQL 在获取数据时作用于数据窗口,可以转换和聚合数据,甚至从数据中实时提取见解,这将是非常强大的。
这就是 Arroyo 的用武之地,我们将把 Arroyo 最好的部分引入管道,并与Workers、 R2 和我们开发人员平台的其他部分深度集成。
Arroyo 的起源故事
(作者:Arroyo 创始人 Micah Wylde)
我们于 2023 年创立了 Arroyo,旨在为每个与数据打交道的人带来实时(流)处理。现代公司依赖数据管道来驱动他们的应用和业务——从用户定制、推荐、反欺诈,到新兴的 AI 代理世界。
但如今,这些管道大多是以批量方式运行,每小时、每天甚至每月运行一次。对于 Lyft 和 Splunk 这样的公司从事流处理工作多年后,其原因显而易见:对于开发人员和数据科学家来说,构建正确、高性能和可靠的管道太难了。大型科技公司会雇佣流媒体专家来构建和运营这些系统,但其他人都在等待批次到达。
我们在创业时,流式管道的主要解决方案是 Apache Flink,这也是我们在 Lyft 和 Splunk 中使用的解决方案。Flink 是第一个成功地将 容错(能够从故障中恢复)、分布式(跨多台机器)有状态(并记住过去事件的数据)数据流与 图构建API结合起来的系统。这些功能的组合意味着我们最终可以构建强大的实时数据应用程序,具有窗口、聚合和连接等功能。尽管 Flink 具有必要的能力,但在实践中,该API对于非专家用户来说过于困难和低级,而且生成的服务的有状态需要无休止的操作。
我们意识到需要打造一个新的流引擎——一个具有 Flink 的强大能力,但为产品工程师和数据科学家而设计,并在现代云基础设施上运行。我们一开始使用 SQL 作为我们的API ,因为它易于使用、众所周知,并且是声明性的。我们用 Rust 构建它,目的是为了提高速度和操作简单性(无需调整 JVM!)。我们构建了一个对象存储的原生状态后端,简化了运行有状态管道的挑战——每个有状态管道都像一个怪异的专用数据库。在 2023 年夏天,我们将其开源。如今,数十家公司正在运行 Arroyo 管道,用例包括数据摄取、反欺诈、IoT可观察性和金融交易。
但我们始终都知道,该引擎只是拼图中的一块。为了使流式传输像批处理一样简单,用户需要能够开发和测试查询逻辑,回填历史数据,以及无服务器部署而不必担心集群大小或持续运营。大众化流媒体最终意味着建立一个完整的数据平台。当我们开始与Cloudflare交谈时,我们意识到他们已经具备了所有要求:R2 为静态状态和数据提供对象存储, Cloudflare Queues用于传输中的数据,以及Workers来安全、高效地运行用户代码。Cloudflare的独特之处在于,允许我们将这些系统一直推到边缘,实现一种本地流处理的新范式,这将是数据主权和 AI 未来的关键。
因此,我们非常高兴能与Cloudflare团队一起将这一愿景变成现实。
大规模摄取
虽然正在为 Pipelines 提供转换和流式 SQL API ,但它已经解决了数据旅程的两个关键部分:全球分布式的高吞吐量摄取和高效加载到对象存储中。
创建管道就像运行命令一样简单:
$ npx wrangler@latest pipelines create my-clickstream-pipeline --r2-bucket my-bucket
🌀 Creating pipeline named "my-clickstream-pipeline"
✅ Successfully created pipeline my-clickstream-pipeline with ID
0e00c5ff09b34d018152af98d06f5a1xvc
Id: 0e00c5ff09b34d018152af98d06f5a1xvc
Name: my-clickstream-pipeline
Sources:
HTTP:
Endpoint: https://0e00c5ff09b34d018152af98d06f5a1xvc.pipelines.cloudflare.com/
Authentication: off
Format: JSON
Worker:
Format: JSON
Destination:
Type: R2
Bucket: my-bucket
Format: newline-delimited JSON
Compression: GZIP
Batch hints:
Max bytes: 100 MB
Max duration: 300 seconds
Max records: 100,000
🎉 You can now send data to your pipeline!
Send data to your pipeline's HTTP endpoint:
curl "https://0e00c5ff09b34d018152af98d06f5a1xvc.pipelines.cloudflare.com/" -d '[{ ...JSON_DATA... }]'默认情况下,管道可以从两个源获取数据( Workers和一个 HTTP端点),并将成批事件加载到一个 R2存储桶中。对于将原始事件数据流式传输到对象存储,这是一个开箱即用的解决方案。如果默认值不起作用,您可以在创建管道时或创建后的任何时间进行配置。选项包括:向 HTTP端点添加 验证 ,配置 CORS 以允许浏览器发出跨源请求,以及指定输出文件压缩和 批量 设置。
我们从第一天就开始构建可用于高摄取量的管道。每个管道可以扩展到约 100,000 条记录/秒(这只是我们的开端)。将记录写入管道后,就会在 R2存储桶中作为文件进行持久存储、批量并写出。批处理在这里至关重要:如果您要对这些数据采取行动,就应该让查询引擎查询数百万(或数千万)个小文件。速度慢(每个文件和请求的开销)、效率低(需要读取更多文件)和成本高昂(更多操作)。相反,您希望在查询引擎的批量大小和延迟(批处理不会等待太长时间)之间找到适当的平衡: Pipelines 允许您进行配置。
为进一步优化查询,使用标准 Hive 分区方案,按日期和时间对输出文件进行分区。这可以进一步优化查询,因为您的查询引擎可以跳过与您正在运行的查询无关的数据。R2存储桶中的输出可能如下所示:

在 R2存储桶中对来自管道的 Hive 分区文件
输出文件存储为换行符分隔的 JSON (NDJSON)——可以很容易地从这些文件中实现流(提示:将来您也可以使用 R2 作为管道源)。最后,文件名是ULID,所以默认是按时间排序的。
首先进行分片,然后进一步分片
Pipeline 的构建方式使其水平可扩展并能够快速确认写入:我们使用 Durable Objects 和每个 Durable Object 中的嵌入式零延迟 SQLite 存储,在数据写入时立即持久化,然后再进行处理和写入连接到 R2。
例如:假设您是一个电子商务或SaaS网站,需要获取网站使用数据(称为 点击流数据),并将其提供给您的数据科学团队进行查询。处理此类工作负载的基础设施必须能够可靠地应对多种故障情景。面对流量激增,摄取服务需要保持高可用性。数据在被摄取后,需要进行缓冲,以最小化下游调用,从而降低成本。最后,缓冲的数据需要传递到接收器,并在接收器不可用时进行适当的重试和故障处理。当过载时,此过程的每一步都需要向上游发出反压信号。还需要扩展规模:在重大销售或活动期间扩大规模,在一天中知名度较低的时期缩小规模。
阅读本文的数据工程师可能会熟悉使用 Kafka 及相关生态系统来处理该问题的现状。但是,如果您是一名应用工程师:您可以使用 Pipelines 构建一个摄取服务,而无需了解有关 Kafka、Zookeeper 和 Kafka Streams 的知识。

管道水平分片
上图显示了 Pipelines 如何拆分控制平面和数据路径,前者负责记账、跟踪分片和生命周期事件,后者是一组可扩展的Durable Objects分片。
将一条记录(或一批记录)写入 Pipelines 时:
- Pipelines Worker 通过提取处理程序或 Worker 绑定接收记录。
- 根据
pipeline_id联系协调者以获取执行计划:后续读取被缓存以减少协调者的压力。 - 执行计划,首先分片为一组执行器,同时主要用于扩展读取请求处理
- 然后,它们重新分片到另一组实际处理写入的执行器,首先是持久存储到 Durable Object 存储,再通过存储中继服务(SRS)复制这些执行器,以提高持久性和可用性。
- 在 SRS 之后,我们会将数据传递给任何已配置的 Transform Workers来自定义数据。
- 数据被批处理,写入输出文件并压缩(如果适用)。
- 文件被压缩,数据被打包到最终批次,并写入配置的 R2存储桶。
管道的每一步都可能发出上游反压信号。通过利用ReadableStreams 并在等待写入的总字节数超过某个阈值时以429错误响应来执行此操作。每个 ReadableStream 都能够通过使用Durable Objects之间的 JSRPC 调用来跨越 Durable Object 边界。为了提高性能,我们使用 RPC 存根在Durable Objects之间进行连接重用。还可以重试每个步骤操作,以处理Durable Objects或 R2 中的任何临时不可用情况。
即使在更新现有的管道时,我们也能保证交付。更新现有管道时,我们会创建一个新的 Deployment,包括上述所有分片和Durable Objects 。请求正常重新路由到新的管道。旧管道继续将数据写入 R2,直到所有 Durable Object 存储耗尽。只有在所有数据都被写出后,我们才会停止运行旧管道。这样,即使在更新管道时,您也不会丢失数据。
您会注意到这里有一个有趣的部分—— Transform Workers我们还没有公开。随着我们努力将 Arroyo 的流式传输引擎与管道集成在一起,这将成为我们将数据移交给 Arroyo 进行处理的关键部分。
那么,它的成本是多少?
在公测的第一阶段,除了标准的 R2 存储以及加载和访问数据产生的操作费用外,不会有任何额外费用。一如既往,直接从 R2 存储桶的出口是免费的,因此您可以从云或地区处理和查询数据,而不必担心增加数据传输成本。
未来,我们计划推出根据引入 Pipelines 和从 Pipelines 交付的数据量定价的做法。
Workers付费 (5 美元/月) | |
|---|---|
摄取 | 包含每月前 50 GB增加每 GB 0.02 美元 |
交付到 R2 | 包含每月前 50 GB增加每 GB 0.02 美元 |
随着测试的进行,我们还计划在Workers Free计划中提供 Pipelines。
我们将分享更多信息,为 Pipelines 带来转换和更多接收器。我们会至少提前 30 天通知,再进行任何更改或开始收费(预计在 2025 年 9 月 15 日前)。
接下来?
这里需要构建很多东西,我们渴望在 Arroyo 已经构建的许多强大组件的基础上再构建:将Workers集成为 UDF(用户定义函数),添加新的来源(如 Kafka 客户端),以及使用新的接收器扩展管道(超越 R2)。
我们还将管道与我们刚刚推出的 R2 数据目录 集成:使您能够将数据流直接摄取到 冰山表 中并立即进行查询,而不需要依赖于其他系统。
同时,您可以:
- 开始并创建您的第一个 Pipeline
- 阅读文档
- 加入我们的 Developer Discord 上的
#pipelines-beta频道
或者直接部署示例项目:
$ npm create cloudflare@latest -- pipelines-starter
--template="cloudflare/pipelines-starter"

