Kafka 到 Iceberg 流水线的性能特征

本页面介绍了 Apache Beam 2.75.0 版中从 Apache Kafka 读取数据并写入 Apache Iceberg 表的 Dataflow 流式作业的性能特征。它评估了直接 Apache Iceberg 写入与通过 Managed BigQuery API 路由的写入之间的性能差异,并将这些结果与 Kafka 到 BigQuery 流水线的基准测试结果进行了比较。由于 Apache Iceberg I/O 的优化工作仍在进行中,因此这些性能指标可能会发生变化。

基准比较涵盖了三种主要的无状态 映射配置(这意味着它们从来源读取数据,将消息转换为记录,并写入接收器,而无需跟踪状态或应用复杂的业务逻辑;在基准测试中称为 map_only 或 mapping):

  1. Kafka 到 BigQuery (map_only)( Kafka 到 BigQuery 性能的基准)
  2. Kafka 到 Iceberg Direct (map_only,autosharding=false)
  3. 使用 Managed BigQuery API 的 Kafka 到 Iceberg (map_only)

此外,本指南还评估了直接 Apache Iceberg 流式处理模式(例如使用 groupbykey 的有状态批处理),并详细介绍了有关 文件大小分布 、自动分片行为 和读取端查询延迟时间 的重要下游注意事项。

测试方法

基准测试使用了以下资源:

  • Managed Service for Apache Kafka 集群: 使用 Dataflow Streaming Data Generator 模板生成流量。
    • 输入吞吐量: 1 GBps
    • 消息速率: 大约每秒 100 万条消息
    • 消息格式: 具有固定架构的 JSON 文本(每条消息大约 1 KB)
    • 分区: 1,000 个 Kafka 分区
  • 目标接收器:
    • BigQuery: 使用 BigQuery Storage Write API 写入的标准表(未分区)。
    • Apache Iceberg: 由 Cloud Storage 提供支持的目录。直接接收器 使用 bucket(id, 64)(在主键上分桶到 64 个分片)进行分区,并使用 hash 分布模式。

横向自动扩缩稳定后,每种流水线配置在稳定状态下运行了 24 小时。每种流水线情况的基准测试都单独执行了 3 次,所有报告的值都代表这些运行的计算平均值,以确保持续可靠的性能指标。

注入性能:映射工作负载

无状态映射流水线从来源读取数据,将消息格式转换为记录,并写入接收器,而无需跨记录跟踪状态。以下部分分析了以 1 GBps 运行的参考架构。

作业配置

设置 Kafka 到 BigQuery (map_only) Kafka 到 Iceberg Direct (autosharding=false) 使用 Managed BigQuery API 的 Kafka 到 Iceberg
工作器机器类型 e2-standard-2 e2-standard-4 e2-standard-4
每个工作器的 vCPU 数 2 4 4
每个工作器的 RAM 8 GB 16 GB 16 GB
Streaming Engine 已启用 已启用 已启用
横向自动扩缩 已启用 已启用 已启用
触发频率 5 秒 60 秒 60 秒

吞吐量和资源用量

与 BigQuery 流式注入相比,直接写入对象存储中的物理 Parquet 文件会产生更高的 I/O 开销。与直接 Iceberg 写入相比,通过 Managed BigQuery API 路由写入可提高工作器 CPU 利用率(大约 70% 对大约 60%),并适度减少 Streaming Engine 消耗量(大约 180 SECU/小时对大约 200 SECU/小时),但总体工作器计算要求保持相似(大约 440 个 vCPU 对大约 450 个 vCPU)。

指标 Kafka 到 BigQuery (map_only) Kafka 到 Iceberg Direct (autosharding=false) 使用 Managed BigQuery API 的 Kafka 到 Iceberg
每个工作器的平均输入吞吐量 大约 15 MBps 大约 9 MBps 大约 9 MBps
平均 CPU 利用率 大约 70% 大约 60% 大约 70%
1 GBps 输入的估计 vCPU 数 大约 126 个 vCPU 大约 450 个 vCPU 大约 440 个 vCPU
1 GBps 输入的估计工作器数 大约 63 个工作器 大约 110 个工作器 大约 110 个工作器
1 GBps 的每小时估计 SECU 数 大约 58 SECU/小时 大约 200 SECU/小时 大约 180 SECU/小时

写入延迟时间配置文件

由于对象存储元数据提交限制,直接 Iceberg 写入表现出严重的尾延迟时间 (P99)。使用 Managed BigQuery API 可消除尾延迟时间峰值,同时保持较低的中位延迟时间。

端到端写入延迟时间 Kafka 到 BigQuery Kafka 到 Iceberg Direct (autosharding=false) 使用 Managed BigQuery API 的 Kafka 到 Iceberg
P50(中位) 大约 1,200 毫秒 大约 1,000 毫秒 大约 1,000 毫秒
P95 大约 3,000 毫秒 大约 7,400 毫秒 大约 1,900 毫秒
P99(尾部) 大约 5,400 毫秒 大约 14,000 毫秒 大约 2,700 毫秒

自动分片注意事项和设计选择

