Kinesis ClickPipes 可以通过 ClickPipes UI 手动部署和管理,也可以使用 OpenAPI 和 Terraform 以编程方式部署和管理。
前置条件
你应已了解 ClickPipes 简介,并已设置好 IAM 凭证 或 IAM 角色。有关如何设置可与 ClickHouse Cloud 配合使用的角色,请参阅 Kinesis 基于角色的访问指南。
创建你的第一个 ClickPipe
- 进入你的 ClickHouse Cloud 服务的 SQL 控制台。

- 在左侧菜单中选择
Data Sources按钮,然后点击“Set up a ClickPipe”

- 选择你的数据源。

- 填写表单,为 ClickPipe 提供名称、描述 (可选) 、IAM 角色 或凭证,以及其他连接详细信息。

- 选择 Kinesis Stream 和起始偏移量。UI 会显示所选源中的一个样本文档 (Kafka topic 等) 。你还可以为 Kinesis 数据流启用 Enhanced Fan-out,以提升 ClickPipe 的性能和稳定性 (有关 Enhanced Fan-out 的更多信息,请参见此处)

- 在下一步中,你可以选择将数据摄取到新的 ClickHouse 表中,或复用现有表。按照界面上的说明修改表名、schema 和设置。你可以在顶部的样本表中实时预览这些更改。

你还可以使用提供的控件自定义高级设置

- 或者,你也可以选择将数据摄取到现有的 ClickHouse 表中。在这种情况下,UI 会允许你将源中的字段映射到所选目标表中的 ClickHouse 字段。

- 最后,你可以为内部 ClickPipes 用户配置权限。
Permissions: ClickPipes 会创建一个专用用户,用于将数据写入目标表。你可以为该内部用户选择一个角色,可使用自定义角色或预定义角色之一:
Full access:对集群具有完全访问权限。如果你在目标表中使用 materialized view 或字典,这可能会很有用。Only destination table:仅具有目标表的INSERT权限。

- 点击“Complete Setup”后,系统将注册你的 ClickPipe,你将能够在摘要表中看到它。


摘要表提供了相关控件,可显示源中的样本数据或 ClickHouse 中目标表的数据

以及用于移除 ClickPipe 并显示摄取作业摘要的控件。

