如何在Kafka Connect Sink中指定Kafka主题的分区
你的理解没有错误,Kafka Connect Sink的底层确实是封装了Kafka Consumer实现,理论上完全支持指定读取特定分区,你找不到对应配置的核心原因是:Kafka Connect框架本身默认没有开放这类手动配置项,主流第三方Sink连接器(包括你使用的DataStax Apache Kafka Connector)也没有额外实现该配置的暴露能力。
Kafka Connect的原生设计逻辑是优先保证自动化运维能力,默认使用消费者组的协作分区分配策略,自动将主题分区均匀分配给连接器的多个Task实例,实现自动负载均衡、故障转移,因此原生通用Sink配置里没有设计手动指定读取分区的参数,大部分第三方连接器为了适配原生的调度逻辑,也不会额外开发手动指定分区的配置项。
如果确实有指定分区读取的需求,可以通过两种方式实现:
方案1:零代码改动的替代方案(优先推荐)
你可以先通过Kafka流处理组件(比如Kafka Streams、Flink)或者专门的转换连接器,把你需要读取的特定分区的数据路由到一个独立的新主题中,再让DataStax Sink连接器只消费这个新主题即可。
该方案完全适配现有连接器的所有能力,不需要修改任何连接器代码,也不会丢失原生的自动负载均衡、故障转移特性,操作成本最低。
方案2:自定义连接器实现指定分区
如果你一定要在连接器层面直接指定分区,需要修改Sink连接器的源码:在SinkTask实现类的open方法中,手动调用消费者的assign()方法,传入你指定的分区列表,覆盖框架默认的分区分配逻辑即可。
Confluent官方开发者文档中提到的动态分区支持,指的就是连接器可以自定义这部分分配逻辑,但是该能力需要连接器开发者主动实现对应的配置项暴露给用户,这也是你目前在各类公开配置文档中找不到对应参数的原因。
注意:如果选择手动指定分区的方案,会丢失Kafka Connect原生的自动故障转移、分区重平衡能力,当负责处理指定分区的Task异常退出时,框架不会自动将该分区调度给其他健康的Task处理,需要额外的人工运维成本,非特殊需求不推荐使用。
内容的提问来源于stack exchange,提问作者RyanQuey