本部分讨论了在写入 Apache Iceberg 时自动分片对文件大小和流水线延迟时间的影响。

为什么选择 autosharding=false 作为基准

在初始测试中,启用自动分片会导致文件大小崩溃为很小的块,并由于局部线程级负载激增触发的动态分片拆分而任意波动,即使在恒定的聚合输入负载下也是如此。

为了保持稳定、可预测的 Parquet 文件布局(平均大约 800 KB)并确保公平的基准而不会过早刷新,我们为直接接收器配置选择了 autosharding=false 。

停用自动分片与保持启用自动分片会发生什么情况?

  • 使用 autosharding=false(基准): 与自动分片相比,您可以获得更大的初始文件大小(平均大约 800 KB)。虽然与理想的 Iceberg 文件大小(128–512 MB)相比,这仍然很小,但它需要的下游压缩要少得多。不过,代价是由于对象存储元数据瓶颈,写入尾延迟时间 (P99) 较高(达到大约 14.0 秒)。
  • 如果启用了自动分片: Dataflow 会动态扩缩写入器线程以吸收局部吞吐量峰值,从而缩短写入尾延迟时间。不过,它会生成大量小的、分散的 Parquet 文件(大约 100 KB 或更小),从而损害存储层。这些文件大小表现出很高的方差,并且在运行之间任意波动(平均范围从大约 39 KB 到大约 100 KB),从而加剧了对积极的下游压缩维护的需求。

分区调整和建议

在评估期间,我们尝试了目标表的各种固定分区值,以找到最佳平衡点。我们发现,对目标表分区使用 64 个存储分区 (例如 bucket(id, 64))可以生成我们目标文件大小,同时保持良好的利用率和吞吐量。这种方法帮助我们实现了自动分片的性能优势,同时避免了与完全动态伸缩相关的任意文件大小碎片问题。

从业者建议: 建议客户使用目标分区设置执行类似的初步测试,以找到最佳点,从而在不影响 Parquet 文件大小的情况下最大限度地提高流水线并行度。

下游读取影响:文件大小和压缩

虽然写入端指标有利于使用 Managed BigQuery API 进行 Iceberg 注入,但流水线的整体效率在很大程度上取决于下游读取性能:

  • Managed BigQuery API 中的小文件生成: Managed BigQuery API 会频繁刷新数据,以确保较低的写入延迟时间。此行为会导致大量小的 Parquet 文件写入目标 Iceberg 目录。
  • 读取查询延迟时间影响: 查询引擎(例如 Starburst/Trino、Apache Spark、BigQuery、Dremio)读取包含数百万个小 Parquet 文件的表时,会产生大量的元数据解析开销和分区扫描惩罚。
  • 压缩要求: 为了在使用 Managed BigQuery API 时防止读取性能下降(或者如果对直接写入启用了自动分片),请运行常规 Iceberg 压缩维护作业(例如 REWRITE DATA FILES)。压缩的计算开销应纳入整体架构设计。
  • 直接写入 (autosharding=false) 文件分布: 具有固定分片的直接 Iceberg 写入会生成较大的平均 Parquet 文件(大约 800 KB),从而生成碎片较少的布局,以便立即查询访问,而无需立即压缩(尽管仍低于理想范围)。

有状态直接 Iceberg 流水线 (groupbykey)

为了评估手动批处理策略,我们针对基准 Kafka 到 Iceberg Direct (map_only,autosharding=false) 流水线测试了有状态键分组 (groupbykey)。这两种配置都将 Parquet 文件直接写入对象存储。

基准比较

指标 / 功能 直接接收器基准 (autosharding=false) 有状态直接接收器 (groupbykey) 性能影响
1 GBps 的估计 vCPU 数 大约 450 个 vCPU 大约 520 个 vCPU 需要大约 +16% 的计算资源
平均 CPU 利用率 大约 60% 大约 50% 工作器效率降低大约 17%
1 GBps 的每小时估计 SECU 数 大约 200 SECU/小时 大约 300 SECU/小时 Streaming Engine 负载增加大约 50%
平均文件大小 大约 800 KB 大约 100 KB 生成较小的文件批次
P50 延迟时间 大约 1,000 毫秒 大约 1,200 毫秒 中位延迟时间慢大约 20%
P95 延迟时间 大约 7,400 毫秒 大约 5,500 毫秒 延迟时间缩短大约 26%
P99 延迟时间 大约 14,000 毫秒 大约 13,000 毫秒 尾延迟时间变化不大

权衡分析

  1. Streaming Engine 开销: 添加有状态 groupbykey 步骤需要 Beam 跨窗口边界存储中间状态。这会将 Streaming Engine 计算单元消耗量增加大约 50% (从大约 200 SECU/小时 增加到大约 300 SECU/小时 )。
  2. 缓冲延迟时间: 手动键聚合会引入强制窗口缓冲,从而将中位写入延迟时间 (P50) 增加到大约 1,200 毫秒 ,并将 P95 延迟时间增加到大约 5.5 秒 。

反向流水线:从 Iceberg 流式传输到 Kafka

