下表列出了 Managed Service for Apache Kafka 支持的 Kafka Connect 连接器类型。您可以使用这些连接器将 Apache Kafka 与您的应用和其他 Google Cloud服务集成。
| 连接器 | 说明 | 使用场景 |
|---|---|---|
| MirrorMaker 2.0 | 将主题和数据从一个 Kafka 集群复制到另一个 Kafka 集群。 | 数据复制、灾难恢复、数据迁移 |
| BigQuery 接收器 | 将数据从 Kafka 主题流式传输到 BigQuery 表。 | 数据仓储、分析 |
| Cloud Storage Sink | 将数据从 Kafka 主题流式传输到 Cloud Storage 存储桶。 | 数据湖注入、数据归档 |
| Pub/Sub Sink | 将数据从 Kafka 主题流式传输到 Pub/Sub 主题。 | 服务集成、实时通知 |
| Pub/Sub Source | 将消息从 Pub/Sub 订阅流式传输到 Kafka 主题。 | 实时数据注入、事件驱动型架构 |
转化者数量
转换器负责 Kafka 记录数据的序列化和反序列化。它们可在 Kafka 主题上找到的原始字节格式与 Kafka Connect 使用的内部结构化数据表示形式之间进行转换。
对于接收器连接器,转换器会将主题的传输格式中的数据反序列化为 Kafka Connect 内部数据格式,连接器使用该格式将数据写入目标系统。
对于源连接器,转换器会将 Kafka Connect 内部数据格式的数据序列化为 Kafka 主题的指定线路格式。
转换器可确保连接器以与外部系统兼容的格式读取或写入 Kafka 记录。
配置连接器时,请设置以下属性:
键转换器 (
key.converter):用于序列化和反序列化 Kafka 记录键的转换器。值转换器 (
value.converter):用于序列化和反序列化 Kafka 记录值的转换器。
如果您未指定转换器,则默认转换器类型为 org.apache.kafka.connect.converters.ByteArrayConverter,它会以原始字节格式传递数据。
支持的转化器
Managed Service for Apache Kafka 支持以下内置转换器:
| 转化者 | 格式 |
|---|---|
io.confluent.connect.avro.AvroConverter |
Apache Avro |
org.apache.kafka.connect.converters.BooleanConverter |
布尔值 |
org.apache.kafka.connect.converters.ByteArrayConverter |
字节数组 默认转换器类型。在两个系统之间保留消息的完整内容。 |
org.apache.kafka.connect.converters.DoubleConverter |
双精度型 |
org.apache.kafka.connect.converters.FloatConverter |
浮点数 |
org.apache.kafka.connect.converters.IntegerConverter |
整数 |
org.apache.kafka.connect.json.JsonConverter |
JSON 对于没有架构的 JSON 数据,也请设置 |
org.apache.kafka.connect.converters.LongConverter |
长 |
org.apache.kafka.connect.converters.ShortConverter |
短 |
org.apache.kafka.connect.storage.StringConverter |
字符串 |
转换器的选择取决于连接器类型以及您存储在 Kafka 中的数据。如需了解详情,请参阅相应连接器的文档。
Tasks
连接器通过创建并行运行的一个或多个任务来转移数据。如需设置连接器创建的任务数量上限,请设置连接器的 tasks.max 配置属性。连接器创建的任务数可能小于此值。
增加 tasks.max 的值可以提高吞吐量,但也会增加资源消耗(CPU 和内存)。最佳值取决于工作负载以及分配给 Connect 集群工作器的资源。对于接收器连接器,Kafka 主题分区的数量也会影响并行性。
任务重启政策
您可以设置连接器的任务重启政策,该政策决定了发生故障时的行为。连接器支持以下政策:
永不重启。连接器不会重启失败的任务。此政策是默认行为。这对于调试或在发生错误后需要人工干预的情况下非常有用。
使用指数退避算法重新启动。连接器会在延迟(称为退避期)后重新启动失败的任务。每次后续失败都会使延迟时间呈指数级增长。建议将此政策用于大多数生产工作负载。
如果您使用指数退避算法,还需为最小退避时间和最大退避时间设置值。退避时长下限应大于 60 秒,退避时长上限应小于 7200 秒。
转换和谓词
Managed Service for Apache Kafka 支持默认的 Kafka Connect 转换和谓词。
借助转换,您可以在将个别消息发送到 Managed Service for Apache Kafka(对于源连接器)或外部系统(对于接收器连接器)之前对其进行修改。您可以使用转换来遮盖敏感数据、添加时间戳或重命名字段。
借助谓词,您可以根据特定条件过滤数据,并根据消息属性确定转换适用于哪些消息。
例如,如需配置接收器连接器以忽略包含 DoNotProcess 标头键的消息,请添加以下配置:
transforms=dropMessage
transforms.dropMessage.type=org.apache.kafka.connect.transforms.Filter
transforms.dropMessage.predicate=hasKey
predicates=hasKey
predicates.hasKey.type=org.apache.kafka.connect.transforms.predicates.HasHeaderKey
predicates.hasKey.name=DoNotProcess
此配置会执行以下操作:
配置一个名为
hasKey且类型为org.apache.kafka.connect.transforms.predicates.HasHeaderKey的谓词。 此谓词可匹配包含键为DoNotProcess的标头的所有消息。配置名为
dropMessage且类型为org.apache.kafka.connect.transforms.Filter的转换。 此转换会舍弃与配置的谓词匹配的所有消息。将转换与谓词
hasKey相关联。这样可确保只有包含DoNotProcess标头键的消息会被转换丢弃。