- 恭喜! 你已成功设置第一个 ClickPipe。如果这是一个流式 ClickPipe,它将持续运行,实时从远程数据源摄取数据。否则,它会摄取该批次的数据并完成。
支持的数据格式
支持的格式如下:
压缩
用于 Kinesis 的 ClickPipes 会自动检测并解压压缩记录。不同于 Kafka 由客户端库透明地完成解压,Kinesis 传递的是原始字节数据——ClickPipes 会为你处理这些,无需任何配置。
支持以下压缩编解码器:
- gzip
- zstd
- lz4
- snappy (帧格式)
系统会根据每条记录中的 magic bytes 自动检测压缩方式。如果未发现已知的压缩签名,则该记录会被视为未压缩。检测到的压缩类型也会在 schema 推断期间显示,因此 UI 中的样本数据预览会正确显示解压后的数据。
受支持的数据类型
标准类型支持
ClickPipes 当前支持以下 ClickHouse 数据类型:
- 基础数值类型 - [U]Int8/16/32/64、Float32/64 和 BFloat16
- 大整数类型 - [U]Int128/256
- Decimal 类型
- Boolean
- String
- FixedString
- Date、Date32
- DateTime、DateTime64 (仅支持 UTC 时区)
- Enum8/Enum16
- UUID
- IPv4
- IPv6
- 所有 ClickHouse LowCardinality 类型
- 键和值可使用上述任意类型的 Map (包括 Nullable)
- 元素可使用上述任意类型的 Tuple 和 Array (包括 Nullable,仅支持一层嵌套)
- SimpleAggregateFunction 类型 (适用于 AggregatingMergeTree 或 SummingMergeTree 目标端)
Variant 类型支持
您可以为源数据流中的任何 JSON 字段手动指定 Variant 类型 (例如 Variant(String, Int64, DateTime)) 。
由于 ClickPipes 确定应使用哪种 Variant 子类型的方式所限,Variant 定义中只能使用一种整数类型或一种 DateTime 类型——例如,不支持 Variant(Int64, UInt32)。
JSON 类型支持
始终为 JSON 对象的 JSON 字段可以映射到 JSON 目标列。您需要手动将目标列调整为所需的 JSON 类型,包括任何固定路径或跳过路径。
Kinesis 虚拟列
Kinesis 数据流支持以下虚拟列。创建新的目标表时,可以使用 Add Column 按钮添加虚拟列。
| 名称 | 描述 | 推荐数据类型 |
|---|---|---|
| _key | Kinesis 分区键 | String |
| _timestamp | Kinesis 近似到达时间戳 (毫秒精度) | DateTime64(3) |
| _stream | Kinesis 数据流名称 | String |
| _sequence_number | Kinesis 序列号 | String |
| _raw_message | 完整 Kinesis 消息 | String |
在仅需完整 Kinesis JSON 记录的场景下,可以使用 _raw_message 字段 (例如使用 ClickHouse JsonExtract* 函数填充下游 materialized view) 。对于这类管道,删除所有“非虚拟”列可能会提升 ClickPipes 性能。
限制
- 不支持 DEFAULT。
- 默认情况下,使用最小 (XS) 副本规格运行时,单条消息大小上限为 16MB (未压缩) ;使用更大副本时,上限为 32MB (未压缩) 。超过此限制的消息将被拒绝并报错。如果您需要更大的消息,请联系支持团队。
性能
批量处理
ClickPipes 以批次方式将数据插入 ClickHouse。这样做是为了避免在数据库中创建过多的 parts,否则可能会导致集群出现性能问题。
当满足以下任一条件时,就会插入批次:
- 批次大小达到上限 (每 1GB 副本内存对应 100,000 行或 32MB)
- 批次的最长保留时间达到上限 (5 秒)
延迟
延迟 (即消息从发送到 Kinesis stream,到在 ClickHouse 中可用之间的时间) 取决于多种因素 (例如 Kinesis 延迟、网络延迟、消息大小/格式) 。上文所述的批处理也会对延迟产生影响。我们始终建议针对您的具体使用场景进行测试,以了解实际可预期的延迟水平。
如果您对低延迟有特定要求,请联系我们。
活跃分片
我们强烈建议将同时处于活跃状态的分片数量限制在满足吞吐量需求的范围内。对于 “On Demand” Kinesis 数据流,AWS 会根据吞吐量自动分配相应数量的分片; 但对于 “Provisioned” 数据流,配置过多分片不仅会导致下文所述的延迟问题,还会增加成本,因为这类 Kinesis 数据流按“每个分片”计费。
如果生产者应用持续向大量活跃分片写入数据,而您的管道规模又不足以高效处理这些分片,就可能产生延迟。根据 Kinesis 的吞吐量限制, ClickPipes 会为每个副本分配固定数量的“工作线程”来读取分片数据。例如,在最小规格下,一个 ClickPipes 副本会有 4 个这样的工作线程。如果生产者同时向 超过 4 个分片写入数据,那么在某个工作线程空闲之前,“额外”分片中的数据不会得到处理。特别是,如果该管道使用了 “enhanced fanout”,每个工作线程都会订阅 单个分片 5 分钟,并且在此期间无法读取任何其他分片。这会导致延迟出现以 5 分钟为倍数的“尖峰”。
扩缩容
ClickPipes for Kinesis 旨在同时支持水平和垂直扩缩容。默认情况下,我们会创建一个仅包含一个消费者的消费者组。这可以在创建 ClickPipe 时配置,也可以在之后的任意时间通过 Settings -> Advanced Settings -> 扩缩容 进行配置。
ClickPipes 通过跨可用区的分布式架构提供高可用性。 这要求将消费者数量至少扩展到两个。
无论当前运行的消费者数量是多少,系统在设计上都具备容错能力。 如果某个消费者或其底层基础设施发生故障, ClickPipe 会自动重启该消费者并继续处理消息。
身份验证
要访问Amazon Kinesis 数据流,您可以使用 IAM 凭证 或 IAM 角色。如需进一步了解如何设置 IAM 角色,您可以参阅本指南,了解如何设置可与 ClickHouse Cloud 配合使用的角色。