ClickPipes 支持与 Schema Registry 集成,以解码采用 Avro 和 Protobuf 编码的记录值以及结构化 Kafka key。
Kafka ClickPipes 支持的 Schema Registry
Kafka ClickPipes 支持两类 schema registry:
- 兼容 Confluent 的 Schema Registry:任何与 Confluent Schema Registry API 兼容的 registry,例如 Confluent Schema Registry 本身和 Redpanda Schema Registry。支持 Avro 和 Protobuf。
- AWS Glue Schema Registry:适用于使用 AWS Glue SerDe 序列化的 Avro 数据,通常来自亚马逊 MSK。
ClickPipes 目前尚不支持 Azure Schema Registry。如需支持,请联系我们的团队。
兼容 Confluent 的 Schema Registry
配置
要在配置 ClickPipes 时集成 Schema Registry,必须使用以下方法之一:
- 提供 schema subject 的完整路径 (例如
https://registry.example.com/subjects/events)- 也可以通过在 URL 后附加
/versions/[version]来指定特定版本 (否则 ClickPipes 会获取最新版本) 。
- 也可以通过在 URL 后附加
- 提供 schema ID 的完整路径 (例如
https://registry.example.com/schemas/ids/1000) - 提供 Schema Registry 的根 URL (例如
https://registry.example.com)
网络连通性
ClickPipes 会通过您提供的 URL 使用 HTTPS 连接到 Schema Registry。Schema Registry 无需可通过公网访问。
如果您的 Kafka 消息代理是通过反向专用终结点 (AWS PrivateLink 或 GCP Private Service Connect) 访问的,Schema Registry 也可以使用相同的私有连接。ClickPipes 会通过反向专用终结点的私有 DNS 解析 registry 主机名,因此,只要其主机名解析到反向专用终结点的私有 IP 地址 (通过该端点的私有 DNS 支持或自定义私有 DNS 映射),与消息代理一起私有托管的 registry 就可以访问。
请注意以下事项:
- Schema Registry URL 必须使用
https://。 - 如果 registry 主机名解析为私有地址,则它必须能通过为 ClickPipe 选择的反向专用终结点访问;否则,设置期间的连通性检查将失败。
工作原理
ClickPipes 会动态从已配置的 Schema Registry 获取并应用 schema。
- 如果记录值中嵌入了 schema ID,则会使用该 ID 获取 schema。
- 如果记录值中未嵌入 schema ID,则会使用 ClickPipe 配置中指定的 schema ID 或 subject 名称来获取 schema。
- 如果记录值写入时未嵌入 schema ID,且 ClickPipe 配置中也未指定 schema ID 或 subject 名称,则不会获取 schema,该消息将被跳过,并在 ClickPipes 错误表中记录
SOURCE_SCHEMA_ERROR。 - 如果记录值不符合 schema,则该消息将被跳过,并在 ClickPipes 错误表中记录
DATA_PARSING_ERROR。 - 仅适用于 Protobuf schema:ClickPipes 会加载定义为依赖项的所有导入 schema。暂不支持带外部引用的 Avro schema。
配置了 _key.id 等字段的映射时,ClickPipes 会独立于记录值解析嵌入在 Kafka 键中的 schema ID。键可以使用不同的 schema ID,但必须与值使用相同的 registry 家族和序列化格式。已解析的键 schema 会被缓存,并会自动检测 schema 变更。
AWS Glue Schema Registry
如果您的 producer 使用 AWS Glue SerDe 序列化 Avro (例如,对亚马逊 MSK topic 使用 AWSKafkaAvroSerializer) ,ClickPipes 可以直接从 AWS Glue Schema Registry 解析这些 schema。Glue 使用的传输格式和 API 与兼容 Confluent 的 registry 不同,因此需要单独配置。
目前,AWS Glue Schema Registry configuration 只能通过 ClickHouse Cloud 控制台进行配置,不支持通过 ClickPipes API 或 Terraform provider 配置。
配置
在 ClickPipe 创建向导的 Kafka 连接步骤中,启用Schema Registry,并将Registry type设为AWS Glue:

| 字段 | 必填 | 说明 | 示例 |
|---|---|---|---|
| Registry type | 是 | 选择AWS Glue | AWS Glue |
| AWS 区域 | 是 | Glue registry 所在的区域。必须与 registry 的区域完全一致。 | us-east-1 |
| Registry name | 是 | Glue registry 的名称。解析到其他 registry 的 schema 会被拒绝,因此 ClickPipes 解析 schema version 时会发现拼写错误。 | my-glue-registry |
| IAM 角色 ARN | 有条件 | 用于访问 registry 的专用角色。如果 broker 使用 IAM 身份验证,则为可选;否则必填。 | arn:aws:iam::123456789012:role/ClickHouseAccessRole-glue |
无需配置 registry URL。Glue SerDe 生成的每条记录都携带自身 schema version 的 ID。ClickPipes 使用 glue:GetSchemaVersion 解析并缓存这些 schema;每个不同的 schema version 仅需一次 API 调用。schema evolution 会自动处理:当记录在 stream 中途切换到新的 schema version 时,会在首次遇到该版本时进行解析。
IAM 设置
请选择以下两种适合您环境的方案之一。对于亚马逊 MSK,通常使用选项 A。
选项 A:复用 broker 的 IAM 身份
如果您的 Kafka ClickPipe 已通过 IAM 向 MSK 进行身份验证,ClickPipes 会使用同一 IAM 身份读取 registry。请将 IAM 角色 ARN 字段留空,并将以下语句添加到该身份的权限中:
- **IAM role:**将该语句添加到为 MSK 配置的角色的权限策略中。
- **IAM credentials:**将该语句添加到与访问密钥关联的 IAM 主体的权限策略中。
{
"Version": "2012-10-17",
"Statement": [
{
"Sid": "ClickPipesGlueSchemaRegistryRead",
"Effect": "Allow",
"Action": ["glue:GetSchemaVersion"],
"Resource": "*"
}
]
}对于基于角色的身份验证,无需更改信任策略;为 MSK 配置的信任关系已涵盖此访问权限。IAM 凭证不使用角色信任策略。
选项 B:使用专用 registry 角色
如果您的 broker 不使用 IAM 身份验证 (SASL/SCRAM、SASL/PLAIN、mTLS) ,或者 registry 与 broker 位于不同的 AWS 账户中,请使用此选项。
获取 ClickHouse 服务 IAM 角色 ARN
打开服务,选择 设置 选项卡,滚动到 Network security 信息 部分,然后复制 服务角色 ID (IAM) 的值。该值是格式类似于 arn:aws:iam::123456789012:role/CH-S3-example-service-Role 的 ARN。下文将其称为 {ClickHouse_IAM_ARN}。部署在 AWS 上的每个 ClickHouse 服务都有各自的角色,因此每个服务的此值都不同。

创建 registry IAM 角色
在您的 AWS 账户中创建 IAM 角色。角色名称必须以 ClickHouseAccessRole- 开头。
配置信任策略
将 {ClickHouse_IAM_ARN} 替换为上一步获取的值。
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Principal": {
"AWS": "{ClickHouse_IAM_ARN}"
},
"Action": "sts:AssumeRole"
}
]
}配置权限策略
{
"Version": "2012-10-17",
"Statement": [
{
"Sid": "ClickPipesGlueSchemaRegistryRead",
"Effect": "Allow",
"Action": ["glue:GetSchemaVersion"],
"Resource": "*"
}
]
}配置 ClickPipe
将新角色的 ARN 粘贴到向导中的 IAM 角色 ARN 字段。
故障排查
| 错误 | 原因和解决方法 |
|---|---|
access denied retrieving schema version …: check the IAM role grants glue:GetSchemaVersion |
用于访问 registry 的 IAM 身份缺少 glue:GetSchemaVersion 权限。对于基于角色的访问,该角色的信任策略中可能也未包含您的服务角色 ID。请重新检查上述 IAM 设置。 |
… is not authorized to perform: sts:AssumeRole on resource: … |
信任策略中指定了错误的主体。该错误包含尝试承担角色的确切角色。请在信任策略中使用该值。 |
schema version … not found in Glue schema registry |
记录引用的 schema version 在配置的账户或区域中不存在。请确认 AWS 区域 与 registry 所在区域一致。 |
schema version … belongs to Glue registry "X", but the pipe is configured for registry "Y" |
您的 producer 在与管道配置不同的 registry 中注册了 schema。请更正 Registry 名称,或将 producer 指向正确的 registry。 |
the AWS Glue schema registry only supports the Avro format |
Glue 管道仅支持 Avro 格式。不支持通过 Glue SerDe 使用 JSON 和 Protobuf。 |
限制
- 仅支持 Avro。不支持通过 Glue SerDe 使用 JSON Schema 或 Protobuf。
- 仅支持 Kafka 源。Kinesis ClickPipes 无法使用 Glue registry。
Schema 映射
以下规则同时适用于与 Confluent 兼容的 registry 和 AWS Glue Schema Registry。它们规定了已获取的值 schema 与 ClickHouse 目标端表之间的映射,也适用于从带有 _key. 前缀的结构化键映射而来的记录或消息字段:
- 如果 schema 包含某个字段,但该字段未包含在 ClickHouse 目标端映射中,则该字段会被忽略。
- 如果 schema 缺少 ClickHouse 目标端映射中定义的某个字段,则 ClickHouse 列将填充为“零”值,例如 0 或空字符串。请注意,不支持
DEFAULT表达式。 - 如果 schema 字段与 ClickHouse 列不兼容,则该行/消息的插入会失败,且失败记录会写入 ClickPipes 错误表。请注意,系统支持一些隐式转换 (例如数值类型之间的转换) ,但并非全部都支持 (例如,Avro record 字段不能插入到
Int32ClickHouse 列中) 。