连接器概览

下表列出了 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 数据,也请设置 value.converter.schemas.enable=false

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

此配置会执行以下操作:

  1. 配置一个名为 hasKey 且类型为 org.apache.kafka.connect.transforms.predicates.HasHeaderKey 的谓词。 此谓词可匹配包含键为 DoNotProcess 的标头的所有消息。

  2. 配置名为 dropMessage 且类型为 org.apache.kafka.connect.transforms.Filter 的转换。 此转换会舍弃与配置的谓词匹配的所有消息。

  3. 将转换与谓词 hasKey 相关联。这样可确保只有包含 DoNotProcess 标头键的消息会被转换丢弃。

接下来怎么做?

Apache Kafka® 是 Apache Software Foundation 或其关联公司在美国和/或其他国家/地区的注册商标。