为了评估双向湖仓功能,我们还针对反向流式传输数据进行了基准测试,即从 Apache Iceberg 表读取仅追加数据流并发布回 Apache Kafka。

作业配置和效率

与必须处理大量对象存储文件写入或元数据提交瓶颈的注入流水线不同,从 Iceberg 读取和流式传输更改的效率很高:

指标 Iceberg 到 Kafka(仅追加,正好一次)
工作器机器类型 e2-standard-4
1 GBps 输入的估计 vCPU 数 大约 30 个 vCPU
1 GBps 输入的估计工作器数 大约 7 个工作器
1 GBps 的每小时估计 SECU 数 大约 0.2 SECU/小时

反向流水线的主要要点

  • 计算开销显著降低: 从 Iceberg 读取和投影 CDC 数据流所需的计算资源要少得多(大约 30 个 vCPU 对大约 450 个 vCPU,用于直接写入),因为它避免了将大量 Parquet 文件分区、编码和提交到对象存储的繁重工作。
  • 资源效率: 与入站注入路径相比,从湖仓格式到流式处理层的下游事件驱动型消耗或复制的效率非常高。

架构建议摘要

架构模式 P99 写入延迟时间 文件布局 下游读取注意事项
Kafka 到 BigQuery (map_only) 大约 5.4 秒 不适用 最佳 (Managed BigQuery 存储引擎)
使用 Managed BigQuery API 的 Kafka 到 Iceberg 大约 2.7 秒 任意小文件 需要定期压缩 以进行大量读取
Kafka 到 Iceberg Direct (autosharding=false) 大约 14.0 秒 大约 800 KB 良好 (初始文件较大,压缩需求较低)
Kafka 到 Iceberg Direct (groupbykey) 大约 13.0 秒 大约 100 KB 中等 (计算和状态开销较高)

估算费用

您可以使用 Google Cloud 价格计算器,估算您自己的、类似流水线的基于资源的结算的基准成本,如下所示:

  1. 打开价格计算器。
  2. 点击添加到估算 。
  3. 选择 Dataflow。
  4. 对于服务类型,选择“Dataflow Classic”。
  5. 选择高级设置以显示完整选项集。
  6. 选择作业的运行位置。
  7. 对于作业类型,请选择“流式处理”。
  8. 选择启用 Streaming Engine 。
  9. 输入作业运行小时数、工作器节点、工作器机器和 Persistent Disk 存储空间的相关信息。
  10. 输入 Streaming Engine 计算单元的估算数量。

资源用量和费用大致随输入吞吐量线性扩展,不过对于只有少量工作器的小型作业,总费用主要由固定费用构成。首先,您可以根据基准测试结果推断工作器节点的数量和资源消耗量。

例如,假设您使用 Kafka 到 Iceberg Direct (autosharding=false) 架构运行流水线,输入数据速率为 100 MBps。根据 1 GBps 流水线的基准测试结果,您可以按如下方式估算资源要求:

  • 缩放比例:(100 MBps) / (1024 MBps) = 大约 0.1
  • 预计工作器节点数:110 个工作器 × 0.1 = 大约 11 个工作器
  • 每小时的 Streaming Engine 计算单元预计数:200 × 0.1 = 大约每小时 20 个单元

此值仅应作为初始估算值使用。实际吞吐量和费用可能会因多种因素而有很大差异,例如机器类型、消息大小分布、用户代码、聚合类型、键并行度和窗口大小。如需了解详情,请参阅 Dataflow 费用优化的最佳实践。

运行测试流水线

如需使用 Dataflow Flex 模板部署 Apache Iceberg 流式作业,请使用 gcloud dataflow flex-template run 命令。

gcloud dataflow flex-template run JOB_NAME \
  --project=PROJECT_ID \
  --region=REGION \
  --template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_To_Iceberg_Yaml \
  --enable-streaming-engine \
  --parameters ^@^bootstrapServers="KAFKA_BOOTSTRAP_ADDRESS"\
@topic="KAFKA_TOPIC"\
@table="ICEBERG_TABLE_IDENTIFIER"\
@catalogName="CATALOG_NAME"\
@catalogProperties='{"type":"CATALOG_TYPE","warehouse":"gs://BUCKET_NAME/warehouse/"}'\
@triggeringFrequencySeconds=60\
@schema='SCHEMA_DEFINITION'

替换以下内容:

  • JOB_NAME:Dataflow 作业的名称
  • PROJECT_ID:您的 Google Cloud 项目 ID
  • REGION: Google Cloud 作业运行所在的区域(例如 us-central1)
  • KAFKA_BOOTSTRAP_ADDRESS:Apache Kafka 集群的引导地址
  • KAFKA_TOPIC:Kafka 主题的名称
  • ICEBERG_TABLE_IDENTIFIER:目标 Iceberg 表的标识符
  • CATALOG_NAME:Iceberg 目录的名称
  • CATALOG_TYPE:要使用的目录的类型(例如 hadoop 或 bigquery)
  • BUCKET_NAME:仓库位置的 Cloud Storage 存储桶的名称
  • SCHEMA_DEFINITION:Kafka 主题数据的架构定义(例如 {"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})