消息压缩
我们强烈建议为 Kafka topic 启用压缩。压缩可以显著节省数据传输成本,同时几乎不会对性能产生影响。 如需进一步了解 Kafka 中的消息压缩,我们建议先阅读这篇指南。
限制
- 不支持
DEFAULT。 - 默认情况下,在使用最小 (XS) 副本大小运行时,单条消息的大小上限为 16MB (未压缩) ;使用更大副本时,上限为 32MB (未压缩) 。超过此限制的消息会因报错而被拒绝。如果您需要更大的消息,请联系支持团队。
交付语义
ClickPipes for Kafka 默认提供至少一次交付保证,并通过 Kafka 消费者组偏移量跟踪摄取进度。它还可选支持精确一次语义:即使发生 pod (容器组) 重启、消费者再平衡和插入失败,每条 Kafka 记录也只会被插入 ClickHouse 一次。
为实现精确一次交付,ClickPipes 会使用两个值在内部状态存储中记录每个分区的进度:
- 高水位标记 — 已确认插入 ClickHouse 的该分区所有记录所对应的最高偏移量。重启时,ClickPipes 会丢弃偏移量等于或低于此标记的记录,因此已落库的数据不会再次发送。
- 待处理范围 — 已发送至 ClickHouse 但尚未确认的插入块所覆盖的偏移量范围。发生故障后,ClickPipes 仅重放这些范围。
每个插入块覆盖连续的偏移量范围,并携带格式为 topic:partition:firstOffset-lastOffset 的确定性去重标记。重放时,ClickPipes 会生成相同的偏移量范围,因而使用相同的标记,ClickHouse 将拒绝重复数据。由于标记仅取决于偏移量范围,即使重建的块并非逐字节完全一致,重放的数据仍会被去重。
主要的权衡在于 part 大小。更大的插入块会在 ClickHouse 中生成更少但更大的 parts,从而降低合并开销。ClickPipes 在构建块时会将一个分区的行保存在内存中,因此可达到的 part 大小取决于管道可用的内存:内存紧张时,它会构建较小的块,表中会累积更多 parts。为管道提供更多内存可使其构建更大的块,从而生成更少的 parts。
当分区数量接近内部插入“工作线程”的数量时,管道的运行效果最佳,因为此时每个工作线程大致处理一个分区,并有足够的内存余量构建大块。工作线程数量和可用内存都会随副本大小和数量而扩缩容,可在 设置 -> 高级设置 -> 扩缩容 下配置。
身份验证
对于 Apache Kafka 协议的数据源,ClickPipes 支持结合 TLS 加密的 SASL/PLAIN 身份验证,以及 SASL/SCRAM-SHA-256 和 SASL/SCRAM-SHA-512。根据流式数据源 (Redpanda、MSK 等) 的兼容性,支持的身份验证机制可能是全部,也可能只是其中一部分。如果您的身份验证需求不同,请向我们提供反馈。
Warpstream 拉取大小
ClickPipes 依赖 Kafka 设置 max.fetch_bytes 来限制单个 ClickPipes 节点在任一时刻处理的数据量。在某些情况下,
Warpstream 可能不会遵循该设置,从而导致管道意外失败。我们强烈建议在配置 WarpStream agent 时,将 Warpstream 专用设置 kafkaMaxFetchPartitionBytesUncompressedOverride
设置为 8MB (或更低) ,以防止 ClickPipes 失败。
IAM
ClickPipes 支持以下 AWS MSK 身份验证方式
- SASL/SCRAM-SHA-512 身份验证
- IAM 凭证或基于角色的访问 身份验证
使用 IAM 身份验证连接到 MSK broker 时,IAM 角色必须具有相应的权限。 下面是 MSK 的 Apache Kafka API 所需 IAM 策略示例:
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"kafka-cluster:Connect"
],
"Resource": [
"arn:aws:kafka:us-west-2:12345678912:cluster/clickpipes-testing-brokers/b194d5ae-5013-4b5b-ad27-3ca9f56299c9-10"
]
},
{
"Effect": "Allow",
"Action": [
"kafka-cluster:DescribeTopic",
"kafka-cluster:ReadData"
],
"Resource": [
"arn:aws:kafka:us-west-2:12345678912:topic/clickpipes-testing-brokers/*"
]
},
{
"Effect": "Allow",
"Action": [
"kafka-cluster:AlterGroup",
"kafka-cluster:DescribeGroup"
],
"Resource": [
"arn:aws:kafka:us-east-1:12345678912:group/clickpipes-testing-brokers/*"
]
}
]
}配置信任关系
如果你使用 IAM 角色 ARN 对 MSK 进行身份验证,则需要在你的 ClickHouse Cloud 实例与该角色之间建立信任关系,以便可以代入该角色。
{
"Version": "2012-10-17",
"Statement": [
...
{
"Effect": "Allow",
"Principal": {
"AWS": "arn:aws:iam::12345678912:role/CH-S3-your-clickhouse-cloud-role"
},
"Action": "sts:AssumeRole"
}
]
}自定义证书
ClickPipes for Kafka 支持为使用非公开服务器证书的 Kafka broker 上传自定义证书。
也支持上传客户端证书和私钥,用于基于双向 TLS (mTLS) 的身份验证。
性能
批处理
ClickPipes 以批次方式将数据插入 ClickHouse。这样做是为了避免在数据库中产生过多的 parts,否则可能会导致集群性能问题。
当满足以下任一条件时,就会插入一个批次:
- 批次大小达到上限 (每 1GB pod (容器组) 内存对应 100,000 行或 28MB)
- 批次保持打开状态达到最长时限 (5 秒)
延迟
延迟 (指 Kafka 消息生成后到该消息在 ClickHouse 中可用之间的时间) 取决于多种因素 (例如 broker 延迟、网络延迟、消息大小/格式) 。上文所述的批处理也会影响延迟。我们始终建议针对您的具体用例,在典型负载下进行测试,以确定预期延迟。
ClickPipes 不对延迟提供任何保证。如果您有特定的低延迟要求,请 联系我们。
扩缩容
ClickPipes for Kafka 设计上支持水平和垂直扩缩容。默认情况下,我们会创建一个仅包含一个消费者的消费者组。你可以在创建 ClickPipe 时进行配置,也可以随时在 设置 -> 高级设置 -> 扩缩容 中调整。
ClickPipes 通过跨可用区分布式架构提供高可用性。 这要求至少扩缩容到两个消费者。
无论当前运行的消费者数量是多少,系统在设计上都具备容错能力。 如果某个消费者或其底层基础设施发生故障, ClickPipe 都会自动重启该消费者并继续处理消息。
基准测试
以下是一些 ClickPipes for Kafka 的非正式基准测试,可用于大致了解其基线性能。需要注意的是,影响性能的因素很多,包括消息大小、数据类型和数据格式。实际表现会因环境而异,这里展示的结果并不构成对实际性能的保证。
基准测试详情:
- 我们使用了生产环境中的 ClickHouse Cloud 服务,并为其配置了充足资源,以确保吞吐量不会受 ClickHouse 侧插入处理瓶颈的影响。
- ClickHouse Cloud 服务、Kafka 集群 (Confluent Cloud) 和 ClickPipe 都运行在同一区域 (
us-east-2) 。 - ClickPipe 配置为单个 L 规格副本 (4 GiB RAM 和 1 个 vCPU) 。
- 样本数据包含嵌套数据,混合使用了
UUID、String和Int数据类型。其他数据类型 (如Float、Decimal和DateTime) 的性能可能会更低。 - 使用压缩和未压缩数据时,性能没有明显差异。
| 副本大小 | 消息大小 | 数据格式 | 吞吐量 |
|---|---|---|---|
| 大型 (L) | 1.6kb | JSON | 63mb/s |
| 大型 (L) | 1.6kb | Avro | 99mb/s |