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

Micah WyldeMatt SilverlockPranshu Maheshwari

阅读时间:9 分钟

本文另有 English日本語.

BLOG-2785 Feature Image

今天,我们推出流式摄取产品 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存储桶中的输出可能如下所示:

BLOG-2785 Image 1

在 R2存储桶中对来自管道的 Hive 分区文件

输出文件存储为换行符分隔的 JSON (NDJSON)——可以很容易地从这些文件中实现流(提示:将来您也可以使用 R2 作为管道源)。最后,文件名是ULID,所以默认是按时间排序的。

首先进行分片,然后进一步分片

Pipeline 的构建方式使其水平可扩展能够快速确认写入:我们使用 Durable Objects 和每个 Durable Object 中的嵌入式零延迟 SQLite 存储,在数据写入时立即持久化,然后再进行处理和写入连接到 R2。

例如:假设您是一个电子商务或SaaS网站,需要获取网站使用数据(称为 点击流数据),并将其提供给您的数据科学团队进行查询。处理此类工作负载的基础设施必须能够可靠地应对多种故障情景。面对流量激增,摄取服务需要保持高可用性。数据在被摄取后,需要进行缓冲,以最小化下游调用,从而降低成本。最后,缓冲的数据需要传递到接收器,并在接收器不可用时进行适当的重试和故障处理。当过载时,此过程的每一步都需要向上游发出反压信号。还需要扩展规模:在重大销售或活动期间扩大规模,在一天中知名度较低的时期缩小规模。

阅读本文的数据工程师可能会熟悉使用 Kafka 及相关生态系统来处理该问题的现状。但是,如果您是一名应用工程师:您可以使用 Pipelines 构建一个摄取服务,而无需了解有关 Kafka、Zookeeper 和 Kafka Streams 的知识。

BLOG-2785 Image 2

管道水平分片

上图显示了 Pipelines 如何拆分控制平面和数据路径,前者负责记账、跟踪分片和生命周期事件,后者是一组可扩展的Durable Objects分片。

将一条记录(或一批记录)写入 Pipelines 时:

  1. Pipelines Worker 通过提取处理程序或 Worker 绑定接收记录。
  2. 根据pipeline_id 联系协调者以获取执行计划:后续读取被缓存以减少协调者的压力。
  3. 执行计划,首先分片为一组执行器,同时主要用于扩展读取请求处理
  4. 然后,它们重新分片到另一组实际处理写入的执行器,首先是持久存储到 Durable Object 存储,再通过存储中继服务(SRS)复制这些执行器,以提高持久性和可用性。
  5. 在 SRS 之后,我们会将数据传递给任何已配置的 Transform Workers来自定义数据。
  6. 数据被批处理,写入输出文件并压缩(如果适用)。
  7. 文件被压缩,数据被打包到最终批次,并写入配置的 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 数据目录 集成:使您能够将数据流直接摄取到 冰山表 中并立即进行查询,而不需要依赖于其他系统。

同时,您可以:

或者直接部署示例项目:

$ npm create cloudflare@latest -- pipelines-starter 
--template="cloudflare/pipelines-starter"