Cloudflare Basin 正式可用:用 Iceberg 与 R2 构建无出口费数据湖

2026-10-01 28 预计阅读时间: 1 分钟
来源: blog.cloudflare.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:9 分钟

Cloudflare Basin 已进入正式可用阶段。它以 Apache Iceberg 作为开放表格式,以 R2 Object Storage 保存底层数据,让开发者能够摄取、管理和查询大规模数据集,同时避免数据出口费用。

这套组合的价值不只是“把文件放进对象存储”。Iceberg 为对象存储上的数据增加了表、快照和元数据层,而 Basin 将这些能力包装成无服务器数据平台,使团队可以把更多精力放在数据模型和查询上,而不是维护存储集群。

Iceberg 与 R2 分别解决什么问题

R2 适合保存体量较大的数据对象,但单纯堆积 Parquet、JSON 或 CSV 文件,很快会遇到几个工程问题:

  • 哪些文件属于当前版本的数据集?
  • 新增字段后,旧文件如何继续参与查询?
  • 写入失败时,如何避免查询读到一半新、一半旧的数据?
  • 文件越来越多后,查询引擎怎样高效定位需要扫描的数据?

Apache Iceberg 在数据文件之上维护表级元数据。查询引擎面对的是一张表及其快照,而不是一组需要自行猜测含义的对象路径。这种设计适合日志、事件、分析明细和持续增长的历史数据。

Basin 的底层存储使用 R2,则进一步改变了数据流动的成本模型。数据可以长期保留在对象存储中,并被不同的数据处理或查询流程使用,而无需为数据出口付费。这里需要注意:免出口费不等于所有操作都免费,存储、请求、摄取和查询所消耗的计算资源仍应纳入成本评估。

无服务器平台改变的是运维边界

传统数据湖往往需要团队分别管理对象存储、目录服务、表元数据、计算集群和访问权限。任何一个组件升级或配置不一致,都可能让数据管道停摆。

Basin 的无服务器定位意味着开发者不必把容量规划和常驻集群维护作为第一优先级。更合适的思路是围绕数据产品建立清晰边界:

  1. 生产端按照稳定模式生成事件或批次数据。
  2. Basin 负责大规模数据的摄取与管理。
  3. Iceberg 表提供可查询的数据组织形式。
  4. 下游任务按需查询,而不是复制一份数据到每个系统。

开放表格式也很重要。采用 Iceberg 后,数据组织不再完全依赖某个查询引擎的私有表结构。不过,“开放格式”不代表任意工具天然兼容;接入前仍要核对引擎支持的 Iceberg 版本、数据类型、分区变换和写入语义。

实践:先生成适合摄取的事件数据

来源摘要没有给出 Basin 的具体上传 API 或命令行接口,因此下面不虚构专有端点,而是演示一个可以直接运行的数据生产步骤:生成带日期分区的 Parquet 事件数据。随后可按照 Basin 官方接入方式,将这些文件或对应事件流摄取到目标表中。

先安装依赖:

python -m venv .venv
source .venv/bin/activate
python -m pip install "pyarrow>=15,<20"

保存下面的代码为 generate_events.py:

from datetime import datetime, timedelta, timezone
from pathlib import Path
import random
import uuid

import pyarrow as pa
import pyarrow.dataset as ds

OUTPUT_DIR = Path("data/events")
EVENT_COUNT = 10_000

now = datetime.now(timezone.utc)
rows = []

for index in range(EVENT_COUNT):
    occurred_at = now - timedelta(minutes=random.randint(0, 60 * 24 * 3))
    rows.append(
        {
            "event_id": str(uuid.uuid4()),
            "account_id": f"acct-{random.randint(1, 500):04d}",
            "event_type": random.choice(["page_view", "signup", "purchase"]),
            "amount": round(random.uniform(5, 250), 2),
            "occurred_at": occurred_at,
            "event_date": occurred_at.date().isoformat(),
        }
    )

table = pa.Table.from_pylist(rows)
OUTPUT_DIR.mkdir(parents=True, exist_ok=True)

ds.write_dataset(
    table,
    base_dir=str(OUTPUT_DIR),
    format="parquet",
    partitioning=["event_date"],
    existing_data_behavior="overwrite_or_ignore",
)

print(f"Wrote {table.num_rows} rows to {OUTPUT_DIR}")
print(table.schema)

运行并检查结果:

python generate_events.py
find data/events -type f -name '*.parquet' | sort

这个例子刻意采用低基数的日期字段进行分区,而没有按 account_id 或 event_id 建立目录。高基数字段容易生成大量小分区和小文件,增加元数据及查询规划开销。

数据进入 Basin 并映射为 Iceberg 表后,查询逻辑可以保持简单。以下 SQL 使用通用 Iceberg 查询语义,表名和时间函数需要按照实际查询接口调整:

SELECT
    event_date,
    event_type,
    COUNT(*) AS event_count,
    ROUND(SUM(amount), 2) AS total_amount
FROM analytics.events
WHERE event_date >= CURRENT_DATE - INTERVAL '7' DAY
GROUP BY event_date, event_type
ORDER BY event_date DESC, event_type;

生产环境还应给摄取任务增加三个字段:来源批次 ID、摄取时间和模式版本。例如,source_batch_id 可以用于重试去重,ingested_at 可以帮助排查延迟,schema_version 则能明确生产端使用的数据契约。

迁移时不要忽略表设计

把现有文件复制到新平台,并不会自动得到高质量的数据湖。迁移或新建 Basin 数据集时,建议重点检查以下内容:

  • 分区策略:优先选择日期、小时或有限类别等常见过滤字段,避免过度分区。
  • 小文件治理:高频微批可能产生大量小文件,需要评估合并或压缩策略。
  • 模式演进:新增字段通常比改变字段含义安全;字段重命名、类型转换和删除必须经过兼容性测试。
  • 幂等写入:网络重试不应产生重复事件,批次标识或业务主键应进入数据契约。
  • 权限边界:区分数据生产者、查询者和管理员,不要让所有任务共享同一组长期凭证。
  • 成本验证:出口费为零只是整体成本的一部分,还要测量存储、请求和查询负载。

适合从哪类工作负载开始

Basin 适合优先验证那些数据量持续增长、需要长期保留,并且会被多个分析流程重复读取的工作负载,例如产品事件、应用日志、审计记录或批量分析数据。

一个稳妥的采用路径是先选择单张追加型事实表,保留现有管道作为回退方案,连续观察摄取延迟、查询耗时、失败重试和月度成本。验证模式演进与权限控制后,再迁移需要更新、删除或跨表一致性的复杂数据集。

Basin 的核心吸引力在于把 Iceberg 的开放表能力、R2 的对象存储和无服务器运维模型放到一起。它能够减少基础设施负担和数据移动成本,但真正决定平台效果的,仍然是表结构、文件大小、写入幂等性以及查询模式这些具体工程决策。


相关推荐